feat(otel v2): let a tenant destination export alongside the operator's own

Override stays the default: a key or team destination replaces the operator's
exporter for that backend. Operators running one org-wide backend across every
team set litellm_settings.otel_tenant_destination_mode to additive, and the
same trace lands in both places. A team that names the operator's own project
is still written once, since the fan-out skips a destination the operator's
exporter is already sending that span to.
This commit is contained in:
Yucheng He 2026-09-04 13:50:10 -07:00
parent 3ba9133fe8
commit e6c5594b8a
6 changed files with 330 additions and 20 deletions

View file

@ -325,6 +325,9 @@ ssl_certificate: Optional[str] = None
user_url_validation: bool = True
user_url_allowed_hosts: List[str] = []
provider_url_destination_allowed_hosts: List[str] = []
#: "override" (default) or "additive": whether a key or team destination replaces
#: the operator's exporter for that backend or exports alongside it.
otel_tenant_destination_mode: Optional[str] = None
ssl_ecdh_curve: Optional[str] = None # Set to 'X25519' to disable PQC and improve performance
disable_streaming_logging: bool = False
disable_token_counter: bool = False

View file

@ -872,7 +872,7 @@ def publish_global_otel_v2_provider(
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)
attach_tenant_fan_out(logger._tracer_provider, logger.config)
set_global_provider(logger._tracer_provider)
return logger

View file

