Files
security-alert-center/backend/app/api/v1/events.py
T
PTah 20fa7e2c27 feat: webhook notification channel with UI and ingest dispatch
Add webhook config (DB/env), JSON POST on events/problems, settings API,
Settings UI section, and notify_dispatch for multi-channel ingest.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-05-29 16:09:40 +10:00

186 lines
6.2 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.problems import maybe_create_problem
from app.services.schema_validate import validate_event_payload
from app.services.notify_dispatch import notify_event, notify_problem
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)
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))
count_stmt = count_stmt.where(Host.hostname.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 = [
EventSummary(
id=e.id,
event_id=e.event_id,
host_id=e.host_id,
hostname=e.host.hostname,
occurred_at=e.occurred_at,
received_at=e.received_at,
category=e.category,
type=e.type,
severity=e.severity,
title=e.title,
summary=e.summary,
)
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,
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,
)