From 61c23a94303a5f0a7bb794220f176a916939ae2c Mon Sep 17 00:00:00 2001 From: yucheng-berriai Date: Thu, 2 Jul 2026 00:27:30 -0700 Subject: [PATCH] refactor(otel/v2): merge fan-out processor into routing; describe abstractions in docstrings Fold TenantFanOutSpanProcessor from plumbing/fan_out.py into plumbing/routing.py so one module owns both span-routing paths: TenantTracerCache for the gen-AI LLM-call span (per-tenant clone providers) and TenantFanOutSpanProcessor for proxy-internal spans on the main provider. Both read the request's destinations from the same server-only contextvar, so there is a single source of truth for routing. Rewrite the OtelDestination and dynamic-params-destination docstrings to describe what the abstraction is (a resolved endpoint plus auth headers the exporter sends) rather than narrating what it is not. --- .../integrations/otel/model/destination.py | 13 +- litellm/integrations/otel/model/metadata.py | 11 +- litellm/integrations/otel/plumbing/fan_out.py | 204 ------------------ .../integrations/otel/plumbing/providers.py | 2 +- litellm/integrations/otel/plumbing/routing.py | 191 +++++++++++++++- .../integrations/otel/test_otel_v2_fan_out.py | 2 +- 6 files changed, 202 insertions(+), 221 deletions(-) delete mode 100644 litellm/integrations/otel/plumbing/fan_out.py diff --git a/litellm/integrations/otel/model/destination.py b/litellm/integrations/otel/model/destination.py index 3341c2cec02..abdc548fbce 100644 --- a/litellm/integrations/otel/model/destination.py +++ b/litellm/integrations/otel/model/destination.py @@ -1,11 +1,10 @@ -"""The resolved, admin-owned OTLP destination. +"""The resolved OTLP destination a request's traces export to. -A trace destination is admin-owned infrastructure config, never request data. -The proxy resolves a key/team's bound named credential into this typed, -backend-agnostic target (an endpoint plus auth headers) server-side, and the v2 -logger exports through it. Every OTEL backend -- Langfuse, Arize, Weave, a -self-hosted collector -- reduces to this shape; the per-backend field mapping -lives in ``litellm.integrations.otel.destinations``. +A destination is a backend-agnostic target: an endpoint plus the auth headers the +exporter sends. The proxy builds it from the named logging credential bound to the +request's identity chain, and the v2 logger exports through it. Every OTEL backend +-- Langfuse, Arize, Weave, a self-hosted collector -- reduces to this shape; the +per-backend field mapping lives in ``litellm.integrations.otel.presets.destinations``. """ from pydantic import BaseModel, ConfigDict, Field diff --git a/litellm/integrations/otel/model/metadata.py b/litellm/integrations/otel/model/metadata.py index 5db88d25bc4..30e286f75ef 100644 --- a/litellm/integrations/otel/model/metadata.py +++ b/litellm/integrations/otel/model/metadata.py @@ -51,13 +51,12 @@ if TYPE_CHECKING: def _otel_destinations(dynamic_params: Any) -> tuple[OtelDestination, ...]: - """The admin-resolved OTLP destinations carried on ``standard_callback_dynamic_params``. + """The admin-resolved OTLP destinations for this call, parsed off + ``standard_callback_dynamic_params``. - Server-set only (the proxy resolves the exporters assigned to the request's - identity chain and strips any client value), so this is the sole source the v2 - router trusts -- request-supplied vendor credentials are never read here. A - request fans out to every destination here; each logger keeps only the ones - tagged with its own backend. + The proxy resolves the destinations assigned to the request's identity chain and + places them here server-side. A request fans out to every destination; each + logger keeps only the ones tagged with its own backend. """ if not isinstance(dynamic_params, Mapping): return () diff --git a/litellm/integrations/otel/plumbing/fan_out.py b/litellm/integrations/otel/plumbing/fan_out.py deleted file mode 100644 index 5a4b35a67fa..00000000000 --- a/litellm/integrations/otel/plumbing/fan_out.py +++ /dev/null @@ -1,204 +0,0 @@ -"""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. - -Each backend often requires backend-specific Resource attributes (Arize rejects -spans missing ``model_id`` / ``arize.project.name``), so the fan-out wraps each -forwarded span with the destination's expected Resource before handing it to -the per-destination exporter. -""" - -from __future__ import annotations - -from collections import OrderedDict -from typing import TYPE_CHECKING - -from opentelemetry.context import Context -from opentelemetry.sdk.resources import Resource -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. Evicted entries are dropped (not shut down; see -# the eviction site) and reclaimed at process exit. -_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 - # The gen-AI LLM-call span (and the MCP tool-call sibling) is already - # routed to per-tenant destinations by the per-backend v2 logger via - # ``TenantTracerCache`` -- the logger picks the right attribute mapper - # (OpenInference for arize, GenAI semconv for langfuse_otel) and ships - # through the clone provider's appended exporter. Forwarding it here too - # would deliver a SECOND copy with the wrong vocabulary and a fresh - # span_id, surfacing in the destination as an orphaned duplicate. Skip. - if _is_genai_span(span): - return - # Proxy-internal spans (FastAPI server, ``auth`` phase, postgres lookups, - # post-call cost ledger) are generic OTel semantic-convention spans with - # no backend-specific vocabulary, so they ship to EVERY admin-resolved - # destination this request fans out to, regardless of the destination's - # ``callback_name``. The owner discriminator only matters for the gen-AI - # span (handled by the skip above). - for destination in destinations: - processor = self._processor_for(destination) - if processor is None: - continue - try: - processor.on_end(_with_destination_resource(span, destination)) - except Exception as exc: # noqa: BLE001 - 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: # noqa: BLE001 - 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: # noqa: BLE001 - 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 as _build_processor, - default_otlp_kind_for_backend, - ) - - try: - spec = ExporterSpec( - kind=default_otlp_kind_for_backend(destination.callback_name), - endpoint=destination.endpoint, - headers=destination.header_string(), - owner=None, - ) - exporter = _exporter_from_spec(spec) - processor = _build_processor(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: - # Evict the LRU entry but do NOT shut it down here: a - # ``BatchSpanProcessor`` may still hold spans queued on its exporter - # thread, and calling ``shutdown`` synchronously can drop or raise on - # those in-flight spans. Dropping the reference lets the worker drain - # naturally and be reclaimed at process exit. The cache is bounded at - # ``_MAX_CACHED_PROCESSORS``, so the un-shut-down working set stays - # bounded. - self._processors.popitem(last=False) - return processor - - -# Attribute set on every gen-AI LLM-call span by the v2 emitter. Used as the -# unambiguous skip signal: only the LLM-call span carries this, and the -# per-backend v2 logger already routes it to per-tenant destinations through -# the TenantTracerCache clone provider's appended exporter. -_GENAI_SPAN_ATTR = "gen_ai.operation.name" - - -def _is_genai_span(span: ReadableSpan) -> bool: - attributes = span.attributes or {} - return _GENAI_SPAN_ATTR in attributes - - -def _with_destination_resource(span: ReadableSpan, destination: OtelDestination) -> ReadableSpan: - """Return ``span`` with its Resource augmented by the destination's required - attributes. The original span object is left untouched; a shallow wrapper - reuses every other field and only swaps the ``resource`` property.""" - from litellm.integrations.otel.plumbing.providers import ( - destination_resource_attrs, - ) - - extra = destination_resource_attrs(destination) - if not extra: - return span - merged = Resource.create({**dict(span.resource.attributes), **extra}) - return _ResourceWrappedReadableSpan(span, merged) - - -class _ResourceWrappedReadableSpan(ReadableSpan): - """A ``ReadableSpan`` view whose ``resource`` is overridden. - - The OTLP exporter reads each span's ``resource`` when serializing; for - backend-specific attributes (Arize's ``model_id``) we substitute the - destination's expected Resource without mutating the underlying span. - """ - - def __init__(self, inner: ReadableSpan, resource: Resource) -> None: - super().__init__( - name=inner.name, - context=inner.context, - parent=inner.parent, - resource=resource, - attributes=inner.attributes, - events=inner.events, - links=inner.links, - kind=inner.kind, - status=inner.status, - start_time=inner.start_time, - end_time=inner.end_time, - instrumentation_scope=inner.instrumentation_scope, - ) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 90fb28a95ab..71eb6aef5aa 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -346,7 +346,7 @@ def build_tracer_provider( provider.add_span_processor(baggage_processor) if attach_tenant_fan_out or tenant_fan_out_owner is not None: - from litellm.integrations.otel.plumbing.fan_out import ( + from litellm.integrations.otel.plumbing.routing import ( TenantFanOutSpanProcessor, ) diff --git a/litellm/integrations/otel/plumbing/routing.py b/litellm/integrations/otel/plumbing/routing.py index 37104331b98..5ffe25bee6b 100644 --- a/litellm/integrations/otel/plumbing/routing.py +++ b/litellm/integrations/otel/plumbing/routing.py @@ -1,4 +1,12 @@ -"""Per-request multi-tenant tracer routing with fan-out. +"""Per-request multi-tenant tracer routing and span fan-out. + +This module owns both halves of routing a request's spans to its admin-owned +destinations. ``TenantTracerCache`` handles the gen-AI LLM-call span by building +per-tenant clone providers (grouped by backend Resource attributes). +``TenantFanOutSpanProcessor`` (at the bottom) handles the proxy-internal spans on +the main provider (server, auth, DB, cost ledger) by forwarding each to every +destination. Both read the request's destinations from the same server-only +contextvar, so there is one source of truth for where a request's traces go. A call's identity chain is assigned a set of admin-owned OTEL destinations (``LLMCallEvent.otel_destinations``, resolved server-side from named credentials). @@ -17,12 +25,15 @@ a caller can neither redirect a trace nor spawn providers. from collections import OrderedDict -from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.context import Context +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import ReadableSpan, Span, SpanProcessor, TracerProvider from opentelemetry.trace import Tracer from litellm._logging import verbose_logger from litellm.integrations.otel.model.config import ExporterSpec, OpenTelemetryV2Config from litellm.integrations.otel.model.destination import OtelDestination +from litellm.integrations.otel.plumbing.context import request_destinations from litellm.integrations.otel.plumbing.providers import ( build_tracer_provider, get_tracer, @@ -224,3 +235,179 @@ class TenantTracerCache: "resource_attributes": merged_resource_attrs, } ) + + +# --- Proxy-internal span fan-out ------------------------------------------- # +# +# ``TenantTracerCache`` above routes the gen-AI LLM-call span (and its MCP-tool +# sibling) to per-tenant destinations through clone providers. The processor below +# handles the OTHER span class: the proxy-internal spans (FastAPI server span, the +# ``auth`` phase, DB lookups, the post-call cost ledger) emitted on the MAIN +# provider. It forwards each to every admin-resolved destination for the request, +# reading them from the same server-only contextvar the cache's callers set, so both +# routing paths share one source of truth for where a request's traces go. + +# 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. Evicted entries are dropped (not shut down; see +# the eviction site) and reclaimed at process exit. +_MAX_CACHED_PROCESSORS = 256 + +# Attribute set on every gen-AI LLM-call span by the v2 emitter. Used as the +# unambiguous skip signal: only the LLM-call span carries this, and the per-backend +# v2 logger already routes it to per-tenant destinations through the +# TenantTracerCache clone provider's appended exporter. +_GENAI_SPAN_ATTR = "gen_ai.operation.name" + + +def _processor_key(destination: OtelDestination) -> tuple: + return (destination.endpoint, tuple(sorted(destination.headers.items()))) + + +def _is_genai_span(span: ReadableSpan) -> bool: + attributes = span.attributes or {} + return _GENAI_SPAN_ATTR in attributes + + +def _with_destination_resource(span: ReadableSpan, destination: OtelDestination) -> ReadableSpan: + """Return ``span`` with its Resource augmented by the destination's required + attributes. The original span object is left untouched; a shallow wrapper reuses + every other field and only swaps the ``resource`` property.""" + from litellm.integrations.otel.plumbing.providers import ( + destination_resource_attrs, + ) + + extra = destination_resource_attrs(destination) + if not extra: + return span + merged = Resource.create({**dict(span.resource.attributes), **extra}) + return _ResourceWrappedReadableSpan(span, merged) + + +class _ResourceWrappedReadableSpan(ReadableSpan): + """A ``ReadableSpan`` view whose ``resource`` is overridden. + + The OTLP exporter reads each span's ``resource`` when serializing; for + backend-specific attributes (Arize's ``model_id``) we substitute the + destination's expected Resource without mutating the underlying span. + """ + + def __init__(self, inner: ReadableSpan, resource: Resource) -> None: + super().__init__( + name=inner.name, + context=inner.context, + parent=inner.parent, + resource=resource, + attributes=inner.attributes, + events=inner.events, + links=inner.links, + kind=inner.kind, + status=inner.status, + start_time=inner.start_time, + end_time=inner.end_time, + instrumentation_scope=inner.instrumentation_scope, + ) + + +class TenantFanOutSpanProcessor(SpanProcessor): + """Forward each finished proxy-internal span to every admin-resolved 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 + 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 + # The gen-AI LLM-call span (and the MCP tool-call sibling) is already routed + # to per-tenant destinations by the per-backend v2 logger via + # ``TenantTracerCache`` -- the logger picks the right attribute mapper + # (OpenInference for arize, GenAI semconv for langfuse_otel) and ships through + # the clone provider's appended exporter. Forwarding it here too would deliver + # a SECOND copy with the wrong vocabulary and a fresh span_id, surfacing in the + # destination as an orphaned duplicate. Skip. + if _is_genai_span(span): + return + # Proxy-internal spans (FastAPI server, ``auth`` phase, postgres lookups, + # post-call cost ledger) are generic OTel semantic-convention spans with no + # backend-specific vocabulary, so they ship to EVERY admin-resolved destination + # this request fans out to, regardless of the destination's ``callback_name``. + for destination in destinations: + processor = self._processor_for(destination) + if processor is None: + continue + try: + processor.on_end(_with_destination_resource(span, destination)) + except Exception as exc: # noqa: BLE001 + 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: # noqa: BLE001 + 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: # noqa: BLE001 + 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 as _build_processor, + default_otlp_kind_for_backend, + ) + + try: + spec = ExporterSpec( + kind=default_otlp_kind_for_backend(destination.callback_name), + endpoint=destination.endpoint, + headers=destination.header_string(), + owner=None, + ) + exporter = _exporter_from_spec(spec) + processor = _build_processor(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: + # Evict the LRU entry but do NOT shut it down here: a + # ``BatchSpanProcessor`` may still hold spans queued on its exporter + # thread, and calling ``shutdown`` synchronously can drop or raise on those + # in-flight spans. Dropping the reference lets the worker drain naturally + # and be reclaimed at process exit. The cache is bounded, so the + # un-shut-down working set stays bounded. + self._processors.popitem(last=False) + return processor 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 index e1af18c2dbc..c014527e32b 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py @@ -33,7 +33,7 @@ from litellm.integrations.otel.plumbing.context import ( # noqa: E402 set_request_destinations, _request_destinations, ) -from litellm.integrations.otel.plumbing.fan_out import ( # noqa: E402 +from litellm.integrations.otel.plumbing.routing import ( # noqa: E402 TenantFanOutSpanProcessor, _processor_key, )