"""HTTP delivery adapters for durable platform events and invitation mail.""" from __future__ import annotations import json from urllib.request import Request, urlopen from user_engine.domain import OutboxEvent _MAIL_EVENTS = {"family_member.invited", "family_invitation.resent"} class HTTPOutboxDeliveryAdapter: """Deliver one outbox event to idempotent platform HTTP lanes.""" def __init__( self, *, event_url: str, bearer_token: str, mail_url: str | None = None, timeout_seconds: float = 5.0, ) -> None: self.event_url = event_url self.mail_url = mail_url self.bearer_token = bearer_token self.timeout_seconds = timeout_seconds def __call__(self, event: OutboxEvent) -> None: envelope = { "id": event.event_id, "type": event.event_type, "source": "user-engine", "subject": event.aggregate_id, "tenant": event.tenant, "correlation_id": event.correlation_id, "occurred_at": event.occurred_at.isoformat(), "data": dict(event.payload), } if self.mail_url and event.event_type in _MAIL_EVENTS: self._post(self.mail_url, envelope) event_envelope = dict(envelope) event_data = dict(envelope["data"]) if "primary_email" in event_data: event_data["recipient_present"] = True del event_data["primary_email"] event_envelope["data"] = event_data self._post(self.event_url, event_envelope) def _post(self, url: str, envelope: dict[str, object]) -> None: with urlopen( Request( url, data=json.dumps(envelope).encode(), headers={ "Authorization": f"Bearer {self.bearer_token}", "Content-Type": "application/json", "Idempotency-Key": str(envelope["id"]), }, method="POST", ), timeout=self.timeout_seconds, ) as response: if response.status < 200 or response.status >= 300: raise RuntimeError(f"delivery rejected with HTTP {response.status}")