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
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,8 @@
<a href="https://github.com/HardMax71/Integr8sCode/actions/workflows/security.yml">
<img src="https://img.shields.io/github/actions/workflow/status/HardMax71/Integr8sCode/security.yml?branch=main&label=security&logo=shieldsdotio&logoColor=white" alt="Security Scan Status" />
</a>
<a href="https://github.com/HardMax71/Integr8sCode/actions/workflows/vulture.yml">
<img src="https://img.shields.io/github/actions/workflow/status/HardMax71/Integr8sCode/vulture.yml?branch=main&label=dead%20code&logo=python&logoColor=white" alt="Dead Code Check" />
<a href="https://github.com/HardMax71/Integr8sCode/actions/workflows/grimp.yml">
<img src="https://img.shields.io/github/actions/workflow/status/HardMax71/Integr8sCode/grimp.yml?branch=main&label=dead%20code&logo=python&logoColor=white" alt="Dead Code Check" />
</a>
<a href="https://github.com/HardMax71/Integr8sCode/actions/workflows/docker.yml">
<img src="https://img.shields.io/github/actions/workflow/status/HardMax71/Integr8sCode/docker.yml?branch=main&label=docker&logo=docker&logoColor=white" alt="Docker Scan Status" />
Expand Down
87 changes: 5 additions & 82 deletions backend/app/core/metrics/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand All @@ -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
Expand All @@ -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:
Expand All @@ -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)
17 changes: 11 additions & 6 deletions backend/app/core/providers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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,
)


Expand Down
11 changes: 7 additions & 4 deletions backend/app/db/repositories/event_repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Comment thread
HardMax71 marked this conversation as resolved.
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)
Expand Down
2 changes: 2 additions & 0 deletions backend/app/events/core/__init__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
from .producer import UnifiedProducer
from .transport import KafkaEventTransport

__all__ = [
"KafkaEventTransport",
"UnifiedProducer",
]
35 changes: 9 additions & 26 deletions backend/app/events/core/producer.py
Original file line number Diff line number Diff line change
@@ -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.
Expand All @@ -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)
Comment thread
HardMax71 marked this conversation as resolved.
except Exception:
await self._event_repository.mark_publish_failed(event_to_produce.event_id)
raise
Comment thread
HardMax71 marked this conversation as resolved.
Comment thread
HardMax71 marked this conversation as resolved.
29 changes: 29 additions & 0 deletions backend/app/events/core/transport.py
Original file line number Diff line number Diff line change
@@ -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)
8 changes: 1 addition & 7 deletions backend/app/events/handlers.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,18 +34,14 @@
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:
metrics.record_events_processing_failed(
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


Expand Down Expand Up @@ -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(**{
Expand Down
4 changes: 1 addition & 3 deletions backend/app/services/k8s_worker/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading