Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 5 additions & 11 deletions backend/app/core/dishka_lifespan.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

import asyncio
import logging
from contextlib import AsyncExitStack, asynccontextmanager
from contextlib import asynccontextmanager
from typing import AsyncGenerator

import redis.asyncio as redis
Expand All @@ -15,7 +15,7 @@
from app.core.startup import initialize_rate_limits
from app.core.tracing import init_tracing
from app.db.docs import ALL_DOCUMENTS
from app.events.event_store_consumer import EventStoreConsumer
from app.events.event_store import EventStore
from app.events.schema.schema_registry import SchemaRegistryManager, initialize_event_schemas
from app.services.notification_scheduler import NotificationScheduler
from app.services.notification_service import NotificationService
Expand Down Expand Up @@ -83,15 +83,15 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
database,
redis_client,
rate_limit_metrics,
event_store_consumer,
_event_store,
_notification_service,
_notification_scheduler,
) = await asyncio.gather(
container.get(SchemaRegistryManager),
container.get(Database),
container.get(redis.Redis),
container.get(RateLimitMetrics),
container.get(EventStoreConsumer),
container.get(EventStore),
container.get(NotificationService),
container.get(NotificationScheduler),
)
Expand All @@ -104,10 +104,4 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
)
logger.info("Infrastructure initialized (schemas, beanie, rate limits)")

# Phase 3: Start lifecycle-managed services
# EventStoreConsumer requires explicit __aenter__; all other services are managed by DI providers
async with AsyncExitStack() as stack:
stack.push_async_callback(event_store_consumer.aclose)
await event_store_consumer.__aenter__()
logger.info("EventStoreConsumer started")
yield
yield
62 changes: 0 additions & 62 deletions backend/app/core/lifecycle.py

This file was deleted.

71 changes: 49 additions & 22 deletions backend/app/core/providers.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,9 +50,15 @@
from app.domain.enums.kafka import CONSUMER_GROUP_SUBSCRIPTIONS, GroupId
from app.domain.idempotency import KeyStrategy
from app.domain.saga.models import SagaConfig
from app.events.core import ConsumerConfig, EventDispatcher, ProducerMetrics, UnifiedConsumer, UnifiedProducer
from app.events.core import (
ConsumerConfig,
EventDispatcher,
ProducerMetrics,
UnifiedConsumer,
UnifiedProducer,
create_dlq_error_handler,
)
from app.events.event_store import EventStore, create_event_store
from app.events.event_store_consumer import EventStoreConsumer, create_event_store_consumer
from app.events.schema.schema_registry import SchemaRegistryManager
from app.infrastructure.kafka.topics import get_all_topics
from app.services.admin import AdminEventsService, AdminSettingsService, AdminUserService
Expand Down Expand Up @@ -236,33 +242,54 @@ def get_schema_registry(self, settings: Settings, logger: logging.Logger) -> Sch

@provide
async def get_event_store(
self, schema_registry: SchemaRegistryManager, logger: logging.Logger, event_metrics: EventMetrics
) -> EventStore:
return create_event_store(
schema_registry=schema_registry, logger=logger, event_metrics=event_metrics, ttl_days=90
)

@provide
async def get_event_store_consumer(
self,
event_store: EventStore,
schema_registry: SchemaRegistryManager,
settings: Settings,
kafka_producer: UnifiedProducer,
logger: logging.Logger,
event_metrics: EventMetrics,
) -> AsyncIterator[EventStoreConsumer]:
) -> AsyncIterator[EventStore]:
event_store = create_event_store(
schema_registry=schema_registry, logger=logger, event_metrics=event_metrics, ttl_days=90
)

dispatcher = EventDispatcher(logger=logger)
for event_type in EventType:
dispatcher.register_handler(event_type, event_store.store_event)

config = ConsumerConfig(
bootstrap_servers=settings.KAFKA_BOOTSTRAP_SERVERS,
group_id=GroupId.EVENT_STORE_CONSUMER,
enable_auto_commit=False,
max_poll_records=100,
session_timeout_ms=settings.KAFKA_SESSION_TIMEOUT_MS,
heartbeat_interval_ms=settings.KAFKA_HEARTBEAT_INTERVAL_MS,
max_poll_interval_ms=settings.KAFKA_MAX_POLL_INTERVAL_MS,
request_timeout_ms=settings.KAFKA_REQUEST_TIMEOUT_MS,
)
kafka_consumer = UnifiedConsumer(
config,
event_dispatcher=dispatcher,
schema_registry=schema_registry,
settings=settings,
logger=logger,
event_metrics=event_metrics,
)

dlq_handler = create_dlq_error_handler(
producer=kafka_producer, logger=logger, max_retries=3,
)
kafka_consumer.register_error_callback(dlq_handler)

