diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 480abe8bf8c..04866c65fd0 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -4,7 +4,6 @@ import queue import threading from collections import OrderedDict from collections.abc import Callable, Iterable -from functools import lru_cache from typing import TYPE_CHECKING, Any, Final, Literal from opentelemetry import _logs, baggage, metrics @@ -215,6 +214,33 @@ _MAX_CACHED_DESTINATION_PROCESSORS: Final = 32 _DRAIN_WORKERS: Final = 2 +class _DrainPool: + """Closes shed destination processors off the span-export path. + + ``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. A fixed set of workers rather than a thread per + processor means a tenant cycling its destination config cannot spawn threads as + fast as it can send requests; slow shutdowns queue behind each other. + + The workers are daemons and belong to the fan-out that sheds the processors, so + neither an unreachable collector nor a lazily built process-wide singleton can + hold the proxy open on the way down. + """ + + def __init__(self, workers: int = _DRAIN_WORKERS) -> None: + self._pending: Final[queue.Queue[SpanProcessor]] = queue.Queue() + for _ in range(workers): + threading.Thread(target=self._drain_forever, daemon=True, name="litellm-otel-destination-drain").start() + + def submit(self, processor: SpanProcessor) -> None: + self._pending.put(processor) + + def _drain_forever(self) -> None: + while True: + _shutdown_quietly(self._pending.get()) + + class _ResourceWrappedReadableSpan(ReadableSpan): """A ``ReadableSpan`` view with an overridden Resource, leaving the original alone.""" @@ -270,6 +296,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): self._processors: OrderedDict[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded LRU self._retired: OrderedDict[int, SpanProcessor] = OrderedDict() # mutable-ok: drains as exports finish self._exporting: dict[int, int] = {} # mutable-ok: per-processor in-flight export count + self._drain: Final = _DrainPool() def on_start(self, span: SDKSpan, parent_context: Context | None = None) -> None: return None @@ -331,7 +358,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): existing: Final = self._processors.get(key) if existing is not None: # Another thread won the race; drop ours rather than leak its thread. - _drain_in_background(built) + self._drain.submit(built) self._exporting[id(existing)] = self._exporting.get(id(existing), 0) + 1 return existing self._processors[key] = built @@ -339,7 +366,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): self._retire_overflow_locked() drained: Final = self._drainable_locked() for processor in drained: - _drain_in_background(processor) + self._drain.submit(processor) return built def _release(self, processor: SpanProcessor) -> None: @@ -351,7 +378,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): self._exporting.pop(id(processor), None) drained: Final = self._drainable_locked() for retired in drained: - _drain_in_background(retired) + self._drain.submit(retired) def _retire_overflow_locked(self) -> None: """Move the LRU processor out of the cache once it is past the cap.""" @@ -386,39 +413,6 @@ def _destination_processor(destination: "OtelDestination") -> SpanProcessor | No return None -def _drain_in_background(processor: SpanProcessor) -> None: - """Close a shed processor off the span-export path. - - ``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 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_queue().put(processor) - - -@lru_cache(maxsize=1) -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: try: processor.shutdown() 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 75ee440e33c..9aee7edf293 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -790,12 +790,13 @@ class TestEvictionSafety: return built[-1] fan_out = TenantFanOutSpanProcessor(processor_factory=factory) + before = self._drain_workers() try: for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 30): fan_out._acquire(self._dest(index)) fan_out._release(built[-1]) - draining = [t for t in threading.enumerate() if t.name.startswith("litellm-otel-destination-drain")] - assert len(draining) <= 2, f"one drain thread per shed processor: {len(draining)}" + grew = self._drain_workers() - before + assert grew == 0, f"one drain thread per shed processor: {grew} new threads" finally: release.set() self._settle(fan_out, built[0]) @@ -806,14 +807,52 @@ class TestEvictionSafety: timeout on the way down.""" import threading - from litellm.integrations.otel.plumbing.providers import _drain_queue - - _drain_queue() + self._fan_out() 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_burst_of_first_evictions_starts_one_set_of_drain_workers(self): + """A drain pool built lazily on first use is not built once: several threads + can each finish the build, and every pool but the winner is left with its + workers blocked on a queue nothing will ever feed again.""" + import threading + + from litellm.integrations.otel.plumbing.providers import ( + _DRAIN_WORKERS, + _MAX_CACHED_DESTINATION_PROCESSORS, + ) + + for _ in range(3): + before = self._drain_workers() + fan_out, built = self._fan_out() + for index in range(_MAX_CACHED_DESTINATION_PROCESSORS): + fan_out._release(fan_out._acquire(self._dest(index))) + barrier = threading.Barrier(16) + + def shed(index, fan_out=fan_out, barrier=barrier): + barrier.wait(timeout=10) + fan_out._release(fan_out._acquire(self._dest(index))) + + threads = [ + threading.Thread(target=shed, args=(_MAX_CACHED_DESTINATION_PROCESSORS + index,)) + for index in range(16) + ] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=10) + self._settle(fan_out) + + assert self._drain_workers() - before == _DRAIN_WORKERS + + @staticmethod + def _drain_workers(): + import threading + + return len([t for t in threading.enumerate() if t.name.startswith("litellm-otel-destination-drain")]) + def test_a_retired_processor_is_still_closed_on_shutdown(self): from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS