from __future__ import annotations import re from dataclasses import dataclass from datetime import datetime, timedelta, timezone from typing import Any, Callable @dataclass(frozen=True) class ClaimedEvent: event_id: str event_type: str payload: dict[str, Any] attempts: int @dataclass(frozen=True) class DispatchResult: status: str attempts: int available_at: datetime | None = None last_error: str | None = None record_consumption: bool = False def retry_delay( attempts: int, *, base_seconds: int = 2, maximum_seconds: int = 300, ) -> timedelta: seconds = min(maximum_seconds, base_seconds * (2 ** max(0, attempts - 1))) return timedelta(seconds=seconds) def _redact_error(error: Exception) -> str: message = str(error) message = re.sub( r"(?i)(secret[-_ ]?token|api[-_ ]?key|authorization|bearer)(?:[=: ]+\S+)?", "[redacted]", message, ) return message[:1000] def dispatch_event( event: ClaimedEvent, handler: Callable[[dict[str, Any]], None], *, already_processed: bool = False, max_attempts: int = 5, now: datetime | None = None, ) -> DispatchResult: if already_processed: return DispatchResult(status="published", attempts=event.attempts) now = now or datetime.now(timezone.utc) try: handler(event.payload) except Exception as exc: attempts = event.attempts + 1 if attempts >= max_attempts: return DispatchResult( status="dead_letter", attempts=attempts, last_error=_redact_error(exc), ) return DispatchResult( status="pending", attempts=attempts, available_at=now + retry_delay(attempts), last_error=_redact_error(exc), ) return DispatchResult( status="published", attempts=event.attempts + 1, record_consumption=True, )