topics = get_all_topics()
async with create_event_store_consumer(
event_store=event_store,
topics=list(topics),
schema_registry_manager=schema_registry,
settings=settings,
producer=kafka_producer,
logger=logger,
event_metrics=event_metrics,
) as consumer:
yield consumer
await kafka_consumer.start(list(topics))
logger.info(f"Event store consumer started for topics: {list(topics)}")
Comment thread
HardMax71 marked this conversation as resolved.

try:
yield event_store
finally:
await kafka_consumer.stop()
logger.info("Event store consumer stopped")



Expand Down
6 changes: 3 additions & 3 deletions backend/app/events/core/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ def __init__(
self._running = False
self._metrics = ConsumerMetrics()
self._event_metrics = event_metrics
self._error_callback: "Callable[[Exception, DomainEvent], Awaitable[None]] | None" = None
self._error_callback: "Callable[[Exception, DomainEvent, str], Awaitable[None]] | None" = None
self._consume_task: asyncio.Task[None] | None = None
self._topic_prefix = settings.KAFKA_TOPIC_PREFIX

Expand Down Expand Up @@ -192,10 +192,10 @@ async def _process_message(self, message: Any) -> None:
topic=topic, consumer_group=self._config.group_id, error_type=type(e).__name__
)
if self._error_callback:
await self._error_callback(e, event)
await self._error_callback(e, event, topic)
raise

def register_error_callback(self, callback: Callable[[Exception, DomainEvent], Awaitable[None]]) -> None:
def register_error_callback(self, callback: Callable[[Exception, DomainEvent, str], Awaitable[None]]) -> None:
self._error_callback = callback

@property
Expand Down
6 changes: 3 additions & 3 deletions backend/app/events/core/dispatcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,13 @@
import logging
from collections import defaultdict
from collections.abc import Awaitable, Callable
from typing import TypeAlias, TypeVar
from typing import Any, TypeAlias, TypeVar

from app.domain.enums.events import EventType
from app.domain.events.typed import DomainEvent

T = TypeVar("T", bound=DomainEvent)
EventHandler: TypeAlias = Callable[[DomainEvent], Awaitable[None]]
EventHandler: TypeAlias = Callable[[DomainEvent], Awaitable[Any]]


class EventDispatcher:
Expand All @@ -26,7 +26,7 @@ def __init__(self, logger: logging.Logger) -> None:
self.logger = logger

# Map event types to their handlers
self._handlers: dict[EventType, list[Callable[[DomainEvent], Awaitable[None]]]] = defaultdict(list)
self._handlers: dict[EventType, list[EventHandler]] = defaultdict(list)

def _wrap_handler(self, handler: EventHandler) -> EventHandler:
"""Hook for subclasses to wrap handlers at registration time."""
Expand Down
16 changes: 8 additions & 8 deletions backend/app/events/core/dlq_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,31 +7,31 @@


def create_dlq_error_handler(
producer: UnifiedProducer, original_topic: str, logger: logging.Logger, max_retries: int = 3
) -> Callable[[Exception, DomainEvent], Awaitable[None]]:
producer: UnifiedProducer, logger: logging.Logger, max_retries: int = 3
) -> Callable[[Exception, DomainEvent, str], Awaitable[None]]:
"""Create an error handler that sends failed events to DLQ after max retries."""
retry_counts: dict[str, int] = {}

async def handle_error_with_dlq(error: Exception, event: DomainEvent) -> None:
async def handle_error_with_dlq(error: Exception, event: DomainEvent, topic: str) -> None:
event_id = event.event_id or "unknown"
retry_count = retry_counts.get(event_id, 0)
retry_counts[event_id] = retry_count + 1
logger.error(f"Error processing {event_id}: {error}. Retry {retry_count + 1}/{max_retries}", exc_info=True)
if retry_count >= max_retries:
logger.warning(f"Event {event_id} exceeded max retries. Sending to DLQ.")
await producer.send_to_dlq(event, original_topic, error, retry_count)
await producer.send_to_dlq(event, topic, error, retry_count)
retry_counts.pop(event_id, None)

return handle_error_with_dlq


def create_immediate_dlq_handler(
producer: UnifiedProducer, original_topic: str, logger: logging.Logger
) -> Callable[[Exception, DomainEvent], Awaitable[None]]:
producer: UnifiedProducer, logger: logging.Logger
) -> Callable[[Exception, DomainEvent, str], Awaitable[None]]:
"""Create an error handler that immediately sends failed events to DLQ."""

async def handle_error_immediate_dlq(error: Exception, event: DomainEvent) -> None:
async def handle_error_immediate_dlq(error: Exception, event: DomainEvent, topic: str) -> None:
logger.error(f"Critical error processing {event.event_id}: {error}. Sending to DLQ.", exc_info=True)
await producer.send_to_dlq(event, original_topic, error, 0)
await producer.send_to_dlq(event, topic, error, 0)

return handle_error_immediate_dlq
Loading
Loading