from collections.abc import Awaitable, Callable from datetime import UTC, datetime, timedelta from typing import Any from sqlalchemy import select, true from sqlalchemy.ext.asyncio import AsyncSession from app.model.platform import DomainEventOutbox, OutboxDelivery class OutboxWorker: """One dispatcher per event type, not a fan-out transport. Handlers must use this session for local writes and must not commit it. External side effects still require their own idempotency protocol. """ def __init__( self, session: AsyncSession, handlers: dict[str, Callable[[dict[str, Any]], Awaitable[None]]], ) -> None: self.session = session self.handlers = handlers async def publish_one(self, *, aggregate_id: str | None = None) -> bool: if not self.handlers: return False now = datetime.now(UTC).replace(tzinfo=None) event = await self.session.scalar( select(DomainEventOutbox) .where( DomainEventOutbox.status.in_({"pending", "failed"}), DomainEventOutbox.event_type.in_(tuple(self.handlers)), DomainEventOutbox.aggregate_id == aggregate_id if aggregate_id else true(), ( DomainEventOutbox.next_retry_at.is_(None) | (DomainEventOutbox.next_retry_at <= now) ), ) .order_by(DomainEventOutbox.id) .limit(1) .with_for_update(skip_locked=True) ) if event is None: await self.session.rollback() return False handler = self.handlers.get(event.event_type) if handler is None: event.status = "dead" event.last_error = "no handler registered" await self._commit() return False try: delivery = await self.session.scalar( select(OutboxDelivery).where( OutboxDelivery.event_id == event.event_id, OutboxDelivery.consumer_name == event.event_type, ) ) except BaseException: await self.session.rollback() raise try: if delivery is None: await handler(event.payload) self.session.add(OutboxDelivery( event_id=event.event_id, consumer_name=event.event_type, delivered_at=datetime.now(UTC).replace(tzinfo=None), )) await self.session.flush() except Exception as exc: event.retry_count += 1 event.last_error = type(exc).__name__ event.status = "dead" if event.retry_count >= 5 else "failed" event.next_retry_at = now + timedelta(seconds=min(300, 2**event.retry_count)) else: event.status = "published" event.published_at = datetime.now(UTC).replace(tzinfo=None) event.last_error = None event.next_retry_at = None event.updated_at = datetime.now(UTC).replace(tzinfo=None) await self._commit() return True async def _commit(self) -> None: try: await self.session.commit() except BaseException: await self.session.rollback() raise