From c5003c88cd77d1263d7c167603d3f6e49241aa6d Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Thu, 3 Sep 2026 22:38:47 -0700 Subject: [PATCH] fix(otel v2): retire the drain workers with the fan-out that started them A proxy that rebuilds its telemetry builds another fan-out, so workers that outlive the one that started them are two more threads per reload. Shutdown now retires them once everything queued is closed, and a processor shed afterwards is closed inline rather than queued to nobody. --- .../integrations/otel/plumbing/providers.py | 31 ++++++++++++++++--- .../otel/test_otel_v2_destinations.py | 27 ++++++++++++++++ 2 files changed, 54 insertions(+), 4 deletions(-) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index b9c842656f4..65d9a01b1a5 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -234,16 +234,38 @@ class _DrainPool: """ def __init__(self, workers: int = _DRAIN_WORKERS) -> None: - self._pending: Final[queue.Queue[SpanProcessor]] = queue.Queue() + self._workers: Final = workers + self._closed: Final = threading.Event() + self._pending: Final[queue.Queue[SpanProcessor | None]] = queue.Queue() for _ in range(workers): - threading.Thread(target=self._drain_forever, daemon=True, name="litellm-otel-destination-drain").start() + threading.Thread( + target=self._drain_until_closed, daemon=True, name="litellm-otel-destination-drain" + ).start() def submit(self, processor: SpanProcessor) -> None: + if self._closed.is_set(): + _shutdown_quietly(processor) + return self._pending.put(processor) - def _drain_forever(self) -> None: + def close(self) -> None: + """Retire the workers once they have closed everything already queued. + + A proxy that rebuilds its telemetry builds another fan-out, so workers that + outlive the one that started them are two more threads per reload, forever. + """ + if self._closed.is_set(): + return + self._closed.set() + for _ in range(self._workers): + self._pending.put(None) + + def _drain_until_closed(self) -> None: while True: - _shutdown_quietly(self._pending.get()) + processor: SpanProcessor | None = self._pending.get() # rebind-ok: loop variable + if processor is None: + return + _shutdown_quietly(processor) class _ResourceWrappedReadableSpan(ReadableSpan): @@ -345,6 +367,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): self._processors.clear() self._retired.clear() self._exporting.clear() + self._drain.close() def force_flush(self, timeout_millis: int = 30000) -> bool: results: Final = tuple(self._flush_one(processor, timeout_millis) for processor in self._snapshot()) 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 25ed0b1af68..e950e5e9f6e 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -894,6 +894,33 @@ class TestEvictionSafety: assert not closed.is_alive(), "shutdown blocked on an export that never finished" + def test_shutdown_retires_the_drain_workers(self): + """A proxy that rebuilds its telemetry builds another fan-out, so workers that + outlive the one that started them are two more threads per reload.""" + from litellm.integrations.otel.plumbing.providers import _DRAIN_WORKERS + + before = self._drain_workers() + fan_out, _ = self._fan_out() + assert self._drain_workers() - before == _DRAIN_WORKERS + + fan_out.shutdown() + for _ in range(500): + if self._drain_workers() == before: + break + time.sleep(0.02) + + assert self._drain_workers() == before, "the drain workers outlived their fan-out" + + def test_a_processor_shed_after_shutdown_is_still_closed(self): + """``close`` retires the workers, so anything handed to the pool afterwards + would sit in a queue nobody reads.""" + fan_out, _ = self._fan_out() + stray = self.Recording() + fan_out.shutdown() + fan_out._drain.submit(stray) + + assert stray.shutdown_calls == 1 + def test_a_retired_processor_is_still_closed_on_shutdown(self): from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS