Files
group_fqcd_jr/app/worker/outbox_worker.py

92 lines
3.3 KiB
Python

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