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.
This commit is contained in:
Yucheng He 2026-09-03 22:38:47 -07:00
parent 5e98a5d5fc
commit c5003c88cd
2 changed files with 54 additions and 4 deletions

View file

@ -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())

View file

@ -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