From fb3b32d22dc5a4d374f6c3391e424434a65504fe Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Thu, 3 Sep 2026 17:38:08 -0700 Subject: [PATCH] fix(otel v2): scope the fan-out to its own backend and close shed processors off the export path Three problems in the fan-out, two of them in the eviction added last round: - Every v2 logger carries its own provider and emits its own copy of a gen-AI span, so a proxy running two of them handed the tenant the same model call twice. A provider now forwards only destinations for the backend it speaks for; the tenant's own backend always has a logger, since naming it in the key or team config is what builds one. Reproduced live against a self-hosted Langfuse on an arize-only proxy and on the bare `otel` callback. - Eviction could close a processor another thread was still exporting through, which drops that span. Exports are now counted, and a retired processor is closed only once its count reaches zero. - That close ran inside `on_end`, where `shutdown` flushes over the network, so one unreachable tenant collector stalled every other tenant's spans. It now runs on a short-lived thread, which also retires the retiree cap: a retiree drains as soon as its export finishes. --- litellm/integrations/otel/logger.py | 2 +- .../integrations/otel/plumbing/providers.py | 84 +++++++--- .../otel/test_otel_v2_destinations.py | 153 ++++++++++++------ 3 files changed, 170 insertions(+), 69 deletions(-) diff --git a/litellm/integrations/otel/logger.py b/litellm/integrations/otel/logger.py index ec3c3eaa004..7348d913b35 100644 --- a/litellm/integrations/otel/logger.py +++ b/litellm/integrations/otel/logger.py @@ -182,7 +182,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) + else build_tracer_provider(self.config, tenant_overrides=True, tenant_callback_name=callback_name) ) self.tracer: Tracer = get_tracer(self._tracer_provider, LITELLM_TRACER_NAME) self._metrics_recorder = self._init_metrics(meter_provider) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index a9dc0798c07..20649332f40 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -207,7 +207,6 @@ def _processor_for(exporter: SpanExporter, use_simple: bool | None) -> SpanProce #: Distinct tenant destinations whose exporters stay alive. Each holds a connection #: pool and a batch thread, so the cache is bounded and evicts least-recently-used. _MAX_CACHED_DESTINATION_PROCESSORS: Final = 32 -_MAX_RETIRED_DESTINATION_PROCESSORS: Final = 8 class _ResourceWrappedReadableSpan(ReadableSpan): @@ -246,29 +245,42 @@ class TenantFanOutSpanProcessor(SpanProcessor): Destinations ride a request-scoped ``ContextVar`` set during auth, so concurrent 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. """ 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[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded drain list + self._retired: OrderedDict[int, SpanProcessor] = OrderedDict() # mutable-ok: drains as exports finish + self._exporting: dict[int, int] = {} # mutable-ok: per-processor in-flight export count def on_start(self, span: SDKSpan, parent_context: Context | None = None) -> None: return None def on_end(self, span: ReadableSpan) -> None: for destination in request_destinations(): - processor = self._processor_for(destination) # rebind-ok: loop variable; pyright forbids Final in a loop + 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 try: processor.on_end(_with_destination_resource(span, destination)) except Exception as exc: # noqa: BLE001 # one destination's failure must not cost the others their span verbose_logger.debug("OTel V2 fan-out: forwarding to %s failed: %s", destination.endpoint, exc) + finally: + self._release(processor) def shutdown(self) -> None: # Snapshot first: ``on_end`` mutates the cache on whichever thread ends a span @@ -282,6 +294,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): with self._lock: self._processors.clear() self._retired.clear() + self._exporting.clear() def force_flush(self, timeout_millis: int = 30000) -> bool: results: Final = tuple(self._flush_one(processor, timeout_millis) for processor in self._snapshot()) @@ -298,12 +311,14 @@ class TenantFanOutSpanProcessor(SpanProcessor): except Exception: # noqa: BLE001 # one exporter's flush failure must not fail the whole flush return False - def _processor_for(self, destination: "OtelDestination") -> SpanProcessor | None: + def _acquire(self, destination: "OtelDestination") -> SpanProcessor | None: + """The processor for ``destination``, marked busy until ``_release``.""" key: Final = destination.cache_key() with self._lock: - cached: Final = self._processors.get(key) + cached = self._processors.get(key) # rebind-ok: reassigned after the build below if cached is not None: self._processors.move_to_end(key) + self._exporting[id(cached)] = self._exporting.get(id(cached), 0) + 1 return cached built: Final = self._build(destination) if built is None: @@ -312,28 +327,44 @@ class TenantFanOutSpanProcessor(SpanProcessor): existing: Final = self._processors.get(key) if existing is not None: # Another thread won the race; drop ours rather than leak its thread. - _shutdown_quietly(built) + _drain_in_background(built) + self._exporting[id(existing)] = self._exporting.get(id(existing), 0) + 1 return existing self._processors[key] = built - overflowed: Final = self._retired_on_overflow_locked() - if overflowed is not None: - _shutdown_quietly(overflowed) + self._exporting[id(built)] = 1 + self._retire_overflow_locked() + drained: Final = self._drainable_locked() + for processor in drained: + _drain_in_background(processor) return built - def _retired_on_overflow_locked(self) -> SpanProcessor | None: - """Drop the LRU processor past the cap; return one only once it is safe to close. + def _release(self, processor: SpanProcessor) -> None: + with self._lock: + remaining: Final = self._exporting.get(id(processor), 1) - 1 + if remaining > 0: + self._exporting[id(processor)] = remaining + else: + self._exporting.pop(id(processor), None) + drained: Final = self._drainable_locked() + for retired in drained: + _drain_in_background(retired) - ``on_end`` hands a processor back and then exports outside the lock, so shutting - an evicted one down there loses that span. Evictions retire to drain instead, and - the retirees are capped so they cannot accumulate a thread each. - """ + def _retire_overflow_locked(self) -> None: + """Move the LRU processor out of the cache once it is past the cap.""" if len(self._processors) <= _MAX_CACHED_DESTINATION_PROCESSORS: - return None + return _, evicted = self._processors.popitem(last=False) self._retired[id(evicted)] = evicted - if len(self._retired) <= _MAX_RETIRED_DESTINATION_PROCESSORS: - return None - return self._retired.popitem(last=False)[1] + + def _drainable_locked(self) -> tuple[SpanProcessor, ...]: + """Retired processors no thread is exporting through, removed from the list. + + ``on_end`` holds a processor across an export, so closing an evicted one there + drops the span it is holding. A retiree is out of the cache and can never be + handed out again, so once its export count reaches zero it stays there. + """ + idle: Final = tuple(key for key in self._retired if self._exporting.get(key, 0) == 0) + return tuple(self._retired.pop(key) for key in idle) def _destination_processor(destination: "OtelDestination") -> SpanProcessor | None: @@ -351,6 +382,18 @@ def _destination_processor(destination: "OtelDestination") -> SpanProcessor | No return None +def _drain_in_background(processor: SpanProcessor) -> None: + """Close a shed processor off the span-export path. + + ``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. + """ + threading.Thread( + target=_shutdown_quietly, args=(processor,), daemon=True, name="litellm-otel-destination-drain" + ).start() + + def _shutdown_quietly(processor: SpanProcessor) -> None: try: processor.shutdown() @@ -629,6 +672,7 @@ 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`. @@ -668,7 +712,7 @@ def build_tracer_provider( _OverriddenBackendFilter(processor, owner) if tenant_overrides and owner is not None else processor ) if tenant_overrides: - provider.add_span_processor(TenantFanOutSpanProcessor()) + provider.add_span_processor(TenantFanOutSpanProcessor(tenant_callback_name)) return provider 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 c9778bbb2ec..71b1164bfc0 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -1,14 +1,15 @@ """Key/team OTLP destinations override the operator's exporters for that backend.""" import contextvars +import time from collections.abc import Mapping -import litellm import pytest from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import SimpleSpanProcessor from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +import litellm from litellm.integrations.otel.model.config import ( ExporterOwner, ExporterSpec, @@ -67,7 +68,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(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) + TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) ) return provider @@ -100,7 +101,7 @@ class TestOverrideSuppression: provider = TracerProvider() provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(arize_exporter), "arize")) provider.add_span_processor( - TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) + TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) ) def run(): @@ -118,7 +119,7 @@ class TestFanOut: dest_exporter = InMemorySpanExporter() provider = TracerProvider() provider.add_span_processor( - TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) + TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) ) tracer = get_tracer(provider, "litellm") @@ -140,12 +141,16 @@ 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_two_destinations_each_receive_their_own_copy(self): - first, second = InMemorySpanExporter(), InMemorySpanExporter() - by_endpoint = {"http://a.local": first, "http://b.local": second} + 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.""" + 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]), ) ) @@ -161,8 +166,25 @@ class TestFanOut: in_fresh_context(run) - assert [s.name for s in first.get_finished_spans()] == ["chat gpt-4"] - assert [s.name for s in second.get_finished_spans()] == ["chat gpt-4"] + assert [s.name for s in langfuse.get_finished_spans()] == ["chat gpt-4"] + assert arize.get_finished_spans() == () + + 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.""" + dest = InMemorySpanExporter() + provider = TracerProvider() + provider.add_span_processor( + TenantFanOutSpanProcessor(None, processor_factory=lambda _d: SimpleSpanProcessor(dest)) + ) + + def run(): + set_request_destinations((LANGFUSE_DEST,)) + emit(provider) + + in_fresh_context(run) + + assert dest.get_finished_spans() == () def test_a_destination_that_cannot_build_a_processor_is_skipped_quietly(self): """An unbuildable destination must not cost the caller its request.""" @@ -173,7 +195,7 @@ class TestFanOut: attempts.append(destination.endpoint) provider = TracerProvider() - provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory)) + provider.add_span_processor(TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory)) def run(): set_request_destinations((LANGFUSE_DEST,)) @@ -194,7 +216,7 @@ class TestFanOut: return processor provider = TracerProvider() - provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory)) + provider.add_span_processor(TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory)) def run(): set_request_destinations((LANGFUSE_DEST,)) @@ -554,66 +576,101 @@ class TestTenantConfigAgreement: class TestEvictionSafety: - def test_an_evicted_processor_is_retired_rather_than_shut_down(self): - """``on_end`` hands a processor back and exports outside the lock, so shutting - an evicted one down there loses that span. Retirees are capped so they cannot - accumulate a thread each.""" - from litellm.integrations.otel.plumbing.providers import ( - _MAX_CACHED_DESTINATION_PROCESSORS, - _MAX_RETIRED_DESTINATION_PROCESSORS, - ) + class Recording(SimpleSpanProcessor): + def __init__(self): + super().__init__(InMemorySpanExporter()) + self.shutdown_calls = 0 - class Recording(SimpleSpanProcessor): - def __init__(self): - super().__init__(InMemorySpanExporter()) - self.shutdown_calls = 0 - - def shutdown(self): - self.shutdown_calls += 1 + def shutdown(self): + self.shutdown_calls += 1 + def _fan_out(self): built = [] def factory(_destination): - built.append(Recording()) + built.append(self.Recording()) return built[-1] - fan_out = TenantFanOutSpanProcessor(processor_factory=factory) - for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + _MAX_RETIRED_DESTINATION_PROCESSORS): - fan_out._processor_for(LANGFUSE_DEST.model_copy(update={"endpoint": f"http://d{index}/otel"})) + return TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory), built - assert [p.shutdown_calls for p in built] == [0] * len(built) - assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS + @staticmethod + def _dest(index): + return LANGFUSE_DEST.model_copy(update={"endpoint": f"http://d{index}/otel"}) - for index in range(2): - fan_out._processor_for(LANGFUSE_DEST.model_copy(update={"endpoint": f"http://late{index}/otel"})) + @staticmethod + def _settle(fan_out): + for _ in range(50): + if not fan_out._retired: + return + time.sleep(0.02) - assert [p.shutdown_calls for p in built[:2]] == [1, 1] - assert built[2].shutdown_calls == 0 - assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS - - def test_a_retired_processor_is_still_flushed_and_closed_on_shutdown(self): + def test_a_processor_still_exporting_a_span_is_not_closed_under_it(self): + """``on_end`` holds a processor across the export, so closing an evicted one + there drops the span it is holding.""" from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS - class Recording(SimpleSpanProcessor): - def __init__(self): - super().__init__(InMemorySpanExporter()) - self.shutdown_calls = 0 + fan_out, built = self._fan_out() + held = fan_out._acquire(self._dest(0)) + for index in range(1, _MAX_CACHED_DESTINATION_PROCESSORS + 1): + fan_out._acquire(self._dest(index)) + fan_out._release(built[-1]) + assert held.shutdown_calls == 0 + assert id(held) in fan_out._retired + + fan_out._release(held) + self._settle(fan_out) + + assert held.shutdown_calls == 1 + + def test_an_idle_evicted_processor_is_closed_off_the_export_path(self): + from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS + + fan_out, built = self._fan_out() + for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 1): + fan_out._acquire(self._dest(index)) + fan_out._release(built[-1]) + self._settle(fan_out) + + assert built[0].shutdown_calls == 1 + assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS + + def test_a_slow_collector_does_not_hold_up_the_export_path(self): + """``shutdown`` flushes over the network and is reached from ``on_end``, so + closing a shed processor inline lets one unreachable tenant collector stall + every other tenant's spans.""" + from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS + + class Slow(self.Recording): def shutdown(self): - self.shutdown_calls += 1 + time.sleep(3) + super().shutdown() built = [] def factory(_destination): - built.append(Recording()) + built.append(Slow()) return built[-1] - fan_out = TenantFanOutSpanProcessor(processor_factory=factory) + fan_out = TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory) + started = time.monotonic() for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 1): - fan_out._processor_for(LANGFUSE_DEST.model_copy(update={"endpoint": f"http://d{index}/otel"})) + fan_out._acquire(self._dest(index)) + fan_out._release(built[-1]) + + assert time.monotonic() - started < 2 + + def test_a_retired_processor_is_still_closed_on_shutdown(self): + from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS + + fan_out, built = self._fan_out() + held = fan_out._acquire(self._dest(0)) + for index in range(1, _MAX_CACHED_DESTINATION_PROCESSORS + 1): + fan_out._acquire(self._dest(index)) + fan_out._release(built[-1]) fan_out.shutdown() - assert built[0].shutdown_calls == 1 + assert held.shutdown_calls == 1 class TestCredentialGatedExporters: