diff --git a/litellm/integrations/opentelemetry.py b/litellm/integrations/opentelemetry.py index c3461c849dc..9363db12385 100644 --- a/litellm/integrations/opentelemetry.py +++ b/litellm/integrations/opentelemetry.py @@ -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( diff --git a/tests/test_litellm/integrations/test_langfuse_otel.py b/tests/test_litellm/integrations/test_langfuse_otel.py index 9392f974570..417921c166b 100644 --- a/tests/test_litellm/integrations/test_langfuse_otel.py +++ b/tests/test_litellm/integrations/test_langfuse_otel.py @@ -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 == { diff --git a/tests/test_litellm/integrations/test_opentelemetry.py b/tests/test_litellm/integrations/test_opentelemetry.py index b300c386326..fa1c9fa8a79 100644 --- a/tests/test_litellm/integrations/test_opentelemetry.py +++ b/tests/test_litellm/integrations/test_opentelemetry.py @@ -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)