from __future__ import annotations import json import threading from fastapi import FastAPI from sqlalchemy import select from services.shared.core import new_id, utc_now_iso from services.shared.db import get_session from services.shared.event_bus import ( consume_one_message, consumer_enabled, event_bus_audit_queue, event_bus_enabled, inbox_seen, poll_forever, record_inbox, ) from services.shared.models import AuditEvent, AuditEventIn, HealthResponse from services.shared.sql_init import init_sql_schema from services.shared.sql_models import AuditEventRow app = FastAPI(title="audit-service", version="1.0.0") init_sql_schema() def _to_out(row: AuditEventRow) -> AuditEvent: return AuditEvent( event_id=row.event_id, actor=row.actor, action=row.action, entity=row.entity, metadata=json.loads(row.metadata_json or "{}"), created_at=row.created_at, ) @app.get("/health", response_model=HealthResponse) def health() -> HealthResponse: return HealthResponse(status="ok", service="audit-service") @app.post("/audit/events", response_model=AuditEvent) def create_event(payload: AuditEventIn) -> AuditEvent: session = get_session() try: row = AuditEventRow( event_id=new_id("aud"), actor=payload.actor, action=payload.action, entity=payload.entity, metadata_json=json.dumps(payload.metadata, ensure_ascii=False), created_at=utc_now_iso(), ) session.add(row) session.commit() session.refresh(row) return _to_out(row) finally: session.close() @app.get("/audit/events", response_model=list[AuditEvent]) def list_events(actor: str | None = None, action: str | None = None, limit: int = 100) -> list[AuditEvent]: session = get_session() try: stmt = select(AuditEventRow).order_by(AuditEventRow.id.asc()) if actor: stmt = stmt.where(AuditEventRow.actor == actor) if action: stmt = stmt.where(AuditEventRow.action == action) rows = session.execute(stmt).scalars().all() rows = rows[-limit:] return [_to_out(r) for r in rows] finally: session.close() def _handle_event(envelope: dict) -> None: session = get_session() try: event_id = str(envelope.get("event_id") or "").strip() event_type = str(envelope.get("event_type") or "").strip() if not event_id or not event_type: raise ValueError("Invalid event envelope") if inbox_seen(session, "audit-service", event_id): return payload = envelope.get("payload") if isinstance(envelope.get("payload"), dict) else {} row = AuditEventRow( event_id=new_id("aud"), actor=str(payload.get("actor") or envelope.get("producer") or "system"), action=event_type, entity=str(envelope.get("entity_type") or "event"), metadata_json=json.dumps(envelope, ensure_ascii=False), created_at=utc_now_iso(), ) session.add(row) record_inbox( session, consumer_name="audit-service", event_id=event_id, event_type=event_type, status="processed", ) session.commit() finally: session.close() def _consume_once() -> None: consume_one_message(event_bus_audit_queue(), _handle_event) @app.on_event("startup") def _startup() -> None: if not event_bus_enabled() or not consumer_enabled(): return threading.Thread(target=lambda: poll_forever(_consume_once), daemon=True).start()