diff --git a/litellm/integrations/otel/logger.py b/litellm/integrations/otel/logger.py index f231a8257fc..6461abaa9eb 100644 --- a/litellm/integrations/otel/logger.py +++ b/litellm/integrations/otel/logger.py @@ -110,7 +110,11 @@ class OpenTelemetryV2(CustomLogger): self.config: OpenTelemetryV2Config = config or OpenTelemetryV2Config(**kwargs) self.callback_name = callback_name self._tracer_provider: TracerProvider = ( - tracer_provider if tracer_provider is not None else build_tracer_provider(self.config) + tracer_provider + if tracer_provider is not None + else build_tracer_provider( + self.config, tenant_fan_out_owner=callback_name + ) ) self.tracer: Tracer = get_tracer(self._tracer_provider, LITELLM_TRACER_NAME) self._metrics_recorder = self._init_metrics(meter_provider) diff --git a/litellm/integrations/otel/plumbing/context.py b/litellm/integrations/otel/plumbing/context.py index 64790da814b..d261244ad79 100644 --- a/litellm/integrations/otel/plumbing/context.py +++ b/litellm/integrations/otel/plumbing/context.py @@ -27,9 +27,28 @@ _PROPAGATOR = TraceContextTextMapPropagator() # and is inherited by ``asyncio.create_task`` children — i.e. the async logging # callbacks that close the span. It is never reset: the contextvar dies with the # request task, so there is nothing to leak. -_request_root_span: "ContextVar[Span | None]" = ContextVar( - "litellm_otel_request_root_span", default=None -) +_request_root_span: "ContextVar[Span | None]" = ContextVar("litellm_otel_request_root_span", default=None) + +# Per-request admin-resolved destinations. Set once at the auth boundary (the +# earliest point a request's identity is known) and read by the global-provider +# fan-out processor at ``on_end`` time, so every span the proxy emits for this +# request -- the FastAPI server span, the ``auth`` phase, DB lookups, the +# batch-write cost ledger -- ships to every per-tenant destination the admin +# assigned. Lives on a ``ContextVar`` so it follows the request task across +# ``asyncio.create_task`` children (the success/failure logging callbacks close +# the LLM span in a worker copied from the request context). Request-scoped: the +# contextvar dies with the request task, so nothing leaks across requests. +_request_destinations: "ContextVar[tuple]" = ContextVar("litellm_otel_request_destinations", default=()) + + +def set_request_destinations(destinations: tuple) -> None: + """Anchor the admin-resolved destinations for this request.""" + _request_destinations.set(tuple(destinations)) + + +def request_destinations() -> tuple: + """Destinations the request fans out to, or empty when none were resolved.""" + return _request_destinations.get() def set_request_root_span(span: Span) -> None: @@ -49,9 +68,7 @@ def request_root_span() -> "Span | None": return span if is_recordable_span(span) else None -def set_request_baggage( - values: Mapping[str, str], context: Context | None = None -) -> Context: +def set_request_baggage(values: Mapping[str, str], context: Context | None = None) -> Context: """Return a context with ``values`` written into Baggage.""" ctx = context for key, value in values.items(): diff --git a/litellm/integrations/otel/plumbing/fan_out.py b/litellm/integrations/otel/plumbing/fan_out.py new file mode 100644 index 00000000000..dbc3f8c4ec2 --- /dev/null +++ b/litellm/integrations/otel/plumbing/fan_out.py @@ -0,0 +1,143 @@ +"""Per-request span fan-out to admin-resolved destinations. + +Attached to the main ``TracerProvider`` so every span on it (the FastAPI server +span, the proxy's ``auth`` phase span, DB lookups, the post-call cost ledger) +ships to every per-tenant destination the admin assigned for this request, on +top of the provider's configured global exporters. + +Spans emitted through the per-tenant ``TenantTracerCache`` clone providers (the +gen-AI LLM-call span and its MCP-tool sibling) reach tenant backends through +the clone's own exporters; the clone's provider has its own processor list and +does NOT carry this fan-out processor, so a span is exported once per backend. + +Identical ``(trace_id, span_id)`` deduplication on the OTLP receiver collapses +the rare overlap into one span per backend (e.g. when the configured global +exporter and a per-tenant destination point at the same vendor account). +""" + +from __future__ import annotations + +from collections import OrderedDict +from typing import TYPE_CHECKING + +from opentelemetry.context import Context +from opentelemetry.sdk.trace import ReadableSpan, Span, SpanProcessor + +from litellm._logging import verbose_logger +from litellm.integrations.otel.model.config import ExporterSpec +from litellm.integrations.otel.plumbing.context import request_destinations + +if TYPE_CHECKING: + from litellm.integrations.otel.model.destination import OtelDestination + +# Bound on cached per-destination processors. One processor per +# ``(endpoint, sorted(headers))`` pair, so the working set is one entry per +# admin-resolved tenant credential -- a real-world deployment with hundreds of +# tenants stays well under this. The LRU shuts down the evicted processor's +# exporter thread so the working set is reclaimed. +_MAX_CACHED_PROCESSORS = 256 + + +def _processor_key(destination: "OtelDestination") -> tuple: + return (destination.endpoint, tuple(sorted(destination.headers.items()))) + + +class TenantFanOutSpanProcessor(SpanProcessor): + """Forward each finished span to every admin-resolved per-tenant destination. + + The destinations are looked up from a request-scoped contextvar set during + auth, so the processor is stateless across requests and isolation across + concurrent requests is guaranteed by Python's contextvars. + """ + + def __init__(self, owner_callback_name: str | None) -> None: + self._owner = owner_callback_name + # Built lazily so we avoid importing the providers module at class + # definition time (which would create a circular import with routing). + self._processors: "OrderedDict[tuple, SpanProcessor]" = OrderedDict() + + def on_start(self, span: "Span", parent_context: Context | None = None) -> None: + return None + + def on_end(self, span: ReadableSpan) -> None: + destinations = request_destinations() + if not destinations: + return + for destination in destinations: + if destination.callback_name != self._owner: + continue + processor = self._processor_for(destination) + if processor is None: + continue + try: + processor.on_end(span) + except Exception as exc: + verbose_logger.debug( + "OTel V2 fan-out: forwarding span to %s failed: %s", + destination.endpoint, + exc, + ) + + def shutdown(self) -> None: + for processor in self._processors.values(): + try: + processor.shutdown() + except Exception as exc: + verbose_logger.debug("OTel V2 fan-out: processor shutdown failed: %s", exc) + self._processors.clear() + + def force_flush(self, timeout_millis: int = 30000) -> bool: + all_ok = True + for processor in self._processors.values(): + try: + if not processor.force_flush(timeout_millis): + all_ok = False + except Exception: + all_ok = False + return all_ok + + def _processor_for(self, destination: "OtelDestination") -> "SpanProcessor | None": + key = _processor_key(destination) + cached = self._processors.get(key) + if cached is not None: + self._processors.move_to_end(key) + return cached + from litellm.integrations.otel.plumbing.providers import ( + _exporter_from_spec, + _processor_for, + ) + + try: + spec = ExporterSpec( + kind=_resolve_kind(destination), + endpoint=destination.endpoint, + headers=destination.header_string(), + owner=None, + ) + exporter = _exporter_from_spec(spec) + processor = _processor_for(exporter, use_simple=False) + except Exception as exc: + verbose_logger.debug( + "OTel V2 fan-out: failed to build processor for %s: %s", + destination.endpoint, + exc, + ) + return None + self._processors[key] = processor + if len(self._processors) > _MAX_CACHED_PROCESSORS: + _, evicted = self._processors.popitem(last=False) + try: + evicted.shutdown() + except Exception: + pass + return processor + + +# Per-backend transport: Arize speaks OTLP/gRPC, every other current preset +# speaks OTLP/HTTP. Kept here so the fan-out picks the same transport the +# preset's own exporter uses, mirroring ``TenantTracerCache._owned_otlp_kind``. +_GRPC_BACKENDS = frozenset({"arize"}) + + +def _resolve_kind(destination: "OtelDestination") -> str: + return "otlp_grpc" if destination.callback_name in _GRPC_BACKENDS else "otlp_http" diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 6d0710397a3..d03d28faa87 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -290,6 +290,7 @@ def build_tracer_provider( exporter: SpanExporter | None = None, baggage_processor: SpanProcessor | None = None, use_simple_processor: bool | None = None, + tenant_fan_out_owner: str | None = None, ) -> TracerProvider: """Build the shared :class:`TracerProvider`. @@ -298,6 +299,13 @@ def build_tracer_provider( ``config.exporters`` entry — this is what fans spans out to multiple backends. ``exporter`` and ``use_simple_processor`` are explicit overrides: pass a single exporter to attach exactly that one (used by tests). + + ``tenant_fan_out_owner`` — when set, attach a ``TenantFanOutSpanProcessor`` + that forwards each finished span to the admin-resolved per-tenant + destinations whose ``callback_name`` matches the owner. Only the MAIN v2 + logger provider opts in; the per-tenant clone providers (built by + ``TenantTracerCache``) intentionally do not, so the LLM-call span exported + through them is not also fanned out here. """ provider = TracerProvider(resource=build_resource(config)) if baggage_processor is None: @@ -306,6 +314,15 @@ def build_tracer_provider( ) provider.add_span_processor(baggage_processor) + if tenant_fan_out_owner is not None: + from litellm.integrations.otel.plumbing.fan_out import ( + TenantFanOutSpanProcessor, + ) + + provider.add_span_processor( + TenantFanOutSpanProcessor(owner_callback_name=tenant_fan_out_owner) + ) + if exporter is not None: provider.add_span_processor(_processor_for(exporter, use_simple_processor)) return provider diff --git a/litellm/proxy/auth/user_api_key_auth.py b/litellm/proxy/auth/user_api_key_auth.py index e439f6a5998..134db25bd30 100644 --- a/litellm/proxy/auth/user_api_key_auth.py +++ b/litellm/proxy/auth/user_api_key_auth.py @@ -949,6 +949,51 @@ async def _resolve_jwt_to_virtual_key( return None +async def _hoist_request_destinations( + request: Request, user_api_key_dict: UserAPIKeyAuth +) -> None: + """Resolve admin-owned OTEL destinations for this request and anchor them. + + Runs after the auth builder, while we are still inside the request task, so + the ``ContextVar`` is visible to every ``SpanProcessor.on_end`` that fires + for spans this request opens. Stashes the same list on ``request.state`` so + ``_apply_admin_logging_exporters`` can reuse it without a second DB pass. + + Best-effort: a resolver failure must not break the request. The contextvar + is left at its default (empty tuple), so the fan-out processor no-ops. + """ + try: + from litellm.integrations.otel.model.destination import OtelDestination + from litellm.integrations.otel.plumbing.context import ( + set_request_destinations, + ) + from litellm.proxy.litellm_pre_call_utils import ( + _resolve_logging_exporters, + ) + + destinations_raw, _backends = await _resolve_logging_exporters( + user_api_key_dict + ) + destinations = tuple( + OtelDestination( + callback_name=item.get("callback_name"), + endpoint=item.get("endpoint", ""), + headers=item.get("headers") or {}, + ) + for item in destinations_raw + if isinstance(item, dict) and item.get("endpoint") + ) + set_request_destinations(destinations) + try: + request.state.otel_destinations = destinations_raw + except Exception: + pass + except Exception as exc: + verbose_proxy_logger.debug( + "OTel V2: hoist destination resolution failed: %s", exc + ) + + def _ensure_parent_otel_span_on_request_state(request: Request) -> None: """Idempotently create the OTEL SERVER span and stash it on ``request.state.parent_otel_span``. Safe to call multiple times. @@ -2598,6 +2643,15 @@ async def user_api_key_auth( ) user_api_key_auth_obj.budget_reservation = None + # Admin-resolved OTEL destinations: anchor them on this request's task + # context BEFORE downstream spans close, so the global-provider fan-out + # processor forwards every span (server, auth, db, batch-write) to the + # admin-assigned per-tenant backends -- not just the gen-AI span the + # ``TenantTracerCache`` already routes. Also stashed on ``request.state`` + # so ``_apply_admin_logging_exporters`` reuses the result instead of + # re-resolving. + await _hoist_request_destinations(request, user_api_key_auth_obj) + ## ENSURE DISABLE ROUTE WORKS ACROSS ALL USER AUTH FLOWS ## RouteChecks.should_call_route( route=route, valid_token=user_api_key_auth_obj, request=request diff --git a/litellm/proxy/litellm_pre_call_utils.py b/litellm/proxy/litellm_pre_call_utils.py index 6f824cd5d2b..298235a6235 100644 --- a/litellm/proxy/litellm_pre_call_utils.py +++ b/litellm/proxy/litellm_pre_call_utils.py @@ -682,7 +682,9 @@ async def _resolve_logging_exporters( async def _apply_admin_logging_exporters( - data: dict, user_api_key_dict: UserAPIKeyAuth + data: dict, + user_api_key_dict: UserAPIKeyAuth, + cached_destinations: "list | None" = None, ) -> None: """Stamp the resolved fan-out destinations onto ``data`` and activate their backends. @@ -691,8 +693,23 @@ async def _apply_admin_logging_exporters( ``all_litellm_params``, so scrubbed from the provider request body), not a top-level key, so an unknown field cannot leak to the provider. Default-deny means an identity with no assignment gets no per-tenant destination here. + + ``cached_destinations`` -- when ``user_api_key_auth`` already resolved the + destinations on this request (the FastAPI path), reuse the result instead of + running the resolver a second time. The SDK path passes ``None`` and the + resolver runs here. """ - destinations, backends = await _resolve_logging_exporters(user_api_key_dict) + if cached_destinations is not None: + destinations = list(cached_destinations) + backends = list( + dict.fromkeys( + str(d["callback_name"]) + for d in destinations + if isinstance(d, dict) and d.get("callback_name") + ) + ) + else: + destinations, backends = await _resolve_logging_exporters(user_api_key_dict) if not destinations: return proxy_metadata = data.get("litellm_metadata") @@ -2012,8 +2029,13 @@ async def add_litellm_data_to_request( # Admin-owned exporter assignment: resolve the union of exporters assigned across # the request's identity chain (key + team + org) into fan-out destinations and - # activate their backends. Default-deny: an unassigned identity gets none. - await _apply_admin_logging_exporters(data, user_api_key_dict) + # activate their backends. Default-deny: an unassigned identity gets none. Reuse + # the result ``user_api_key_auth`` already cached on ``request.state`` so the + # resolver runs once per request, not twice. + cached = getattr(getattr(request, "state", None), "otel_destinations", None) + await _apply_admin_logging_exporters( + data, user_api_key_dict, cached_destinations=cached + ) # Add disabled callbacks from key metadata if ( diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py b/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py new file mode 100644 index 00000000000..7d44522d2c1 --- /dev/null +++ b/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py @@ -0,0 +1,268 @@ +"""Tests for the per-request tenant fan-out SpanProcessor (Hoist). + +The fan-out processor lives on the main v2 ``TracerProvider`` and forwards every +finished span to the admin-resolved per-tenant destinations carried on the +request's ``ContextVar``. The contract under test: + +- Spans on the main provider land at every per-tenant destination matching this + backend (the ``owner_callback_name``). +- When destinations is empty (an unassigned identity, the SDK path, a request + before the resolver ran), the processor is a no-op. +- Concurrent requests with different destinations stay isolated -- contextvars + scope per task, so one request's tenant doesn't receive another's spans. +- The processor caches one BatchSpanProcessor per ``(endpoint, headers)`` pair + and skips destinations whose ``callback_name`` doesn't match its owner. +""" + +import asyncio +from typing import Any + +import pytest + +pytest.importorskip("opentelemetry") + +from opentelemetry.sdk.trace import TracerProvider # noqa: E402 +from opentelemetry.sdk.trace.export import SimpleSpanProcessor # noqa: E402 +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( # noqa: E402 + InMemorySpanExporter, +) +from opentelemetry.sdk.trace.export import SpanExporter # noqa: E402 + +from litellm.integrations.otel.model.destination import OtelDestination # noqa: E402 +from litellm.integrations.otel.plumbing.context import ( # noqa: E402 + set_request_destinations, + _request_destinations, +) +from litellm.integrations.otel.plumbing.fan_out import ( # noqa: E402 + TenantFanOutSpanProcessor, + _processor_key, +) +from litellm.integrations.otel.plumbing import providers # noqa: E402 +from litellm.integrations.otel.model.config import OpenTelemetryV2Config # noqa: E402 + + +@pytest.fixture(autouse=True) +def _reset_request_destinations(): + _request_destinations.set(()) + yield + _request_destinations.set(()) + + +def _build_provider_with_fan_out( + owner: str, exporter: SpanExporter +) -> tuple[TracerProvider, TenantFanOutSpanProcessor]: + """Build a real TracerProvider with the in-memory exporter as the configured + backend, plus the fan-out processor wired in. The fan-out processor's + per-destination processors are real BatchSpanProcessors, replaced in the + test by an injected SimpleSpanProcessor against an in-memory exporter via + ``monkeypatch`` so the test can read what was forwarded.""" + cfg = OpenTelemetryV2Config(exporter="in_memory") + provider = providers.build_tracer_provider( + cfg, exporter=exporter, tenant_fan_out_owner=owner + ) + fan_out = next( + p + for p in provider._active_span_processor._span_processors + if isinstance(p, TenantFanOutSpanProcessor) + ) + return provider, fan_out + + +def test_fan_out_forwards_span_to_matching_destination(monkeypatch): + """A span on the main provider lands at the per-tenant destination whose + callback_name matches the fan-out owner. Removing the fan-out branch makes + this test fail (the tenant exporter never sees the span).""" + main_exporter = InMemorySpanExporter() + provider, fan_out = _build_provider_with_fan_out("langfuse_otel", main_exporter) + + tenant_exporter = InMemorySpanExporter() + # Swap the lazy per-destination processor build for an in-memory one so we + # don't actually hit the network. + def _stub_processor_for(self, destination): + return SimpleSpanProcessor(tenant_exporter) + + monkeypatch.setattr( + TenantFanOutSpanProcessor, "_processor_for", _stub_processor_for + ) + + set_request_destinations( + ( + OtelDestination( + callback_name="langfuse_otel", + endpoint="https://cloud.langfuse.com/api/public/otel", + headers={"Authorization": "Bearer pk:sk"}, + ), + ) + ) + tracer = provider.get_tracer("test") + with tracer.start_as_current_span("auth /chat/completions"): + pass + + main_names = [s.name for s in main_exporter.get_finished_spans()] + tenant_names = [s.name for s in tenant_exporter.get_finished_spans()] + assert "auth /chat/completions" in main_names + assert "auth /chat/completions" in tenant_names + + +def test_fan_out_skips_destinations_for_other_backends(monkeypatch): + """Each fan-out processor owns ONE backend; a destination tagged for a + different backend is skipped, so the wrong exporter never sees the span.""" + main_exporter = InMemorySpanExporter() + _provider, _fan_out = _build_provider_with_fan_out( + "langfuse_otel", main_exporter + ) + + tenant_exporter = InMemorySpanExporter() + monkeypatch.setattr( + TenantFanOutSpanProcessor, + "_processor_for", + lambda self, destination: SimpleSpanProcessor(tenant_exporter), + ) + + set_request_destinations( + ( + OtelDestination( + callback_name="arize", + endpoint="https://otlp.arize.com/v1", + headers={"api_key": "k", "space_id": "s"}, + ), + ) + ) + tracer = _provider.get_tracer("test") + with tracer.start_as_current_span("auth"): + pass + + assert tenant_exporter.get_finished_spans() == () + + +def test_fan_out_noop_when_no_destinations(monkeypatch): + """An unassigned identity / pre-auth / SDK path leaves the contextvar at + its empty-tuple default; the processor must short-circuit, NOT crash, NOT + forward.""" + main_exporter = InMemorySpanExporter() + provider, _ = _build_provider_with_fan_out("langfuse_otel", main_exporter) + forwarded: list[Any] = [] + monkeypatch.setattr( + TenantFanOutSpanProcessor, + "_processor_for", + lambda self, destination: forwarded.append(destination) or None, + ) + tracer = provider.get_tracer("test") + with tracer.start_as_current_span("auth"): + pass + assert forwarded == [] + + +def test_fan_out_caches_processor_per_destination_key(monkeypatch): + """Two requests with the same destination must share one cached + BatchSpanProcessor; two different destinations must build two. Otherwise + every request rebuilds the OTLP exporter (and its background thread).""" + main_exporter = InMemorySpanExporter() + _provider, fan_out = _build_provider_with_fan_out("langfuse_otel", main_exporter) + + built: list[OtelDestination] = [] + + real_processor_for = TenantFanOutSpanProcessor._processor_for + + def _spy(self, destination): + result = real_processor_for(self, destination) + if result is not None: + built.append(destination) + return result + + # Build always returns None to keep the test offline, but we tracked + # invocations via the spy above. Easier: stub the inner builder. + def _stub(self, destination): + built.append(destination) + key = _processor_key(destination) + if key in self._processors: + return self._processors[key] + proc = SimpleSpanProcessor(InMemorySpanExporter()) + self._processors[key] = proc + return proc + + monkeypatch.setattr(TenantFanOutSpanProcessor, "_processor_for", _stub) + + dest_a = OtelDestination( + callback_name="langfuse_otel", endpoint="https://a/", headers={"k": "1"} + ) + dest_b = OtelDestination( + callback_name="langfuse_otel", endpoint="https://b/", headers={"k": "1"} + ) + set_request_destinations((dest_a,)) + tracer = _provider.get_tracer("test") + with tracer.start_as_current_span("s1"): + pass + set_request_destinations((dest_a,)) + with tracer.start_as_current_span("s2"): + pass + set_request_destinations((dest_b,)) + with tracer.start_as_current_span("s3"): + pass + + # _stub appends every call; the cache size is what matters: two unique + # ``(endpoint, headers)`` pairs -> two cached processors. + assert len(fan_out._processors) == 2 + + +def test_fan_out_per_request_isolation_with_concurrent_tasks(monkeypatch): + """Two requests running concurrently with different destinations must each + see only THEIR destination's spans. The contextvar scopes per asyncio task, + so this is a pin on the contextvar approach: switching to a global variable + would break this test.""" + main_exporter = InMemorySpanExporter() + provider, _ = _build_provider_with_fan_out("langfuse_otel", main_exporter) + + exporter_a = InMemorySpanExporter() + exporter_b = InMemorySpanExporter() + + def _stub(self, destination): + if destination.endpoint == "https://a/": + return SimpleSpanProcessor(exporter_a) + return SimpleSpanProcessor(exporter_b) + + monkeypatch.setattr(TenantFanOutSpanProcessor, "_processor_for", _stub) + tracer = provider.get_tracer("test") + + async def fire(label: str, endpoint: str): + set_request_destinations( + ( + OtelDestination( + callback_name="langfuse_otel", + endpoint=endpoint, + headers={}, + ), + ) + ) + with tracer.start_as_current_span(label): + await asyncio.sleep(0) + + async def driver(): + await asyncio.gather( + fire("req_a", "https://a/"), + fire("req_b", "https://b/"), + ) + + asyncio.run(driver()) + + names_a = {s.name for s in exporter_a.get_finished_spans()} + names_b = {s.name for s in exporter_b.get_finished_spans()} + assert "req_a" in names_a and "req_b" not in names_a + assert "req_b" in names_b and "req_a" not in names_b + + +def test_fan_out_owner_set_by_build_tracer_provider(): + """The v2 logger opts the MAIN provider into fan-out by passing its + callback_name; the TenantTracerCache clone providers do NOT, so the gen-AI + span emitted through them is not also forwarded by the fan-out processor. + Regressing this (adding the processor to clones) would double-export every + gen-AI span.""" + from litellm.integrations.otel.logger import OpenTelemetryV2 + + logger = OpenTelemetryV2(callback_name="arize") + main_procs = logger._tracer_provider._active_span_processor._span_processors + assert any(isinstance(p, TenantFanOutSpanProcessor) for p in main_procs) + # Build a tenant clone provider and confirm it has no fan-out processor. + clone = providers.build_tracer_provider(logger.config) + clone_procs = clone._active_span_processor._span_processors + assert not any(isinstance(p, TenantFanOutSpanProcessor) for p in clone_procs)