From 161a4e7e11cbba4fcf58657d730a7c662ba09ade Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Sun, 1 Mar 2026 23:08:15 +0100 Subject: [PATCH 1/4] feat: removed calculation/acquisition of kafka-internal metrics (using faststream middleware for that) --- README.md | 4 +- backend/app/core/metrics/events.py | 87 ++----------------- backend/app/core/providers.py | 17 ++-- backend/app/events/core/__init__.py | 2 + backend/app/events/core/producer.py | 35 ++------ backend/app/events/core/transport.py | 33 +++++++ backend/app/events/handlers.py | 8 +- backend/app/services/k8s_worker/worker.py | 4 +- .../dashboards/http-middleware.json | 8 +- .../dashboards/kafka-events-monitoring.json | 22 ++--- .../test_execution_and_events_metrics.py | 13 --- .../unit/core/metrics/test_metrics_classes.py | 12 --- 12 files changed, 79 insertions(+), 166 deletions(-) create mode 100644 backend/app/events/core/transport.py diff --git a/README.md b/README.md index acfecf08..dfcf8041 100644 --- a/README.md +++ b/README.md @@ -12,8 +12,8 @@ Security Scan Status - - Dead Code Check + + Dead Code Check Docker Scan Status diff --git a/backend/app/core/metrics/events.py b/backend/app/core/metrics/events.py index f5dbdf49..6c6dd74f 100644 --- a/backend/app/core/metrics/events.py +++ b/backend/app/core/metrics/events.py @@ -2,10 +2,13 @@ class EventMetrics(BaseMetrics): - """Metrics for event processing and Kafka.""" + """Metrics for domain-level event processing. + + Transport-level Kafka metrics (produced/consumed/errors) are handled + automatically by KafkaTelemetryMiddleware on the broker. + """ def _create_instruments(self) -> None: - # Core event metrics self.event_published = self._meter.create_counter( name="events.published.total", description="Total number of events published", unit="1" ) @@ -18,33 +21,6 @@ def _create_instruments(self) -> None: name="event.processing.errors.total", description="Total number of event processing errors", unit="1" ) - # Event bus metrics - self.event_bus_queue_size = self._meter.create_up_down_counter( - name="event.bus.queue.size", description="Size of event bus message queue", unit="1" - ) - - # Event replay metrics - self.event_replay_operations = self._meter.create_counter( - name="event.replay.operations.total", description="Total number of event replay operations", unit="1" - ) - - # Kafka-specific metrics - self.kafka_messages_produced = self._meter.create_counter( - name="kafka.messages.produced.total", description="Total number of messages produced to Kafka", unit="1" - ) - - self.kafka_messages_consumed = self._meter.create_counter( - name="kafka.messages.consumed.total", description="Total number of messages consumed from Kafka", unit="1" - ) - - self.kafka_production_errors = self._meter.create_counter( - name="kafka.production.errors.total", description="Total number of Kafka production errors", unit="1" - ) - - self.kafka_consumption_errors = self._meter.create_counter( - name="kafka.consumption.errors.total", description="Total number of Kafka consumption errors", unit="1" - ) - def record_event_published(self, event_type: str, event_category: str | None = None) -> None: if event_category is None: event_category = event_type.split(".")[0] if "." in event_type else event_type @@ -54,12 +30,6 @@ def record_event_published(self, event_type: str, event_category: str | None = N def record_event_processing_duration(self, duration_seconds: float, event_type: str) -> None: self.event_processing_duration.record(duration_seconds, attributes={"event_type": event_type}) - def record_event_replay_operation(self, operation: str, status: str) -> None: - self.event_replay_operations.add(1, attributes={"operation": operation, "status": status}) - - def record_event_stored(self, event_type: str, collection: str) -> None: - self.event_published.add(1, attributes={"event_type": event_type, "aggregate_type": collection}) - def record_events_processing_failed( self, topic: str, event_type: str, consumer_group: str, error_type: str ) -> None: @@ -72,50 +42,3 @@ def record_events_processing_failed( "error_type": error_type, }, ) - - def record_event_store_duration(self, duration: float, operation: str, collection: str) -> None: - self.event_processing_duration.record(duration, attributes={"operation": operation, "collection": collection}) - - def record_event_store_failed(self, event_type: str, error_type: str) -> None: - self.event_processing_errors.add( - 1, attributes={"event_type": event_type, "error_type": error_type, "operation": "store"} - ) - - def record_event_query_duration(self, duration: float, query_type: str, collection: str) -> None: - self.event_processing_duration.record( - duration, attributes={"operation": f"query_{query_type}", "collection": collection} - ) - - def record_processing_duration( - self, duration_seconds: float, topic: str, event_type: str, consumer_group: str - ) -> None: - self.event_processing_duration.record( - duration_seconds, attributes={"topic": topic, "event_type": event_type, "consumer_group": consumer_group} - ) - - def record_kafka_message_produced(self, topic: str, partition: int = -1) -> None: - self.kafka_messages_produced.add( - 1, attributes={"topic": topic, "partition": str(partition) if partition >= 0 else "auto"} - ) - - def record_kafka_message_consumed(self, topic: str, consumer_group: str) -> None: - self.kafka_messages_consumed.add(1, attributes={"topic": topic, "consumer_group": consumer_group}) - - def record_kafka_production_error(self, topic: str, error_type: str) -> None: - self.kafka_production_errors.add(1, attributes={"topic": topic, "error_type": error_type}) - - def record_kafka_consumption_error(self, topic: str, consumer_group: str, error_type: str) -> None: - self.kafka_consumption_errors.add( - 1, attributes={"topic": topic, "consumer_group": consumer_group, "error_type": error_type} - ) - - def update_event_bus_queue_size(self, delta: int, queue_name: str = "default") -> None: - self.event_bus_queue_size.add(delta, attributes={"queue": queue_name}) - - def set_event_bus_queue_size(self, size: int, queue_name: str = "default") -> None: - key = f"_event_bus_size_{queue_name}" - current_val = getattr(self, key, 0) - delta = size - current_val - if delta != 0: - self.event_bus_queue_size.add(delta, attributes={"queue": queue_name}) - setattr(self, key, size) diff --git a/backend/app/core/providers.py b/backend/app/core/providers.py index 1acdd9df..b4219689 100644 --- a/backend/app/core/providers.py +++ b/backend/app/core/providers.py @@ -43,6 +43,7 @@ from app.dlq.manager import DLQManager from app.domain.saga import SagaConfig from app.events import UnifiedProducer +from app.events.core.transport import KafkaEventTransport from app.services.admin import AdminEventsService, AdminSettingsService, AdminUserService from app.services.admin.admin_execution_service import AdminExecutionService from app.services.auth_service import AuthService @@ -167,14 +168,20 @@ class MessagingProvider(Provider): scope = Scope.APP @provide - def get_unified_producer( + def get_kafka_event_transport( self, broker: KafkaBroker, - event_repository: EventRepository, logger: structlog.stdlib.BoundLogger, - event_metrics: EventMetrics, + ) -> KafkaEventTransport: + return KafkaEventTransport(broker, logger) + + @provide + def get_unified_producer( + self, + event_repository: EventRepository, + transport: KafkaEventTransport, ) -> UnifiedProducer: - return UnifiedProducer(broker, event_repository, logger, event_metrics) + return UnifiedProducer(event_repository, transport) @provide def get_idempotency_repository(self, redis_client: redis.Redis) -> RedisIdempotencyRepository: @@ -622,14 +629,12 @@ def get_kubernetes_worker( kafka_producer: UnifiedProducer, settings: Settings, logger: structlog.stdlib.BoundLogger, - event_metrics: EventMetrics, ) -> KubernetesWorker: return KubernetesWorker( api_client=api_client, producer=kafka_producer, settings=settings, logger=logger, - event_metrics=event_metrics, ) diff --git a/backend/app/events/core/__init__.py b/backend/app/events/core/__init__.py index 555a3e77..ca050d03 100644 --- a/backend/app/events/core/__init__.py +++ b/backend/app/events/core/__init__.py @@ -1,5 +1,7 @@ from .producer import UnifiedProducer +from .transport import KafkaEventTransport __all__ = [ + "KafkaEventTransport", "UnifiedProducer", ] diff --git a/backend/app/events/core/producer.py b/backend/app/events/core/producer.py index f067b7e5..e94ee57d 100644 --- a/backend/app/events/core/producer.py +++ b/backend/app/events/core/producer.py @@ -1,29 +1,23 @@ -import structlog -from faststream.kafka import KafkaBroker - -from app.core.metrics import EventMetrics from app.db import EventRepository from app.domain.events import DomainEvent +from app.events.core.transport import KafkaEventTransport class UnifiedProducer: - """Kafka producer backed by FastStream KafkaBroker. + """Orchestrates the store-then-publish outbox pattern. - FastStream handles Pydantic JSON serialization natively. - The broker's lifecycle is managed externally (FastStream app or FastAPI lifespan). + Persists the event to MongoDB, then delegates to KafkaEventTransport + for Kafka delivery. On transport failure the event is marked as + failed-to-publish before the exception propagates. """ def __init__( self, - broker: KafkaBroker, event_repository: EventRepository, - logger: structlog.stdlib.BoundLogger, - event_metrics: EventMetrics, + transport: KafkaEventTransport, ): - self._broker = broker self._event_repository = event_repository - self.logger = logger - self._event_metrics = event_metrics + self._transport = transport async def produce(self, event_to_produce: DomainEvent, key: str) -> None: """Persist event to MongoDB, then publish to Kafka. @@ -32,19 +26,8 @@ async def produce(self, event_to_produce: DomainEvent, key: str) -> None: in MongoDB before the exception propagates. """ await self._event_repository.store_event(event_to_produce) - topic = event_to_produce.event_type try: - await self._broker.publish( - message=event_to_produce, - topic=topic, - key=key.encode(), - ) - - self._event_metrics.record_kafka_message_produced(topic) - self.logger.debug("Event sent to topic", event_type=event_to_produce.event_type, topic=topic) - - except Exception as e: - self._event_metrics.record_kafka_production_error(topic=topic, error_type=type(e).__name__) - self.logger.error("Failed to produce message", topic=topic, error=str(e)) + await self._transport.publish(event_to_produce, event_to_produce.event_type, key) + except Exception: await self._event_repository.mark_publish_failed(event_to_produce.event_id) raise diff --git a/backend/app/events/core/transport.py b/backend/app/events/core/transport.py new file mode 100644 index 00000000..c05dde7d --- /dev/null +++ b/backend/app/events/core/transport.py @@ -0,0 +1,33 @@ +import structlog +from faststream.kafka import KafkaBroker + +from app.domain.events import DomainEvent + + +class KafkaEventTransport: + """Publishes events to Kafka. + + Transport-level metrics (produced count, error count, latency) are + recorded automatically by KafkaTelemetryMiddleware on the broker. + """ + + def __init__( + self, + broker: KafkaBroker, + logger: structlog.stdlib.BoundLogger, + ): + self._broker = broker + self._logger = logger + + async def publish(self, event: DomainEvent, topic: str, key: str) -> None: + """Publish event to Kafka.""" + try: + await self._broker.publish( + message=event, + topic=topic, + key=key.encode(), + ) + self._logger.debug("Event sent to topic", event_type=event.event_type, topic=topic) + except Exception: + self._logger.error("Failed to produce message", topic=topic) + raise diff --git a/backend/app/events/handlers.py b/backend/app/events/handlers.py index b864f0cb..b2af375b 100644 --- a/backend/app/events/handlers.py +++ b/backend/app/events/handlers.py @@ -34,8 +34,7 @@ async def _track_consumed( metrics: EventMetrics, event: DomainEvent, consumer_group: str, coro: Awaitable[None], ) -> None: - """Record consumption metric, await *coro*, and record failure metric on error.""" - metrics.record_kafka_message_consumed(topic=event.event_type, consumer_group=consumer_group) + """Await *coro* and record domain-level failure metric on error.""" try: await coro except Exception as e: @@ -43,9 +42,6 @@ async def _track_consumed( topic=event.event_type, event_type=event.event_type, consumer_group=consumer_group, error_type=type(e).__name__, ) - metrics.record_kafka_consumption_error( - topic=event.event_type, consumer_group=consumer_group, error_type=type(e).__name__, - ) raise @@ -266,9 +262,7 @@ def register_sse_subscriber(broker: KafkaBroker, settings: Settings) -> None: async def on_sse_event( body: DomainEvent, sse_bus: FromDishka[SSERedisBus], - event_metrics: FromDishka[EventMetrics], ) -> None: - event_metrics.record_kafka_message_consumed(topic=body.event_type, consumer_group=group_id) execution_id = getattr(body, "execution_id", None) if execution_id: sse_data = SSEExecutionEventData(**{ diff --git a/backend/app/services/k8s_worker/worker.py b/backend/app/services/k8s_worker/worker.py index d78084af..14fbff23 100644 --- a/backend/app/services/k8s_worker/worker.py +++ b/backend/app/services/k8s_worker/worker.py @@ -7,7 +7,7 @@ from kubernetes_asyncio import client as k8s_client from kubernetes_asyncio.client.rest import ApiException -from app.core.metrics import EventMetrics, ExecutionMetrics, KubernetesMetrics +from app.core.metrics import ExecutionMetrics, KubernetesMetrics from app.domain.enums import ExecutionErrorType from app.domain.events import ( CreatePodCommandEvent, @@ -41,9 +41,7 @@ def __init__( producer: UnifiedProducer, settings: Settings, logger: structlog.stdlib.BoundLogger, - event_metrics: EventMetrics, ): - self._event_metrics = event_metrics self.logger = logger self.metrics = KubernetesMetrics(settings) self.execution_metrics = ExecutionMetrics(settings) diff --git a/backend/grafana/provisioning/dashboards/http-middleware.json b/backend/grafana/provisioning/dashboards/http-middleware.json index f5305117..5e2bd2f6 100644 --- a/backend/grafana/provisioning/dashboards/http-middleware.json +++ b/backend/grafana/provisioning/dashboards/http-middleware.json @@ -222,14 +222,14 @@ "id": 9, "targets": [ { - "expr": "rate(kafka_production_errors_total[5m]) or vector(0)", - "legendFormat": "Production", + "expr": "sum(rate(event_processing_errors_total[5m])) by (error_type) or vector(0)", + "legendFormat": "{{error_type}}", "refId": "A", "datasource": "Victoria Metrics" }, { - "expr": "rate(kafka_consumption_errors_total[5m]) or vector(0)", - "legendFormat": "Consumption", + "expr": "sum(rate(event_processing_errors_total[5m])) by (consumer_group) or vector(0)", + "legendFormat": "{{consumer_group}}", "refId": "B", "datasource": "Victoria Metrics" } diff --git a/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json b/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json index f16a07e5..bb3daf20 100644 --- a/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json +++ b/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json @@ -130,7 +130,7 @@ { "datasource": "Victoria Metrics", "editorMode": "code", - "expr": "sum(rate(kafka_messages_produced_total[1m])) or vector(0)", + "expr": "sum(rate(messaging_publish_messages_total[1m])) or vector(0)", "instant": false, "legendFormat": "Messages/sec", "range": true, @@ -324,7 +324,7 @@ { "datasource": "Victoria Metrics", "editorMode": "code", - "expr": "sum(rate(kafka_messages_produced_total[1m])) or vector(0)", + "expr": "sum(rate(messaging_publish_messages_total[1m])) or vector(0)", "instant": false, "legendFormat": "Produced", "range": true, @@ -333,7 +333,7 @@ { "datasource": "Victoria Metrics", "editorMode": "code", - "expr": "sum(rate(kafka_messages_consumed_total[1m])) or vector(0)", + "expr": "sum(rate(messaging_process_messages_total[1m])) or vector(0)", "instant": false, "legendFormat": "Consumed", "range": true, @@ -436,9 +436,9 @@ { "datasource": "Victoria Metrics", "editorMode": "code", - "expr": "sum(rate(kafka_messages_produced_total[1m])) by (topic) or vector(0)", + "expr": "sum(rate(messaging_publish_messages_total[1m])) by (messaging_destination_name) or vector(0)", "instant": false, - "legendFormat": "{{topic}}", + "legendFormat": "{{messaging_destination_name}}", "range": true, "refId": "A" } @@ -654,9 +654,9 @@ { "datasource": "Victoria Metrics", "editorMode": "code", - "expr": "sum by (consumer_group) (rate(kafka_messages_consumed_total[1m])) or vector(0)", + "expr": "sum(rate(messaging_process_messages_total[1m])) by (messaging_destination_publish_name) or vector(0)", "instant": false, - "legendFormat": "{{consumer_group}}", + "legendFormat": "{{messaging_destination_publish_name}}", "range": true, "refId": "A" } @@ -956,9 +956,9 @@ { "datasource": "Victoria Metrics", "editorMode": "code", - "expr": "rate(kafka_messages_produced_total[1m]) or vector(0)", + "expr": "rate(messaging_publish_messages_total[1m]) or vector(0)", "instant": false, - "legendFormat": "{{topic}}", + "legendFormat": "{{messaging_destination_name}}", "range": true, "refId": "A" } @@ -1044,9 +1044,9 @@ { "datasource": "Victoria Metrics", "editorMode": "code", - "expr": "sum(kafka_messages_produced_total) by (topic) or vector(0)", + "expr": "sum(messaging_publish_messages_total) by (messaging_destination_name) or vector(0)", "instant": false, - "legendFormat": "{{topic}}-{{partition}}", + "legendFormat": "{{messaging_destination_name}}", "range": true, "refId": "A" } diff --git a/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py b/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py index a61b3f5c..8eb620d8 100644 --- a/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py +++ b/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py @@ -28,17 +28,4 @@ def test_event_metrics_methods(test_settings: Settings) -> None: m = EventMetrics(test_settings) m.record_event_published("execution.requested", None) m.record_event_processing_duration(0.05, "execution.requested") - m.record_event_replay_operation("prepare", "success") - m.record_event_stored("execution.requested", "events") m.record_events_processing_failed("topic", "etype", "group", "error") - m.record_event_store_duration(0.1, "insert", "events") - m.record_event_store_failed("etype", "fail") - m.record_event_query_duration(0.2, "by_type", "events") - m.record_processing_duration(0.3, "topic", "etype", "group") - m.record_kafka_message_produced("t") - m.record_kafka_message_consumed("t", "g") - m.record_kafka_production_error("t", "e") - m.record_kafka_consumption_error("t", "g", "e") - m.update_event_bus_queue_size(1, "default") - m.set_event_bus_queue_size(5, "default") - m.set_event_bus_queue_size(2, "default") diff --git a/backend/tests/unit/core/metrics/test_metrics_classes.py b/backend/tests/unit/core/metrics/test_metrics_classes.py index 5fcf7b6a..a8600153 100644 --- a/backend/tests/unit/core/metrics/test_metrics_classes.py +++ b/backend/tests/unit/core/metrics/test_metrics_classes.py @@ -33,19 +33,7 @@ def test_event_metrics_smoke(test_settings: Settings) -> None: m = EventMetrics(test_settings) m.record_event_published("execution.requested") m.record_event_processing_duration(0.01, "execution.requested") - m.record_event_replay_operation("replay", "success") - m.record_event_stored("x", "events") m.record_events_processing_failed("t", "x", "g", "ValueError") - m.record_event_store_duration(0.01, "store", "events") - m.record_event_store_failed("x", "RuntimeError") - m.record_event_query_duration(0.02, "by_id", "events") - m.record_processing_duration(0.03, "t", "x", "g") - m.record_kafka_message_produced("t") - m.record_kafka_message_consumed("t", "g") - m.record_kafka_production_error("t", "E") - m.record_kafka_consumption_error("t", "g", "E") - m.update_event_bus_queue_size(1) - m.set_event_bus_queue_size(5) def test_other_metrics_classes_smoke(test_settings: Settings) -> None: From a2640ac4e65e772d92abc7638abae870deced183 Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Sun, 1 Mar 2026 23:09:42 +0100 Subject: [PATCH 2/4] fix: removed try-catch-raise --- backend/app/events/core/transport.py | 16 ++++++---------- 1 file changed, 6 insertions(+), 10 deletions(-) diff --git a/backend/app/events/core/transport.py b/backend/app/events/core/transport.py index c05dde7d..0e89f076 100644 --- a/backend/app/events/core/transport.py +++ b/backend/app/events/core/transport.py @@ -21,13 +21,9 @@ def __init__( async def publish(self, event: DomainEvent, topic: str, key: str) -> None: """Publish event to Kafka.""" - try: - await self._broker.publish( - message=event, - topic=topic, - key=key.encode(), - ) - self._logger.debug("Event sent to topic", event_type=event.event_type, topic=topic) - except Exception: - self._logger.error("Failed to produce message", topic=topic) - raise + await self._broker.publish( + message=event, + topic=topic, + key=key.encode(), + ) + self._logger.debug("Event sent to topic", event_type=event.event_type, topic=topic) From ca1508ce89d83b8f15b08fbc2e39799d06d2cedc Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Sun, 1 Mar 2026 23:32:06 +0100 Subject: [PATCH 3/4] fix: detected issues --- .../app/db/repositories/event_repository.py | 11 ++++--- .../dashboards/http-middleware.json | 30 ++----------------- .../dashboards/kafka-events-monitoring.json | 4 +-- 3 files changed, 12 insertions(+), 33 deletions(-) diff --git a/backend/app/db/repositories/event_repository.py b/backend/app/db/repositories/event_repository.py index f04043cc..50d7cd2f 100644 --- a/backend/app/db/repositories/event_repository.py +++ b/backend/app/db/repositories/event_repository.py @@ -58,10 +58,13 @@ async def store_event(self, event: DomainEvent) -> str: return event.event_id async def mark_publish_failed(self, event_id: str) -> None: - """Mark an event as failed to publish to Kafka for later retry.""" - await EventDocument.find_one( - EventDocument.event_id == event_id, - ).update({"$set": {"publish_failed": True, "publish_failed_at": datetime.now(timezone.utc)}}) + """Best-effort mark of an event as failed to publish. Never raises.""" + try: + await EventDocument.find_one( + EventDocument.event_id == event_id, + ).update({"$set": {"publish_failed": True, "publish_failed_at": datetime.now(timezone.utc)}}) + except Exception: + self.logger.warning("Could not mark event as publish-failed", event_id=event_id) async def get_event(self, event_id: str) -> DomainEvent | None: doc = await EventDocument.find_one(EventDocument.event_id == event_id) diff --git a/backend/grafana/provisioning/dashboards/http-middleware.json b/backend/grafana/provisioning/dashboards/http-middleware.json index 5e2bd2f6..182f57bc 100644 --- a/backend/grafana/provisioning/dashboards/http-middleware.json +++ b/backend/grafana/provisioning/dashboards/http-middleware.json @@ -203,7 +203,7 @@ }, "id": 102, "panels": [], - "title": "Kafka Errors & Event Bus", + "title": "Event Processing Errors", "type": "row" }, { @@ -215,7 +215,7 @@ }, "gridPos": { "h": 8, - "w": 12, + "w": 24, "x": 0, "y": 27 }, @@ -234,33 +234,9 @@ "datasource": "Victoria Metrics" } ], - "title": "Kafka Errors", + "title": "Event Processing Errors", "type": "timeseries" }, - { - "datasource": "Victoria Metrics", - "fieldConfig": { - "defaults": { - "unit": "short" - } - }, - "gridPos": { - "h": 8, - "w": 12, - "x": 12, - "y": 27 - }, - "id": 10, - "targets": [ - { - "expr": "event_bus_queue_size or vector(0)", - "refId": "A", - "datasource": "Victoria Metrics" - } - ], - "title": "Event Bus Queue Size", - "type": "stat" - }, { "collapsed": false, "datasource": null, diff --git a/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json b/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json index bb3daf20..68be9782 100644 --- a/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json +++ b/backend/grafana/provisioning/dashboards/kafka-events-monitoring.json @@ -577,7 +577,7 @@ }, { "datasource": "Victoria Metrics", - "description": "Message consumption rate by consumer group", + "description": "Message consumption rate by topic", "fieldConfig": { "defaults": { "color": { @@ -661,7 +661,7 @@ "refId": "A" } ], - "title": "Consumer Rate by Group", + "title": "Consumer Rate by Topic", "type": "timeseries" }, { From b358cc9622cc6113ec6992fd171c2fa99c3377a1 Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Sun, 1 Mar 2026 23:41:48 +0100 Subject: [PATCH 4/4] fix: logging caught error in event repo --- backend/app/db/repositories/event_repository.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/backend/app/db/repositories/event_repository.py b/backend/app/db/repositories/event_repository.py index 50d7cd2f..f9b9d54a 100644 --- a/backend/app/db/repositories/event_repository.py +++ b/backend/app/db/repositories/event_repository.py @@ -63,8 +63,8 @@ async def mark_publish_failed(self, event_id: str) -> None: await EventDocument.find_one( EventDocument.event_id == event_id, ).update({"$set": {"publish_failed": True, "publish_failed_at": datetime.now(timezone.utc)}}) - except Exception: - self.logger.warning("Could not mark event as publish-failed", event_id=event_id) + except Exception as exc: + self.logger.warning("Could not mark event as publish-failed", event_id=event_id, error=str(exc)) async def get_event(self, event_id: str) -> DomainEvent | None: doc = await EventDocument.find_one(EventDocument.event_id == event_id)