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 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.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) @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, )