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
9 changes: 0 additions & 9 deletions backend/app/core/container.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@
DatabaseProvider,
DLQProvider,
DLQWorkerProvider,
EventProvider,
EventReplayProvider,
EventReplayWorkerProvider,
K8sWorkerProvider,
Expand Down Expand Up @@ -52,7 +51,6 @@ def create_app_container(settings: Settings, broker: KafkaBroker) -> AsyncContai
RepositoryProvider(),
MessagingProvider(),
DLQProvider(),
EventProvider(),
SagaOrchestratorProvider(),
KafkaServicesProvider(),
SSEProvider(),
Expand All @@ -79,7 +77,6 @@ def create_result_processor_container(settings: Settings, broker: KafkaBroker) -
CoreServicesProvider(),
MetricsProvider(),
RepositoryProvider(),
EventProvider(),
MessagingProvider(),
DLQProvider(),
ResultProcessorProvider(),
Expand All @@ -99,7 +96,6 @@ def create_coordinator_container(settings: Settings, broker: KafkaBroker) -> Asy
RepositoryProvider(),
MessagingProvider(),
DLQProvider(),
EventProvider(),
CoordinatorProvider(),
context={Settings: settings, KafkaBroker: broker},
)
Expand All @@ -117,7 +113,6 @@ def create_k8s_worker_container(settings: Settings, broker: KafkaBroker) -> Asyn
RepositoryProvider(),
MessagingProvider(),
DLQProvider(),
EventProvider(),
KubernetesProvider(),
K8sWorkerProvider(),
context={Settings: settings, KafkaBroker: broker},
Expand All @@ -136,7 +131,6 @@ def create_pod_monitor_container(settings: Settings, broker: KafkaBroker) -> Asy
RepositoryProvider(),
MessagingProvider(),
DLQProvider(),
EventProvider(),
KafkaServicesProvider(),
KubernetesProvider(),
PodMonitorProvider(),
Expand All @@ -159,7 +153,6 @@ def create_saga_orchestrator_container(settings: Settings, broker: KafkaBroker)
RepositoryProvider(),
MessagingProvider(),
DLQProvider(),
EventProvider(),
SagaWorkerProvider(),
context={Settings: settings, KafkaBroker: broker},
)
Expand All @@ -180,7 +173,6 @@ def create_event_replay_container(settings: Settings, broker: KafkaBroker) -> As
RepositoryProvider(),
MessagingProvider(),
DLQProvider(),
EventProvider(),
EventReplayWorkerProvider(),
context={Settings: settings, KafkaBroker: broker},
)
Expand All @@ -202,6 +194,5 @@ def create_dlq_processor_container(settings: Settings, broker: KafkaBroker) -> A
RepositoryProvider(),
MessagingProvider(),
DLQWorkerProvider(),
EventProvider(),
context={Settings: settings, KafkaBroker: broker},
)
80 changes: 63 additions & 17 deletions backend/app/core/providers.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,10 +49,11 @@
from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository
from app.db.repositories.user_settings_repository import UserSettingsRepository
from app.dlq.manager import DLQManager
from app.dlq.models import RetryPolicy, RetryStrategy
from app.domain.enums.kafka import KafkaTopic
from app.domain.rate_limit import RateLimitConfig
from app.domain.saga.models import SagaConfig
from app.events.core import UnifiedProducer
from app.events.schema.schema_registry import SchemaRegistryManager
from app.services.admin import AdminEventsService, AdminSettingsService, AdminUserService
from app.services.auth_service import AuthService
from app.services.coordinator.coordinator import ExecutionCoordinator
Expand Down Expand Up @@ -180,13 +181,12 @@ class MessagingProvider(Provider):
def get_unified_producer(
self,
broker: KafkaBroker,
schema_registry: SchemaRegistryManager,
event_repository: EventRepository,
logger: logging.Logger,
settings: Settings,
event_metrics: EventMetrics,
) -> UnifiedProducer:
return UnifiedProducer(broker, schema_registry, event_repository, logger, settings, event_metrics)
return UnifiedProducer(broker, event_repository, logger, settings, event_metrics)

