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
2 changes: 2 additions & 0 deletions backend/app/core/providers.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from dishka import Provider, Scope, from_context, provide
from faststream.kafka import KafkaBroker
from faststream.kafka.opentelemetry import KafkaTelemetryMiddleware
from kubernetes_asyncio import client as k8s_client
from kubernetes_asyncio import config as k8s_config
from kubernetes_asyncio.client.rest import ApiException
Expand Down Expand Up @@ -96,6 +97,7 @@ async def get_broker(
logger=logger,
client_id=f"integr8scode-{settings.SERVICE_NAME}",
request_timeout_ms=settings.KAFKA_REQUEST_TIMEOUT_MS,
middlewares=(KafkaTelemetryMiddleware(),),
)
logger.info("Kafka broker created")
try:
Expand Down
4 changes: 0 additions & 4 deletions backend/app/core/tracing/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,7 @@
# Import utilities and decorators
from app.core.tracing.utils import (
add_span_attributes,
extract_trace_context,
get_tracer,
inject_trace_context,
trace_span,
)

Expand All @@ -38,9 +36,7 @@
"init_tracing",
# Utilities and decorators
"add_span_attributes",
"extract_trace_context",
"get_tracer",
"inject_trace_context",
"trace_span",
# OpenTelemetry types
"context",
Expand Down
29 changes: 1 addition & 28 deletions backend/app/core/tracing/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
from contextlib import contextmanager
from typing import Any

from opentelemetry import context, propagate, trace
from opentelemetry import trace
from opentelemetry.trace import SpanKind, Status, StatusCode


Expand Down Expand Up @@ -45,33 +45,6 @@ def trace_span(
raise


def inject_trace_context(headers: dict[str, str]) -> dict[str, str]:
"""
Inject current trace context into headers for propagation.

Args:
headers: Existing headers dictionary

Returns:
Headers with trace context injected
"""
propagation_headers = headers.copy()
propagate.inject(propagation_headers)
return propagation_headers


def extract_trace_context(headers: dict[str, str]) -> context.Context:
"""
Extract trace context from headers.

Args:
headers: Headers containing trace context

Returns:
Extracted OpenTelemetry context
"""
return propagate.extract(headers)


def add_span_attributes(**attributes: Any) -> None:
"""
Expand Down
8 changes: 0 additions & 8 deletions backend/app/dlq/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@
from faststream.kafka import KafkaBroker

from app.core.metrics import DLQMetrics
from app.core.tracing import inject_trace_context
from app.db.repositories import DLQRepository
from app.dlq.models import (
DLQBatchRetryResult,
Expand Down Expand Up @@ -127,17 +126,10 @@ async def retry_message(self, message: DLQMessage) -> None:

FastStream handles JSON serialization of Pydantic models natively.
"""
hdrs: dict[str, str] = {
"event_type": message.event.event_type,
}
hdrs = inject_trace_context(hdrs)

# Publish directly to original topic - FastStream serializes Pydantic to JSON
await self._broker.publish(
message=message.event,
topic=message.original_topic,
key=message.event.event_id.encode(),
headers=hdrs,
)

self.metrics.record_dlq_message_retried(message.original_topic, message.event.event_type, "success")
Expand Down
1 change: 0 additions & 1 deletion backend/app/dlq/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,6 @@ class DLQMessage(BaseModel):
dlq_offset: int | None = None
dlq_partition: int | None = None
last_error: str | None = None
headers: dict[str, str] = Field(default_factory=dict)


@dataclass
Expand Down
32 changes: 11 additions & 21 deletions backend/app/events/core/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,8 @@
from faststream.kafka import KafkaBroker

from app.core.metrics import EventMetrics
from app.core.tracing import inject_trace_context
from app.db.repositories import EventRepository
from app.dlq.models import DLQMessageStatus
from app.dlq.models import DLQMessage, DLQMessageStatus
from app.domain.enums import KafkaTopic
from app.domain.events import DomainEvent
from app.infrastructure.kafka.mappings import EVENT_TYPE_TO_TOPIC
Expand Down Expand Up @@ -41,16 +40,10 @@ 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:
headers = inject_trace_context({
"event_type": event_to_produce.event_type,
})

# FastStream handles Pydantic → JSON serialization natively
await self._broker.publish(
message=event_to_produce,
topic=topic,
key=key.encode(),
headers=headers,
)

self._event_metrics.record_kafka_message_produced(topic)
Expand All @@ -72,23 +65,20 @@ async def send_to_dlq(

dlq_topic = f"{self._topic_prefix}{KafkaTopic.DEAD_LETTER_QUEUE}"

headers = inject_trace_context({
"event_type": original_event.event_type,
"original_topic": original_topic,
"error_type": type(error).__name__,
"error": str(error),
"retry_count": str(retry_count),
"failed_at": datetime.now(timezone.utc).isoformat(),
"status": DLQMessageStatus.PENDING,
"producer_id": producer_id,
})
dlq_msg = DLQMessage(
event=original_event,
original_topic=original_topic,
error=str(error),
retry_count=retry_count,
failed_at=datetime.now(timezone.utc),
status=DLQMessageStatus.PENDING,
producer_id=producer_id,
)

# FastStream handles Pydantic → JSON serialization natively
await self._broker.publish(
message=original_event,
message=dlq_msg,
topic=dlq_topic,
key=original_event.event_id.encode(),
headers=headers,
)

self._event_metrics.record_kafka_message_produced(dlq_topic)
Expand Down
Loading
Loading