fix(otel): serialize first fan-out attach

This commit is contained in:
Yucheng He 2026-09-05 03:18:52 -07:00
parent 836d31babc
commit 0cce313796
3 changed files with 41 additions and 7 deletions

View file

@ -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

View file

@ -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(

View file

@ -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