From 0cce313796d6f876e46928791e91f87549a2d42f Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Sat, 5 Sep 2026 03:18:52 -0700 Subject: [PATCH] fix(otel): serialize first fan-out attach --- litellm/__init__.py | 2 +- .../integrations/otel/plumbing/providers.py | 17 +++++++---- .../otel/test_otel_v2_destinations.py | 29 +++++++++++++++++++ 3 files changed, 41 insertions(+), 7 deletions(-) diff --git a/litellm/__init__.py b/litellm/__init__.py index dfa72d2aa68..14327d54897 100644 --- a/litellm/__init__.py +++ b/litellm/__init__.py @@ -327,7 +327,7 @@ 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 +otel_tenant_destination_mode: str | None = 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/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 350566d1985..ff1da6c5f4b 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -887,17 +887,22 @@ def build_tracer_provider( return provider +_FAN_OUT_ATTACH_LOCK: Final = threading.Lock() + + 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. ``config`` - names the operator's own exporters so an additive destination pointing at one of - them is delivered once rather than twice. + second publish (a test, a re-initialized proxy) cannot double-export. Concurrent + first calls (requests racing to anchor before any publish) serialize on one lock + so exactly one fan-out lands. ``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(operator_sinks=operator_sink_keys(config))) + with _FAN_OUT_ATTACH_LOCK: + if any(isinstance(processor, TenantFanOutSpanProcessor) for processor in _attached_processors(provider)): + return + provider.add_span_processor(TenantFanOutSpanProcessor(operator_sinks=operator_sink_keys(config))) def deliverable_destinations( 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 fbda890da88..a603990f5e7 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -716,6 +716,35 @@ class TestProviderWiring: assert fan_out_provider() is logger.tracer_provider assert deliverable_destinations((LANGFUSE_DEST,), fan_out_provider()) == (LANGFUSE_DEST,) + def test_concurrent_anchoring_attaches_exactly_one_fan_out(self): + """Requests race to anchor when the startup publish never ran, and a fan-out + attached twice delivers every tenant span twice.""" + import threading + + from litellm.integrations.otel.plumbing.providers import attach_tenant_fan_out + + class SlowAttachProvider(TracerProvider): + def add_span_processor(self, span_processor): + time.sleep(0.05) + super().add_span_processor(span_processor) + + provider = SlowAttachProvider() + config = OpenTelemetryV2Config(exporters=[ExporterSpec(kind="in_memory", owner=ExporterOwner.LANGFUSE_OTEL)]) + barrier = threading.Barrier(8) + + def anchor(): + barrier.wait(timeout=10) + attach_tenant_fan_out(provider, config) + + threads = [threading.Thread(target=anchor) for _ in range(8)] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=10) + + kinds = [type(p).__name__ for p in provider._active_span_processor._span_processors] + assert kinds.count("TenantFanOutSpanProcessor") == 1, f"one fan-out per provider, got {kinds}" + def test_without_a_publish_anchoring_falls_back_to_the_otel_global(self, monkeypatch): from opentelemetry import trace