Add the Healthcare platform backend (phases 13.0–13.7) with Alembic migration revision fix, production deploy script, validation reports, shared UI exports, and loyalty list page client directives so Next.js build succeeds. Co-authored-by: Cursor <cursoragent@cursor.com>
122 lines
3.6 KiB
Python
122 lines
3.6 KiB
Python
"""Delivery event publisher — in-memory + transactional outbox (ADR-006)."""
|
|
from __future__ import annotations
|
|
|
|
from typing import Any, Protocol
|
|
from uuid import UUID, uuid4
|
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from shared.events import EventEnvelope, EventStatus
|
|
|
|
from app.core.config import settings
|
|
from app.events.types import HealthcareEventType
|
|
from app.models.outbox import OutboxEvent
|
|
|
|
|
|
class EventPublisher(Protocol):
|
|
def publish(
|
|
self,
|
|
*,
|
|
event_type: HealthcareEventType,
|
|
aggregate_type: str,
|
|
aggregate_id: UUID,
|
|
tenant_id: UUID,
|
|
payload: dict[str, Any] | None = None,
|
|
) -> EventEnvelope: ...
|
|
|
|
|
|
class InMemoryEventPublisher:
|
|
"""Records published envelopes for tests and local verification."""
|
|
|
|
def __init__(self) -> None:
|
|
self.published: list[EventEnvelope] = []
|
|
|
|
def publish(
|
|
self,
|
|
*,
|
|
event_type: HealthcareEventType,
|
|
aggregate_type: str,
|
|
aggregate_id: UUID,
|
|
tenant_id: UUID,
|
|
payload: dict[str, Any] | None = None,
|
|
) -> EventEnvelope:
|
|
envelope = EventEnvelope(
|
|
event_id=uuid4(),
|
|
event_type=event_type.value,
|
|
aggregate_type=aggregate_type,
|
|
aggregate_id=str(aggregate_id),
|
|
tenant_id=tenant_id,
|
|
source_service=settings.service_name,
|
|
payload=payload or {},
|
|
)
|
|
self.published.append(envelope)
|
|
return envelope
|
|
|
|
def record(self, envelope: EventEnvelope) -> EventEnvelope:
|
|
self.published.append(envelope)
|
|
return envelope
|
|
|
|
|
|
class TransactionalEventPublisher:
|
|
"""Persist outbox row in the current transaction; mirror to memory in tests."""
|
|
|
|
def __init__(
|
|
self, session: AsyncSession, memory: InMemoryEventPublisher | None = None
|
|
) -> None:
|
|
self.session = session
|
|
if memory is not None:
|
|
self.memory = memory
|
|
elif settings.environment == "test":
|
|
self.memory = get_event_publisher()
|
|
else:
|
|
self.memory = None
|
|
|
|
async def publish(
|
|
self,
|
|
*,
|
|
event_type: HealthcareEventType,
|
|
aggregate_type: str,
|
|
aggregate_id: UUID,
|
|
tenant_id: UUID,
|
|
payload: dict[str, Any] | None = None,
|
|
) -> EventEnvelope:
|
|
envelope = EventEnvelope(
|
|
event_id=uuid4(),
|
|
event_type=event_type.value,
|
|
aggregate_type=aggregate_type,
|
|
aggregate_id=str(aggregate_id),
|
|
tenant_id=tenant_id,
|
|
source_service=settings.service_name,
|
|
payload=payload or {},
|
|
)
|
|
row = OutboxEvent(
|
|
tenant_id=tenant_id,
|
|
event_type=envelope.event_type,
|
|
aggregate_type=aggregate_type,
|
|
aggregate_id=str(aggregate_id),
|
|
payload={
|
|
"event_id": str(envelope.event_id),
|
|
"source_service": settings.service_name,
|
|
**(payload or {}),
|
|
},
|
|
status=EventStatus.PENDING,
|
|
)
|
|
self.session.add(row)
|
|
await self.session.flush()
|
|
if self.memory is not None:
|
|
self.memory.record(envelope)
|
|
return envelope
|
|
|
|
|
|
_default_publisher = InMemoryEventPublisher()
|
|
|
|
|
|
def get_event_publisher() -> InMemoryEventPublisher:
|
|
return _default_publisher
|
|
|
|
|
|
def reset_event_publisher() -> InMemoryEventPublisher:
|
|
global _default_publisher
|
|
_default_publisher = InMemoryEventPublisher()
|
|
return _default_publisher
|