From b143605890365a0bd09f69ec362465632d5676de Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Thu, 5 Feb 2026 19:00:45 +0100 Subject: [PATCH 1/2] dlq manager direct calls of broker --- backend/app/core/providers.py | 59 +++++++- backend/app/dlq/manager.py | 155 +++++++--------------- backend/tests/e2e/dlq/test_dlq_manager.py | 13 +- 3 files changed, 105 insertions(+), 122 deletions(-) diff --git a/backend/app/core/providers.py b/backend/app/core/providers.py index ae6579fa..bd04f55b 100644 --- a/backend/app/core/providers.py +++ b/backend/app/core/providers.py @@ -49,6 +49,7 @@ 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.rate_limit import RateLimitConfig from app.domain.saga.models import SagaConfig from app.events.core import UnifiedProducer @@ -199,6 +200,53 @@ 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() -> dict[str, RetryPolicy]: + """Topic-specific retry policies for DLQ.""" + return { + "execution_requested": RetryPolicy( + topic="execution_requested", + strategy=RetryStrategy.EXPONENTIAL_BACKOFF, + max_retries=5, + base_delay_seconds=30, + max_delay_seconds=300, + retry_multiplier=2.0, + ), + "pod_created": RetryPolicy( + topic="pod_created", + strategy=RetryStrategy.EXPONENTIAL_BACKOFF, + max_retries=3, + base_delay_seconds=60, + max_delay_seconds=600, + retry_multiplier=3.0, + ), + "execution_completed": RetryPolicy( + topic="execution_completed", + strategy=RetryStrategy.EXPONENTIAL_BACKOFF, + max_retries=5, + base_delay_seconds=30, + max_delay_seconds=300, + retry_multiplier=2.0, + ), + "execution_failed": RetryPolicy( + topic="execution_failed", + strategy=RetryStrategy.IMMEDIATE, + max_retries=3, + ), + } + + class DLQProvider(Provider): """Provides DLQManager without scheduling. Used by all containers except the DLQ worker.""" @@ -209,7 +257,6 @@ def get_dlq_manager( self, broker: KafkaBroker, settings: Settings, - schema_registry: SchemaRegistryManager, logger: logging.Logger, dlq_metrics: DLQMetrics, repository: DLQRepository, @@ -217,18 +264,18 @@ def get_dlq_manager( return 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(), ) 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 @@ -238,7 +285,6 @@ async def get_dlq_manager( self, broker: KafkaBroker, settings: Settings, - schema_registry: SchemaRegistryManager, logger: logging.Logger, dlq_metrics: DLQMetrics, repository: DLQRepository, @@ -247,10 +293,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(), ) scheduler = AsyncIOScheduler() diff --git a/backend/app/dlq/manager.py b/backend/app/dlq/manager.py index 5cf21e36..008c6cb5 100644 --- a/backend/app/dlq/manager.py +++ b/backend/app/dlq/manager.py @@ -23,7 +23,6 @@ DLQMessageRetriedEvent, EventMetadata, ) -from app.events.schema.schema_registry import SchemaRegistryManager from app.settings import Settings @@ -39,63 +38,22 @@ def __init__( self, settings: Settings, broker: KafkaBroker, - schema_registry: SchemaRegistryManager, logger: logging.Logger, dlq_metrics: DLQMetrics, repository: DLQRepository, + default_retry_policy: RetryPolicy, + retry_policies: dict[str, RetryPolicy], dlq_topic: KafkaTopic = KafkaTopic.DEAD_LETTER_QUEUE, - retry_topic_suffix: str = "-retry", - default_retry_policy: RetryPolicy | None = None, - retry_policies: dict[str, RetryPolicy] | None = None, filters: list[Callable[[DLQMessage], bool]] | None = None, ): self.settings = settings self._broker = broker - self.schema_registry = schema_registry self.logger = logger self.metrics = dlq_metrics self.repository = repository self.dlq_topic = dlq_topic - self.retry_topic_suffix = retry_topic_suffix - - self.default_retry_policy = default_retry_policy or RetryPolicy( - topic="default", - strategy=RetryStrategy.EXPONENTIAL_BACKOFF, - max_retries=4, - base_delay_seconds=60, - max_delay_seconds=1800, - retry_multiplier=2.5, - ) - - self._retry_policies: dict[str, RetryPolicy] = retry_policies if retry_policies is not None else { - "execution-requests": RetryPolicy( - topic="execution-requests", - 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, - ), - "resource-allocation": RetryPolicy( - topic="resource-allocation", - strategy=RetryStrategy.IMMEDIATE, - max_retries=3, - ), - "websocket-events": RetryPolicy( - topic="websocket-events", - strategy=RetryStrategy.FIXED_INTERVAL, - max_retries=10, - base_delay_seconds=10, - ), - } + self.default_retry_policy = default_retry_policy + self._retry_policies = retry_policies self._filters: list[Callable[[DLQMessage], bool]] = filters if filters is not None else [ f for f in [ @@ -131,7 +89,20 @@ async def handle_message(self, message: DLQMessage) -> None: message.status = DLQMessageStatus.PENDING message.last_updated = datetime.now(timezone.utc) await self.repository.save_message(message) - await self._emit_message_received_event(message) + + await self._broker.publish( + DLQMessageReceivedEvent( + dlq_event_id=message.event.event_id, + original_topic=message.original_topic, + original_event_type=str(message.event.event_type), + error=message.error, + retry_count=message.retry_count, + producer_id=message.producer_id, + failed_at=message.failed_at, + metadata=self._event_metadata, + ), + topic=self._dlq_events_topic, + ) retry_policy = self._retry_policies.get(message.original_topic, self.default_retry_policy) @@ -149,9 +120,10 @@ async def handle_message(self, message: DLQMessage) -> None: await self.retry_message(message) async def retry_message(self, message: DLQMessage) -> None: - """Retry a DLQ message by republishing to the retry topic and original topic.""" - retry_topic = f"{message.original_topic}{self.retry_topic_suffix}" + """Retry a DLQ message by republishing to the original topic. + FastStream handles JSON serialization of Pydantic models natively. + """ hdrs: dict[str, str] = { "event_type": message.event.event_type, "dlq_retry_count": str(message.retry_count + 1), @@ -160,18 +132,11 @@ async def retry_message(self, message: DLQMessage) -> None: } hdrs = inject_trace_context(hdrs) - serialized = await self.schema_registry.serialize_event(message.event) - + # Publish directly to original topic - FastStream serializes Pydantic to JSON await self._broker.publish( - message=serialized, - topic=retry_topic, - key=message.event.event_id.encode(), - headers=hdrs, - ) - await self._broker.publish( - message=serialized, + message=message.event, topic=message.original_topic, - key=message.event.event_id.encode(), + key=message.event.event_id.encode() if message.event.event_id else None, headers=hdrs, ) @@ -187,7 +152,17 @@ async def retry_message(self, message: DLQMessage) -> None: ), ) - await self._emit_message_retried_event(message, retry_topic, new_retry_count) + await self._broker.publish( + DLQMessageRetriedEvent( + dlq_event_id=message.event.event_id, + original_topic=message.original_topic, + original_event_type=str(message.event.event_type), + retry_count=new_retry_count, + retry_topic=message.original_topic, + metadata=self._event_metadata, + ), + topic=self._dlq_events_topic, + ) self.logger.info("Successfully retried message", extra={"event_id": message.event.event_id}) async def discard_message(self, message: DLQMessage, reason: str) -> None: @@ -203,7 +178,17 @@ async def discard_message(self, message: DLQMessage, reason: str) -> None: ), ) - await self._emit_message_discarded_event(message, reason) + await self._broker.publish( + DLQMessageDiscardedEvent( + dlq_event_id=message.event.event_id, + original_topic=message.original_topic, + original_event_type=str(message.event.event_type), + reason=reason, + retry_count=message.retry_count, + metadata=self._event_metadata, + ), + topic=self._dlq_events_topic, + ) self.logger.warning("Discarded message", extra={"event_id": message.event.event_id, "reason": reason}) async def process_due_retries(self) -> int: @@ -291,51 +276,3 @@ async def discard_message_manually(self, event_id: str, reason: str) -> bool: await self.discard_message(message, reason) return True - - async def _emit_message_received_event(self, message: DLQMessage) -> None: - event = DLQMessageReceivedEvent( - dlq_event_id=message.event.event_id, - original_topic=message.original_topic, - original_event_type=str(message.event.event_type), - error=message.error, - retry_count=message.retry_count, - producer_id=message.producer_id, - failed_at=message.failed_at, - metadata=self._event_metadata, - ) - await self._produce_dlq_event(event) - - async def _emit_message_retried_event(self, message: DLQMessage, retry_topic: str, new_retry_count: int) -> None: - event = DLQMessageRetriedEvent( - dlq_event_id=message.event.event_id, - original_topic=message.original_topic, - original_event_type=str(message.event.event_type), - retry_count=new_retry_count, - retry_topic=retry_topic, - metadata=self._event_metadata, - ) - await self._produce_dlq_event(event) - - async def _emit_message_discarded_event(self, message: DLQMessage, reason: str) -> None: - event = DLQMessageDiscardedEvent( - dlq_event_id=message.event.event_id, - original_topic=message.original_topic, - original_event_type=str(message.event.event_type), - reason=reason, - retry_count=message.retry_count, - metadata=self._event_metadata, - ) - await self._produce_dlq_event(event) - - async def _produce_dlq_event( - self, event: DLQMessageReceivedEvent | DLQMessageRetriedEvent | DLQMessageDiscardedEvent - ) -> None: - try: - serialized = await self.schema_registry.serialize_event(event) - await self._broker.publish( - message=serialized, - topic=self._dlq_events_topic, - key=event.event_id.encode(), - ) - except Exception as e: - self.logger.error(f"Failed to emit DLQ event {event.event_type}: {e}") diff --git a/backend/tests/e2e/dlq/test_dlq_manager.py b/backend/tests/e2e/dlq/test_dlq_manager.py index e19e4528..c7914ab1 100644 --- a/backend/tests/e2e/dlq/test_dlq_manager.py +++ b/backend/tests/e2e/dlq/test_dlq_manager.py @@ -1,21 +1,22 @@ import asyncio +import json import logging import uuid from datetime import datetime, timezone import pytest from aiokafka import AIOKafkaConsumer -from faststream.kafka import KafkaBroker from app.core.metrics import DLQMetrics +from app.core.providers import _default_retry_policies, _default_retry_policy from app.db.repositories.dlq_repository import DLQRepository from app.dlq.manager import DLQManager from app.dlq.models import DLQMessage from app.domain.enums.events import EventType from app.domain.enums.kafka import KafkaTopic from app.domain.events.typed import DLQMessageReceivedEvent, DomainEventAdapter -from app.events.schema.schema_registry import SchemaRegistryManager from app.settings import Settings from dishka import AsyncContainer +from faststream.kafka import KafkaBroker from tests.conftest import make_execution_requested_event @@ -30,7 +31,6 @@ @pytest.mark.asyncio async def test_dlq_manager_persists_and_emits_event(scope: AsyncContainer, test_settings: Settings) -> None: """Test that DLQ manager persists messages and emits DLQMessageReceivedEvent.""" - schema_registry = SchemaRegistryManager(test_settings, _test_logger) dlq_metrics: DLQMetrics = await scope.get(DLQMetrics) prefix = test_settings.KAFKA_TOPIC_PREFIX @@ -53,9 +53,7 @@ async def consume_dlq_events() -> None: """Consume DLQ events and set future when our event is received.""" async for msg in events_consumer: try: - payload = await schema_registry.serializer.decode_message(msg.value) - if payload is None: - continue + payload = json.loads(msg.value.decode()) event = DomainEventAdapter.validate_python(payload) if ( isinstance(event, DLQMessageReceivedEvent) @@ -79,10 +77,11 @@ async def consume_dlq_events() -> None: manager = DLQManager( settings=test_settings, broker=broker, - schema_registry=schema_registry, logger=_test_logger, dlq_metrics=dlq_metrics, repository=repository, + default_retry_policy=_default_retry_policy(), + retry_policies=_default_retry_policies(), ) # Build a DLQMessage directly and call handle_message (no internal consumer loop) From 0a31a6cfab880b727f1a93fa6746684d7c08fb4f Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Thu, 5 Feb 2026 19:20:43 +0100 Subject: [PATCH 2/2] using faststream broker directly --- backend/app/core/container.py | 9 ---- backend/app/core/providers.py | 45 +++++++++---------- backend/app/events/broker.py | 27 ----------- backend/app/events/core/producer.py | 24 +++------- backend/app/events/schema/__init__.py | 1 - backend/app/events/schema/schema_registry.py | 28 ------------ backend/app/main.py | 6 +-- backend/tests/e2e/app/test_main_app.py | 6 --- backend/tests/e2e/core/test_container.py | 10 ----- .../tests/e2e/core/test_dishka_lifespan.py | 10 ----- backend/tests/e2e/dlq/test_dlq_manager.py | 2 +- .../e2e/events/test_schema_registry_real.py | 29 ------------ .../events/test_schema_registry_roundtrip.py | 25 ----------- .../result_processor/test_result_processor.py | 3 -- .../events/test_schema_registry_manager.py | 34 -------------- backend/workers/dlq_processor.py | 6 +-- backend/workers/run_coordinator.py | 6 +-- backend/workers/run_event_replay.py | 6 +-- backend/workers/run_k8s_worker.py | 6 +-- backend/workers/run_pod_monitor.py | 6 +-- backend/workers/run_result_processor.py | 6 +-- backend/workers/run_saga_orchestrator.py | 6 +-- 22 files changed, 46 insertions(+), 255 deletions(-) delete mode 100644 backend/app/events/broker.py delete mode 100644 backend/app/events/schema/__init__.py delete mode 100644 backend/app/events/schema/schema_registry.py delete mode 100644 backend/tests/e2e/events/test_schema_registry_real.py delete mode 100644 backend/tests/e2e/events/test_schema_registry_roundtrip.py delete mode 100644 backend/tests/unit/events/test_schema_registry_manager.py diff --git a/backend/app/core/container.py b/backend/app/core/container.py index 44c281d4..5aca7e45 100644 --- a/backend/app/core/container.py +++ b/backend/app/core/container.py @@ -11,7 +11,6 @@ DatabaseProvider, DLQProvider, DLQWorkerProvider, - EventProvider, EventReplayProvider, EventReplayWorkerProvider, K8sWorkerProvider, @@ -52,7 +51,6 @@ def create_app_container(settings: Settings, broker: KafkaBroker) -> AsyncContai RepositoryProvider(), MessagingProvider(), DLQProvider(), - EventProvider(), SagaOrchestratorProvider(), KafkaServicesProvider(), SSEProvider(), @@ -79,7 +77,6 @@ def create_result_processor_container(settings: Settings, broker: KafkaBroker) - CoreServicesProvider(), MetricsProvider(), RepositoryProvider(), - EventProvider(), MessagingProvider(), DLQProvider(), ResultProcessorProvider(), @@ -99,7 +96,6 @@ def create_coordinator_container(settings: Settings, broker: KafkaBroker) -> Asy RepositoryProvider(), MessagingProvider(), DLQProvider(), - EventProvider(), CoordinatorProvider(), context={Settings: settings, KafkaBroker: broker}, ) @@ -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}, @@ -136,7 +131,6 @@ def create_pod_monitor_container(settings: Settings, broker: KafkaBroker) -> Asy RepositoryProvider(), MessagingProvider(), DLQProvider(), - EventProvider(), KafkaServicesProvider(), KubernetesProvider(), PodMonitorProvider(), @@ -159,7 +153,6 @@ def create_saga_orchestrator_container(settings: Settings, broker: KafkaBroker) RepositoryProvider(), MessagingProvider(), DLQProvider(), - EventProvider(), SagaWorkerProvider(), context={Settings: settings, KafkaBroker: broker}, ) @@ -180,7 +173,6 @@ def create_event_replay_container(settings: Settings, broker: KafkaBroker) -> As RepositoryProvider(), MessagingProvider(), DLQProvider(), - EventProvider(), EventReplayWorkerProvider(), context={Settings: settings, KafkaBroker: broker}, ) @@ -202,6 +194,5 @@ def create_dlq_processor_container(settings: Settings, broker: KafkaBroker) -> A RepositoryProvider(), MessagingProvider(), DLQWorkerProvider(), - EventProvider(), context={Settings: settings, KafkaBroker: broker}, ) diff --git a/backend/app/core/providers.py b/backend/app/core/providers.py index bd04f55b..df62a876 100644 --- a/backend/app/core/providers.py +++ b/backend/app/core/providers.py @@ -50,10 +50,10 @@ 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 @@ -181,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: @@ -212,35 +211,43 @@ def _default_retry_policy() -> RetryPolicy: ) -def _default_retry_policies() -> dict[str, RetryPolicy]: - """Topic-specific retry policies for DLQ.""" +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_requested": RetryPolicy( - topic="execution_requested", + 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_created": RetryPolicy( - topic="pod_created", + 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, ), - "execution_completed": RetryPolicy( - topic="execution_completed", + 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_failed": RetryPolicy( - topic="execution_failed", + execution_results: RetryPolicy( + topic=execution_results, strategy=RetryStrategy.IMMEDIATE, max_retries=3, ), @@ -268,7 +275,7 @@ def get_dlq_manager( dlq_metrics=dlq_metrics, repository=repository, default_retry_policy=_default_retry_policy(), - retry_policies=_default_retry_policies(), + retry_policies=_default_retry_policies(settings.KAFKA_TOPIC_PREFIX), ) @@ -297,7 +304,7 @@ async def get_dlq_manager( dlq_metrics=dlq_metrics, repository=repository, default_retry_policy=_default_retry_policy(), - retry_policies=_default_retry_policies(), + retry_policies=_default_retry_policies(settings.KAFKA_TOPIC_PREFIX), ) scheduler = AsyncIOScheduler() @@ -319,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 diff --git a/backend/app/events/broker.py b/backend/app/events/broker.py deleted file mode 100644 index 01e70da4..00000000 --- a/backend/app/events/broker.py +++ /dev/null @@ -1,27 +0,0 @@ -import logging -from typing import Any - -from faststream import StreamMessage -from faststream.kafka import KafkaBroker - -from app.domain.events.typed import DomainEvent, DomainEventAdapter -from app.events.schema.schema_registry import SchemaRegistryManager -from app.settings import Settings - - -def create_broker( - settings: Settings, - schema_registry: SchemaRegistryManager, - logger: logging.Logger, -) -> KafkaBroker: - """Create a KafkaBroker with Avro decoder for standalone workers.""" - - async def avro_decoder(msg: StreamMessage[Any]) -> DomainEvent: - payload = await schema_registry.serializer.decode_message(msg.body) - return DomainEventAdapter.validate_python(payload) - - return KafkaBroker( - settings.KAFKA_BOOTSTRAP_SERVERS, - decoder=avro_decoder, - logger=logger, - ) diff --git a/backend/app/events/core/producer.py b/backend/app/events/core/producer.py index daad87e2..e3b5fa57 100644 --- a/backend/app/events/core/producer.py +++ b/backend/app/events/core/producer.py @@ -11,29 +11,26 @@ from app.dlq.models import DLQMessageStatus from app.domain.enums.kafka import KafkaTopic from app.domain.events.typed import DomainEvent -from app.events.schema.schema_registry import SchemaRegistryManager from app.infrastructure.kafka.mappings import EVENT_TYPE_TO_TOPIC from app.settings import Settings class UnifiedProducer: - """Fully async Kafka producer backed by FastStream KafkaBroker. + """Kafka producer backed by FastStream KafkaBroker. - The broker's lifecycle (start/stop) is managed externally — either by - the FastStream app (worker entry points) or by the FastAPI lifespan. + FastStream handles Pydantic JSON serialization natively. + The broker's lifecycle is managed externally (FastStream app or FastAPI lifespan). """ def __init__( self, broker: KafkaBroker, - schema_registry_manager: SchemaRegistryManager, event_repository: EventRepository, logger: logging.Logger, settings: Settings, event_metrics: EventMetrics, ): self._broker = broker - self._schema_registry = schema_registry_manager self._event_repository = event_repository self.logger = logger self._event_metrics = event_metrics @@ -44,8 +41,6 @@ async def produce(self, event_to_produce: DomainEvent, key: str) -> None: await self._event_repository.store_event(event_to_produce) topic = f"{self._topic_prefix}{EVENT_TYPE_TO_TOPIC[event_to_produce.event_type]}" try: - serialized_value = await self._schema_registry.serialize_event(event_to_produce) - headers = inject_trace_context({ "event_type": event_to_produce.event_type, "correlation_id": event_to_produce.metadata.correlation_id or "", @@ -53,14 +48,14 @@ async def produce(self, event_to_produce: DomainEvent, key: str) -> None: }) await self._broker.publish( - message=serialized_value, + message=event_to_produce, topic=topic, key=key.encode(), headers=headers, ) self._event_metrics.record_kafka_message_produced(topic) - self.logger.debug(f"Message [{event_to_produce}] sent to topic: {topic}") + self.logger.debug(f"Event {event_to_produce.event_type} sent to topic: {topic}") except Exception as e: self._event_metrics.record_kafka_production_error(topic=topic, error_type=type(e).__name__) @@ -70,17 +65,12 @@ async def produce(self, event_to_produce: DomainEvent, key: str) -> None: async def send_to_dlq( self, original_event: DomainEvent, original_topic: str, error: Exception, retry_count: int = 0 ) -> None: - """Send a failed event to the Dead Letter Queue. - - The event body is Avro-encoded (same as every other topic). - DLQ metadata is carried in Kafka headers. - """ + """Send a failed event to the Dead Letter Queue.""" try: current_task = asyncio.current_task() task_name = current_task.get_name() if current_task else "main" producer_id = f"{socket.gethostname()}-{task_name}" - serialized_value = await self._schema_registry.serialize_event(original_event) dlq_topic = f"{self._topic_prefix}{KafkaTopic.DEAD_LETTER_QUEUE}" headers = inject_trace_context({ @@ -95,7 +85,7 @@ async def send_to_dlq( }) await self._broker.publish( - message=serialized_value, + message=original_event, topic=dlq_topic, key=original_event.event_id.encode() if original_event.event_id else None, headers=headers, diff --git a/backend/app/events/schema/__init__.py b/backend/app/events/schema/__init__.py deleted file mode 100644 index 80bb333a..00000000 --- a/backend/app/events/schema/__init__.py +++ /dev/null @@ -1 +0,0 @@ -"""Event schema management, validation, and serialization.""" diff --git a/backend/app/events/schema/schema_registry.py b/backend/app/events/schema/schema_registry.py deleted file mode 100644 index af67d943..00000000 --- a/backend/app/events/schema/schema_registry.py +++ /dev/null @@ -1,28 +0,0 @@ -import logging - -from schema_registry.client import AsyncSchemaRegistryClient, schema -from schema_registry.serializers import AsyncAvroMessageSerializer # type: ignore[attr-defined] - -from app.domain.events.typed import DomainEvent -from app.settings import Settings - - -class SchemaRegistryManager: - """Avro serialization via Confluent Schema Registry. - - Schemas are registered lazily by the underlying serializer on first - produce — no eager bootstrap needed. - """ - - def __init__(self, settings: Settings, logger: logging.Logger): - self.logger = logger - self.namespace = "com.integr8scode.events" - self.subject_prefix = settings.SCHEMA_SUBJECT_PREFIX - self._client = AsyncSchemaRegistryClient(url=settings.SCHEMA_REGISTRY_URL) - self.serializer = AsyncAvroMessageSerializer(self._client) - - async def serialize_event(self, event: DomainEvent) -> bytes: - """Serialize event to Confluent wire format: [0x00][4-byte schema id][Avro binary].""" - avro = schema.AvroSchema(event.avro_schema(namespace=self.namespace)) - subject = f"{self.subject_prefix}{avro.name}-value" - return await self.serializer.encode_record_with_schema(subject, avro, event.model_dump()) diff --git a/backend/app/main.py b/backend/app/main.py index 607ba1b7..8f67b2a4 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -3,6 +3,7 @@ from dishka.integrations.faststream import setup_dishka as setup_dishka_faststream from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware +from faststream.kafka import KafkaBroker from app.api.routes import ( auth, @@ -40,12 +41,10 @@ RequestSizeLimitMiddleware, setup_metrics, ) -from app.events.broker import create_broker from app.events.handlers import ( register_notification_subscriber, register_sse_subscriber, ) -from app.events.schema.schema_registry import SchemaRegistryManager from app.settings import Settings @@ -61,8 +60,7 @@ def create_app(settings: Settings | None = None) -> FastAPI: logger = setup_logger(settings.LOG_LEVEL) # Create Kafka broker and register in-app subscribers - schema_registry = SchemaRegistryManager(settings, logger) - broker = create_broker(settings, schema_registry, logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=logger) register_sse_subscriber(broker, settings) register_notification_subscriber(broker, settings) diff --git a/backend/tests/e2e/app/test_main_app.py b/backend/tests/e2e/app/test_main_app.py index 160f5a83..219f087f 100644 --- a/backend/tests/e2e/app/test_main_app.py +++ b/backend/tests/e2e/app/test_main_app.py @@ -6,7 +6,6 @@ import redis.asyncio as aioredis from app.core.database_context import Database from app.domain.exceptions import DomainError -from app.events.schema.schema_registry import SchemaRegistryManager from app.settings import Settings from dishka import AsyncContainer from fastapi import FastAPI @@ -294,11 +293,6 @@ async def test_redis_connected(self, scope: AsyncContainer) -> None: pong = await redis_client.ping() # type: ignore[misc] assert pong is True - @pytest.mark.asyncio - async def test_schema_registry_initialized(self, scope: AsyncContainer) -> None: - """Schema registry manager is initialized.""" - schema_registry = await scope.get(SchemaRegistryManager) - assert schema_registry is not None class TestCreateAppFunction: diff --git a/backend/tests/e2e/core/test_container.py b/backend/tests/e2e/core/test_container.py index 45ac8ae5..8456a0f0 100644 --- a/backend/tests/e2e/core/test_container.py +++ b/backend/tests/e2e/core/test_container.py @@ -4,7 +4,6 @@ import redis.asyncio as aioredis from app.core.database_context import Database from app.core.security import SecurityService -from app.events.schema.schema_registry import SchemaRegistryManager from app.services.event_service import EventService from app.services.execution_service import ExecutionService from app.services.notification_service import NotificationService @@ -82,15 +81,6 @@ async def test_resolves_event_service(self, scope: AsyncContainer) -> None: assert isinstance(service, EventService) - @pytest.mark.asyncio - async def test_resolves_schema_registry( - self, scope: AsyncContainer - ) -> None: - """Container resolves SchemaRegistryManager.""" - registry = await scope.get(SchemaRegistryManager) - - assert isinstance(registry, SchemaRegistryManager) - class TestBusinessServices: """Tests for business service resolution.""" diff --git a/backend/tests/e2e/core/test_dishka_lifespan.py b/backend/tests/e2e/core/test_dishka_lifespan.py index e1e387a8..466e88b5 100644 --- a/backend/tests/e2e/core/test_dishka_lifespan.py +++ b/backend/tests/e2e/core/test_dishka_lifespan.py @@ -70,16 +70,6 @@ async def test_redis_connected(self, scope: AsyncContainer) -> None: pong = await redis_client.ping() # type: ignore[misc] assert pong is True - @pytest.mark.asyncio - async def test_schema_registry_initialized( - self, scope: AsyncContainer - ) -> None: - """Schema registry is initialized during lifespan.""" - from app.events.schema.schema_registry import SchemaRegistryManager - - registry = await scope.get(SchemaRegistryManager) - assert registry is not None - @pytest.mark.asyncio async def test_sse_redis_bus_available(self, scope: AsyncContainer) -> None: """SSE Redis bus is available after lifespan.""" diff --git a/backend/tests/e2e/dlq/test_dlq_manager.py b/backend/tests/e2e/dlq/test_dlq_manager.py index c7914ab1..631faaa7 100644 --- a/backend/tests/e2e/dlq/test_dlq_manager.py +++ b/backend/tests/e2e/dlq/test_dlq_manager.py @@ -81,7 +81,7 @@ async def consume_dlq_events() -> None: dlq_metrics=dlq_metrics, repository=repository, default_retry_policy=_default_retry_policy(), - retry_policies=_default_retry_policies(), + retry_policies=_default_retry_policies(test_settings.KAFKA_TOPIC_PREFIX), ) # Build a DLQMessage directly and call handle_message (no internal consumer loop) diff --git a/backend/tests/e2e/events/test_schema_registry_real.py b/backend/tests/e2e/events/test_schema_registry_real.py deleted file mode 100644 index 58e4900d..00000000 --- a/backend/tests/e2e/events/test_schema_registry_real.py +++ /dev/null @@ -1,29 +0,0 @@ -import logging - -import pytest - -from app.domain.events.typed import DomainEventAdapter, EventMetadata, PodCreatedEvent -from app.events.schema.schema_registry import SchemaRegistryManager -from app.settings import Settings - -pytestmark = [pytest.mark.e2e, pytest.mark.kafka] - -_test_logger = logging.getLogger("test.events.schema_registry_real") - - -@pytest.mark.asyncio -async def test_serialize_and_deserialize_event_real_registry(test_settings: Settings) -> None: - # Uses real Schema Registry configured via env (SCHEMA_REGISTRY_URL) - m = SchemaRegistryManager(settings=test_settings, logger=_test_logger) - ev = PodCreatedEvent( - execution_id="e1", - pod_name="p", - namespace="n", - metadata=EventMetadata(service_name="s", service_version="1"), - ) - data = await m.serialize_event(ev) - payload = await m.serializer.decode_message(data) - assert payload is not None - obj = DomainEventAdapter.validate_python(payload) - assert isinstance(obj, PodCreatedEvent) - assert obj.namespace == "n" diff --git a/backend/tests/e2e/events/test_schema_registry_roundtrip.py b/backend/tests/e2e/events/test_schema_registry_roundtrip.py deleted file mode 100644 index 1fc83467..00000000 --- a/backend/tests/e2e/events/test_schema_registry_roundtrip.py +++ /dev/null @@ -1,25 +0,0 @@ -import logging - -import pytest - -from app.domain.events.typed import DomainEventAdapter -from app.events.schema.schema_registry import SchemaRegistryManager -from dishka import AsyncContainer - -from tests.conftest import make_execution_requested_event - -pytestmark = [pytest.mark.e2e] - -_test_logger = logging.getLogger("test.events.schema_registry_roundtrip") - - -@pytest.mark.asyncio -async def test_schema_registry_serialize_deserialize_roundtrip(scope: AsyncContainer) -> None: - reg: SchemaRegistryManager = await scope.get(SchemaRegistryManager) - ev = make_execution_requested_event(execution_id="e-rt") - data = await reg.serialize_event(ev) - assert data[:1] == b"\x00" # Confluent wire format magic byte - payload = await reg.serializer.decode_message(data) - assert payload is not None - back = DomainEventAdapter.validate_python(payload) - assert back.event_id == ev.event_id and getattr(back, "execution_id", None) == ev.execution_id diff --git a/backend/tests/e2e/result_processor/test_result_processor.py b/backend/tests/e2e/result_processor/test_result_processor.py index 7976d1b0..0d318648 100644 --- a/backend/tests/e2e/result_processor/test_result_processor.py +++ b/backend/tests/e2e/result_processor/test_result_processor.py @@ -13,7 +13,6 @@ ) from app.domain.execution import DomainExecutionCreate from app.events.core import UnifiedProducer -from app.events.schema.schema_registry import SchemaRegistryManager from app.services.result_processor.processor import ResultProcessor from app.settings import Settings from dishka import AsyncContainer @@ -30,8 +29,6 @@ @pytest.mark.asyncio async def test_result_processor_persists_and_emits(scope: AsyncContainer) -> None: - # Schemas are initialized inside the SchemaRegistryManager DI provider - registry: SchemaRegistryManager = await scope.get(SchemaRegistryManager) settings: Settings = await scope.get(Settings) execution_metrics: ExecutionMetrics = await scope.get(ExecutionMetrics) diff --git a/backend/tests/unit/events/test_schema_registry_manager.py b/backend/tests/unit/events/test_schema_registry_manager.py deleted file mode 100644 index cdc3159b..00000000 --- a/backend/tests/unit/events/test_schema_registry_manager.py +++ /dev/null @@ -1,34 +0,0 @@ -import pytest -from pydantic import ValidationError - -from app.domain.enums.execution import QueuePriority -from app.domain.events.typed import DomainEventAdapter, ExecutionRequestedEvent - - -def test_domain_event_adapter_execution_requested() -> None: - data = { - "event_type": "execution_requested", - "execution_id": "e1", - "script": "print('ok')", - "language": "python", - "language_version": "3.11", - "runtime_image": "python:3.11-slim", - "runtime_command": ["python"], - "runtime_filename": "main.py", - "timeout_seconds": 30, - "cpu_limit": "100m", - "memory_limit": "128Mi", - "cpu_request": "50m", - "memory_request": "64Mi", - "priority": QueuePriority.NORMAL, - "metadata": {"service_name": "t", "service_version": "1.0"}, - } - ev = DomainEventAdapter.validate_python(data) - assert isinstance(ev, ExecutionRequestedEvent) - assert ev.execution_id == "e1" - assert ev.language == "python" - - -def test_domain_event_adapter_missing_type_raises() -> None: - with pytest.raises(ValidationError): - DomainEventAdapter.validate_python({}) diff --git a/backend/workers/dlq_processor.py b/backend/workers/dlq_processor.py index 6f2561cc..ef195b80 100644 --- a/backend/workers/dlq_processor.py +++ b/backend/workers/dlq_processor.py @@ -5,12 +5,11 @@ from app.core.tracing import init_tracing from app.dlq.manager import DLQManager from app.domain.enums.kafka import GroupId -from app.events.broker import create_broker from app.events.handlers import register_dlq_subscriber -from app.events.schema.schema_registry import SchemaRegistryManager from app.settings import Settings from dishka.integrations.faststream import setup_dishka from faststream import FastStream +from faststream.kafka import KafkaBroker def main() -> None: @@ -33,8 +32,7 @@ def main() -> None: logger.info("Tracing initialized for DLQ Processor") # Create Kafka broker and register DLQ subscriber - schema_registry = SchemaRegistryManager(settings, logger) - broker = create_broker(settings, schema_registry, logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=logger) register_dlq_subscriber(broker, settings) # Create DI container with broker in context diff --git a/backend/workers/run_coordinator.py b/backend/workers/run_coordinator.py index d2d5ae30..40b8876f 100644 --- a/backend/workers/run_coordinator.py +++ b/backend/workers/run_coordinator.py @@ -5,12 +5,11 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.domain.enums.kafka import GroupId -from app.events.broker import create_broker from app.events.handlers import register_coordinator_subscriber -from app.events.schema.schema_registry import SchemaRegistryManager from app.settings import Settings from dishka.integrations.faststream import setup_dishka from faststream import FastStream +from faststream.kafka import KafkaBroker def main() -> None: @@ -33,8 +32,7 @@ def main() -> None: logger.info("Tracing initialized for ExecutionCoordinator") # Create Kafka broker and register subscriber - schema_registry = SchemaRegistryManager(settings, logger) - broker = create_broker(settings, schema_registry, logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=logger) register_coordinator_subscriber(broker, settings) # Create DI container with broker in context diff --git a/backend/workers/run_event_replay.py b/backend/workers/run_event_replay.py index 81aac922..ec31c92e 100644 --- a/backend/workers/run_event_replay.py +++ b/backend/workers/run_event_replay.py @@ -5,17 +5,15 @@ from app.core.container import create_event_replay_container from app.core.logging import setup_logger from app.core.tracing import init_tracing -from app.events.broker import create_broker -from app.events.schema.schema_registry import SchemaRegistryManager from app.services.event_replay.replay_service import EventReplayService from app.settings import Settings +from faststream.kafka import KafkaBroker async def run_replay_service(settings: Settings) -> None: """Run the event replay service with DI-managed cleanup scheduler.""" tmp_logger = setup_logger(settings.LOG_LEVEL) - schema_registry = SchemaRegistryManager(settings, tmp_logger) - broker = create_broker(settings, schema_registry, tmp_logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=tmp_logger) container = create_event_replay_container(settings, broker) logger = await container.get(logging.Logger) diff --git a/backend/workers/run_k8s_worker.py b/backend/workers/run_k8s_worker.py index 3044625b..173fd7b4 100644 --- a/backend/workers/run_k8s_worker.py +++ b/backend/workers/run_k8s_worker.py @@ -5,13 +5,12 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.domain.enums.kafka import GroupId -from app.events.broker import create_broker from app.events.handlers import register_k8s_worker_subscriber -from app.events.schema.schema_registry import SchemaRegistryManager from app.services.k8s_worker import KubernetesWorker from app.settings import Settings from dishka.integrations.faststream import setup_dishka from faststream import FastStream +from faststream.kafka import KafkaBroker def main() -> None: @@ -34,8 +33,7 @@ def main() -> None: logger.info("Tracing initialized for KubernetesWorker") # Create Kafka broker and register subscriber - schema_registry = SchemaRegistryManager(settings, logger) - broker = create_broker(settings, schema_registry, logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=logger) register_k8s_worker_subscriber(broker, settings) # Create DI container with broker in context diff --git a/backend/workers/run_pod_monitor.py b/backend/workers/run_pod_monitor.py index 1f72277e..9be78252 100644 --- a/backend/workers/run_pod_monitor.py +++ b/backend/workers/run_pod_monitor.py @@ -4,12 +4,11 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.domain.enums.kafka import GroupId -from app.events.broker import create_broker -from app.events.schema.schema_registry import SchemaRegistryManager from app.services.pod_monitor.monitor import PodMonitor from app.settings import Settings from dishka.integrations.faststream import setup_dishka from faststream import FastStream +from faststream.kafka import KafkaBroker def main() -> None: @@ -32,8 +31,7 @@ def main() -> None: logger.info("Tracing initialized for PodMonitor Service") # Create Kafka broker (PodMonitor publishes events via KafkaEventService) - schema_registry = SchemaRegistryManager(settings, logger) - broker = create_broker(settings, schema_registry, logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=logger) # Create DI container with broker in context container = create_pod_monitor_container(settings, broker) diff --git a/backend/workers/run_result_processor.py b/backend/workers/run_result_processor.py index 2a1061a2..ee430315 100644 --- a/backend/workers/run_result_processor.py +++ b/backend/workers/run_result_processor.py @@ -5,12 +5,11 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.domain.enums.kafka import GroupId -from app.events.broker import create_broker from app.events.handlers import register_result_processor_subscriber -from app.events.schema.schema_registry import SchemaRegistryManager from app.settings import Settings from dishka.integrations.faststream import setup_dishka from faststream import FastStream +from faststream.kafka import KafkaBroker def main() -> None: @@ -33,8 +32,7 @@ def main() -> None: logger.info("Tracing initialized for ResultProcessor Service") # Create Kafka broker and register subscriber - schema_registry = SchemaRegistryManager(settings, logger) - broker = create_broker(settings, schema_registry, logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=logger) register_result_processor_subscriber(broker, settings) # Create DI container with broker in context diff --git a/backend/workers/run_saga_orchestrator.py b/backend/workers/run_saga_orchestrator.py index 4f355f86..c25f8ced 100644 --- a/backend/workers/run_saga_orchestrator.py +++ b/backend/workers/run_saga_orchestrator.py @@ -4,13 +4,12 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.domain.enums.kafka import GroupId -from app.events.broker import create_broker from app.events.handlers import register_saga_subscriber -from app.events.schema.schema_registry import SchemaRegistryManager from app.services.saga import SagaOrchestrator from app.settings import Settings from dishka.integrations.faststream import setup_dishka from faststream import FastStream +from faststream.kafka import KafkaBroker def main() -> None: @@ -33,8 +32,7 @@ def main() -> None: logger.info("Tracing initialized for Saga Orchestrator Service") # Create Kafka broker and register subscriber - schema_registry = SchemaRegistryManager(settings, logger) - broker = create_broker(settings, schema_registry, logger) + broker = KafkaBroker(settings.KAFKA_BOOTSTRAP_SERVERS, logger=logger) register_saga_subscriber(broker, settings) # Create DI container with broker in context