fix(otel v2): deliver tenant destinations from the published global provider

Scoping the fan-out by callback name in the previous commit left every backend
that is not the canonical logger with a one-span trace: only the published
global provider sees the FastAPI server span, the auth span and the post-call
database spans, so an arize-only proxy handed a team's Langfuse just the model
call. Attach the fan-out once, to that provider, and let it forward every
destination.

An overridden backend now skips per-request tracer routing outright rather than
only clearing its credential headers, since a key or team otel_service_name was
still enough to detach the model call onto a second provider. The destination
carries that service name as a resource attribute instead.

Shed processors drain on a two-thread pool rather than a thread each, so a
tenant cycling its destination config cannot spawn threads as fast as it sends
requests.
This commit is contained in:
Yucheng He 2026-09-03 18:21:34 -07:00
parent fb3b32d22d
commit f1cab32449
6 changed files with 252 additions and 60 deletions

View file

@ -63,6 +63,7 @@ from litellm.integrations.otel.plumbing.metrics import (
create_genai_metrics,
)
from litellm.integrations.otel.plumbing.providers import (
attach_tenant_fan_out,
build_tracer_provider,
get_event_logger,
get_meter,
@ -182,7 +183,7 @@ class OpenTelemetryV2(CustomLogger):
self._tracer_provider: TracerProvider = (
tracer_provider
if tracer_provider is not None
else build_tracer_provider(self.config, tenant_overrides=True, tenant_callback_name=callback_name)
else build_tracer_provider(self.config, tenant_overrides=True)
)
self.tracer: Tracer = get_tracer(self._tracer_provider, LITELLM_TRACER_NAME)
self._metrics_recorder = self._init_metrics(meter_provider)
@ -865,8 +866,13 @@ def publish_global_otel_v2_provider(
``opentelemetry.trace.set_tracer_provider``) are injected so the publish step is
unit-testable without reading or mutating real global OTel state. Returns the
logger whose provider was published.
The published provider is also the one that fans spans out to key/team
destinations, because it is the only provider the whole request tree passes
through; see :func:`attach_tenant_fan_out`.
"""
logger: Final = select_global_otel_v2_logger(in_memory_loggers, registered=registered)
attach_tenant_fan_out(logger._tracer_provider)
set_global_provider(logger._tracer_provider)
return logger

View file

@ -3,6 +3,8 @@
import threading
from collections import OrderedDict
from collections.abc import Callable, Iterable
from concurrent.futures import ThreadPoolExecutor
from functools import lru_cache
from typing import TYPE_CHECKING, Any, Final, Literal
from opentelemetry import _logs, baggage, metrics
@ -246,20 +248,20 @@ class TenantFanOutSpanProcessor(SpanProcessor):
requests stay isolated. The forwarded view keeps the original trace and parent
ids, so the tenant gets the same tree the operator would have received.
``callback_name`` scopes the fan-out to the backend this provider speaks for. Every
v2 logger carries its own provider and emits its own copy of a gen-AI span, so a
proxy running two of them would otherwise deliver the tenant two copies of the same
model call. The tenant's own backend always has a logger, since naming it in the
key or team config is what builds one.
Exactly one provider carries this processor, the one published as the OTel global
(see :func:`attach_tenant_fan_out`). That provider is the only one every span
passes through: the FastAPI server span, the auth span and the post-call database
spans are emitted on the global, while a second v2 logger's provider sees only
that logger's own gen-AI span. Attaching the fan-out per logger would hand a
tenant a one-span trace whenever its backend is not the global one, and two
copies of the model call whenever it is.
"""
def __init__(
self,
callback_name: str | None = None,
processor_factory: 'Callable[["OtelDestination"], SpanProcessor | None] | None' = None,
) -> None:
self._lock: Final = threading.Lock()
self._callback_name: Final = callback_name
self._build: Final = processor_factory if processor_factory is not None else _destination_processor
self._processors: OrderedDict[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded LRU
self._retired: OrderedDict[int, SpanProcessor] = OrderedDict() # mutable-ok: drains as exports finish
@ -270,8 +272,6 @@ class TenantFanOutSpanProcessor(SpanProcessor):
def on_end(self, span: ReadableSpan) -> None:
for destination in request_destinations():
if destination.callback_name != self._callback_name:
continue
processor = self._acquire(destination) # rebind-ok: loop variable; pyright forbids Final in a loop
if processor is None:
continue
@ -387,11 +387,16 @@ def _drain_in_background(processor: SpanProcessor) -> None:
``shutdown`` flushes over the network and is reached from ``on_end``, so closing
one inline would let a single unreachable tenant collector stall every other
tenant's spans behind it.
tenant's spans behind it. The work goes to a two-thread pool rather than a thread
per processor, so a tenant cycling its destination config cannot spawn threads as
fast as it can send requests; slow shutdowns queue behind each other instead.
"""
threading.Thread(
target=_shutdown_quietly, args=(processor,), daemon=True, name="litellm-otel-destination-drain"
).start()
_drain_pool().submit(_shutdown_quietly, processor)
@lru_cache(maxsize=1)
def _drain_pool() -> ThreadPoolExecutor:
return ThreadPoolExecutor(max_workers=2, thread_name_prefix="litellm-otel-destination-drain")
def _shutdown_quietly(processor: SpanProcessor) -> None:
@ -672,7 +677,6 @@ def build_tracer_provider(
baggage_processor: SpanProcessor | None = None,
use_simple_processor: bool | None = None,
tenant_overrides: bool = False,
tenant_callback_name: str | None = None,
) -> TracerProvider:
"""Build the shared :class:`TracerProvider`.
@ -682,11 +686,12 @@ def build_tracer_provider(
backends. ``exporter`` and ``use_simple_processor`` are explicit overrides:
pass a single exporter to attach exactly that one (used by tests).
``tenant_overrides`` belongs to the operator-level provider alone: it wraps each
owned exporter so a request that pointed that backend at a key's or team's own
account skips it, and adds the fan-out processor that delivers to that account
instead. The per-tenant providers this same function builds must leave it off,
or they would filter out the very spans they exist to carry.
``tenant_overrides`` wraps each owned exporter so a request that pointed that
backend at a key's or team's own account skips it. Every v2 logger's provider
wants it, since any of them may own the overridden backend; delivering to the
tenant is a separate job, done once by :func:`attach_tenant_fan_out`. The
per-tenant providers this same function builds must leave it off, or they would
filter out the very spans they exist to carry.
"""
provider: Final = TracerProvider(resource=build_resource(config))
if baggage_processor is None:
@ -711,11 +716,26 @@ def build_tracer_provider(
provider.add_span_processor(
_OverriddenBackendFilter(processor, owner) if tenant_overrides and owner is not None else processor
)
if tenant_overrides:
provider.add_span_processor(TenantFanOutSpanProcessor(tenant_callback_name))
return provider
def attach_tenant_fan_out(provider: TracerProvider) -> None:
"""Give ``provider`` the fan-out that delivers spans to key/team destinations.
Called on the one provider published as the OTel global, and idempotent so a
second publish (a test, a re-initialized proxy) cannot double-export.
"""
if any(isinstance(processor, TenantFanOutSpanProcessor) for processor in _attached_processors(provider)):
return
provider.add_span_processor(TenantFanOutSpanProcessor())
def _attached_processors(provider: TracerProvider) -> "tuple[SpanProcessor, ...]":
"""The processors already on ``provider``, or empty when the SDK hides them."""
multi: Final = getattr(provider, "_active_span_processor", None)
return tuple(getattr(multi, "_span_processors", ()))
def get_tracer(provider: TracerProvider, name: str = "litellm") -> Tracer:
# Stamp the instrumentation scope with the LiteLLM package version so every
# emitted span carries a deterministic ``scope.version`` (the standard OTel

View file

@ -233,13 +233,12 @@ class TenantTracerCache:
the caller's span start. The caller must ``release`` it exactly once.
"""
# An overridden backend is delivered by the fan-out processor, which carries the
# whole trace. Routing here too would detach this span onto a second provider,
# whole trace and already carries this tenant's credentials and service name.
# Routing here too would detach this span onto a second provider,
# so the tenant would get the request tree plus a stray one-span trace.
credential_headers: Final = (
_NO_HEADERS
if self._callback_name is not None and self._callback_name in overridden_backends()
else self._credential_headers(dynamic_params)
)
if self._callback_name is not None and self._callback_name in overridden_backends():
return TenantRoute(tracer=default, detached=False)
credential_headers: Final = self._credential_headers(dynamic_params)
project_headers: Final = self._project_headers(auth_metadata)
service_name: Final = tenant_service_name(auth_metadata)
if not credential_headers and not project_headers and service_name is None:

View file

@ -118,11 +118,17 @@ def destination_capable_backends() -> frozenset[str]:
return frozenset(_DESTINATION_BY_CALLBACK) & frozenset(DYNAMIC_HEADERS_BY_CALLBACK)
def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -> OtelDestination | None:
def destination_for(
callback_name: str,
params: StandardCallbackDynamicParams,
service_name: str | None = None,
) -> OtelDestination | None:
"""The destination ``params`` names for ``callback_name``, or ``None``.
``None`` means the caller configured nothing usable for this backend, so the
request keeps the operator's global exporters.
request keeps the operator's global exporters. ``service_name`` is the key's or
team's ``otel_service_name``, which the per-request tracer route applies when the
backend is not overridden and the destination has to apply once it is.
"""
from litellm.integrations.otel.presets import DYNAMIC_HEADERS_BY_CALLBACK
@ -140,7 +146,7 @@ def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -
return OtelDestination(
endpoint=endpoint,
headers=MappingProxyType(dict(headers)), # mutable-ok: MappingProxyType needs a concrete mapping to wrap
resource_attributes=_NO_ATTRS,
resource_attributes=MappingProxyType({"service.name": service_name}) if service_name else _NO_ATTRS,
callback_name=callback_name,
protocol=protocol,
)

View file

@ -1044,12 +1044,33 @@ def resolve_tenant_otel_destinations(
}
)
),
_tenant_service_name(user_api_key_dict),
)
)
is not None
)
def _tenant_service_name(user_api_key_dict: UserAPIKeyAuth) -> str | None:
"""The ``service.name`` this key or team configured, the key winning over its team.
Same fields and same precedence the request-metadata build applies, read straight
off the auth object because destinations resolve during auth, before that metadata
is assembled.
"""
sources: Final = (user_api_key_dict.metadata, user_api_key_dict.team_metadata)
return next(
(
stripped
for source in sources
if source
for field in OTEL_SERVICE_NAME_METADATA_KEYS
if isinstance(value := source.get(field), str) and (stripped := value.strip())
),
None,
)
def clean_headers(
headers: Headers,
litellm_key_header_name: str | None = None,

View file

@ -17,6 +17,10 @@ from litellm.integrations.otel.model.config import (
is_otel_v2_enabled,
)
from litellm.integrations.otel.model.destination import OtelDestination
from litellm.integrations.otel.logger import (
OpenTelemetryV2,
publish_global_otel_v2_provider,
)
from litellm.integrations.otel.plumbing.context import (
overridden_backends,
request_destinations,
@ -68,7 +72,7 @@ def wired_provider(dest_exporter: InMemorySpanExporter, global_exporter: InMemor
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(global_exporter), "langfuse_otel"))
provider.add_span_processor(
TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
)
return provider
@ -101,7 +105,7 @@ class TestOverrideSuppression:
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(arize_exporter), "arize"))
provider.add_span_processor(
TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
)
def run():
@ -119,7 +123,7 @@ class TestFanOut:
dest_exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(
TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
)
tracer = get_tracer(provider, "litellm")
@ -141,18 +145,14 @@ class TestFanOut:
for child in ("auth /v1/chat/completions", "chat gpt-4"):
assert by_name[child].parent.span_id == root.context.span_id
def test_a_destination_for_another_backend_is_left_to_that_backends_provider(self):
"""Every v2 logger has its own provider and emits its own copy of a gen-AI
span, so a proxy running two of them would hand the tenant the same model call
twice if each provider forwarded every destination."""
def test_a_team_naming_two_backends_gets_the_trace_at_both(self):
"""The fan-out rides one provider, so it cannot skip a destination on the
grounds that some other backend owns it: nothing else would deliver it."""
langfuse, arize = InMemorySpanExporter(), InMemorySpanExporter()
by_endpoint = {"http://a.local": langfuse, "http://b.local": arize}
provider = TracerProvider()
provider.add_span_processor(
TenantFanOutSpanProcessor(
"langfuse_otel",
processor_factory=lambda d: SimpleSpanProcessor(by_endpoint[d.endpoint]),
)
TenantFanOutSpanProcessor(processor_factory=lambda d: SimpleSpanProcessor(by_endpoint[d.endpoint]))
)
def run():
@ -167,24 +167,30 @@ class TestFanOut:
in_fresh_context(run)
assert [s.name for s in langfuse.get_finished_spans()] == ["chat gpt-4"]
assert arize.get_finished_spans() == ()
assert [s.name for s in arize.get_finished_spans()] == ["chat gpt-4"]
def test_a_provider_that_speaks_for_no_backend_forwards_nothing(self):
"""The bare ``otel`` callback has no backend of its own; forwarding from it
would duplicate whatever the tenant's own backend provider already sent."""
def test_a_destination_carries_the_tenants_service_name(self):
"""An overridden backend skips per-request tracer routing, so the service name
that route used to apply has to travel on the destination instead."""
dest = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(
TenantFanOutSpanProcessor(None, processor_factory=lambda _d: SimpleSpanProcessor(dest))
)
provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest)))
def run():
set_request_destinations((LANGFUSE_DEST,))
set_request_destinations(
(
OtelDestination(
endpoint="http://a.local",
callback_name="langfuse_otel",
resource_attributes={"service.name": "team-checkout"},
),
)
)
emit(provider)
in_fresh_context(run)
assert dest.get_finished_spans() == ()
assert {s.resource.attributes["service.name"] for s in dest.get_finished_spans()} == {"team-checkout"}
def test_a_destination_that_cannot_build_a_processor_is_skipped_quietly(self):
"""An unbuildable destination must not cost the caller its request."""
@ -195,7 +201,7 @@ class TestFanOut:
attempts.append(destination.endpoint)
provider = TracerProvider()
provider.add_span_processor(TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory))
provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory))
def run():
set_request_destinations((LANGFUSE_DEST,))
@ -216,7 +222,7 @@ class TestFanOut:
return processor
provider = TracerProvider()
provider.add_span_processor(TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory))
provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory))
def run():
set_request_destinations((LANGFUSE_DEST,))
@ -238,10 +244,35 @@ class TestProviderWiring:
return [type(p).__name__ for p in provider._active_span_processor._span_processors]
assert "_OverriddenBackendFilter" in kinds(operator)
assert "TenantFanOutSpanProcessor" in kinds(operator)
assert "_OverriddenBackendFilter" not in kinds(tenant), "a per-tenant provider must not filter itself out"
assert "TenantFanOutSpanProcessor" not in kinds(operator), "delivery belongs to the published global alone"
assert "TenantFanOutSpanProcessor" not in kinds(tenant)
def test_only_the_published_global_provider_delivers_to_tenants(self):
"""A second v2 logger's provider never sees the server, auth or database spans,
so fanning out from it would hand the tenant a one-span trace. Publishing is
what picks the one provider the whole request tree passes through."""
config = OpenTelemetryV2Config(exporters=[ExporterSpec(kind="in_memory", owner=ExporterOwner.ARIZE_AX)])
published, other = OpenTelemetryV2(config=config, callback_name="arize"), OpenTelemetryV2(config=config)
publish_global_otel_v2_provider([other], lambda _p: None, registered=published)
def kinds(logger):
return [type(p).__name__ for p in logger._tracer_provider._active_span_processor._span_processors]
assert kinds(published).count("TenantFanOutSpanProcessor") == 1
assert "TenantFanOutSpanProcessor" not in kinds(other)
def test_publishing_twice_does_not_double_export(self):
config = OpenTelemetryV2Config(exporters=[ExporterSpec(kind="in_memory", owner=ExporterOwner.ARIZE_AX)])
logger = OpenTelemetryV2(config=config, callback_name="arize")
publish_global_otel_v2_provider([], lambda _p: None, registered=logger)
publish_global_otel_v2_provider([], lambda _p: None, registered=logger)
kinds = [type(p).__name__ for p in logger._tracer_provider._active_span_processor._span_processors]
assert kinds.count("TenantFanOutSpanProcessor") == 1
class TestRouting:
def test_an_overridden_backend_is_not_detached_onto_a_second_provider(self):
@ -263,6 +294,27 @@ class TestRouting:
assert route.tracer is default
assert route.provider is None
def test_an_overridden_backend_does_not_detach_on_a_service_name_either(self):
"""A key or team service name is its own reason to build a second provider, so
clearing only the credentials would still take the model call out of the tree."""
config = OpenTelemetryV2Config(
exporters=[ExporterSpec(kind="otlp_http", endpoint="http://op.local", owner=ExporterOwner.LANGFUSE_OTEL)]
)
cache = TenantTracerCache(config, "langfuse_otel", "litellm")
default = get_tracer(TracerProvider(), "litellm")
auth_metadata = {"otel_service_name": "team-checkout"}
assert cache.route_for(default, None, auth_metadata).detached is False
assert cache.route_for(default, None, auth_metadata).tracer is not default
def run():
set_request_destinations((LANGFUSE_DEST,))
return cache.route_for(default, None, auth_metadata)
route = in_fresh_context(run)
assert route.tracer is default, "the fan-out carries the service name on the destination instead"
assert route.provider is None
@pytest.mark.usefixtures("allow_test_hosts")
class TestDestinationResolution:
@ -290,6 +342,57 @@ class TestDestinationResolution:
assert [d.endpoint for d in destinations] == ["http://team.local/api/public/otel"]
assert destinations[0].callback_name == "langfuse_otel"
def test_a_keys_service_name_outranks_its_teams_on_the_destination(self, monkeypatch):
"""The key/team ``otel_service_name`` used to reach the backend through
per-request tracer routing, which an overridden backend skips."""
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
is_otel_v2_enabled.cache_clear()
auth = UserAPIKeyAuth(
metadata={"otel_service_name": "key-svc"},
team_metadata={
"otel_service_name": "team-svc",
"logging": [
{
"callback_name": "langfuse_otel",
"callback_type": "success",
"callback_vars": {
"langfuse_public_key": "pk-team",
"langfuse_secret_key": "sk-team",
"langfuse_host": "http://team.local",
},
}
],
},
)
destinations = resolve_tenant_otel_destinations(auth)
assert dict(destinations[0].resource_attributes) == {"service.name": "key-svc"}
def test_a_team_that_named_no_service_name_gets_no_resource_override(self, monkeypatch):
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
is_otel_v2_enabled.cache_clear()
auth = UserAPIKeyAuth(
team_metadata={
"otel_service_name": " ",
"logging": [
{
"callback_name": "langfuse_otel",
"callback_type": "success",
"callback_vars": {
"langfuse_public_key": "pk-team",
"langfuse_secret_key": "sk-team",
"langfuse_host": "http://team.local",
},
}
],
}
)
destinations = resolve_tenant_otel_destinations(auth)
assert dict(destinations[0].resource_attributes) == {}
def test_the_key_wins_over_the_team_for_the_same_backend(self, monkeypatch):
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
is_otel_v2_enabled.cache_clear()
@ -591,16 +694,21 @@ class TestEvictionSafety:
built.append(self.Recording())
return built[-1]
return TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory), built
return TenantFanOutSpanProcessor(processor_factory=factory), built
@staticmethod
def _dest(index):
return LANGFUSE_DEST.model_copy(update={"endpoint": f"http://d{index}/otel"})
@staticmethod
def _settle(fan_out):
for _ in range(50):
if not fan_out._retired:
def _settle(fan_out, processor=None):
"""Wait for retirement to clear and, when given, for the drain to run.
The drain pool is shared and bounded, so a shed processor is closed once a
worker picks it up rather than the moment it is handed over.
"""
for _ in range(500):
if not fan_out._retired and (processor is None or processor.shutdown_calls):
return
time.sleep(0.02)
@ -630,7 +738,7 @@ class TestEvictionSafety:
for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 1):
fan_out._acquire(self._dest(index))
fan_out._release(built[-1])
self._settle(fan_out)
self._settle(fan_out, built[0])
assert built[0].shutdown_calls == 1
assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS
@ -652,7 +760,7 @@ class TestEvictionSafety:
built.append(Slow())
return built[-1]
fan_out = TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory)
fan_out = TenantFanOutSpanProcessor(processor_factory=factory)
started = time.monotonic()
for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 1):
fan_out._acquire(self._dest(index))
@ -660,6 +768,38 @@ class TestEvictionSafety:
assert time.monotonic() - started < 2
def test_shedding_many_processors_does_not_spawn_a_thread_each(self):
"""A tenant that cycles its destination config sheds a processor per request,
so a thread per shed processor is a thread per request against a slow
collector."""
import threading
from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS
release = threading.Event()
class Blocking(self.Recording):
def shutdown(self):
release.wait(timeout=10)
super().shutdown()
built = []
def factory(_destination):
built.append(Blocking())
return built[-1]
fan_out = TenantFanOutSpanProcessor(processor_factory=factory)
try:
for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 30):
fan_out._acquire(self._dest(index))
fan_out._release(built[-1])
draining = [t for t in threading.enumerate() if t.name.startswith("litellm-otel-destination-drain")]
assert len(draining) <= 2, f"one drain thread per shed processor: {len(draining)}"
finally:
release.set()
self._settle(fan_out, built[0])
def test_a_retired_processor_is_still_closed_on_shutdown(self):
from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS