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