@ -1,5 +1,6 @@
"""Trace-context + Baggage helpers."""
import os
from collections.abc import Mapping
from contextvars import ContextVar, Token
from typing import TYPE_CHECKING, Final
@ -329,12 +330,38 @@ def request_destinations() -> 'tuple["OtelDestination", ...]':
return _request_destinations.get()
def overridden_backends() -> frozenset[str]:
"""Backends whose global exporters this request must NOT reach.
#: ``litellm_settings: otel_tenant_destination_mode`` and its env equivalent.
ADDITIVE_DESTINATION_MODE: Final = "additive"
OTEL_TENANT_DESTINATION_MODE_ENV: Final = "LITELLM_OTEL_TENANT_DESTINATION_MODE"
A team destination is an override, not an addition: once the request resolved a
destination for a backend, that backend's operator-level exporters are suppressed
for every span of the request, so the tenant's traffic reaches the tenant's
account and nowhere else.
def tenant_destinations_are_additive() -> bool:
"""Whether a tenant destination exports alongside the operator's own exporter.
Override is the default: the tenant's traffic reaches the tenant's account and
nowhere else. Operators running one org-wide backend across every team set this
to ``additive`` so the same trace lands in both places.
"""
import litellm
configured: Final = litellm.otel_tenant_destination_mode or os.environ.get(OTEL_TENANT_DESTINATION_MODE_ENV)
return isinstance(configured, str) and configured.strip().lower() == ADDITIVE_DESTINATION_MODE
def destination_backends() -> frozenset[str]:
"""Backends this request resolved a tenant destination for.
The fan-out already carries the whole trace to those destinations, so the
per-request tracer route must never send a second copy, in either mode.
"""
return frozenset(d.callback_name for d in _request_destinations.get() if d.callback_name)
def suppressed_backends() -> frozenset[str]:
"""Backends whose operator-level exporters this request must NOT reach.
Empty under ``additive``, where the operator keeps its copy of every span.
"""
if tenant_destinations_are_additive():
return frozenset()
return destination_backends()

View file

@ -3,7 +3,7 @@
import queue
import threading
from collections import OrderedDict
from collections.abc import Callable, Iterable
from collections.abc import Callable, Iterable, Mapping
from typing import TYPE_CHECKING, Any, Final, Literal
from opentelemetry import _logs, baggage, metrics
@ -42,8 +42,8 @@ from litellm.integrations.otel.model.config import ExporterSpec, OpenTelemetryV2
from litellm.integrations.otel.model.semconv import LiteLLM
from litellm.integrations.otel.model.spans import LiteLLMSpanKind
from litellm.integrations.otel.plumbing.context import (
overridden_backends,
request_destinations,
suppressed_backends,
)
if TYPE_CHECKING:
@ -218,6 +218,9 @@ _DRAIN_WORKERS: Final = 2
#: proxy open.
_SHUTDOWN_DRAIN_SECONDS: Final = 5.0
#: An exporter's account: its normalized endpoint and its credentials.
_SinkKey = tuple[str, tuple[tuple[str, str], ...]]
class _DrainPool:
"""Closes shed destination processors off the span-export path.
@ -331,7 +334,9 @@ class TenantFanOutSpanProcessor(SpanProcessor):
self,
processor_factory: 'Callable[["OtelDestination"], SpanProcessor | None] | None' = None,
shutdown_drain_seconds: float = _SHUTDOWN_DRAIN_SECONDS,
operator_sinks: frozenset[_SinkKey] = frozenset(),
) -> None:
self._operator_sinks: Final = operator_sinks
self._drain_seconds: Final = shutdown_drain_seconds
self._lock: Final = threading.Condition()
self._closed = False # guarded by ``_lock``: an unlocked read races the teardown it gates
@ -345,7 +350,10 @@ class TenantFanOutSpanProcessor(SpanProcessor):
return None
def on_end(self, span: ReadableSpan) -> None:
suppressed: Final = suppressed_backends()
for destination in request_destinations():
if self._operator_already_writes(destination, suppressed):
continue
processor = self._acquire(destination) # rebind-ok: loop variable; pyright forbids Final in a loop
if processor is None:
continue
@ -356,6 +364,18 @@ class TenantFanOutSpanProcessor(SpanProcessor):
finally:
self._release(processor)
def _operator_already_writes(self, destination: "OtelDestination", suppressed: frozenset[str]) -> bool:
"""Whether the operator's own exporter is sending this span to the same account.
Only reachable under ``additive``, where nothing is suppressed: a team that
names the operator's own project would otherwise have every span written
there twice, once by the operator's exporter and once by the fan-out.
"""
return (
destination.callback_name not in suppressed
and _sink_key(destination.endpoint, destination.headers) in self._operator_sinks
)
def shutdown(self) -> None:
"""Close every destination processor, once the spans in flight have landed.
@ -489,6 +509,9 @@ class _OverriddenBackendFilter(SpanProcessor):
Wrapping is the only place this works: ``SynchronousMultiSpanProcessor.on_end``
ignores return values, so a sibling processor can never veto the export.
Under ``additive`` mode nothing is suppressed, so the wrapper passes every span
straight through and the operator keeps its copy.
"""
def __init__(self, inner: SpanProcessor, owner: str) -> None:
@ -499,7 +522,7 @@ class _OverriddenBackendFilter(SpanProcessor):
self._inner.on_start(span, parent_context)
def on_end(self, span: ReadableSpan) -> None:
if self._owner in overridden_backends():
if self._owner in suppressed_backends():
return
self._inner.on_end(span)
@ -796,15 +819,43 @@ def build_tracer_provider(
return provider
def attach_tenant_fan_out(provider: TracerProvider) -> None:
def attach_tenant_fan_out(provider: TracerProvider, config: OpenTelemetryV2Config | None = None) -> 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.
second publish (a test, a re-initialized proxy) cannot double-export. ``config``
names the operator's own exporters so an additive destination pointing at one of
them is delivered once rather than twice.
"""
if any(isinstance(processor, TenantFanOutSpanProcessor) for processor in _attached_processors(provider)):
return
provider.add_span_processor(TenantFanOutSpanProcessor())
provider.add_span_processor(TenantFanOutSpanProcessor(operator_sinks=operator_sink_keys(config)))
def operator_sink_keys(config: OpenTelemetryV2Config | None) -> frozenset[_SinkKey]:
"""The accounts the operator's own exporters write to, in destination terms.
An exporter with no endpoint of its own resolves one from the environment at
export time, so it has no comparable identity and is left out.
"""
if config is None:
return frozenset()
return frozenset(
key for spec in config.exporters if (key := _sink_key(spec.endpoint, parse_headers(spec.headers))) is not None
)
def _sink_key(endpoint: str | None, headers: Mapping[str, str]) -> "_SinkKey | None":
"""The account an exporter writes to, or ``None`` when it has no fixed one.
Normalized on both counts that make the same account look like two: the operator's
spec carries the signal path a tenant destination leaves for the exporter to
append, and header names survive one round trip lowercased and the other not.
"""
normalized: Final = _otlp_traces_endpoint(endpoint)
if normalized is None:
return None
return (normalized, tuple(sorted((name.lower(), value) for name, value in headers.items())))
def _attached_processors(provider: TracerProvider) -> "tuple[SpanProcessor, ...]":

View file

@ -25,7 +25,7 @@ from opentelemetry.trace import Tracer
from litellm._logging import verbose_logger
from litellm.constants import OTEL_SERVICE_NAME_METADATA_KEYS
from litellm.integrations.otel.model.config import ExporterSpec, OpenTelemetryV2Config
from litellm.integrations.otel.plumbing.context import overridden_backends
from litellm.integrations.otel.plumbing.context import destination_backends
from litellm.integrations.otel.plumbing.providers import (
build_tracer_provider,
exporter_transport,
@ -232,11 +232,11 @@ class TenantTracerCache:
concurrent overflow eviction can't shut it down between selection and
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 and already carries this tenant's credentials and service name.
# Routing here too would detach this span onto a second provider,
# A backend with a destination is delivered by the fan-out processor, which
# carries the 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.
if self._callback_name is not None and self._callback_name in overridden_backends():
if self._callback_name is not None and self._callback_name in destination_backends():
return TenantRoute(tracer=default, detached=False)
credential_headers: Final = self._credential_headers(dynamic_params)
project_headers: Final = self._project_headers(auth_metadata)

View file

@ -3,6 +3,7 @@
import contextvars
import time
from collections.abc import Mapping
from types import MappingProxyType
import pytest
from opentelemetry.sdk.trace import TracerProvider
@ -22,14 +23,16 @@ from litellm.integrations.otel.logger import (
publish_global_otel_v2_provider,
)
from litellm.integrations.otel.plumbing.context import (
overridden_backends,
destination_backends,
request_destinations,
set_request_destinations,
)
from litellm.integrations.otel.plumbing.providers import (
TenantFanOutSpanProcessor,
_OverriddenBackendFilter,
_sink_key,
build_tracer_provider,
operator_sink_keys,
)
from litellm.integrations.otel.plumbing.routing import TenantTracerCache, get_tracer
from litellm.integrations.otel.presets.destinations import (
@ -38,6 +41,7 @@ from litellm.integrations.otel.presets.destinations import (
)
from litellm.integrations.otel.presets.langfuse import langfuse_preset
from litellm.proxy._types import UserAPIKeyAuth
from litellm.types.utils import StandardCallbackDynamicParams
from litellm.proxy.litellm_pre_call_utils import resolve_tenant_otel_destinations
LANGFUSE_DEST = OtelDestination(
@ -117,6 +121,231 @@ class TestOverrideSuppression:
assert [s.name for s in arize_exporter.get_finished_spans()] == ["chat gpt-4"]
class TestRoutingMode:
"""The operator's choice between replacing its own exporter and exporting alongside it.
One org-wide backend across every team is a real deployment, and losing it the
moment a team configures its own is what ``additive`` exists to prevent.
"""
OPERATOR_SINK = ("https://cloud.langfuse.com/api/public/otel/v1/traces", (("authorization", "Basic op"),))
#: What a tenant destination for that same project looks like before normalizing:
#: no signal path yet, and the header name cased the way the backend writes it.
SAME_ACCOUNT_ENDPOINT = "https://cloud.langfuse.com/api/public/otel"
@staticmethod
def _additive(monkeypatch):
monkeypatch.setattr(litellm, "otel_tenant_destination_mode", "additive", raising=False)
@staticmethod
def _tree(provider):
tracer = get_tracer(provider, "litellm")
with tracer.start_as_current_span("POST /v1/chat/completions"):
with tracer.start_as_current_span("auth /v1/chat/completions"):
pass
with tracer.start_as_current_span("chat gpt-4"):
pass
def _run(self, provider, destinations=(LANGFUSE_DEST,)):
def run():
set_request_destinations(destinations)
self._tree(provider)
in_fresh_context(run)
def test_global_only_keeps_every_span_and_delivers_to_nobody(self):
"""No team destination resolved, so the operator's backbone is untouched."""
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
self._run(provider, destinations=())
assert len(global_exporter.get_finished_spans()) == 3
assert dest_exporter.get_finished_spans() == ()
def test_team_only_gets_the_whole_tree_with_no_operator_exporter(self):
"""A deployment with no operator credentials still gives the team its trace."""
dest_exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(
TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
)
self._run(provider)
assert {s.name for s in dest_exporter.get_finished_spans()} == {
"POST /v1/chat/completions",
"auth /v1/chat/completions",
"chat gpt-4",
}
def test_additive_gives_the_operator_and_the_team_the_same_tree(self, monkeypatch):
self._additive(monkeypatch)
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
self._run(provider)
names = {"POST /v1/chat/completions", "auth /v1/chat/completions", "chat gpt-4"}
assert {s.name for s in global_exporter.get_finished_spans()} == names
assert {s.name for s in dest_exporter.get_finished_spans()} == names
assert len(global_exporter.get_finished_spans()) == 3, "the operator must not get a span twice"
def test_override_moves_the_tree_off_the_operator(self):
"""The default, unchanged: the tenant's traffic reaches the tenant and nowhere else."""
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
self._run(provider)
assert global_exporter.get_finished_spans() == ()
assert len(dest_exporter.get_finished_spans()) == 3
def test_a_team_naming_the_operators_own_project_is_written_once(self, monkeypatch):
"""Fanning out to two accounts is the point. Writing the same account twice
is a duplicate the operator would see in their own project."""
self._additive(monkeypatch)
shared = InMemorySpanExporter()
same = OtelDestination(
endpoint=self.SAME_ACCOUNT_ENDPOINT,
headers=MappingProxyType({"Authorization": "Basic op"}),
callback_name="langfuse_otel",
)
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(shared), "langfuse_otel"))
provider.add_span_processor(
TenantFanOutSpanProcessor(
processor_factory=lambda _d: SimpleSpanProcessor(shared),
operator_sinks=frozenset({self.OPERATOR_SINK}),
)
)
self._run(provider, destinations=(same,))
assert len(shared.get_finished_spans()) == 3, "the same account received the trace twice"
def test_in_override_a_team_naming_the_operators_project_still_gets_the_trace(self):
"""Override suppresses the operator's own exporter, so the fan-out is the only
thing left delivering. Skipping it on a matching account leaves the team with
nothing at all."""
shared = InMemorySpanExporter()
same = OtelDestination(
endpoint=self.SAME_ACCOUNT_ENDPOINT,
headers=MappingProxyType({"Authorization": "Basic op"}),
callback_name="langfuse_otel",
)
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(shared), "langfuse_otel"))
provider.add_span_processor(
TenantFanOutSpanProcessor(
processor_factory=lambda _d: SimpleSpanProcessor(shared),
operator_sinks=frozenset({self.OPERATOR_SINK}),
)
)
self._run(provider, destinations=(same,))
assert len(shared.get_finished_spans()) == 3, "the team's own destination received nothing"
def test_a_team_naming_a_different_project_still_gets_its_copy(self, monkeypatch):
"""The dedup keys on the account, so a second project is still a second copy."""
self._additive(monkeypatch)
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(global_exporter), "langfuse_otel"))
provider.add_span_processor(
TenantFanOutSpanProcessor(
processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter),
operator_sinks=frozenset({self.OPERATOR_SINK}),
)
)
self._run(provider)
assert len(global_exporter.get_finished_spans()) == 3
assert len(dest_exporter.get_finished_spans()) == 3
@pytest.mark.parametrize("additive", [True, False])
def test_a_failing_team_destination_leaves_the_operator_alone(self, monkeypatch, additive):
"""A tenant collector that raises on every span must not cost the operator
its own telemetry, nor take the request down with it."""
if additive:
self._additive(monkeypatch)
global_exporter, arize_exporter = InMemorySpanExporter(), InMemorySpanExporter()
class Exploding(SimpleSpanProcessor):
def on_end(self, span):
raise RuntimeError("tenant collector is down")
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(global_exporter), "langfuse_otel"))
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(arize_exporter), "arize"))
provider.add_span_processor(
TenantFanOutSpanProcessor(processor_factory=lambda _d: Exploding(InMemorySpanExporter()))
)
self._run(provider)
assert len(arize_exporter.get_finished_spans()) == 3, "an unrelated backend lost spans"
assert len(global_exporter.get_finished_spans()) == (3 if additive else 0)
def test_the_env_var_turns_additive_on_without_a_config_file(self, monkeypatch):
monkeypatch.setattr(litellm, "otel_tenant_destination_mode", None, raising=False)
monkeypatch.setenv("LITELLM_OTEL_TENANT_DESTINATION_MODE", "Additive")
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
self._run(provider)
assert len(global_exporter.get_finished_spans()) == 3
assert len(dest_exporter.get_finished_spans()) == 3
def test_an_unrecognized_mode_stays_on_override(self, monkeypatch):
monkeypatch.setattr(litellm, "otel_tenant_destination_mode", "both", raising=False)
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
self._run(provider)
assert global_exporter.get_finished_spans() == ()
def test_operator_sink_keys_skips_an_exporter_with_no_endpoint_of_its_own(self):
"""Such an exporter resolves its endpoint from the environment at export
time, so it has no identity to compare a destination against."""
config = OpenTelemetryV2Config(
exporters=(
ExporterSpec(kind="otlp_http", endpoint=self.OPERATOR_SINK[0], headers="authorization=Basic op"),
ExporterSpec(kind="otlp_http", endpoint=None, headers="authorization=Basic other"),
)
)
assert operator_sink_keys(config) == frozenset({self.OPERATOR_SINK})
def test_the_operators_own_langfuse_and_a_team_naming_it_are_one_account(self, monkeypatch):
"""The two sides are built by different code that writes the endpoint and the
header names differently, so comparing them raw silently never matches."""
monkeypatch.setenv("LANGFUSE_HOST", "https://lf.internal")
monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-op")
monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-op")
monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", ["lf.internal"], raising=False)
operator = operator_sink_keys(langfuse_preset())
def sink(public_key, secret_key):
destination = destination_for(
"langfuse_otel",
StandardCallbackDynamicParams(
langfuse_public_key=public_key,
langfuse_secret_key=secret_key,
langfuse_host="https://lf.internal",
),
)
assert destination is not None
return _sink_key(destination.endpoint, destination.headers)
assert sink("pk-op", "sk-op") in operator, "a team naming the operator's own project"
assert sink("pk-team", "sk-team") not in operator, "a different project on the same server"
class TestFanOut:
def test_every_span_of_the_request_reaches_the_destination_in_one_trace(self):
"""The whole tree, gen-AI span included, parented as the operator would see it."""
@ -508,7 +737,7 @@ class TestContextIsolation:
def test_destinations_do_not_leak_between_requests(self):
def first():
set_request_destinations((LANGFUSE_DEST,))
return overridden_backends()
return destination_backends()
assert in_fresh_context(first) == frozenset({"langfuse_otel"})
assert in_fresh_context(request_destinations) == ()