diff --git a/litellm/integrations/otel/logger.py b/litellm/integrations/otel/logger.py index 7348d913b35..3014a334eeb 100644 --- a/litellm/integrations/otel/logger.py +++ b/litellm/integrations/otel/logger.py @@ -63,6 +63,7 @@ from litellm.integrations.otel.plumbing.metrics import ( create_genai_metrics, ) from litellm.integrations.otel.plumbing.providers import ( + attach_tenant_fan_out, build_tracer_provider, get_event_logger, get_meter, @@ -182,7 +183,7 @@ class OpenTelemetryV2(CustomLogger): self._tracer_provider: TracerProvider = ( tracer_provider if tracer_provider is not None - else build_tracer_provider(self.config, tenant_overrides=True, tenant_callback_name=callback_name) + else build_tracer_provider(self.config, tenant_overrides=True) ) self.tracer: Tracer = get_tracer(self._tracer_provider, LITELLM_TRACER_NAME) self._metrics_recorder = self._init_metrics(meter_provider) @@ -865,8 +866,13 @@ def publish_global_otel_v2_provider( ``opentelemetry.trace.set_tracer_provider``) are injected so the publish step is unit-testable without reading or mutating real global OTel state. Returns the logger whose provider was published. + + The published provider is also the one that fans spans out to key/team + destinations, because it is the only provider the whole request tree passes + through; see :func:`attach_tenant_fan_out`. """ logger: Final = select_global_otel_v2_logger(in_memory_loggers, registered=registered) + attach_tenant_fan_out(logger._tracer_provider) set_global_provider(logger._tracer_provider) return logger diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 20649332f40..0ec570eb9e3 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -3,6 +3,8 @@ import threading from collections import OrderedDict from collections.abc import Callable, Iterable +from concurrent.futures import ThreadPoolExecutor +from functools import lru_cache from typing import TYPE_CHECKING, Any, Final, Literal from opentelemetry import _logs, baggage, metrics @@ -246,20 +248,20 @@ class TenantFanOutSpanProcessor(SpanProcessor): requests stay isolated. The forwarded view keeps the original trace and parent ids, so the tenant gets the same tree the operator would have received. - ``callback_name`` scopes the fan-out to the backend this provider speaks for. Every - v2 logger carries its own provider and emits its own copy of a gen-AI span, so a - proxy running two of them would otherwise deliver the tenant two copies of the same - model call. The tenant's own backend always has a logger, since naming it in the - key or team config is what builds one. + Exactly one provider carries this processor, the one published as the OTel global + (see :func:`attach_tenant_fan_out`). That provider is the only one every span + passes through: the FastAPI server span, the auth span and the post-call database + spans are emitted on the global, while a second v2 logger's provider sees only + that logger's own gen-AI span. Attaching the fan-out per logger would hand a + tenant a one-span trace whenever its backend is not the global one, and two + copies of the model call whenever it is. """ def __init__( self, - callback_name: str | None = None, processor_factory: 'Callable[["OtelDestination"], SpanProcessor | None] | None' = None, ) -> None: self._lock: Final = threading.Lock() - self._callback_name: Final = callback_name self._build: Final = processor_factory if processor_factory is not None else _destination_processor self._processors: OrderedDict[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded LRU self._retired: OrderedDict[int, SpanProcessor] = OrderedDict() # mutable-ok: drains as exports finish @@ -270,8 +272,6 @@ class TenantFanOutSpanProcessor(SpanProcessor): def on_end(self, span: ReadableSpan) -> None: for destination in request_destinations(): - if destination.callback_name != self._callback_name: - continue processor = self._acquire(destination) # rebind-ok: loop variable; pyright forbids Final in a loop if processor is None: continue @@ -387,11 +387,16 @@ def _drain_in_background(processor: SpanProcessor) -> None: ``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. + tenant's spans behind it. The work goes to a two-thread pool 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 instead. """ - threading.Thread( - target=_shutdown_quietly, args=(processor,), daemon=True, name="litellm-otel-destination-drain" - ).start() + _drain_pool().submit(_shutdown_quietly, processor) + + +@lru_cache(maxsize=1) +def _drain_pool() -> ThreadPoolExecutor: + return ThreadPoolExecutor(max_workers=2, thread_name_prefix="litellm-otel-destination-drain") def _shutdown_quietly(processor: SpanProcessor) -> None: @@ -672,7 +677,6 @@ def build_tracer_provider( baggage_processor: SpanProcessor | None = None, use_simple_processor: bool | None = None, tenant_overrides: bool = False, - tenant_callback_name: str | None = None, ) -> TracerProvider: """Build the shared :class:`TracerProvider`. @@ -682,11 +686,12 @@ def build_tracer_provider( backends. ``exporter`` and ``use_simple_processor`` are explicit overrides: pass a single exporter to attach exactly that one (used by tests). - ``tenant_overrides`` belongs to the operator-level provider alone: it wraps each - owned exporter so a request that pointed that backend at a key's or team's own - account skips it, and adds the fan-out processor that delivers to that account - instead. The per-tenant providers this same function builds must leave it off, - or they would filter out the very spans they exist to carry. + ``tenant_overrides`` wraps each owned exporter so a request that pointed that + backend at a key's or team's own account skips it. Every v2 logger's provider + wants it, since any of them may own the overridden backend; delivering to the + tenant is a separate job, done once by :func:`attach_tenant_fan_out`. The + per-tenant providers this same function builds must leave it off, or they would + filter out the very spans they exist to carry. """ provider: Final = TracerProvider(resource=build_resource(config)) if baggage_processor is None: @@ -711,11 +716,26 @@ def build_tracer_provider( provider.add_span_processor( _OverriddenBackendFilter(processor, owner) if tenant_overrides and owner is not None else processor ) - if tenant_overrides: - provider.add_span_processor(TenantFanOutSpanProcessor(tenant_callback_name)) return provider +def attach_tenant_fan_out(provider: TracerProvider) -> None: + """Give ``provider`` the fan-out that delivers spans to key/team destinations. + + Called on the one provider published as the OTel global, and idempotent so a + second publish (a test, a re-initialized proxy) cannot double-export. + """ + if any(isinstance(processor, TenantFanOutSpanProcessor) for processor in _attached_processors(provider)): + return + provider.add_span_processor(TenantFanOutSpanProcessor()) + + +def _attached_processors(provider: TracerProvider) -> "tuple[SpanProcessor, ...]": + """The processors already on ``provider``, or empty when the SDK hides them.""" + multi: Final = getattr(provider, "_active_span_processor", None) + return tuple(getattr(multi, "_span_processors", ())) + + def get_tracer(provider: TracerProvider, name: str = "litellm") -> Tracer: # Stamp the instrumentation scope with the LiteLLM package version so every # emitted span carries a deterministic ``scope.version`` (the standard OTel diff --git a/litellm/integrations/otel/plumbing/routing.py b/litellm/integrations/otel/plumbing/routing.py index 710c51d4942..a0831db2eea 100644 --- a/litellm/integrations/otel/plumbing/routing.py +++ b/litellm/integrations/otel/plumbing/routing.py @@ -233,13 +233,12 @@ class TenantTracerCache: the caller's span start. The caller must ``release`` it exactly once. """ # An overridden backend is delivered by the fan-out processor, which carries the - # whole trace. Routing here too would detach this span onto a second provider, + # whole trace and already carries this tenant's credentials and service name. + # Routing here too would detach this span onto a second provider, # so the tenant would get the request tree plus a stray one-span trace. - credential_headers: Final = ( - _NO_HEADERS - if self._callback_name is not None and self._callback_name in overridden_backends() - else self._credential_headers(dynamic_params) - ) + if self._callback_name is not None and self._callback_name in overridden_backends(): + return TenantRoute(tracer=default, detached=False) + credential_headers: Final = self._credential_headers(dynamic_params) project_headers: Final = self._project_headers(auth_metadata) service_name: Final = tenant_service_name(auth_metadata) if not credential_headers and not project_headers and service_name is None: diff --git a/litellm/integrations/otel/presets/destinations.py b/litellm/integrations/otel/presets/destinations.py index 815e8575a76..2bf9bfa5261 100644 --- a/litellm/integrations/otel/presets/destinations.py +++ b/litellm/integrations/otel/presets/destinations.py @@ -118,11 +118,17 @@ def destination_capable_backends() -> frozenset[str]: return frozenset(_DESTINATION_BY_CALLBACK) & frozenset(DYNAMIC_HEADERS_BY_CALLBACK) -def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -> OtelDestination | None: +def destination_for( + callback_name: str, + params: StandardCallbackDynamicParams, + service_name: str | None = None, +) -> OtelDestination | None: """The destination ``params`` names for ``callback_name``, or ``None``. ``None`` means the caller configured nothing usable for this backend, so the - request keeps the operator's global exporters. + request keeps the operator's global exporters. ``service_name`` is the key's or + team's ``otel_service_name``, which the per-request tracer route applies when the + backend is not overridden and the destination has to apply once it is. """ from litellm.integrations.otel.presets import DYNAMIC_HEADERS_BY_CALLBACK @@ -140,7 +146,7 @@ def destination_for(callback_name: str, params: StandardCallbackDynamicParams) - return OtelDestination( endpoint=endpoint, headers=MappingProxyType(dict(headers)), # mutable-ok: MappingProxyType needs a concrete mapping to wrap - resource_attributes=_NO_ATTRS, + resource_attributes=MappingProxyType({"service.name": service_name}) if service_name else _NO_ATTRS, callback_name=callback_name, protocol=protocol, ) diff --git a/litellm/proxy/litellm_pre_call_utils.py b/litellm/proxy/litellm_pre_call_utils.py index aed9933424b..5c6af78aad7 100644 --- a/litellm/proxy/litellm_pre_call_utils.py +++ b/litellm/proxy/litellm_pre_call_utils.py @@ -1044,12 +1044,33 @@ def resolve_tenant_otel_destinations( } ) ), + _tenant_service_name(user_api_key_dict), ) ) is not None ) +def _tenant_service_name(user_api_key_dict: UserAPIKeyAuth) -> str | None: + """The ``service.name`` this key or team configured, the key winning over its team. + + Same fields and same precedence the request-metadata build applies, read straight + off the auth object because destinations resolve during auth, before that metadata + is assembled. + """ + sources: Final = (user_api_key_dict.metadata, user_api_key_dict.team_metadata) + return next( + ( + stripped + for source in sources + if source + for field in OTEL_SERVICE_NAME_METADATA_KEYS + if isinstance(value := source.get(field), str) and (stripped := value.strip()) + ), + None, + ) + + def clean_headers( headers: Headers, litellm_key_header_name: str | None = 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 71b1164bfc0..9da2e25312b 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -17,6 +17,10 @@ from litellm.integrations.otel.model.config import ( is_otel_v2_enabled, ) from litellm.integrations.otel.model.destination import OtelDestination +from litellm.integrations.otel.logger import ( + OpenTelemetryV2, + publish_global_otel_v2_provider, +) from litellm.integrations.otel.plumbing.context import ( overridden_backends, request_destinations, @@ -68,7 +72,7 @@ def wired_provider(dest_exporter: InMemorySpanExporter, global_exporter: InMemor provider = TracerProvider() provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(global_exporter), "langfuse_otel")) provider.add_span_processor( - TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) + TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) ) return provider @@ -101,7 +105,7 @@ class TestOverrideSuppression: provider = TracerProvider() provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(arize_exporter), "arize")) provider.add_span_processor( - TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) + TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) ) def run(): @@ -119,7 +123,7 @@ class TestFanOut: dest_exporter = InMemorySpanExporter() provider = TracerProvider() provider.add_span_processor( - TenantFanOutSpanProcessor("langfuse_otel", processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) + TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter)) ) tracer = get_tracer(provider, "litellm") @@ -141,18 +145,14 @@ class TestFanOut: for child in ("auth /v1/chat/completions", "chat gpt-4"): assert by_name[child].parent.span_id == root.context.span_id - def test_a_destination_for_another_backend_is_left_to_that_backends_provider(self): - """Every v2 logger has its own provider and emits its own copy of a gen-AI - span, so a proxy running two of them would hand the tenant the same model call - twice if each provider forwarded every destination.""" + def test_a_team_naming_two_backends_gets_the_trace_at_both(self): + """The fan-out rides one provider, so it cannot skip a destination on the + grounds that some other backend owns it: nothing else would deliver it.""" langfuse, arize = InMemorySpanExporter(), InMemorySpanExporter() by_endpoint = {"http://a.local": langfuse, "http://b.local": arize} provider = TracerProvider() provider.add_span_processor( - TenantFanOutSpanProcessor( - "langfuse_otel", - processor_factory=lambda d: SimpleSpanProcessor(by_endpoint[d.endpoint]), - ) + TenantFanOutSpanProcessor(processor_factory=lambda d: SimpleSpanProcessor(by_endpoint[d.endpoint])) ) def run(): @@ -167,24 +167,30 @@ class TestFanOut: in_fresh_context(run) assert [s.name for s in langfuse.get_finished_spans()] == ["chat gpt-4"] - assert arize.get_finished_spans() == () + assert [s.name for s in arize.get_finished_spans()] == ["chat gpt-4"] - def test_a_provider_that_speaks_for_no_backend_forwards_nothing(self): - """The bare ``otel`` callback has no backend of its own; forwarding from it - would duplicate whatever the tenant's own backend provider already sent.""" + def test_a_destination_carries_the_tenants_service_name(self): + """An overridden backend skips per-request tracer routing, so the service name + that route used to apply has to travel on the destination instead.""" dest = InMemorySpanExporter() provider = TracerProvider() - provider.add_span_processor( - TenantFanOutSpanProcessor(None, processor_factory=lambda _d: SimpleSpanProcessor(dest)) - ) + provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest))) def run(): - set_request_destinations((LANGFUSE_DEST,)) + set_request_destinations( + ( + OtelDestination( + endpoint="http://a.local", + callback_name="langfuse_otel", + resource_attributes={"service.name": "team-checkout"}, + ), + ) + ) emit(provider) in_fresh_context(run) - assert dest.get_finished_spans() == () + assert {s.resource.attributes["service.name"] for s in dest.get_finished_spans()} == {"team-checkout"} def test_a_destination_that_cannot_build_a_processor_is_skipped_quietly(self): """An unbuildable destination must not cost the caller its request.""" @@ -195,7 +201,7 @@ class TestFanOut: attempts.append(destination.endpoint) provider = TracerProvider() - provider.add_span_processor(TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory)) + provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory)) def run(): set_request_destinations((LANGFUSE_DEST,)) @@ -216,7 +222,7 @@ class TestFanOut: return processor provider = TracerProvider() - provider.add_span_processor(TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory)) + provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory)) def run(): set_request_destinations((LANGFUSE_DEST,)) @@ -238,10 +244,35 @@ class TestProviderWiring: return [type(p).__name__ for p in provider._active_span_processor._span_processors] assert "_OverriddenBackendFilter" in kinds(operator) - assert "TenantFanOutSpanProcessor" in kinds(operator) assert "_OverriddenBackendFilter" not in kinds(tenant), "a per-tenant provider must not filter itself out" + assert "TenantFanOutSpanProcessor" not in kinds(operator), "delivery belongs to the published global alone" assert "TenantFanOutSpanProcessor" not in kinds(tenant) + def test_only_the_published_global_provider_delivers_to_tenants(self): + """A second v2 logger's provider never sees the server, auth or database spans, + so fanning out from it would hand the tenant a one-span trace. Publishing is + what picks the one provider the whole request tree passes through.""" + config = OpenTelemetryV2Config(exporters=[ExporterSpec(kind="in_memory", owner=ExporterOwner.ARIZE_AX)]) + published, other = OpenTelemetryV2(config=config, callback_name="arize"), OpenTelemetryV2(config=config) + + publish_global_otel_v2_provider([other], lambda _p: None, registered=published) + + def kinds(logger): + return [type(p).__name__ for p in logger._tracer_provider._active_span_processor._span_processors] + + assert kinds(published).count("TenantFanOutSpanProcessor") == 1 + assert "TenantFanOutSpanProcessor" not in kinds(other) + + def test_publishing_twice_does_not_double_export(self): + config = OpenTelemetryV2Config(exporters=[ExporterSpec(kind="in_memory", owner=ExporterOwner.ARIZE_AX)]) + logger = OpenTelemetryV2(config=config, callback_name="arize") + + publish_global_otel_v2_provider([], lambda _p: None, registered=logger) + publish_global_otel_v2_provider([], lambda _p: None, registered=logger) + + kinds = [type(p).__name__ for p in logger._tracer_provider._active_span_processor._span_processors] + assert kinds.count("TenantFanOutSpanProcessor") == 1 + class TestRouting: def test_an_overridden_backend_is_not_detached_onto_a_second_provider(self): @@ -263,6 +294,27 @@ class TestRouting: assert route.tracer is default assert route.provider is None + def test_an_overridden_backend_does_not_detach_on_a_service_name_either(self): + """A key or team service name is its own reason to build a second provider, so + clearing only the credentials would still take the model call out of the tree.""" + config = OpenTelemetryV2Config( + exporters=[ExporterSpec(kind="otlp_http", endpoint="http://op.local", owner=ExporterOwner.LANGFUSE_OTEL)] + ) + cache = TenantTracerCache(config, "langfuse_otel", "litellm") + default = get_tracer(TracerProvider(), "litellm") + auth_metadata = {"otel_service_name": "team-checkout"} + + assert cache.route_for(default, None, auth_metadata).detached is False + assert cache.route_for(default, None, auth_metadata).tracer is not default + + def run(): + set_request_destinations((LANGFUSE_DEST,)) + return cache.route_for(default, None, auth_metadata) + + route = in_fresh_context(run) + assert route.tracer is default, "the fan-out carries the service name on the destination instead" + assert route.provider is None + @pytest.mark.usefixtures("allow_test_hosts") class TestDestinationResolution: @@ -290,6 +342,57 @@ class TestDestinationResolution: assert [d.endpoint for d in destinations] == ["http://team.local/api/public/otel"] assert destinations[0].callback_name == "langfuse_otel" + def test_a_keys_service_name_outranks_its_teams_on_the_destination(self, monkeypatch): + """The key/team ``otel_service_name`` used to reach the backend through + per-request tracer routing, which an overridden backend skips.""" + monkeypatch.setenv("LITELLM_OTEL_V2", "true") + is_otel_v2_enabled.cache_clear() + auth = UserAPIKeyAuth( + metadata={"otel_service_name": "key-svc"}, + team_metadata={ + "otel_service_name": "team-svc", + "logging": [ + { + "callback_name": "langfuse_otel", + "callback_type": "success", + "callback_vars": { + "langfuse_public_key": "pk-team", + "langfuse_secret_key": "sk-team", + "langfuse_host": "http://team.local", + }, + } + ], + }, + ) + + destinations = resolve_tenant_otel_destinations(auth) + + assert dict(destinations[0].resource_attributes) == {"service.name": "key-svc"} + + def test_a_team_that_named_no_service_name_gets_no_resource_override(self, monkeypatch): + monkeypatch.setenv("LITELLM_OTEL_V2", "true") + is_otel_v2_enabled.cache_clear() + auth = UserAPIKeyAuth( + team_metadata={ + "otel_service_name": " ", + "logging": [ + { + "callback_name": "langfuse_otel", + "callback_type": "success", + "callback_vars": { + "langfuse_public_key": "pk-team", + "langfuse_secret_key": "sk-team", + "langfuse_host": "http://team.local", + }, + } + ], + } + ) + + destinations = resolve_tenant_otel_destinations(auth) + + assert dict(destinations[0].resource_attributes) == {} + def test_the_key_wins_over_the_team_for_the_same_backend(self, monkeypatch): monkeypatch.setenv("LITELLM_OTEL_V2", "true") is_otel_v2_enabled.cache_clear() @@ -591,16 +694,21 @@ class TestEvictionSafety: built.append(self.Recording()) return built[-1] - return TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory), built + return TenantFanOutSpanProcessor(processor_factory=factory), built @staticmethod def _dest(index): return LANGFUSE_DEST.model_copy(update={"endpoint": f"http://d{index}/otel"}) @staticmethod - def _settle(fan_out): - for _ in range(50): - if not fan_out._retired: + def _settle(fan_out, processor=None): + """Wait for retirement to clear and, when given, for the drain to run. + + The drain pool is shared and bounded, so a shed processor is closed once a + worker picks it up rather than the moment it is handed over. + """ + for _ in range(500): + if not fan_out._retired and (processor is None or processor.shutdown_calls): return time.sleep(0.02) @@ -630,7 +738,7 @@ class TestEvictionSafety: for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 1): fan_out._acquire(self._dest(index)) fan_out._release(built[-1]) - self._settle(fan_out) + self._settle(fan_out, built[0]) assert built[0].shutdown_calls == 1 assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS @@ -652,7 +760,7 @@ class TestEvictionSafety: built.append(Slow()) return built[-1] - fan_out = TenantFanOutSpanProcessor("langfuse_otel", processor_factory=factory) + fan_out = TenantFanOutSpanProcessor(processor_factory=factory) started = time.monotonic() for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 1): fan_out._acquire(self._dest(index)) @@ -660,6 +768,38 @@ class TestEvictionSafety: assert time.monotonic() - started < 2 + def test_shedding_many_processors_does_not_spawn_a_thread_each(self): + """A tenant that cycles its destination config sheds a processor per request, + so a thread per shed processor is a thread per request against a slow + collector.""" + 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) + 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)}" + finally: + release.set() + self._settle(fan_out, built[0]) + def test_a_retired_processor_is_still_closed_on_shutdown(self): from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS