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.
This commit is contained in:
Yucheng He 2026-09-04 15:57:07 -07:00
parent e6c5594b8a
commit 3c36082c62
2 changed files with 56 additions and 10 deletions

View file

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

View file

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