From 1338433b223602ca013c92d48ef699fcae31f95b Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Thu, 3 Sep 2026 21:31:26 -0700 Subject: [PATCH] fix(otel v2): drain shed destination processors on daemon workers A ThreadPoolExecutor joins its workers at interpreter exit, so one unreachable tenant collector would hold the whole proxy open for its export timeout on the way down. Two long-lived daemon workers off a queue keep the thread count bounded without blocking shutdown. --- .../integrations/otel/plumbing/providers.py | 34 +++++++++++++++---- .../otel/test_otel_v2_destinations.py | 16 ++++++++- 2 files changed, 42 insertions(+), 8 deletions(-) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 0ec570eb9e3..480abe8bf8c 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -1,9 +1,9 @@ """Provider / exporter factory + the Baggage span processor.""" +import queue import threading from collections import OrderedDict from collections.abc import Callable, Iterable -from concurrent.futures import ThreadPoolExecutor from functools import lru_cache from typing import TYPE_CHECKING, Any, Final, Literal @@ -210,6 +210,10 @@ def _processor_for(exporter: SpanExporter, use_simple: bool | None) -> SpanProce #: pool and a batch thread, so the cache is bounded and evicts least-recently-used. _MAX_CACHED_DESTINATION_PROCESSORS: Final = 32 +#: Workers closing shed destination processors, bounding the threads a tenant can +#: create by cycling its destination config. +_DRAIN_WORKERS: Final = 2 + class _ResourceWrappedReadableSpan(ReadableSpan): """A ``ReadableSpan`` view with an overridden Resource, leaving the original alone.""" @@ -387,16 +391,32 @@ def _drain_in_background(processor: SpanProcessor) -> None: ``shutdown`` flushes over the network and is reached from ``on_end``, so closing one inline would let a single unreachable tenant collector stall every other - tenant's spans behind it. The work goes to a two-thread pool rather than a thread - per processor, so a tenant cycling its destination config cannot spawn threads as - fast as it can send requests; slow shutdowns queue behind each other instead. + tenant's spans behind it. The work goes to two long-lived workers rather than a + thread per processor, so a tenant cycling its destination config cannot spawn + threads as fast as it can send requests; slow shutdowns queue behind each other. """ - _drain_pool().submit(_shutdown_quietly, processor) + _drain_queue().put(processor) @lru_cache(maxsize=1) -def _drain_pool() -> ThreadPoolExecutor: - return ThreadPoolExecutor(max_workers=2, thread_name_prefix="litellm-otel-destination-drain") +def _drain_queue() -> "queue.Queue[SpanProcessor]": + """The shed-processor queue, with its daemon workers started on first use. + + Daemon on purpose. ``ThreadPoolExecutor`` joins its workers at interpreter exit, + so a single unreachable tenant collector would hold the whole proxy open for its + export timeout on the way down. + """ + pending: queue.Queue[SpanProcessor] = queue.Queue() + for _ in range(_DRAIN_WORKERS): + threading.Thread( + target=_drain_forever, args=(pending,), daemon=True, name="litellm-otel-destination-drain" + ).start() + return pending + + +def _drain_forever(pending: "queue.Queue[SpanProcessor]") -> None: + while True: + _shutdown_quietly(pending.get()) def _shutdown_quietly(processor: SpanProcessor) -> None: 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 9da2e25312b..75ee440e33c 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -727,7 +727,7 @@ class TestEvictionSafety: assert id(held) in fan_out._retired fan_out._release(held) - self._settle(fan_out) + self._settle(fan_out, held) assert held.shutdown_calls == 1 @@ -800,6 +800,20 @@ class TestEvictionSafety: release.set() self._settle(fan_out, built[0]) + def test_the_drain_workers_do_not_hold_the_process_open(self): + """Python joins a ThreadPoolExecutor's workers at interpreter exit, so one + unreachable tenant collector would hold the proxy open for its export + timeout on the way down.""" + import threading + + from litellm.integrations.otel.plumbing.providers import _drain_queue + + _drain_queue() + workers = [t for t in threading.enumerate() if t.name.startswith("litellm-otel-destination-drain")] + + assert workers, "no drain worker was started" + assert all(t.daemon for t in workers), "a non-daemon drain worker blocks interpreter exit" + def test_a_retired_processor_is_still_closed_on_shutdown(self): from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS