From 5e98a5d5fcf61a22913159afb9e932ef282f38d7 Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Thu, 3 Sep 2026 22:20:34 -0700 Subject: [PATCH] fix(otel v2): close no destination processor under a span still in flight The fan-out now refuses new work once shutdown starts and waits out the spans already being forwarded, so teardown neither drops a trace mid-forward nor hands the next caller an exporter nothing will ever close. The wait is bounded so a dead collector cannot hold the proxy open. --- .../integrations/otel/plumbing/providers.py | 31 ++++++++++++-- .../otel/test_otel_v2_destinations.py | 41 +++++++++++++++++++ 2 files changed, 68 insertions(+), 4 deletions(-) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 04866c65fd0..b9c842656f4 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -213,6 +213,11 @@ _MAX_CACHED_DESTINATION_PROCESSORS: Final = 32 #: create by cycling its destination config. _DRAIN_WORKERS: Final = 2 +#: How long ``shutdown`` waits for spans already being forwarded, so teardown closes +#: no processor under one. Bounded: an exporter that never returns must not hold the +#: proxy open. +_SHUTDOWN_DRAIN_SECONDS: Final = 5.0 + class _DrainPool: """Closes shed destination processors off the span-export path. @@ -290,8 +295,11 @@ class TenantFanOutSpanProcessor(SpanProcessor): def __init__( self, processor_factory: 'Callable[["OtelDestination"], SpanProcessor | None] | None' = None, + shutdown_drain_seconds: float = _SHUTDOWN_DRAIN_SECONDS, ) -> None: - self._lock: Final = threading.Lock() + self._drain_seconds: Final = shutdown_drain_seconds + self._lock: Final = threading.Condition() + self._closed: Final = threading.Event() 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[int, SpanProcessor] = OrderedDict() # mutable-ok: drains as exports finish @@ -314,9 +322,20 @@ class TenantFanOutSpanProcessor(SpanProcessor): self._release(processor) def shutdown(self) -> None: - # Snapshot first: ``on_end`` mutates the cache on whichever thread ends a span - # and can run concurrently with this SDK-driven shutdown, so iterating the live - # mapping risks a "mutated during iteration" the per-item except cannot catch. + """Close every destination processor, once the spans in flight have landed. + + ``on_end`` runs on whichever thread ends a span and can reach this fan-out + 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. + """ + self._closed.set() + with self._lock: + 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(): try: processor.shutdown() @@ -344,6 +363,8 @@ class TenantFanOutSpanProcessor(SpanProcessor): def _acquire(self, destination: "OtelDestination") -> SpanProcessor | None: """The processor for ``destination``, marked busy until ``_release``.""" + if self._closed.is_set(): + return None key: Final = destination.cache_key() with self._lock: cached = self._processors.get(key) # rebind-ok: reassigned after the build below @@ -376,6 +397,8 @@ class TenantFanOutSpanProcessor(SpanProcessor): self._exporting[id(processor)] = remaining else: self._exporting.pop(id(processor), None) + if not self._exporting: + self._lock.notify_all() drained: Final = self._drainable_locked() for retired in drained: self._drain.submit(retired) 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 9aee7edf293..25ed0b1af68 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -853,6 +853,47 @@ class TestEvictionSafety: return len([t for t in threading.enumerate() if t.name.startswith("litellm-otel-destination-drain")]) + def test_shutdown_does_not_close_a_processor_under_an_in_flight_export(self): + """``on_end`` runs on whichever thread ends a span, so it reaches the fan-out + while the SDK tears the provider down.""" + import threading + + fan_out, _ = self._fan_out() + held = fan_out._acquire(self._dest(0)) + closed = threading.Thread(target=fan_out.shutdown) + closed.start() + try: + time.sleep(0.3) + + assert held.shutdown_calls == 0, "closed a processor with a span still being forwarded" + finally: + fan_out._release(held) + closed.join(timeout=10) + + assert held.shutdown_calls == 1 + + def test_a_closed_fan_out_builds_no_new_processor(self): + """A processor built after shutdown is one nothing will ever close, and it + exports to a tenant on a provider the SDK has already torn down.""" + fan_out, built = self._fan_out() + fan_out.shutdown() + + assert fan_out._acquire(self._dest(0)) is None + assert built == [] + + def test_shutdown_gives_up_on_an_export_that_never_finishes(self): + """The wait is bounded: an exporter stuck on a dead collector must not hold + the proxy open on the way down.""" + import threading + + fan_out = TenantFanOutSpanProcessor(processor_factory=lambda _d: self.Recording(), shutdown_drain_seconds=0.2) + fan_out._acquire(self._dest(0)) + closed = threading.Thread(target=fan_out.shutdown) + closed.start() + closed.join(timeout=5) + + assert not closed.is_alive(), "shutdown blocked on an export that never finished" + def test_a_retired_processor_is_still_closed_on_shutdown(self): from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS