From e6c5594b8a6444879e510dace74de623589db74d Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Fri, 4 Sep 2026 13:50:10 -0700 Subject: [PATCH] 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. --- litellm/__init__.py | 3 + litellm/integrations/otel/logger.py | 2 +- litellm/integrations/otel/plumbing/context.py | 39 ++- .../integrations/otel/plumbing/providers.py | 63 ++++- litellm/integrations/otel/plumbing/routing.py | 10 +- .../otel/test_otel_v2_destinations.py | 233 +++++++++++++++++- 6 files changed, 330 insertions(+), 20 deletions(-) diff --git a/litellm/__init__.py b/litellm/__init__.py index 42c0ea881fd..dfa72d2aa68 100644 --- a/litellm/__init__.py +++ b/litellm/__init__.py @@ -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 diff --git a/litellm/integrations/otel/logger.py b/litellm/integrations/otel/logger.py index 3014a334eeb..e893fec508c 100644 --- a/litellm/integrations/otel/logger.py +++ b/litellm/integrations/otel/logger.py @@ -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 diff --git a/litellm/integrations/otel/plumbing/context.py b/litellm/integrations/otel/plumbing/context.py index 3d8d8b53388..11b2b66d642 100644 --- a/litellm/integrations/otel/plumbing/context.py +++ b/litellm/integrations/otel/plumbing/context.py @@ -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() diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 40aa0d2404b..f113876b477 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -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, ...]": diff --git a/litellm/integrations/otel/plumbing/routing.py b/litellm/integrations/otel/plumbing/routing.py index a0831db2eea..8fd24d8706c 100644 --- a/litellm/integrations/otel/plumbing/routing.py +++ b/litellm/integrations/otel/plumbing/routing.py @@ -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) diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py index e47532fcde5..e7e0f716f41 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -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) == ()