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_response_fields, command_to_dict, get_command_by_uuid, ) from app.services.rdg_winrm_actions import execute_logoff_via_winrm, execute_qwinsta_via_winrm 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, db) 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 target: str | None = None client_hostname: str | None = None internal_ip: str | None = None def _agent_command_response(cmd) -> AgentCommandResponse: data = command_to_dict(cmd) extra = command_response_fields(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"), target=extra.get("target"), client_hostname=extra.get("client_hostname"), internal_ip=extra.get("internal_ip"), ) 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 = execute_qwinsta_via_winrm(db, event, requested_by=str(user)) db.commit() return _agent_command_response(cmd) @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 = execute_logoff_via_winrm( db, event, session_id=body.session_id, requested_by=str(user), ) db.commit() return _agent_command_response(cmd) @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") return _agent_command_response(cmd) @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, db) return EventDetail( **base.model_dump(), details=event.details, raw=event.raw, dedup_key=event.dedup_key, correlation_id=event.correlation_id, payload=event.payload, )