fix(otel): bound and shut down credential-scoped tracer providers (#36591)

* fix(otel): bound and shut down credential-scoped tracer providers

Each credential-scoped TracerProvider owns a BatchSpanProcessor worker thread that
only stops on shutdown, and the v1 cache holding them was an unbounded, unsynchronized
dict that never shut anything down. Every distinct team/key credential set therefore
added a thread for the life of the process, and concurrent first-requests for the same
credential set orphaned duplicate providers outright.

Make the cache a lock-guarded bounded LRU that shuts down whatever it drops, matching
the v2 TenantTracerCache. Providers wrapping a caller-supplied SpanExporter instance
share that exporter with the logger's own provider, so they are dropped without
shutdown; those use SimpleSpanProcessor and own no thread.

* fix(otel): reclaim dropped providers on a dedicated executor

Sustained credential churn queues one blocking shutdown per eviction, so using the
shared logging executor let an unreachable tenant endpoint stall unrelated logging
work behind the OTLP retry budget. Give provider shutdown its own bounded pool; its
threads spawn lazily, so a proxy that never evicts still pays nothing.

* fix(otel): decide provider shutdown from the victim, not the evicting request

Both dynamic entry points share one provider cache, so it can hold providers of
mixed exporter ownership. Reading the ownership flag from the evicting request
therefore stopped a shared caller-supplied exporter in one direction, silencing
telemetry process-wide, and leaked a BatchSpanProcessor thread in the other.

Cache ownership alongside the provider so the drop decision reads the victim's
own flag.

* fix(otel): honor the widened header mapping type instead of dict only

Widening the header parameter to Mapping left the isinstance check on dict, so a
non-dict Mapping silently returned no headers at all, which for the OTLP path means
an unauthenticated exporter and no traces with nothing raised. The dict branch also
returned the caller's own object, and dropping the defensive copy at the call site
let that alias reach a long-lived exporter. Match on Mapping and copy.

* fix(otel): do not give a provider we may never stop an interpreter-exit hook

Every TracerProvider registers an atexit hook by default, and that hook holds a strong
reference. Providers wrapping a caller-supplied exporter are dropped without shutdown,
so they stayed pinned for the life of the process and then stopped the shared exporter
at exit. Tie shutdown_on_exit to ownership: those providers use SimpleSpanProcessor and
buffer nothing, so they lose no flush, while providers that own their exporter keep the
hook and their exit flush.

Also stop the victim the eviction test leaves behind, and trim the added comments.
This commit is contained in:
yucheng-berri 2026-08-18 16:21:58 -07:00 • committed by GitHub
parent 087d82ffca
commit 55ec491d03
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 336 additions and 31 deletions

View file

@ -1,5 +1,8 @@
import os
from collections.abc import Mapping
import threading
from collections import OrderedDict
from collections.abc import Callable, Mapping
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass, field
from datetime import datetime
from typing import TYPE_CHECKING, Any, Final, TypedDict, cast
@ -83,6 +86,12 @@ class _ResponseWithUsageView(TypedDict, total=False):
usage: "_UsageCompletionTokensView | None"
# Cap on credential-scoped providers held at once; each one owns an exporter thread.
_MAX_DYNAMIC_TRACER_PROVIDERS: Final = 256
# Dedicated so a slow exporter shutdown cannot starve the shared logging executor.
_PROVIDER_SHUTDOWN_EXECUTOR: Final = ThreadPoolExecutor(max_workers=4, thread_name_prefix="OtelProviderShutdown")
LITELLM_TRACER_NAME: Final = os.getenv("OTEL_TRACER_NAME", "litellm")
LITELLM_METER_NAME: Final = os.getenv("LITELLM_METER_NAME", "litellm")
LITELLM_LOGGER_NAME: Final = os.getenv("LITELLM_LOGGER_NAME", "litellm")
@ -227,6 +236,34 @@ def _freeze_for_dedupe(value: object, _depth: int = 0) -> HashableScope:
return repr(value)
def _shutdown_tracer_provider(provider: "_SDKTracerProvider") -> None:
"""Flush and stop a dropped provider so its exporter thread is reclaimed."""
try:
provider.shutdown()
except Exception as e: # noqa: BLE001 # exporter shutdown must not fail the request that dropped it
verbose_logger.debug("OpenTelemetry: error shutting down dropped tracer provider: %s", e)
@dataclass(frozen=True, slots=True)
class _CachedTracerProvider:
"""A cached credential-scoped provider plus whether it may be shut down when dropped."""
provider: "_SDKTracerProvider"
owns_exporter: bool
def _provider_owns_exporter(exporter: "str | _SpanExporter") -> bool:
"""Whether a provider built for ``exporter`` may be shut down when it is dropped.
``_get_span_processor`` builds a fresh exporter for a named kind, but wraps a
caller-supplied ``SpanExporter`` instance as-is, and that instance is shared with the
logger's own provider. Shutting a dropped provider down would then stop exporting for
the whole process. The shared case also uses ``SimpleSpanProcessor``, so it owns no
thread and there is nothing to reclaim.
"""
return not hasattr(exporter, "export")
@dataclass
class OpenTelemetryConfig:
exporter: str | SpanExporter = "console"
@ -322,6 +359,7 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
tracer_provider: object | None = None,
logger_provider: object | None = None,
meter_provider: object | None = None,
max_dynamic_tracer_providers: int = _MAX_DYNAMIC_TRACER_PROVIDERS,
**kwargs,
):
team_metadata_keys_override: Final = kwargs.pop("baggage_team_metadata_keys", None)
@ -347,7 +385,9 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
self.OTEL_EXPORTER = self.config.exporter
self.OTEL_ENDPOINT = self.config.endpoint
self.OTEL_HEADERS = self.config.headers
self._tracer_provider_cache: dict[str, _SDKTracerProvider] = {}
self._tracer_provider_cache: OrderedDict[str, _CachedTracerProvider] = OrderedDict()
self._tracer_provider_cache_lock: Final = threading.Lock()
self._max_dynamic_tracer_providers: Final = max(1, max_dynamic_tracer_providers)
self._init_tracing(tracer_provider)
_debug_otel: Final = str(os.getenv("DEBUG_OTEL", "False")).lower()
@ -1027,38 +1067,98 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
return self.construct_dynamic_otel_config(standard_callback_dynamic_params=standard_callback_dynamic_params)
def _get_tracer_with_dynamic_config(self, dynamic_config: OpenTelemetryConfig):
def _insert_or_drop(
self, cache_key: str, built: "_CachedTracerProvider"
) -> "tuple[_CachedTracerProvider, _CachedTracerProvider | None]":
"""Cache ``built`` under ``cache_key``, returning the entry to use and what to drop.
Caller holds ``_tracer_provider_cache_lock``. The drop is either the loser of a
concurrent build for this key or the LRU victim its insertion pushed out.
"""
raced: Final = self._tracer_provider_cache.get(cache_key)
if raced is not None:
self._tracer_provider_cache.move_to_end(cache_key)
return raced, built
self._tracer_provider_cache[cache_key] = built
if len(self._tracer_provider_cache) > self._max_dynamic_tracer_providers:
return built, self._tracer_provider_cache.popitem(last=False)[1]
return built, None
def _cached_dynamic_tracer(
self,
cache_key: str,
build: Callable[[], "_SDKTracerProvider"],
owns_exporter: bool,
) -> "_Tracer":
"""Return the tracer for ``cache_key``, building and caching a provider on miss.
A provider that owns its exporter also owns a ``BatchSpanProcessor`` worker thread
that only stops on ``shutdown()``, so the cache is a bounded LRU and whatever it
drops is shut down. Without both, a proxy serving key-scoped credentials accumulates
one live thread per credential set for the life of the process.
``owns_exporter`` also decides ``shutdown_on_exit`` at build time: a provider we may
never shut down must not hold an interpreter-exit hook, which would both pin it in
memory for the life of the process and stop the shared exporter at exit. Those
providers use ``SimpleSpanProcessor``, which buffers nothing, so the hook costs them
no flush.
``owns_exporter`` describes the provider being built, and is cached with it, because
the two dynamic entry points share this cache and can disagree: whether the LRU
victim may be shut down is a property of the victim, never of the request that
happened to evict it.
"""
with self._tracer_provider_cache_lock:
cached: Final = self._tracer_provider_cache.get(cache_key)
if cached is not None:
self._tracer_provider_cache.move_to_end(cache_key)
return cached.provider.get_tracer(LITELLM_TRACER_NAME)
# Built outside the lock: exporter construction can block on DNS/TLS.
built: Final = _CachedTracerProvider(provider=build(), owns_exporter=owns_exporter)
with self._tracer_provider_cache_lock:
winner, dropped = self._insert_or_drop(cache_key, built)
if dropped is not None and dropped.owns_exporter:
# Off the caller's thread: shutdown joins the exporter worker.
_PROVIDER_SHUTDOWN_EXECUTOR.submit(_shutdown_tracer_provider, dropped.provider)
return winner.provider.get_tracer(LITELLM_TRACER_NAME)
def _get_tracer_with_dynamic_config(self, dynamic_config: OpenTelemetryConfig) -> "_Tracer":
"""Create (or reuse) a tracer whose exporter target comes from a per-request config."""
from opentelemetry.sdk.trace import TracerProvider
cache_key = f"dynamic_config:{dynamic_config.exporter}:{dynamic_config.endpoint}:{dynamic_config.headers}"
if cache_key in self._tracer_provider_cache:
return self._tracer_provider_cache[cache_key].get_tracer(LITELLM_TRACER_NAME)
owns_exporter: Final = _provider_owns_exporter(dynamic_config.exporter)
temp_provider: Final = TracerProvider(resource=self._get_litellm_resource(self.config))
temp_provider.add_span_processor(self._get_span_processor(config_override=dynamic_config))
def _build() -> "_SDKTracerProvider":
provider: Final = TracerProvider(
resource=self._get_litellm_resource(self.config), shutdown_on_exit=owns_exporter
)
provider.add_span_processor(self._get_span_processor(config_override=dynamic_config))
return provider
self._tracer_provider_cache[cache_key] = temp_provider
cache_key: Final = (
f"dynamic_config:{dynamic_config.exporter}:{dynamic_config.endpoint}:{dynamic_config.headers}"
)
return self._cached_dynamic_tracer(cache_key, _build, owns_exporter)
return temp_provider.get_tracer(LITELLM_TRACER_NAME)
def _get_tracer_with_dynamic_headers(self, dynamic_headers: dict):
"""Create a temporary tracer with dynamic headers for this request only."""
def _get_tracer_with_dynamic_headers(self, dynamic_headers: Mapping[str, str]) -> "_Tracer":
"""Create (or reuse) a tracer whose OTLP headers come from a per-request credential set."""
from opentelemetry.sdk.trace import TracerProvider
# Prevents thread exhaustion by reusing providers for the same credential sets (e.g. per-team keys)
owns_exporter: Final = _provider_owns_exporter(self.OTEL_EXPORTER)
def _build() -> "_SDKTracerProvider":
provider: Final = TracerProvider(
resource=self._get_litellm_resource(self.config), shutdown_on_exit=owns_exporter
)
provider.add_span_processor(self._get_span_processor(dynamic_headers=dynamic_headers))
return provider
cache_key: Final = str(sorted(dynamic_headers.items()))
if cache_key in self._tracer_provider_cache:
return self._tracer_provider_cache[cache_key].get_tracer(LITELLM_TRACER_NAME)
# Create a temporary tracer provider with dynamic headers
temp_provider: Final = TracerProvider(resource=self._get_litellm_resource(self.config))
temp_provider.add_span_processor(self._get_span_processor(dynamic_headers=dynamic_headers))
# Store in cache for reuse
self._tracer_provider_cache[cache_key] = temp_provider
return temp_provider.get_tracer(LITELLM_TRACER_NAME)
return self._cached_dynamic_tracer(cache_key, _build, owns_exporter)
def construct_dynamic_otel_headers(
self, standard_callback_dynamic_params: StandardCallbackDynamicParams
@ -2832,7 +2932,7 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
def _get_span_processor(
self,
dynamic_headers: dict | None = None,
dynamic_headers: Mapping[str, str] | None = None,
config_override: OpenTelemetryConfig | None = None,
):
from opentelemetry.sdk.trace.export import (
@ -3144,7 +3244,7 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
@staticmethod
def _get_headers_dictionary(
headers: str | dict | None,
headers: "str | Mapping[str, str] | None",
) -> dict[str, str]:
"""
Convert a string or dictionary of headers into a dictionary of headers.
@ -3158,8 +3258,8 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
for part in parts:
key, value = part.split("=", 1)
_split_otel_headers[key] = value
elif isinstance(headers, dict):
_split_otel_headers = headers
elif isinstance(headers, Mapping):
_split_otel_headers.update(headers)
return _split_otel_headers
async def async_management_endpoint_success_hook(

View file

@ -554,7 +554,7 @@ class TestLangfuseOtelKeyDynamicConfig:
assert tracer is not logger.tracer
assert len(logger._tracer_provider_cache) == 1
provider = next(iter(logger._tracer_provider_cache.values()))
provider = next(iter(logger._tracer_provider_cache.values())).provider
span_processors = provider._active_span_processor._span_processors
assert len(span_processors) == 1
assert isinstance(span_processors[0], BatchSpanProcessor)
@ -619,7 +619,7 @@ class TestLangfuseOtelKeyDynamicConfig:
assert secret not in logged
assert f"Basic {secret}" not in logged
provider = next(iter(logger._tracer_provider_cache.values()))
provider = next(iter(logger._tracer_provider_cache.values())).provider
exporter = provider._active_span_processor._span_processors[0].span_exporter
assert isinstance(exporter, OTLPSpanExporter)
assert exporter._headers == {

View file

@ -1,10 +1,15 @@
import asyncio
import concurrent.futures
import gc
import json
import os
import sys
import threading
import time
import unittest
import weakref
from datetime import datetime, timedelta, timezone
from types import MappingProxyType
from parameterized import parameterized
from unittest.mock import MagicMock, patch
@ -20,6 +25,7 @@ from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
import litellm
from litellm.integrations import opentelemetry as otel_module
from litellm.integrations.opentelemetry import (
OpenTelemetry,
OpenTelemetryConfig,
@ -1840,6 +1846,22 @@ class TestOpenTelemetryHeaderSplitting(unittest.TestCase):
result, {"api-key": "value1=part2", "config": "setting=enabled"}
)
def test_accepts_any_mapping_not_only_dict(self):
"""The parameter is typed Mapping, so a non-dict Mapping must not silently drop
every header and leave the exporter unauthenticated."""
otel = OpenTelemetry()
headers = MappingProxyType({"authorization": "Basic abc"})
self.assertEqual(otel._get_headers_dictionary(headers), {"authorization": "Basic abc"})
def test_returns_a_copy_so_the_exporter_never_aliases_the_caller(self):
"""The result is handed to a long-lived exporter, so it must not be the caller's
own dict."""
otel = OpenTelemetry()
headers = {"authorization": "Basic abc"}
result = otel._get_headers_dictionary(headers)
self.assertIsNot(result, headers)
self.assertEqual(result, headers)
class TestOpenTelemetryEndpointNormalization(unittest.TestCase):
"""Test suite for the unified _normalize_otel_endpoint method"""
@ -6007,3 +6029,186 @@ class TestOTELServiceTierAttributes(unittest.TestCase):
response_obj,
)
self.assertEqual(attributes[self.RESPONSE_KEY], "tier-added-by-provider-later")
class TestDynamicTracerProviderCache(unittest.TestCase):
"""Every credential-scoped TracerProvider that owns its exporter also owns a
BatchSpanProcessor worker thread that only stops on shutdown, so the cache holding them
must be bounded and must shut down whatever it drops (LIT-5437: threads accumulated
until pods were OOMKilled)."""
BSP_THREAD_NAME = "OtelBatchSpanProcessor"
def _logger(self, cap=3, exporter="console"):
logger = OpenTelemetry(
config=OpenTelemetryConfig(exporter=exporter, skip_set_global=True),
max_dynamic_tracer_providers=cap,
)
self.addCleanup(logger._tracer_provider.shutdown)
self.addCleanup(self._drain, logger)
return logger
def _drain(self, logger):
for entry in list(logger._tracer_provider_cache.values()):
entry.provider.shutdown()
logger._tracer_provider_cache.clear()
def _live_exporter_threads(self):
return [t for t in threading.enumerate() if t.name == self.BSP_THREAD_NAME]
def _wait_for_exporter_threads(self, expected, timeout=10.0):
"""Dropped providers are shut down off-thread, so poll instead of sleeping."""
deadline = time.time() + timeout
while time.time() < deadline:
count = len(self._live_exporter_threads())
if count <= expected:
return count
time.sleep(0.05)
return len(self._live_exporter_threads())
def test_distinct_credential_sets_stay_bounded(self):
"""One tenant per credential set must not mean one live thread per credential set."""
logger = self._logger(cap=3)
before = len(self._live_exporter_threads())
for i in range(25):
logger._get_tracer_with_dynamic_headers({"authorization": f"Basic tenant-{i}"})
self.assertEqual(len(logger._tracer_provider_cache), 3)
# Guards the thread-name constant: a rename upstream would make this read 0 and the
# bound assertion below would pass while measuring nothing.
self.assertGreaterEqual(len(self._live_exporter_threads()), 1)
self.assertLessEqual(self._wait_for_exporter_threads(before + 3) - before, 3)
def test_evicted_provider_is_shut_down(self):
"""An evicted provider is stopped, not silently dropped with its thread running."""
logger = self._logger(cap=3)
before = len(self._live_exporter_threads())
with patch.object(otel_module, "_shutdown_tracer_provider") as mock_shutdown:
logger._get_tracer_with_dynamic_headers({"authorization": "Basic evict-me"})
evicted = next(iter(logger._tracer_provider_cache.values()))
for i in range(3):
logger._get_tracer_with_dynamic_headers({"authorization": f"Basic keep-{i}"})
self.assertNotIn(evicted, logger._tracer_provider_cache.values())
self._wait_for_call(mock_shutdown)
mock_shutdown.assert_called_once_with(evicted.provider)
# The patch stopped the real shutdown, so stop the victim here; leaving its
# exporter thread alive would perturb the thread-census assertions elsewhere.
evicted.provider.shutdown()
still_cached = len(logger._tracer_provider_cache)
self.assertEqual(self._wait_for_exporter_threads(before + still_cached), before + still_cached)
def _wait_for_call(self, mock_fn, timeout=10.0):
"""The shutdown runs on a worker thread, so give it a moment to land."""
deadline = time.time() + timeout
while time.time() < deadline and not mock_fn.call_args_list:
time.sleep(0.05)
def test_concurrent_first_requests_build_one_provider(self):
"""Concurrent misses on one credential set race to build; only the winner may survive,
and the losers must be shut down rather than orphaned with their threads running."""
logger = self._logger(cap=3)
before = len(self._live_exporter_threads())
headers = {"authorization": "Basic same-tenant"}
barrier = threading.Barrier(16)
def _request_tracer(_):
barrier.wait()
return logger._get_tracer_with_dynamic_headers(headers)
with concurrent.futures.ThreadPoolExecutor(max_workers=16) as pool:
list(pool.map(_request_tracer, range(16)))
self.assertEqual(len(logger._tracer_provider_cache), 1)
self.assertEqual(self._wait_for_exporter_threads(before + 1) - before, 1)
def test_shared_exporter_instance_survives_dropped_providers(self):
"""A caller-supplied SpanExporter is shared with the logger's own provider, so a
dropped provider must not shut it down and silence the whole process."""
shared = InMemorySpanExporter()
logger = self._logger(cap=1, exporter=shared)
with logger.tracer.start_as_current_span("before"):
pass
for i in range(4):
logger._get_tracer_with_dynamic_headers({"authorization": f"Basic tenant-{i}"})
with logger.tracer.start_as_current_span("after"):
pass
self.assertEqual(
[span.name for span in shared.get_finished_spans()], ["before", "after"]
)
def test_mixed_ownership_cache_shuts_down_only_the_victims_that_own_their_exporter(self):
"""Both dynamic entry points share one cache, so it can hold providers of mixed
ownership. Whether an evicted provider may be shut down is a property of that
provider, not of the request that evicted it."""
shared = InMemorySpanExporter()
logger = self._logger(cap=1, exporter=shared)
with logger.tracer.start_as_current_span("before"):
pass
# Cached by the headers path, so its processor wraps the SHARED exporter.
logger._get_tracer_with_dynamic_headers({"authorization": "Basic shared-owner"})
# Evicted by the config path, which builds its OWN exporter from a named kind.
logger._get_tracer_with_dynamic_config(
OpenTelemetryConfig(exporter="console", skip_set_global=True)
)
with logger.tracer.start_as_current_span("after"):
pass
self.assertFalse(shared._stopped)
self.assertEqual(
[span.name for span in shared.get_finished_spans()], ["before", "after"]
)
def test_mixed_ownership_cache_still_reclaims_a_thread_owning_victim(self):
"""The other direction of the same defect: a victim that owns a real exporter
thread must still be shut down even when the evicting request does not."""
shared = InMemorySpanExporter()
logger = self._logger(cap=1, exporter=shared)
before = len(self._live_exporter_threads())
# Cached by the config path with a named kind, so it owns a BatchSpanProcessor thread.
logger._get_tracer_with_dynamic_config(
OpenTelemetryConfig(exporter="console", skip_set_global=True)
)
self.assertEqual(len(self._live_exporter_threads()) - before, 1)
# Evicted by the headers path, whose own exporter is the shared instance.
logger._get_tracer_with_dynamic_headers({"authorization": "Basic shared-owner"})
self.assertEqual(self._wait_for_exporter_threads(before) - before, 0)
def test_dropped_shared_exporter_provider_is_not_retained_by_an_exit_hook(self):
"""A provider we may never shut down must not register an interpreter-exit hook.
The hook holds a strong reference, so the provider would be pinned for the life of
the process (the very leak this fixes) and would stop the shared exporter at exit."""
shared = InMemorySpanExporter()
logger = self._logger(cap=1, exporter=shared)
logger._get_tracer_with_dynamic_headers({"authorization": "Basic a"})
entry = next(iter(logger._tracer_provider_cache.values()))
self.assertFalse(entry.owns_exporter)
victim = weakref.ref(entry.provider)
logger._get_tracer_with_dynamic_headers({"authorization": "Basic b"})
del entry
gc.collect()
self.assertIsNone(victim(), "evicted shared-exporter provider is still referenced")
def test_provider_that_owns_its_exporter_keeps_its_exit_flush(self):
"""The counterpart: a provider that owns a buffering processor must keep its exit
hook so its last batch still flushes when the process stops."""
logger = self._logger(cap=3)
logger._get_tracer_with_dynamic_headers({"authorization": "Basic owned"})
entry = next(iter(logger._tracer_provider_cache.values()))
self.assertTrue(entry.owns_exporter)
self.assertIsNotNone(entry.provider._atexit_handler)