feat: RDG session flap rule, qwinsta/logoff queue and dashboard UI (0.9.13)

Detect 302→303 within 1–10s, Problem with 30s dedup, rdg_flap flag.
Agent command poll API, admin qwinsta/logoff from Overview. Migration 016.
This commit is contained in:
PTah
2026-06-19 23:34:01 +10:00
parent b328d32f97
commit 2aec488226
24 changed files with 1012 additions and 25 deletions
+61
View File
@@ -0,0 +1,61 @@
from typing import Any
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel, Field
from sqlalchemy.orm import Session
from app.auth.api_key import get_api_key_auth
from app.database import get_db
from app.services.agent_commands import (
command_to_dict,
complete_command,
get_command_for_host_poll,
resolve_host_by_agent_instance_id,
)
router = APIRouter(prefix="/agent", tags=["agent"])
class AgentCommandResultBody(BaseModel):
status: str = Field(pattern="^(completed|failed)$")
stdout: str | None = None
stderr: str | None = None
class AgentCommandsPollResponse(BaseModel):
commands: list[dict[str, Any]]
@router.get("/commands", response_model=AgentCommandsPollResponse)
def poll_agent_commands(
agent_instance_id: str = Query(..., min_length=1),
db: Session = Depends(get_db),
_api_key: str = Depends(get_api_key_auth),
) -> AgentCommandsPollResponse:
host = resolve_host_by_agent_instance_id(db, agent_instance_id)
if host is None:
raise HTTPException(status_code=404, detail="Unknown agent_instance_id")
pending = get_command_for_host_poll(db, host)
return AgentCommandsPollResponse(
commands=[command_to_dict(c, include_run_as=True) for c in pending]
)
@router.post("/commands/{command_uuid}/result")
def post_agent_command_result(
command_uuid: str,
body: AgentCommandResultBody,
db: Session = Depends(get_db),
_api_key: str = Depends(get_api_key_auth),
) -> dict[str, str]:
cmd = complete_command(
db,
command_uuid,
status=body.status,
stdout=body.stdout,
stderr=body.stderr,
)
if cmd is None:
raise HTTPException(status_code=404, detail="Command not found")
db.commit()
return {"status": cmd.status, "command_uuid": cmd.command_uuid}
+95 -1
View File
@@ -9,11 +9,17 @@ from sqlalchemy import func, select
from sqlalchemy.orm import Session, joinedload
from app.auth.api_key import get_api_key_auth
from app.auth.jwt_auth import get_current_user
from app.auth.jwt_auth import get_current_user, require_admin
from app.config import get_settings
from app.database import get_db
from app.models import Event, Host
from app.schemas.list_models import EventDetail, EventListResponse, EventSummary
from app.services.agent_commands import (
command_to_dict,
get_command_by_uuid,
queue_logoff,
queue_qwinsta,
)
from app.services.ingest import ingest_event
from app.services.event_summary import event_to_summary
from app.services.problems import maybe_create_problem
@@ -168,6 +174,94 @@ def list_events(
return EventListResponse(items=items, total=total, page=page, page_size=page_size)
class AgentCommandResponse(BaseModel):
command_uuid: str
command_type: str
status: str
result_stdout: str | None = None
result_stderr: str | None = None
created_at: str | None = None
completed_at: str | None = None
class LogoffActionBody(BaseModel):
session_id: int
@router.post("/{event_db_id}/actions/qwinsta", response_model=AgentCommandResponse)
def post_event_qwinsta(
event_db_id: int,
db: Session = Depends(get_db),
user=Depends(require_admin),
) -> AgentCommandResponse:
event = db.get(Event, event_db_id)
if event is None:
raise HTTPException(status_code=404, detail="Event not found")
cmd = queue_qwinsta(db, event, requested_by=str(user))
db.commit()
data = command_to_dict(cmd)
return AgentCommandResponse(
command_uuid=data["id"],
command_type=data["type"],
status=data["status"],
result_stdout=data.get("result_stdout"),
result_stderr=data.get("result_stderr"),
created_at=data.get("created_at"),
completed_at=data.get("completed_at"),
)
@router.post("/{event_db_id}/actions/logoff", response_model=AgentCommandResponse)
def post_event_logoff(
event_db_id: int,
body: LogoffActionBody,
db: Session = Depends(get_db),
user=Depends(require_admin),
) -> AgentCommandResponse:
event = db.get(Event, event_db_id)
if event is None:
raise HTTPException(status_code=404, detail="Event not found")
cmd = queue_logoff(
db,
event,
session_id=body.session_id,
requested_by=str(user),
)
db.commit()
data = command_to_dict(cmd)
return AgentCommandResponse(
command_uuid=data["id"],
command_type=data["type"],
status=data["status"],
result_stdout=data.get("result_stdout"),
result_stderr=data.get("result_stderr"),
created_at=data.get("created_at"),
completed_at=data.get("completed_at"),
)
@router.get("/{event_db_id}/actions/{command_uuid}", response_model=AgentCommandResponse)
def get_event_action_status(
event_db_id: int,
command_uuid: str,
db: Session = Depends(get_db),
_user=Depends(get_current_user),
) -> AgentCommandResponse:
cmd = get_command_by_uuid(db, command_uuid)
if cmd is None or cmd.event_id != event_db_id:
raise HTTPException(status_code=404, detail="Command not found")
data = command_to_dict(cmd)
return AgentCommandResponse(
command_uuid=data["id"],
command_type=data["type"],
status=data["status"],
result_stdout=data.get("result_stdout"),
result_stderr=data.get("result_stderr"),
created_at=data.get("created_at"),
completed_at=data.get("completed_at"),
)
@router.get("/{event_db_id}", response_model=EventDetail)
def get_event(
event_db_id: int,
+2 -1
View File
@@ -1,10 +1,11 @@
from fastapi import APIRouter
from app.api.v1 import auth, dashboards, events, health, hosts, mobile, problems, settings, stream, system, users
from app.api.v1 import agent, auth, dashboards, events, health, hosts, mobile, problems, settings, stream, system, users
api_router = APIRouter()
api_router.include_router(health.router)
api_router.include_router(auth.router)
api_router.include_router(agent.router)
api_router.include_router(events.router)
api_router.include_router(hosts.router)
api_router.include_router(problems.router)
+9
View File
@@ -102,6 +102,15 @@ class Settings(BaseSettings):
sac_privilege_spike_window_minutes: int = 10
sac_privilege_spike_threshold: int = 10
# rule:rdg_session_flap — RD Gateway 302→303 within window
sac_rdg_flap_window_min_sec: int = 1
sac_rdg_flap_window_max_sec: int = 10
sac_rdg_flap_dedup_sec: int = 30
# Windows admin for agent qwinsta/logoff (domain-wide)
sac_win_admin_user: str = ""
sac_win_admin_password: str = ""
# Retention (app.jobs.retention / systemd timer)
sac_events_retention_days: int = 90
sac_problems_retention_days: int = 180
+2
View File
@@ -1,3 +1,4 @@
from app.models.agent_command import AgentCommand
from app.models.api_key import ApiKey
from app.models.event import Event
from app.models.host import Host
@@ -16,6 +17,7 @@ from app.models.mobile_refresh_token import MobileRefreshToken
__all__ = [
"ApiKey",
"AgentCommand",
"Event",
"Host",
"NotificationChannel",
+27
View File
@@ -0,0 +1,27 @@
from datetime import datetime
from sqlalchemy import DateTime, ForeignKey, String, Text, func
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.orm import Mapped, mapped_column, relationship
from app.database import Base
class AgentCommand(Base):
__tablename__ = "agent_commands"
id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
command_uuid: Mapped[str] = mapped_column(String(36), unique=True, index=True)
host_id: Mapped[int] = mapped_column(ForeignKey("hosts.id", ondelete="CASCADE"), index=True)
event_id: Mapped[int | None] = mapped_column(ForeignKey("events.id", ondelete="SET NULL"), nullable=True)
command_type: Mapped[str] = mapped_column(String(32))
params: Mapped[dict | None] = mapped_column(JSONB, default=dict)
status: Mapped[str] = mapped_column(String(16), default="pending", index=True)
result_stdout: Mapped[str | None] = mapped_column(Text)
result_stderr: Mapped[str | None] = mapped_column(Text)
requested_by: Mapped[str | None] = mapped_column(String(128))
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())
completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
host: Mapped["Host"] = relationship(back_populates="agent_commands")
event: Mapped["Event | None"] = relationship()
+1
View File
@@ -28,3 +28,4 @@ class Host(Base):
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())
events: Mapped[list["Event"]] = relationship(back_populates="host")
agent_commands: Mapped[list["AgentCommand"]] = relationship(back_populates="host")
+1
View File
@@ -55,6 +55,7 @@ class EventSummary(BaseModel):
title: str
summary: str
actor_user: str | None = None
rdg_flap: bool = False
model_config = {"from_attributes": True}
+144
View File
@@ -0,0 +1,144 @@
"""Queue and resolve agent commands (qwinsta, logoff)."""
from __future__ import annotations
import uuid
from datetime import datetime, timezone
from fastapi import HTTPException
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.config import get_settings
from app.models import AgentCommand, Event, Host
from app.services.rdg_session_flap import event_has_rdg_flap
def _win_admin_configured() -> bool:
settings = get_settings()
return bool(settings.sac_win_admin_user.strip() and settings.sac_win_admin_password.strip())
def win_admin_run_as() -> dict[str, str] | None:
if not _win_admin_configured():
return None
settings = get_settings()
return {
"user": settings.sac_win_admin_user.strip(),
"password": settings.sac_win_admin_password,
}
def command_to_dict(cmd: AgentCommand, *, include_run_as: bool = False) -> dict:
out = {
"id": cmd.command_uuid,
"type": cmd.command_type,
"params": cmd.params if isinstance(cmd.params, dict) else {},
"status": cmd.status,
"result_stdout": cmd.result_stdout,
"result_stderr": cmd.result_stderr,
"created_at": cmd.created_at.isoformat() if cmd.created_at else None,
"completed_at": cmd.completed_at.isoformat() if cmd.completed_at else None,
}
if include_run_as and cmd.status == "pending":
run_as = win_admin_run_as()
if run_as:
out["run_as"] = run_as
return out
def queue_qwinsta(db: Session, event: Event, *, requested_by: str) -> AgentCommand:
if not event_has_rdg_flap(event):
raise HTTPException(status_code=400, detail="Event is not flagged as RDG session flap")
if not _win_admin_configured():
raise HTTPException(
status_code=503,
detail="SAC_WIN_ADMIN_USER / SAC_WIN_ADMIN_PASSWORD not configured",
)
details = event.details if isinstance(event.details, dict) else {}
user = details.get("user")
cmd = AgentCommand(
command_uuid=str(uuid.uuid4()),
host_id=event.host_id,
event_id=event.id,
command_type="qwinsta",
params={"user": user} if user else {},
status="pending",
requested_by=requested_by,
)
db.add(cmd)
db.flush()
return cmd
def queue_logoff(
db: Session,
event: Event,
*,
session_id: int,
requested_by: str,
) -> AgentCommand:
if not event_has_rdg_flap(event):
raise HTTPException(status_code=400, detail="Event is not flagged as RDG session flap")
if not _win_admin_configured():
raise HTTPException(
status_code=503,
detail="SAC_WIN_ADMIN_USER / SAC_WIN_ADMIN_PASSWORD not configured",
)
cmd = AgentCommand(
command_uuid=str(uuid.uuid4()),
host_id=event.host_id,
event_id=event.id,
command_type="logoff",
params={"session_id": session_id},
status="pending",
requested_by=requested_by,
)
db.add(cmd)
db.flush()
return cmd
def get_command_for_host_poll(db: Session, host: Host, *, limit: int = 5) -> list[AgentCommand]:
return list(
db.scalars(
select(AgentCommand)
.where(
AgentCommand.host_id == host.id,
AgentCommand.status == "pending",
)
.order_by(AgentCommand.created_at.asc())
.limit(limit)
).all()
)
def resolve_host_by_agent_instance_id(db: Session, agent_instance_id: str) -> Host | None:
if not agent_instance_id.strip():
return None
return db.scalar(select(Host).where(Host.agent_instance_id == agent_instance_id.strip()))
def complete_command(
db: Session,
command_uuid: str,
*,
status: str,
stdout: str | None = None,
stderr: str | None = None,
) -> AgentCommand | None:
cmd = db.scalar(select(AgentCommand).where(AgentCommand.command_uuid == command_uuid))
if cmd is None:
return None
if cmd.status != "pending":
return cmd
cmd.status = status
cmd.result_stdout = stdout
cmd.result_stderr = stderr
cmd.completed_at = datetime.now(timezone.utc)
db.flush()
return cmd
def get_command_by_uuid(db: Session, command_uuid: str) -> AgentCommand | None:
return db.scalar(select(AgentCommand).where(AgentCommand.command_uuid == command_uuid))
+2
View File
@@ -1,6 +1,7 @@
from app.models.event import Event
from app.schemas.list_models import EventSummary
from app.services.event_actor_user import extract_event_actor_user
from app.services.rdg_session_flap import event_has_rdg_flap
def event_to_summary(event: Event) -> EventSummary:
@@ -20,4 +21,5 @@ def event_to_summary(event: Event) -> EventSummary:
title=event.title,
summary=event.summary,
actor_user=extract_event_actor_user(event.type, event.details),
rdg_flap=event_has_rdg_flap(event),
)
+5 -1
View File
@@ -18,6 +18,7 @@ PRIVILEGE_SUDO_TYPE = "privilege.sudo.command"
RULE_BRUTE_FORCE = "rule:brute_force_burst"
RULE_PRIVILEGE_SPIKE = "rule:privilege_spike"
RULE_HOST_SILENCE = "rule:host_silence"
RULE_RDG_SESSION_FLAP = "rule:rdg_session_flap"
@dataclass(frozen=True)
@@ -204,8 +205,11 @@ def build_fingerprint(host_id: int, correlation_type: str, rule_id: str, suffix:
def pick_rule_match(db: Session, event: Event) -> RuleMatch | None:
"""Первое сработавшее правило (приоритет: burst → spike → silence → immediate)."""
"""Первое сработавшее правило (приоритет: rdg flap → burst → spike → silence → immediate)."""
from app.services.rdg_session_flap import evaluate_rdg_session_flap
for evaluator in (
evaluate_rdg_session_flap,
evaluate_brute_force_burst,
evaluate_privilege_spike,
evaluate_host_silence,
+6 -4
View File
@@ -1,6 +1,6 @@
"""Auto-create Problems from ingested events (rules + correlation)."""
from datetime import datetime, timezone
from datetime import datetime, timedelta, timezone
from sqlalchemy import select
from sqlalchemy.orm import Session
@@ -11,6 +11,7 @@ from app.services.host_health import HEARTBEAT_TYPE
from app.services.problem_rules import (
RuleMatch,
RULE_HOST_SILENCE,
RULE_RDG_SESSION_FLAP,
build_fingerprint,
pick_rule_match,
resolve_host_silence_problems,
@@ -18,8 +19,6 @@ from app.services.problem_rules import (
def _correlation_cutoff(now: datetime) -> datetime:
from datetime import timedelta
minutes = get_settings().sac_problem_correlation_window_minutes
return now - timedelta(minutes=minutes)
@@ -66,7 +65,10 @@ def open_or_append_problem(
match.fingerprint_suffix,
)
now = datetime.now(timezone.utc)
cutoff = _correlation_cutoff(now)
if match.rule_id == RULE_RDG_SESSION_FLAP:
cutoff = now - timedelta(seconds=get_settings().sac_rdg_flap_dedup_sec)
else:
cutoff = _correlation_cutoff(now)
open_problem = db.scalar(
select(Problem)
+120
View File
@@ -0,0 +1,120 @@
"""RDG 302→303 session flap: detect, flag event, build Problem match."""
from __future__ import annotations
from datetime import datetime, timedelta, timezone
from sqlalchemy import select
from sqlalchemy.orm import Session
from sqlalchemy.orm.attributes import flag_modified
from app.config import get_settings
from app.models import Event
from app.services.problem_rules import RULE_RDG_SESSION_FLAP, RuleMatch
RDG_SUCCESS_TYPE = "rdg.connection.success"
RDG_END_TYPES = frozenset({"rdg.connection.disconnected", "rdg.connection.failed"})
def _event_user(event: Event) -> str:
details = event.details if isinstance(event.details, dict) else {}
user = details.get("user")
return str(user).strip() if user else ""
def _event_internal_ip(event: Event) -> str:
details = event.details if isinstance(event.details, dict) else {}
for key in ("internal_ip", "client_ip", "ip_address"):
val = details.get(key)
if val:
return str(val).strip()
return ""
def _users_match(end_event: Event, success_event: Event) -> bool:
return _event_user(end_event) != "" and _event_user(end_event) == _event_user(success_event)
def _internal_ips_compatible(end_event: Event, success_event: Event) -> bool:
end_ip = _event_internal_ip(end_event)
success_ip = _event_internal_ip(success_event)
if not end_ip or not success_ip:
return True
return end_ip == success_ip
def _as_utc(dt: datetime) -> datetime:
if dt.tzinfo is None:
return dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
def find_rdg_success_before_end(db: Session, end_event: Event) -> Event | None:
if end_event.type not in RDG_END_TYPES:
return None
settings = get_settings()
min_sec = settings.sac_rdg_flap_window_min_sec
max_sec = settings.sac_rdg_flap_window_max_sec
end_at = _as_utc(end_event.occurred_at)
window_start = end_at - timedelta(seconds=max_sec)
window_end = end_at - timedelta(seconds=min_sec)
candidates = db.scalars(
select(Event)
.where(
Event.host_id == end_event.host_id,
Event.type == RDG_SUCCESS_TYPE,
Event.occurred_at >= window_start,
Event.occurred_at <= window_end,
Event.id != end_event.id,
)
.order_by(Event.occurred_at.desc())
).all()
for prior in candidates:
if not _users_match(end_event, prior):
continue
if not _internal_ips_compatible(end_event, prior):
continue
prior_at = _as_utc(prior.occurred_at)
delta = (end_at - prior_at).total_seconds()
if min_sec <= delta <= max_sec:
return prior
return None
def mark_rdg_flap(event: Event, *, pair_event: Event) -> None:
details = dict(event.details) if isinstance(event.details, dict) else {}
details["rdg_flap"] = True
details["rdg_flap_pair_event_id"] = pair_event.id
event.details = details
flag_modified(event, "details")
def evaluate_rdg_session_flap(db: Session, event: Event) -> RuleMatch | None:
prior = find_rdg_success_before_end(db, event)
if prior is None:
return None
mark_rdg_flap(event, pair_event=prior)
user = _event_user(event)
internal_ip = _event_internal_ip(event)
ip_note = f", client {internal_ip}" if internal_ip else ""
delta_sec = int((_as_utc(event.occurred_at) - _as_utc(prior.occurred_at)).total_seconds())
return RuleMatch(
rule_id=RULE_RDG_SESSION_FLAP,
correlation_type="rdg.session.flap",
fingerprint_suffix=f"u{user}:ip{internal_ip or 'any'}",
title=f"RDG session flap: {user}",
summary=(
f"302→303 за {delta_sec} с "
f"({user}{ip_note}). Возможна зависшая сессия на ПК пользователя — qwinsta/logoff."
),
severity="warning",
)
def event_has_rdg_flap(event: Event) -> bool:
details = event.details if isinstance(event.details, dict) else {}
return details.get("rdg_flap") is True
+1 -1
View File
@@ -1,5 +1,5 @@
"""Единый источник версии SAC (API, health, логи, OpenAPI)."""
APP_NAME = "Security Alert Center"
APP_VERSION = "0.9.12"
APP_VERSION = "0.9.13"
APP_VERSION_LABEL = f"{APP_NAME} v.{APP_VERSION}"