92 lines
3.3 KiB
Python
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
|