From 742fde2330e8de036ca1d7c519a7e16f01bd13b8 Mon Sep 17 00:00:00 2001 From: yucheng-berriai Date: Tue, 23 Jun 2026 11:12:51 -0700 Subject: [PATCH] fix(otel/v2): forward proxy-internal spans to every destination; wrap with backend-specific Resource The first cut of the fan-out processor had two bugs caught in live testing against real Arize and Langfuse Cloud: 1. The owner-match filter rejected destinations whose callback_name didn't match the v2 logger's own backend. For a YAML configured with only langfuse_otel, a Team A request whose only assigned destination is arize would have every proxy-internal span dropped at the filter, so Arize only received the gen-AI span (via the existing TenantTracerCache clone provider) and rendered it as an orphan. Proxy-internal spans (FastAPI server, auth phase, postgres lookups, batch-write cost ledger) are generic OTel semconv with no backend-specific vocabulary. They ship to every admin-resolved destination this request fans out to, regardless of callback_name. The owner discriminator is now used only to skip the gen-AI span, which the per-backend v2 logger already routes via the clone provider with its own attribute mapper -- forwarding it here too would deliver a duplicate with the wrong vocabulary. 2. Arize rejected the forwarded proxy-internal spans with 'InvalidArgument: model_id span resource attribute or arize.project.name span attribute is required'. The arize preset's own Resource carries model_id, but my fan-out built a fresh OTLPSpanExporter whose spans inherited the main provider's Resource (no model_id). Added _with_destination_resource which wraps each forwarded span with a destination-specific Resource (for arize: model_id + arize.project.name from ARIZE_PROJECT_NAME env). The wrapper is a shallow view; the underlying span object is unchanged. End-to-end verified: a single Team B request now delivers the full proxy-internal tree (server + auth + postgres + chat + batch_write) under one trace_id to both Arize Cloud and Langfuse Cloud, parented correctly, no orphan warnings. Tests updated: the prior owner-rejection test was inverted (it pinned the wrong behavior); added a new test that pins the gen-AI skip. All 7 fan-out tests pass; mutation-checked by removing the wiring. --- litellm/integrations/otel/plumbing/fan_out.py | 88 +++++++++++++++++-- .../integrations/otel/test_otel_v2_fan_out.py | 78 +++++++++++----- 2 files changed, 138 insertions(+), 28 deletions(-) diff --git a/litellm/integrations/otel/plumbing/fan_out.py b/litellm/integrations/otel/plumbing/fan_out.py index dbc3f8c4ec2..91973ff8b7c 100644 --- a/litellm/integrations/otel/plumbing/fan_out.py +++ b/litellm/integrations/otel/plumbing/fan_out.py @@ -10,9 +10,10 @@ gen-AI LLM-call span and its MCP-tool sibling) reach tenant backends through the clone's own exporters; the clone's provider has its own processor list and does NOT carry this fan-out processor, so a span is exported once per backend. -Identical ``(trace_id, span_id)`` deduplication on the OTLP receiver collapses -the rare overlap into one span per backend (e.g. when the configured global -exporter and a per-tenant destination point at the same vendor account). +Each backend often requires backend-specific Resource attributes (Arize rejects +spans missing ``model_id`` / ``arize.project.name``), so the fan-out wraps each +forwarded span with the destination's expected Resource before handing it to +the per-destination exporter. """ from __future__ import annotations @@ -20,7 +21,10 @@ from __future__ import annotations from collections import OrderedDict from typing import TYPE_CHECKING +import os + from opentelemetry.context import Context +from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import ReadableSpan, Span, SpanProcessor from litellm._logging import verbose_logger @@ -63,14 +67,27 @@ class TenantFanOutSpanProcessor(SpanProcessor): destinations = request_destinations() if not destinations: return + # The gen-AI LLM-call span (and the MCP tool-call sibling) is already + # routed to per-tenant destinations by the per-backend v2 logger via + # ``TenantTracerCache`` -- the logger picks the right attribute mapper + # (OpenInference for arize, GenAI semconv for langfuse_otel) and ships + # through the clone provider's appended exporter. Forwarding it here too + # would deliver a SECOND copy with the wrong vocabulary and a fresh + # span_id, surfacing in the destination as an orphaned duplicate. Skip. + if _is_genai_span(span): + return + # Proxy-internal spans (FastAPI server, ``auth`` phase, postgres lookups, + # post-call cost ledger) are generic OTel semantic-convention spans with + # no backend-specific vocabulary, so they ship to EVERY admin-resolved + # destination this request fans out to, regardless of the destination's + # ``callback_name``. The owner discriminator only matters for the gen-AI + # span (handled by the skip above). for destination in destinations: - if destination.callback_name != self._owner: - continue processor = self._processor_for(destination) if processor is None: continue try: - processor.on_end(span) + processor.on_end(_with_destination_resource(span, destination)) except Exception as exc: verbose_logger.debug( "OTel V2 fan-out: forwarding span to %s failed: %s", @@ -141,3 +158,62 @@ _GRPC_BACKENDS = frozenset({"arize"}) def _resolve_kind(destination: "OtelDestination") -> str: return "otlp_grpc" if destination.callback_name in _GRPC_BACKENDS else "otlp_http" + + +# Attribute set on every gen-AI LLM-call span by the v2 emitter. Used as the +# unambiguous skip signal: only the LLM-call span carries this, and the +# per-backend v2 logger already routes it to per-tenant destinations through +# the TenantTracerCache clone provider's appended exporter. +_GENAI_SPAN_ATTR = "gen_ai.operation.name" + + +def _is_genai_span(span: ReadableSpan) -> bool: + attributes = span.attributes or {} + return _GENAI_SPAN_ATTR in attributes + + +# Backend-specific Resource attributes required by the destination. Arize +# rejects spans missing ``model_id`` (or the alternative ``arize.project.name`` +# span attribute); other backends accept the proxy's default Resource. +def _destination_resource_attrs(destination: "OtelDestination") -> dict[str, str]: + if destination.callback_name == "arize": + project = os.environ.get("ARIZE_PROJECT_NAME") + if project: + return {"model_id": project, "arize.project.name": project} + return {} + + +def _with_destination_resource(span: ReadableSpan, destination: "OtelDestination") -> ReadableSpan: + """Return ``span`` with its Resource augmented by the destination's required + attributes. The original span object is left untouched; a shallow wrapper + reuses every other field and only swaps the ``resource`` property.""" + extra = _destination_resource_attrs(destination) + if not extra: + return span + merged = Resource.create({**dict(span.resource.attributes), **extra}) + return _ResourceWrappedReadableSpan(span, merged) + + +class _ResourceWrappedReadableSpan(ReadableSpan): + """A ``ReadableSpan`` view whose ``resource`` is overridden. + + The OTLP exporter reads each span's ``resource`` when serializing; for + backend-specific attributes (Arize's ``model_id``) we substitute the + destination's expected Resource without mutating the underlying span. + """ + + def __init__(self, inner: ReadableSpan, resource: Resource) -> None: + super().__init__( + name=inner.name, + context=inner.context, + parent=inner.parent, + resource=resource, + attributes=inner.attributes, + events=inner.events, + links=inner.links, + kind=inner.kind, + status=inner.status, + start_time=inner.start_time, + end_time=inner.end_time, + instrumentation_scope=inner.instrumentation_scope, + ) diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py b/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py index 7d44522d2c1..36610fee18a 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py @@ -57,13 +57,9 @@ def _build_provider_with_fan_out( test by an injected SimpleSpanProcessor against an in-memory exporter via ``monkeypatch`` so the test can read what was forwarded.""" cfg = OpenTelemetryV2Config(exporter="in_memory") - provider = providers.build_tracer_provider( - cfg, exporter=exporter, tenant_fan_out_owner=owner - ) + provider = providers.build_tracer_provider(cfg, exporter=exporter, tenant_fan_out_owner=owner) fan_out = next( - p - for p in provider._active_span_processor._span_processors - if isinstance(p, TenantFanOutSpanProcessor) + p for p in provider._active_span_processor._span_processors if isinstance(p, TenantFanOutSpanProcessor) ) return provider, fan_out @@ -76,14 +72,13 @@ def test_fan_out_forwards_span_to_matching_destination(monkeypatch): provider, fan_out = _build_provider_with_fan_out("langfuse_otel", main_exporter) tenant_exporter = InMemorySpanExporter() + # Swap the lazy per-destination processor build for an in-memory one so we # don't actually hit the network. def _stub_processor_for(self, destination): return SimpleSpanProcessor(tenant_exporter) - monkeypatch.setattr( - TenantFanOutSpanProcessor, "_processor_for", _stub_processor_for - ) + monkeypatch.setattr(TenantFanOutSpanProcessor, "_processor_for", _stub_processor_for) set_request_destinations( ( @@ -104,13 +99,15 @@ def test_fan_out_forwards_span_to_matching_destination(monkeypatch): assert "auth /chat/completions" in tenant_names -def test_fan_out_skips_destinations_for_other_backends(monkeypatch): - """Each fan-out processor owns ONE backend; a destination tagged for a - different backend is skipped, so the wrong exporter never sees the span.""" +def test_fan_out_forwards_proxy_internal_spans_to_every_destination(monkeypatch): + """Proxy-internal spans (server, auth phase, postgres, post-call ledger) + are generic OTel semconv with no backend-specific vocabulary, so they ship + to every admin-resolved destination regardless of its ``callback_name``. + The owner discriminator is only used to skip the gen-AI span (the v2 logger + routes that itself through TenantTracerCache to avoid wrong-vocabulary + duplicates).""" main_exporter = InMemorySpanExporter() - _provider, _fan_out = _build_provider_with_fan_out( - "langfuse_otel", main_exporter - ) + _provider, _ = _build_provider_with_fan_out("langfuse_otel", main_exporter) tenant_exporter = InMemorySpanExporter() monkeypatch.setattr( @@ -119,6 +116,8 @@ def test_fan_out_skips_destinations_for_other_backends(monkeypatch): lambda self, destination: SimpleSpanProcessor(tenant_exporter), ) + # Destination is tagged for a DIFFERENT backend (``arize``) than the + # owner (``langfuse_otel``). Proxy-internal spans must still forward. set_request_destinations( ( OtelDestination( @@ -132,7 +131,46 @@ def test_fan_out_skips_destinations_for_other_backends(monkeypatch): with tracer.start_as_current_span("auth"): pass - assert tenant_exporter.get_finished_spans() == () + names = [s.name for s in tenant_exporter.get_finished_spans()] + assert "auth" in names + + +def test_fan_out_skips_genai_span_to_avoid_double_export(monkeypatch): + """The gen-AI LLM-call span is routed by the v2 logger through + ``TenantTracerCache`` to per-tenant exporters with the right attribute + mapper. Forwarding it here too would deliver a duplicate with the wrong + vocabulary, surfacing in the destination as an orphaned second span.""" + main_exporter = InMemorySpanExporter() + _provider, _ = _build_provider_with_fan_out("langfuse_otel", main_exporter) + + tenant_exporter = InMemorySpanExporter() + monkeypatch.setattr( + TenantFanOutSpanProcessor, + "_processor_for", + lambda self, destination: SimpleSpanProcessor(tenant_exporter), + ) + + set_request_destinations( + ( + OtelDestination( + callback_name="arize", + endpoint="https://otlp.arize.com/v1", + headers={}, + ), + ) + ) + tracer = _provider.get_tracer("test") + # A gen-AI span sets ``gen_ai.operation.name`` (the v2 emitter does this + # via the SpanRole.LLM_CALL mapper). Fake it explicitly here. + with tracer.start_as_current_span("chat gpt-4o") as span: + span.set_attribute("gen_ai.operation.name", "chat") + # A proxy-internal span also fires. + with tracer.start_as_current_span("auth /v1/chat/completions"): + pass + + names = [s.name for s in tenant_exporter.get_finished_spans()] + assert "auth /v1/chat/completions" in names + assert "chat gpt-4o" not in names # gen-AI skipped def test_fan_out_noop_when_no_destinations(monkeypatch): @@ -183,12 +221,8 @@ def test_fan_out_caches_processor_per_destination_key(monkeypatch): monkeypatch.setattr(TenantFanOutSpanProcessor, "_processor_for", _stub) - dest_a = OtelDestination( - callback_name="langfuse_otel", endpoint="https://a/", headers={"k": "1"} - ) - dest_b = OtelDestination( - callback_name="langfuse_otel", endpoint="https://b/", headers={"k": "1"} - ) + dest_a = OtelDestination(callback_name="langfuse_otel", endpoint="https://a/", headers={"k": "1"}) + dest_b = OtelDestination(callback_name="langfuse_otel", endpoint="https://b/", headers={"k": "1"}) set_request_destinations((dest_a,)) tracer = _provider.get_tracer("test") with tracer.start_as_current_span("s1"):