from __future__ import annotations import threading from fastapi import Depends, FastAPI, HTTPException from sqlalchemy import select from services.shared.core import Role from services.shared.db import get_session from services.shared.event_bus import ( dispatch_outbox_batch, event_bus_enabled, outbox_item_from_row, poll_forever, retry_outbox_event, ) from services.shared.models import EventOutboxItem, HealthResponse from services.shared.security import require_roles from services.shared.sql_init import init_sql_schema from services.shared.sql_models import EventOutboxRow app = FastAPI(title="event-bus-service", version="1.0.0") init_sql_schema() @app.get("/health", response_model=HealthResponse) def health() -> HealthResponse: suffix = "enabled" if event_bus_enabled() else "disabled" return HealthResponse(status="ok", service=f"event-bus-service ({suffix})") @app.get("/bus/outbox", response_model=list[EventOutboxItem]) def list_outbox( status: str | None = None, limit: int = 100, _: dict = Depends(require_roles(Role.ADMIN)), ) -> list[EventOutboxItem]: session = get_session() try: stmt = select(EventOutboxRow).order_by(EventOutboxRow.id.desc()) if status: stmt = stmt.where(EventOutboxRow.status == status) rows = session.execute(stmt.limit(max(limit, 1))).scalars().all() return [outbox_item_from_row(row) for row in rows] finally: session.close() @app.get("/bus/outbox/{event_id}", response_model=EventOutboxItem) def get_outbox_event( event_id: str, _: dict = Depends(require_roles(Role.ADMIN)), ) -> EventOutboxItem: session = get_session() try: row = session.execute(select(EventOutboxRow).where(EventOutboxRow.event_id == event_id)).scalar_one_or_none() if not row: raise HTTPException(status_code=404, detail="Outbox event not found") return outbox_item_from_row(row) finally: session.close() @app.post("/bus/outbox/{event_id}/retry", response_model=EventOutboxItem) def retry_outbox( event_id: str, _: dict = Depends(require_roles(Role.ADMIN)), ) -> EventOutboxItem: session = get_session() try: row = session.execute(select(EventOutboxRow).where(EventOutboxRow.event_id == event_id)).scalar_one_or_none() if not row: raise HTTPException(status_code=404, detail="Outbox event not found") retry_outbox_event(session, row) session.commit() session.refresh(row) return outbox_item_from_row(row) finally: session.close() def _dispatch_once() -> None: session = get_session() try: dispatch_outbox_batch(session) finally: session.close() @app.on_event("startup") def _startup() -> None: if not event_bus_enabled(): return threading.Thread(target=lambda: poll_forever(_dispatch_once), daemon=True).start()