diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 64a45151ab5..b60bcde2b5a 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -216,6 +216,12 @@ _MAX_CACHED_DESTINATION_PROCESSORS: Final = 32 #: create by cycling its destination config. _DRAIN_WORKERS: Final = 2 +#: Shed processors waiting to be closed before the fan-out stops building new ones. +#: Each still owns a batch thread until its close returns, and a collector that never +#: answers makes every close take the exporter's full timeout, so past this many the +#: operator's exporter keeps the span instead (see ``deliverable``). +_MAX_PENDING_DRAINS: Final = 64 + #: 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. @@ -247,10 +253,13 @@ class _DrainPool: self, workers: int = _DRAIN_WORKERS, pending: "queue.Queue[SpanProcessor | None] | None" = None, + capacity: int = _MAX_PENDING_DRAINS, ) -> None: self._workers: Final = workers + self._capacity: Final = capacity self._lock: Final = threading.Lock() self._closed = False + self._backlog = 0 # guarded by ``_lock``: submitted processors whose close has not returned self._pending: Final[queue.Queue[SpanProcessor | None]] = pending if pending is not None else queue.Queue() self._threads: Final = tuple( threading.Thread(target=self._drain_until_closed, daemon=True, name="litellm-otel-destination-drain") @@ -274,6 +283,7 @@ class _DrainPool: """ with self._lock: if not self._closed: + self._backlog += 1 self._pending.put(processor) return threading.Thread( @@ -283,6 +293,18 @@ class _DrainPool: name="litellm-otel-destination-drain-straggler", ).start() + def saturated(self) -> bool: + """Whether enough closes are outstanding that building another processor must wait. + + The workers close in order and each close blocks for as long as its exporter + does, so a collector that stopped answering would otherwise turn every new + destination into one more batch thread parked behind them, for as long as the + tenants keep rotating. Holding the count here rather than reading the queue + keeps the two processors a worker is mid-close on in the total. + """ + with self._lock: + return self._backlog >= self._capacity + def close(self, timeout: float | None = None) -> None: """Retire the workers once they have closed everything already queued. @@ -311,6 +333,8 @@ class _DrainPool: if processor is None: return _shutdown_quietly(processor) + with self._lock: + self._backlog -= 1 _NO_ATTRIBUTES: Final[Mapping[str, AttributeValue]] = MappingProxyType({}) @@ -405,6 +429,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): processor_factory: 'Callable[["OtelDestination"], SpanProcessor | None] | None' = None, shutdown_drain_seconds: float = _SHUTDOWN_DRAIN_SECONDS, operator_sinks: frozenset[_SinkKey] = frozenset(), + pending_drains: int = _MAX_PENDING_DRAINS, ) -> None: self._operator_sinks: Final = operator_sinks self._drain_seconds: Final = shutdown_drain_seconds @@ -414,7 +439,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() + self._drain: Final = _DrainPool(capacity=pending_drains) def on_start(self, span: SDKSpan, parent_context: Context | None = None) -> None: return None @@ -544,6 +569,16 @@ class TenantFanOutSpanProcessor(SpanProcessor): return self._build_locked(destination, key) def _build_locked(self, destination: "OtelDestination", key: object) -> SpanProcessor | None: + """Build and cache a processor for ``destination``, unless the drain is saturated. + + Every build past the cache cap sheds one processor into the drain, so while the + shed ones are stuck closing against a collector that stopped answering, a new + destination is refused rather than parked behind them: ``deliverable`` then + leaves its spans with the operator's exporter until the drain catches up. + """ + if self._drain.saturated(): + verbose_logger.debug("OTel V2 fan-out: drain saturated, not building for %s", destination.endpoint) + return None built: Final = self._build(destination) if built is None: return 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 f051c6dfcb0..1787dd68ad8 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -1676,7 +1676,45 @@ class TestEvictionSafety: release.set() self._settle(fan_out, built[0]) - def test_the_drain_workers_do_not_hold_the_process_open(self): + def test_a_saturated_drain_leaves_new_destinations_with_the_operator(self): + """A shed processor keeps its batch thread until its close returns, and against + a collector that never answers every close waits out the exporter's timeout. + Tenants rotating past the cache cap would otherwise queue one more processor, + and one more thread, per request for as long as the outage lasts.""" + import threading + + from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS + + release = threading.Event() + + class Blocking(self.Recording): + def shutdown(self): + release.wait(timeout=10) + super().shutdown() + + built = [] + + def factory(_destination): + built.append(Blocking()) + return built[-1] + + fan_out = TenantFanOutSpanProcessor(processor_factory=factory, pending_drains=3) + try: + for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 40): + processor = fan_out._acquire(self._dest(index)) + if processor is not None: + fan_out._release(processor) + + assert len(built) == _MAX_CACHED_DESTINATION_PROCESSORS + 3, "a processor per request during the outage" + assert fan_out.deliverable((self._dest(999),)) == (), "the span would vanish instead of staying with the operator" + finally: + release.set() + for _ in range(500): + if not fan_out._drain.saturated(): + break + time.sleep(0.02) + + assert fan_out.deliverable((self._dest(999),)) == (self._dest(999),), "the fan-out never recovered" """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."""