from __future__ import annotations from datetime import datetime, timezone from typing import Any from sqlalchemy import text from app.core.common.identifiers import new_governance_uid from app.core.events.consumer import ClaimedEvent, DispatchResult def enqueue_outbox( session: Any, *, aggregate_type: str, aggregate_id: str, event_type: str, payload: dict[str, Any], correlation_id: str | None = None, event_id: str | None = None, ) -> str: event_id = event_id or new_governance_uid() session.execute( text( """ INSERT INTO public.outbox_events ( event_id, aggregate_type, aggregate_id, event_type, payload, correlation_id ) VALUES ( CAST(:event_id AS uuid), :aggregate_type, :aggregate_id, :event_type, CAST(:payload AS jsonb), :correlation_id ) """ ), { "event_id": event_id, "aggregate_type": aggregate_type, "aggregate_id": aggregate_id, "event_type": event_type, "payload": __import__("json").dumps(payload, ensure_ascii=False), "correlation_id": correlation_id, }, ) return event_id def claim_outbox(session: Any, limit: int = 50) -> list[ClaimedEvent]: rows = session.execute( text( """ WITH selected AS ( SELECT event_id FROM public.outbox_events WHERE status = 'pending' AND available_at <= CURRENT_TIMESTAMP ORDER BY available_at, created_at FOR UPDATE SKIP LOCKED LIMIT :limit ) UPDATE public.outbox_events AS event SET status = 'processing' FROM selected WHERE event.event_id = selected.event_id RETURNING event.event_id::text, event.event_type, event.payload, event.attempts """ ), {"limit": int(limit)}, ) return [ ClaimedEvent( event_id=row[0], event_type=row[1], payload=dict(row[2]), attempts=int(row[3]), ) for row in rows ] def consumed(session: Any, consumer_name: str, event_id: str) -> bool: return bool( session.execute( text( "SELECT 1 FROM public.outbox_consumptions " "WHERE consumer_name = :consumer_name " "AND event_id = CAST(:event_id AS uuid)" ), {"consumer_name": consumer_name, "event_id": event_id}, ).scalar() ) def apply_dispatch_result( session: Any, *, consumer_name: str, event: ClaimedEvent, result: DispatchResult, ) -> None: if result.record_consumption: session.execute( text( """ INSERT INTO public.outbox_consumptions (consumer_name, event_id) VALUES (:consumer_name, CAST(:event_id AS uuid)) ON CONFLICT (consumer_name, event_id) DO NOTHING """ ), {"consumer_name": consumer_name, "event_id": event.event_id}, ) session.execute( text( """ UPDATE public.outbox_events SET status = :status, attempts = :attempts, available_at = COALESCE(:available_at, available_at), published_at = CASE WHEN :status = 'published' THEN CURRENT_TIMESTAMP ELSE published_at END, last_error = :last_error WHERE event_id = CAST(:event_id AS uuid) """ ), { "status": result.status, "attempts": result.attempts, "available_at": result.available_at, "last_error": result.last_error, "event_id": event.event_id, }, ) def outbox_health(session: Any) -> dict[str, Any]: row = session.execute( text( """ SELECT COUNT(*) FILTER (WHERE status IN ('pending', 'processing')) AS backlog, COUNT(*) FILTER (WHERE status IN ('failed', 'dead_letter')) AS failed, EXTRACT(EPOCH FROM ( CURRENT_TIMESTAMP - MIN(created_at) FILTER (WHERE status IN ('pending', 'processing')) )) AS oldest_age_seconds FROM public.outbox_events """ ) ).one() return { "backlog": int(row[0] or 0), "failed": int(row[1] or 0), "oldest_age_seconds": float(row[2] or 0), "checked_at": datetime.now(timezone.utc).isoformat(), } def reset_stale_processing(session: Any, *, older_than_seconds: int = 300) -> int: result = session.execute( text( """ UPDATE public.outbox_events SET status = 'pending', available_at = CURRENT_TIMESTAMP WHERE status = 'processing' AND updated_at < CURRENT_TIMESTAMP - (:older_than_seconds * INTERVAL '1 second') """ ), {"older_than_seconds": int(older_than_seconds)}, ) return int(result.rowcount or 0)