Files
security-alert-center/backend/app/api/v1/events.py
T
2026-06-01 10:21:43 +10:00

193 lines
6.6 KiB
Python

import logging
from datetime import datetime
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) -> datetime | None:
if not value:
return None
return datetime.fromisoformat(value.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)
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")
return EventDetail(
id=event.id,
event_id=event.event_id,
host_id=event.host_id,
hostname=event.host.hostname,
display_name=event.host.display_name,
occurred_at=event.occurred_at,
received_at=event.received_at,
category=event.category,
type=event.type,
severity=event.severity,
title=event.title,
summary=event.summary,
details=event.details,
raw=event.raw,
dedup_key=event.dedup_key,
correlation_id=event.correlation_id,
payload=event.payload,
)