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.
This commit is contained in:
Yucheng He 2026-09-03 17:38:08 -07:00
parent e9df4458c7
commit fb3b32d22d
3 changed files with 170 additions and 69 deletions

View file

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

View file

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

View file

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