TorbatYar/backend/services/beauty_business/app/events/publisher.py
Mortezakoohjani d579d0b142 feat(loyalty): add Loyalty Platform Frontend module
Add frontend/modules/loyalty with types, API client, design system, feature pages and thin App Router routes under app/loyalty/. BFF proxy at app/api/loyalty/. Include loyalty frontend docs.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-27 10:50:55 +03:30

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 BeautyBusinessEventType
from app.models.outbox import OutboxEvent
class EventPublisher(Protocol):
def publish(
self,
*,
event_type: BeautyBusinessEventType,
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: BeautyBusinessEventType,
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: BeautyBusinessEventType,
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