mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
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.
This commit is contained in:
parent
a42efec324
commit
742fde2330
2 changed files with 138 additions and 28 deletions
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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"):
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue