Files
security-alert-center/backend/app/api/v1/events.py
T
PTah 2aec488226 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.
2026-06-19 23:34:01 +10:00

285 lines
9.6 KiB
Python

import logging
from datetime import datetime, timezone
from typing import Any
from fastapi import APIRouter, Body, Depends, HTTPException, Query
from fastapi.responses import JSONResponse
from pydantic import BaseModel
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, 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
from app.services.schema_validate import validate_event_payload
from app.services.notify_dispatch import (
AUTH_LOGIN_SUCCESS_TYPES,
DAILY_REPORT_EVENT_TYPES,
LIFECYCLE_EVENT_TYPE,
RDG_CONNECTION_TYPES,
notify_auth_login,
notify_daily_report,
notify_event,
notify_lifecycle,
notify_problem,
notify_rdg_connection,
)
router = APIRouter(prefix="/events", tags=["events"])
logger = logging.getLogger("sac.ingest")
class IngestResponse(BaseModel):
status: str
event_id: str
created: bool
sac_event_url: str
problem_id: int | None = None
def _ingest_response(event: Event, *, created: bool) -> IngestResponse:
settings = get_settings()
base = settings.sac_public_url.rstrip("/")
return IngestResponse(
status="created" if created else "duplicate",
event_id=event.event_id,
created=created,
sac_event_url=f"{base}/api/v1/events/{event.id}",
problem_id=None,
)
@router.post("", response_model=IngestResponse)
def post_event(
payload: dict[str, Any] = Body(...),
db: Session = Depends(get_db),
_api_key: str = Depends(get_api_key_auth),
) -> JSONResponse:
errors = validate_event_payload(payload)
if errors:
event_id_hint = payload.get("event_id") if isinstance(payload.get("event_id"), str) else None
logger.warning(
"ingest rejected status=422 event_id=%s errors=%s",
event_id_hint,
errors[:5],
)
raise HTTPException(status_code=422, detail={"schema_errors": errors[:20]})
event, created = ingest_event(db, payload)
problem = None
problem_created = False
if created:
problem, problem_created = maybe_create_problem(db, event)
if event.type in DAILY_REPORT_EVENT_TYPES:
notify_daily_report(event, db=db)
elif event.type == LIFECYCLE_EVENT_TYPE:
notify_lifecycle(event, db=db)
elif event.type in AUTH_LOGIN_SUCCESS_TYPES:
notify_auth_login(event, db=db)
elif event.type in RDG_CONNECTION_TYPES:
notify_rdg_connection(event, db=db)
else:
notify_event(event, db=db)
if problem is not None and problem_created:
notify_problem(problem, event, db=db)
logger.info("ingest created event_id=%s type=%s host_id=%s", event.event_id, event.type, event.host_id)
else:
logger.info("ingest duplicate event_id=%s", event.event_id)
db.commit()
body = _ingest_response(event, created=created)
if problem is not None:
body.problem_id = problem.id
status_code = 201 if created else 409
return JSONResponse(status_code=status_code, content=body.model_dump(mode="json"))
def _parse_optional_dt(value: str | None, *, end_of_day: bool = False) -> datetime | None:
if not value:
return None
raw = value.strip()
if len(raw) == 10 and raw[4] == "-" and raw[7] == "-":
dt = datetime.fromisoformat(raw)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
if end_of_day:
return dt.replace(hour=23, minute=59, second=59, microsecond=999999)
return dt.replace(hour=0, minute=0, second=0, microsecond=0)
return datetime.fromisoformat(raw.replace("Z", "+00:00"))
@router.get("", response_model=EventListResponse)
def list_events(
page: int = Query(1, ge=1),
page_size: int = Query(50, ge=1, le=200),
severity: str | None = None,
type: str | None = None,
host_id: int | None = None,
hostname: str | None = None,
from_time: str | None = Query(None, alias="from"),
to_time: str | None = Query(None, alias="to"),
q: str | None = None,
db: Session = Depends(get_db),
_user: str = Depends(get_current_user),
) -> EventListResponse:
stmt = select(Event).join(Host).options(joinedload(Event.host))
count_stmt = select(func.count()).select_from(Event).join(Host)
if severity:
stmt = stmt.where(Event.severity == severity)
count_stmt = count_stmt.where(Event.severity == severity)
if type:
stmt = stmt.where(Event.type == type)
count_stmt = count_stmt.where(Event.type == type)
if host_id is not None:
stmt = stmt.where(Event.host_id == host_id)
count_stmt = count_stmt.where(Event.host_id == host_id)
if hostname:
like = f"%{hostname}%"
stmt = stmt.where(Host.hostname.ilike(like) | Host.display_name.ilike(like))
count_stmt = count_stmt.where(Host.hostname.ilike(like) | Host.display_name.ilike(like))
dt_from = _parse_optional_dt(from_time)
dt_to = _parse_optional_dt(to_time, end_of_day=True)
if dt_from:
stmt = stmt.where(Event.occurred_at >= dt_from)
count_stmt = count_stmt.where(Event.occurred_at >= dt_from)
if dt_to:
stmt = stmt.where(Event.occurred_at <= dt_to)
count_stmt = count_stmt.where(Event.occurred_at <= dt_to)
if q:
like = f"%{q}%"
stmt = stmt.where(Event.summary.ilike(like) | Event.title.ilike(like))
count_stmt = count_stmt.where(Event.summary.ilike(like) | Event.title.ilike(like))
total = db.scalar(count_stmt) or 0
rows = db.scalars(
stmt.order_by(Event.occurred_at.desc())
.offset((page - 1) * page_size)
.limit(page_size)
).all()
items = [event_to_summary(e) for e in rows]
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,
db: Session = Depends(get_db),
_user: str = Depends(get_current_user),
) -> EventDetail:
event = db.scalar(
select(Event).where(Event.id == event_db_id).options(joinedload(Event.host))
)
if event is None:
raise HTTPException(status_code=404, detail="Event not found")
base = event_to_summary(event)
return EventDetail(
**base.model_dump(),
details=event.details,
raw=event.raw,
dedup_key=event.dedup_key,
correlation_id=event.correlation_id,
payload=event.payload,
)