From 3c36082c620718a06d497add8824afce3e6c6269 Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Fri, 4 Sep 2026 15:57:07 -0700 Subject: [PATCH] fix(otel v2): let a straggling export close its own destination processor Shutdown waits out the exports in flight, but the wait has to be bounded or a tenant collector that stops answering holds the proxy open on the way down. Past the bound it closed everything anyway, which is the case it was written to avoid: a processor closed under the span it is carrying loses that span. Keep the bound and retire the stragglers instead. The thread still exporting one closes it through the drain as soon as its export returns, so teardown stays bounded and no span is dropped mid-forward. --- .../integrations/otel/plumbing/providers.py | 21 +++++---- .../otel/test_otel_v2_destinations.py | 45 ++++++++++++++++++- 2 files changed, 56 insertions(+), 10 deletions(-) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index f113876b477..858cf8a1121 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -383,23 +383,26 @@ class TenantFanOutSpanProcessor(SpanProcessor): while the SDK is tearing the provider down, so closing blind would drop a trace mid-forward and would hand the next caller a fresh exporter nothing will ever close. Refusing new work and then waiting out the in-flight ones - keeps both from happening. + keeps both from happening. The wait has to be bounded, or one destination + whose collector stopped answering would hold the proxy open on the way down, + so a straggler past the bound is retired instead of closed: the thread still + exporting it closes it through the drain as soon as its export returns, and + no span is ever dropped mid-forward. """ with self._lock: self._closed = True self._lock.wait_for(lambda: not self._exporting, timeout=self._drain_seconds) - # Snapshot: ``on_end`` mutates the cache on whichever thread ends a span, so - # iterating the live mapping risks a "mutated during iteration" the per-item - # except cannot catch. - for processor in self._snapshot(): + live: Final = tuple((id(p), p) for p in (*self._processors.values(), *self._retired.values())) + closing: Final = tuple(p for ident, p in live if ident not in self._exporting) + self._processors.clear() + self._retired = OrderedDict( # mutable-ok: the same bounded map, keeping only what is still exporting + (ident, p) for ident, p in live if ident in self._exporting + ) + for processor in closing: try: processor.shutdown() except Exception as exc: # noqa: BLE001 # one processor's shutdown must not abort the rest verbose_logger.debug("OTel V2 fan-out: processor shutdown failed: %s", exc) - with self._lock: - self._processors.clear() - self._retired.clear() - self._exporting.clear() self._drain.close() def force_flush(self, timeout_millis: int = 30000) -> bool: 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 e7e0f716f41..04051e70cb9 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -1150,6 +1150,42 @@ class TestEvictionSafety: assert stray.shutdown_calls == 1 + def test_shutdown_waits_out_an_export_that_lands_inside_the_bound(self): + """Without the wait the closing is left to a daemon thread, which the + interpreter can retire before it runs, so the last spans never reach the + tenant.""" + import threading + + fan_out, built = self._fan_out() + held = fan_out._acquire(self._dest(0)) + threading.Timer(0.2, lambda: fan_out._release(held)).start() + + fan_out.shutdown() + + assert held.shutdown_calls == 1, "shutdown returned before the export it should have waited out" + + def test_a_straggler_past_the_drain_bound_is_closed_by_its_own_thread(self): + """The wait is bounded so one dead collector cannot hold the proxy open, which + means a processor still exporting when it expires has to be left to the thread + holding it rather than closed under the span it is carrying.""" + built = [] + + def factory(_destination): + built.append(self.Recording()) + return built[-1] + + fan_out = TenantFanOutSpanProcessor(processor_factory=factory, shutdown_drain_seconds=0.05) + held = fan_out._acquire(self._dest(0)) + + fan_out.shutdown() + + assert held.shutdown_calls == 0 + + fan_out._release(held) + self._settle(fan_out, held) + + assert held.shutdown_calls == 1 + def test_a_processor_built_during_shutdown_is_not_left_in_a_cleared_cache(self): """The build runs outside the lock, so shutdown can finish inside it. Inserting afterwards leaves a live exporter, with its batch thread and its connection @@ -1210,7 +1246,9 @@ class TestEvictionSafety: assert submitted.shutdown_calls == 1, "a processor was queued behind the sentinels and never closed" - def test_a_retired_processor_is_still_closed_on_shutdown(self): + def test_a_retired_processor_is_still_closed_after_shutdown(self): + """Eviction and shutdown can both land while a span is being forwarded, and the + evicted processor still has to be closed once that export returns.""" from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS fan_out, built = self._fan_out() @@ -1220,6 +1258,11 @@ class TestEvictionSafety: fan_out._release(built[-1]) fan_out.shutdown() + assert held.shutdown_calls == 0 + + fan_out._release(held) + self._settle(fan_out, held) + assert held.shutdown_calls == 1