@provide
def get_idempotency_repository(self, redis_client: redis.Redis) -> RedisIdempotencyRepository:
Expand All @@ -199,6 +199,61 @@ def get_idempotency_manager(
return IdempotencyManager(IdempotencyConfig(), repo, logger, database_metrics)


def _default_retry_policy() -> RetryPolicy:
"""Default retry policy for DLQ messages."""
return RetryPolicy(
topic="default",
strategy=RetryStrategy.EXPONENTIAL_BACKOFF,
max_retries=4,
base_delay_seconds=60,
max_delay_seconds=1800,
retry_multiplier=2.5,
)


def _default_retry_policies(prefix: str) -> dict[str, RetryPolicy]:
"""Topic-specific retry policies for DLQ.

Keys must match message.original_topic (full prefixed topic name).
"""
execution_events = f"{prefix}{KafkaTopic.EXECUTION_EVENTS}"
pod_events = f"{prefix}{KafkaTopic.POD_EVENTS}"
saga_commands = f"{prefix}{KafkaTopic.SAGA_COMMANDS}"
execution_results = f"{prefix}{KafkaTopic.EXECUTION_RESULTS}"

return {
execution_events: RetryPolicy(
topic=execution_events,
strategy=RetryStrategy.EXPONENTIAL_BACKOFF,
max_retries=5,
base_delay_seconds=30,
max_delay_seconds=300,
retry_multiplier=2.0,
),
pod_events: RetryPolicy(
topic=pod_events,
strategy=RetryStrategy.EXPONENTIAL_BACKOFF,
max_retries=3,
base_delay_seconds=60,
max_delay_seconds=600,
retry_multiplier=3.0,
),
saga_commands: RetryPolicy(
topic=saga_commands,
strategy=RetryStrategy.EXPONENTIAL_BACKOFF,
max_retries=5,
base_delay_seconds=30,
max_delay_seconds=300,
retry_multiplier=2.0,
),
execution_results: RetryPolicy(
topic=execution_results,
strategy=RetryStrategy.IMMEDIATE,
max_retries=3,
),
}

Comment thread
HardMax71 marked this conversation as resolved.

class DLQProvider(Provider):
"""Provides DLQManager without scheduling. Used by all containers except the DLQ worker."""

Expand All @@ -209,26 +264,25 @@ def get_dlq_manager(
self,
broker: KafkaBroker,
settings: Settings,
schema_registry: SchemaRegistryManager,
logger: logging.Logger,
dlq_metrics: DLQMetrics,
repository: DLQRepository,
) -> DLQManager:
return DLQManager(
settings=settings,
broker=broker,
schema_registry=schema_registry,
logger=logger,
dlq_metrics=dlq_metrics,
repository=repository,
default_retry_policy=_default_retry_policy(),
Comment thread
HardMax71 marked this conversation as resolved.
retry_policies=_default_retry_policies(settings.KAFKA_TOPIC_PREFIX),
)


class DLQWorkerProvider(Provider):
"""Provides DLQManager with APScheduler-managed retry monitoring.

Used by the DLQ worker container only. DLQManager configures its own
retry policies and filters; the provider only handles scheduling.
Used by the DLQ worker container only.
"""

scope = Scope.APP
Expand All @@ -238,7 +292,6 @@ async def get_dlq_manager(
self,
broker: KafkaBroker,
settings: Settings,
schema_registry: SchemaRegistryManager,
logger: logging.Logger,
dlq_metrics: DLQMetrics,
repository: DLQRepository,
Expand All @@ -247,10 +300,11 @@ async def get_dlq_manager(
manager = DLQManager(
settings=settings,
broker=broker,
schema_registry=schema_registry,
logger=logger,
dlq_metrics=dlq_metrics,
repository=repository,
default_retry_policy=_default_retry_policy(),
retry_policies=_default_retry_policies(settings.KAFKA_TOPIC_PREFIX),
)

scheduler = AsyncIOScheduler()
Expand All @@ -272,14 +326,6 @@ async def get_dlq_manager(
logger.info("DLQManager retry monitor stopped")


class EventProvider(Provider):
scope = Scope.APP

@provide
def get_schema_registry(self, settings: Settings, logger: logging.Logger) -> SchemaRegistryManager:
return SchemaRegistryManager(settings, logger)


class KubernetesProvider(Provider):
scope = Scope.APP

Expand Down
Loading
Loading