| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879 |
- 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,
- )
|