diff --git a/README.md b/README.md
index acfecf08..dfcf8041 100644
--- a/README.md
+++ b/README.md
@@ -12,8 +12,8 @@
-
-
+
+
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/db/repositories/event_repository.py b/backend/app/db/repositories/event_repository.py
index f04043cc..f9b9d54a 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 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)
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..0e89f076
--- /dev/null
+++ b/backend/app/events/core/transport.py
@@ -0,0 +1,29 @@
+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."""
+ 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)
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..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,52 +215,28 @@
},
"gridPos": {
"h": 8,
- "w": 12,
+ "w": 24,
"x": 0,
"y": 27
},
"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"
}
],
- "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 f16a07e5..68be9782 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"
}
@@ -577,7 +577,7 @@
},
{
"datasource": "Victoria Metrics",
- "description": "Message consumption rate by consumer group",
+ "description": "Message consumption rate by topic",
"fieldConfig": {
"defaults": {
"color": {
@@ -654,14 +654,14 @@
{
"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"
}
],
- "title": "Consumer Rate by Group",
+ "title": "Consumer Rate by Topic",
"type": "timeseries"
},
{
@@ -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: