From da4583a110ddb19b8800c5e1584d82d83942f14c Mon Sep 17 00:00:00 2001 From: Yucheng Zhu Date: Tue, 28 Jul 2026 23:47:06 -0700 Subject: [PATCH] docs(otel/v2): tighten verbose docstrings and comments Several docstrings and comment blocks in the otel v2 plumbing ran 10-35 lines of rationale where a sentence or two states what the code does; the how and why belong in the PR and design notes. Trims them across routing, context, providers, logger, and the metadata model, keeping the load-bearing invariants (no-op-when-empty, the MCP anti-spoof rule, post-auth ordering, gen-AI-span skipping) in compressed form. No code changes: the executable tokens are byte-identical, only docstrings and comments moved. --- litellm/integrations/otel/logger.py | 274 ++++++------------ litellm/integrations/otel/model/metadata.py | 142 +++------ litellm/integrations/otel/plumbing/context.py | 149 +++------- .../integrations/otel/plumbing/providers.py | 110 +++---- litellm/integrations/otel/plumbing/routing.py | 113 +++----- 5 files changed, 230 insertions(+), 558 deletions(-) diff --git a/litellm/integrations/otel/logger.py b/litellm/integrations/otel/logger.py index 43a4ede1fe5..a76ee61ef9c 100644 --- a/litellm/integrations/otel/logger.py +++ b/litellm/integrations/otel/logger.py @@ -76,11 +76,9 @@ def _span_error_from_exception( status_code: int | None = None, traceback_str: str | None = None, ) -> SpanError: - """A ``SpanError`` for a proxy-level failure that never produced a - ``StandardLoggingPayload`` (auth / validation / malformed-body rejections), - mirroring ``_parse_error``'s field mapping so it stamps the same v2 keys a - failed LLM call does. ``status_code`` pins ``error.code`` to the real response - status, matching v1's SERVER-span behavior.""" + """A ``SpanError`` for a proxy-level failure that never produced a ``StandardLoggingPayload`` + (auth/validation/malformed-body rejections). ``status_code`` pins ``error.code`` to the real + response status.""" from litellm.litellm_core_utils.litellm_logging import StandardLoggingPayloadSetup info = StandardLoggingPayloadSetup.get_error_information( @@ -104,24 +102,17 @@ _OTEL_MODULES = ( ) -# Cap on the open-call carrier map. A span opened at ``pre_call`` that never -# reaches a success/failure callback (e.g. a stream that only fires stream -# events) would otherwise linger; bounding the map evicts the oldest so memory -# stays flat on a long-running proxy while covering every concurrent in-flight -# call. +# Cap on the open-call carrier map: a span opened at ``pre_call`` that never reaches a close +# callback (a stream that only fires stream events) would otherwise linger. Evicts the oldest. _OPEN_CALLS_MAX = 10_000 class _LLMCallSpan: """The state carried from the ``pre_call`` boundary to span close. - ``spans`` is one live span per destination Resource group (a backend like Arize - routing two projects yields two), opened at the boundary when the server span was - ambient. It is empty when creation was deferred because no ambient parent was visible - — in which case the async callback creates the span(s) against its own - (worker-copied) ambient context using ``start_time_ns``. The presence of a carrier - for a call at all is the proof that ``pre_call`` ran, i.e. that an upstream call was - actually attempted. + ``spans`` is one live span per destination Resource group, opened at the boundary; empty when + creation was deferred (no ambient parent visible), where the async callback creates them from + ``start_time_ns``. A carrier existing at all proves ``pre_call`` ran (an upstream call was attempted). """ __slots__ = ("spans", "start_time_ns") @@ -172,10 +163,8 @@ class OpenTelemetryV2(CustomLogger): def _init_metrics(self, meter_provider: Any | None) -> "GenAIMetricRecorder | None": """Create the six GenAI histograms when metrics are enabled, else ``None``. - ``meter_provider`` is an explicit override (tests inject one); otherwise the - provider is resolved from the OTel global so the operator's configured - readers/exporters receive the metrics, building and registering one only - when no global provider is set. + ``meter_provider`` is an explicit override (tests inject one); otherwise it is resolved + from the OTel global so the operator's readers/exporters receive the metrics. """ if not self.config.enable_metrics: return None @@ -186,11 +175,9 @@ class OpenTelemetryV2(CustomLogger): def _init_events(self, logger_provider: LoggerProvider | None) -> "GenAIEventRecorder | None": """Create the GenAI event recorder when events are enabled, else ``None``. - ``logger_provider`` is an explicit override (tests inject one); otherwise the - provider is resolved from the OTel global so an operator-configured logs - pipeline receives the events, building and registering one only when no - global provider is set. A ``None`` resolution means the operator opted out - of the logs signal, so no recorder is built. + ``logger_provider`` is an explicit override (tests inject one); otherwise it is resolved + from the OTel global. A ``None`` resolution means the operator opted out of logs, so no + recorder is built. """ if not self.config.enable_events: return None @@ -226,12 +213,8 @@ class OpenTelemetryV2(CustomLogger): setattr(proxy_server, "open_telemetry_logger", self) def _destinations_for_backend(self, call: "LLMCallEvent") -> tuple: - """The call's admin-resolved destinations that belong to THIS logger's backend. - - A request fans out across whatever exporters its identity chain is assigned; - each logger exports only the destinations tagged with its own callback_name, - so each backend's span keeps its own attribute vocabulary. - """ + """The call's admin-resolved destinations tagged with THIS logger's callback_name, + so each backend's span keeps its own attribute vocabulary.""" return tuple(d for d in call.otel_destinations if d.callback_name == self.callback_name) # ====================================================================== # @@ -242,21 +225,11 @@ class OpenTelemetryV2(CustomLogger): def log_pre_api_call(self, model, messages, kwargs): """Open the LLM-call span at the call boundary. - Runs synchronously inside the request task, before the upstream call — - the one place where the live server span is genuinely the ambient OTel - context — so the span parents to it natively, with no span threaded - through a metadata dict. The open span is stashed on the per-request - ``LiteLLMLoggingObj`` (a typed object) and closed in the async callback. - - When no recordable parent is visible (``pre_call`` was driven from a thread - pool for a sync-only provider, where contextvars — and so the anchor — - don't follow), creation is deferred: only the start time is recorded, and - the async callback — whose worker context was copied from the request task - and so still carries the anchor — creates the span then. - - Synthetic proxy-gate error logs (auth/rate-limit rejections) also fire this - hook but never made an upstream call; they are tagged and skipped so no - phantom LLM-call span is produced. + Runs synchronously in the request task, where the live server span is the ambient OTel + context, so the span parents to it natively; it is stashed and closed in the async + callback. When no recordable parent is visible (a thread-pool sync-only provider, where + the anchor doesn't follow), creation is deferred to the close callback. Synthetic + proxy-gate logs (auth/rate-limit rejections) made no upstream call and are skipped. """ call = LLMCallEvent.from_dict(kwargs) if call.is_no_upstream_call: @@ -270,10 +243,8 @@ class OpenTelemetryV2(CustomLogger): return start_time_ns = to_ns(datetime.now()) spans: tuple[Span, ...] = () - # Parent to the request's anchored root span (stable across the request), - # falling back to ambient on the SDK path. Open the span live only when - # that resolves to a recordable parent; otherwise defer to the close - # callback (the thread-pool case, where the anchor isn't visible here). + # Parent to the anchored root span (ambient on the SDK path); open live only when that + # resolves to a recordable parent, else defer to the close callback (the thread-pool case). parent_context = resolve_request_span_context() if is_recordable_span(get_current_span(parent_context)): spans = tuple( @@ -287,9 +258,8 @@ class OpenTelemetryV2(CustomLogger): for tracer in self._tenant_tracers.tracers_for(self.tracer, self._destinations_for_backend(call)) ) self._open_llm_calls[call_id] = _LLMCallSpan(spans=spans, start_time_ns=start_time_ns) - # Evict the oldest open call if the map is over budget. A call that opens - # but never closes (a stream that only fires stream events) would linger - # otherwise; the evicted span is simply dropped (never exported). + # Evict the oldest open call if over budget; a call that opens but never closes would + # linger otherwise (the evicted span is dropped, never exported). if len(self._open_llm_calls) > _OPEN_CALLS_MAX: self._open_llm_calls.popitem(last=False) @@ -327,10 +297,9 @@ class OpenTelemetryV2(CustomLogger): self._close_llm_call(kwargs, start_time, end_time) def _seed_identity_baggage(self, identity: RequestIdentity, model: str | None, context: Context) -> Context: - """Seed authenticated request-identity Baggage onto ``context`` so the Baggage - processor stamps team/key/metadata onto the span. Identity is read from the - parsed payload, never the client's ``params._meta`` carrier, so it can't be - spoofed.""" + """Seed authenticated request-identity Baggage onto ``context`` so the Baggage processor + stamps team/key/metadata onto the span. Read from the parsed payload, never the client's + ``params._meta`` carrier, so it can't be spoofed.""" bag = promoted_baggage( identity, model, @@ -348,14 +317,10 @@ class OpenTelemetryV2(CustomLogger): ) -> bool: """Emit an MCP tool-call span when the closed request was a tool call. - MCP tool calls reach the success/failure callbacks like any other request - (with ``call_type`` ``call_mcp_tool``), but they are not LLM calls and have - no ``pre_call`` carrier — so they get their own CLIENT span here. Per the MCP - semconv it parents to the trace context the client propagated in - ``params._meta`` (or starts a new root) and links the transport span, rather - than nesting under the HTTP/session span. Returns whether it handled the - event, so the caller skips the LLM-call path. The whole span is emitted at - once (there is no boundary to open it at), deduped on the call id. + MCP tool calls reach the callbacks with no ``pre_call`` carrier, so they get their own + CLIENT span here, emitted at once and deduped on the call id. Per the MCP semconv it + parents to the ``params._meta`` trace context (or a new root) and links the transport + span. Returns whether it handled the event, so the caller skips the LLM-call path. """ raw_payload = kwargs.get("standard_logging_object") if not raw_payload or not is_mcp_tool_call(cast(Mapping[str, object], raw_payload)): @@ -364,9 +329,7 @@ class OpenTelemetryV2(CustomLogger): data = MCPToolCallSpanData.from_standard_logging_payload( payload, capture_content=self.config.capture_span_content ) - # A stray LLM carrier from a ``pre_call`` that mis-fired for this id would - # otherwise linger until evicted; drop it so it's neither leaked nor closed - # as a phantom LLM span. + # Drop any stray LLM carrier for this id so it's neither leaked nor closed as a phantom span. if data.identity.call_id: self._open_llm_calls.pop(data.identity.call_id, None) parent_context, links = resolve_mcp_span_context() @@ -389,12 +352,10 @@ class OpenTelemetryV2(CustomLogger): ) -> bool: """Emit an MCP ``tools/list`` span when the closed request was a discovery call. - Like a tool call, listing reaches the success/failure callbacks (here with - ``call_type`` ``list_mcp_tools``) with no ``pre_call`` carrier, so it gets its - own CLIENT span. Per the MCP semconv it parents to the ``params._meta`` trace - context (or starts a new root) and links the transport span, rather than - nesting under the HTTP/session span. Returns whether it handled the event so - the caller skips the LLM-call path. + Like a tool call, listing has no ``pre_call`` carrier, so it gets its own CLIENT span. + Per the MCP semconv it parents to the ``params._meta`` trace context (or a new root) and + links the transport span. Returns whether it handled the event so the caller skips the + LLM-call path. """ raw_payload = kwargs.get("standard_logging_object") if not raw_payload or not is_mcp_list_tools(cast(Mapping[str, object], raw_payload)): @@ -425,15 +386,10 @@ class OpenTelemetryV2(CustomLogger): ) -> Span | None: """Finish the LLM-call span opened at ``pre_call`` (or create it deferred). - Missing carrier has two shapes. ``pre_call`` genuinely never ran -- the - request was rejected at the gate or blocked by a pre-call guardrail before - any upstream call, so no payload exists and dropping is correct. OR this - v2 instance was lazily activated AFTER ``pre_call`` iterated the callback - list (the destination-resolver path: a credential resolved a backend the - YAML didn't pre-list), so the upstream call DID happen, the payload IS - set, and the admin-resolved destinations name this backend -- emit a - deferred span with the success event's start time so the per-tenant - exporter ships it. + A missing carrier means either ``pre_call`` never ran (rejected at the gate or by a + pre-call guardrail, no payload, so dropping is correct) or this v2 instance was lazily + activated after ``pre_call`` (the destination-resolver path), where the payload IS set and + names this backend, so a deferred span is emitted with the success event's start time. """ call = LLMCallEvent.from_dict(kwargs) call_id = call.call_id @@ -482,11 +438,7 @@ class OpenTelemetryV2(CustomLogger): ) def _mark_closed(self, call_id: str | None) -> None: - """Remember a call_id has been closed so a duplicate callback no-ops. - - Bounded by the same ceiling as ``_open_llm_calls`` to prevent unbounded - growth; oldest entries are evicted FIFO. - """ + """Remember a call_id has been closed so a duplicate callback no-ops (bounded FIFO).""" if not call_id: return self._closed_call_ids[call_id] = None @@ -503,11 +455,8 @@ class OpenTelemetryV2(CustomLogger): ) -> Span | None: """Emit an LLM-call span outside the ``pre_call`` boundary. - Two callers: the SDK thread-pool path (carrier existed but ``pre_call`` - saw no recordable parent) and the destination-resolver path (this v2 - instance was born after ``pre_call`` ran, so no carrier was ever opened). - Both anchor to the request's root span via the worker-copied context and - seed identity Baggage so the span is labeled consistently. + Two callers: the SDK thread-pool path and the destination-resolver path. Both anchor to + the request's root span via the worker-copied context and seed identity Baggage. """ data = LLMCallSpanData.from_standard_logging_payload( payload, @@ -574,19 +523,13 @@ class OpenTelemetryV2(CustomLogger): error_override: str | None, ) -> Span | None: data = ServiceSpanData.from_payload(payload, event_metadata=event_metadata) - # Decide whether this service call is a span at all, and of what kind. - # ``None`` means metrics-only (framework instrumentation that duplicates a - # gen-AI span — ``self``/``router``/``proxy_pre_call`` — or ``auth``, which - # gets a live phase span instead). Those still feed Prometheus/Datadog via - # their own hooks; they just never enter the trace. + # ``None`` role means metrics-only (framework instrumentation that duplicates a gen-AI + # span, or ``auth`` which gets a live phase span instead); it never enters the trace. role = span_role_for_service(data.service_name) if role is None: return None - # A metrics-only ping with neither timing nor a parent (in-memory queue - # gauges) is not a traceable operation; a span for it would be a - # zero-duration root with no context, so skip it. Real background work - # (budget/reset jobs, spend flush) passes start/end times and still emits - # as a root; anything with a parent emits regardless. + # Skip a ping with neither timing nor a parent (in-memory queue gauges): a span would be a + # zero-duration root. Real background work passes start/end times; anything with a parent emits. if error_override is None and start_time is None and end_time is None and parent_otel_span is None: return None if error_override is not None and data.error is None: @@ -596,11 +539,9 @@ class OpenTelemetryV2(CustomLogger): error=SpanError(message=error_override), event_metadata=data.event_metadata, ) - # Parent like every other span: ambient context first (so identity Baggage - # rides along and the call nests under whatever request phase is active — - # e.g. a DB lookup under the live ``auth`` span), falling back to the - # server span the proxy threaded as ``parent_otel_span``. A background - # service call has neither, so it starts its own root trace. + # Parent ambient-first (so the call nests under the active request phase, e.g. a DB lookup + # under ``auth``), falling back to the threaded ``parent_otel_span``; a background call has + # neither and starts its own root trace. parent_context = resolve_parent_context(threaded=parent_otel_span) return self._emitter.emit( role, @@ -618,13 +559,9 @@ class OpenTelemetryV2(CustomLogger): def seed_request_identity(self, user_api_key_dict: Any, model: Any = None) -> None: """Attach request-identity Baggage to the current context + server span. - Seeding identity into Baggage makes **every** span emitted afterwards for - this request — LLM call, guardrail, DB call — inherit it via - ``LiteLLMBaggageSpanProcessor``. Called once at the auth boundary (as soon - as the key resolves) so post-auth spans are labeled consistently; the - Baggage rides the request task's contextvar from there on. Auth-internal - DB lookups that run before the key is known stay unlabeled — identity - isn't determined yet, which is correct. + Called once at the auth boundary so every span emitted afterwards inherits identity via + ``LiteLLMBaggageSpanProcessor``. Auth-internal DB lookups before the key resolves stay + unlabeled, which is correct. """ try: identity = RequestIdentity.from_user_api_key_auth(user_api_key_dict) @@ -639,19 +576,14 @@ class OpenTelemetryV2(CustomLogger): # Attach (no detach): the contextvar is scoped to this request's # asyncio task and is reclaimed when the task ends. attach(set_request_baggage(bag, context=get_current())) - # The server span was started by the instrumentor before this ran, - # so the Baggage processor (which only fires at span start) won't - # backfill it — stamp identity on it directly. Prefer the anchored - # root span over the ambient one so identity still lands on the - # server span when seeding from inside the live ``auth`` phase span - # (the auth-failure path), where ``get_current_span`` is the phase - # span, not the request's root. + # The instrumentor started the server span before this ran, so the Baggage + # processor (fires only at span start) won't backfill it; stamp identity directly. + # Prefer the anchored root over ambient so identity lands on the server span even + # when seeding from inside the live ``auth`` phase span. server_span = request_root_span() or get_current_span() if is_recordable_span(server_span): - # Re-capture the anchor here too: this runs post-auth with the - # server span active and covers entrypoints that bypass - # ``create_litellm_proxy_request_started_span`` (e.g. the SDK - # path's ``async_pre_call_hook``). Idempotent. + # Re-capture the anchor here too, covering entrypoints that bypass + # ``create_litellm_proxy_request_started_span`` (e.g. the SDK path). Idempotent. set_request_root_span(server_span) for key, value in bag.items(): server_span.set_attribute(key, value) @@ -688,14 +620,12 @@ class OpenTelemetryV2(CustomLogger): exception: "Exception | None", status_code: int, ) -> None: - """Stamp the v2 error.* attributes on the FastAPI-owned SERVER span for a - failure that dies before any LLM-call span exists (malformed body, auth / - validation rejection). Called from the proxy's global exception handler via - ``_close_dangling_otel_server_span``. The instrumentor still owns the span's - status and lifecycle, so this only decorates it — never sets status, never - ends it — and emits no exception event, matching v1's SERVER-span behavior - and avoiding a duplicate of the event ``async_post_call_failure_hook`` or - the ``auth`` phase span already records.""" + """Stamp v2 error.* attributes on the FastAPI-owned SERVER span for a failure that dies + before any LLM-call span exists (malformed body, auth/validation rejection). + + Only decorates the span (never sets status, ends it, or emits an exception event); the + instrumentor still owns the span's lifecycle, matching v1's SERVER-span behavior. + """ if span is None or not is_recordable_span(span): return stamp_error( @@ -712,18 +642,13 @@ class OpenTelemetryV2(CustomLogger): user_api_key_dict: "UserAPIKeyAuth", traceback_str: "str | None" = None, ) -> None: - """Stamp error.* on the request's root SERVER span for a proxy-level - failure that never reached an LLM call (empty body rejected in the - endpoint, auth failure), so the failed request carries the same error keys - a failed LLM call does. v1's ``OpenTelemetry`` implemented this same hook; - v2 lost it when it stopped subclassing ``OpenTelemetry``, which is the - LIT-4179 regression for pre-call failures. + """Stamp error.* on the request's root SERVER span for a proxy-level failure that never + reached an LLM call (empty body rejected, auth failure), so it carries the same error keys + a failed LLM call does. - An MCP message is handled on the session's task, where the request-root - anchor is whatever request opened the session, so prefer the transport the - gateway published for this specific message. Without that, a failed tool - call aimed its error at the ``initialize`` request's finished span and the - SDK dropped it, leaving the POST that actually failed unmarked.""" + For an MCP message the session-task anchor is whatever request opened the session, so + prefer the transport the gateway published for this specific message. + """ span = mcp_message_transport_span() or request_root_span() or user_api_key_dict.parent_otel_span if span is None or not is_recordable_span(span): return None @@ -731,18 +656,10 @@ class OpenTelemetryV2(CustomLogger): return None def emit_guardrail_span(self, entry: "StandardLoggingGuardrailInformation") -> None: - # Emitted by the guardrail-recording code the moment a guardrail finishes, - # not from a post-call hook — that hook does not fire on every path (a - # pass-through request that passes its guardrails never reaches it), which - # left passing guardrails without a span. - # - # A guardrail is a sibling of the LLM call under the request's root span, - # so parent it to the explicit anchor — never the active span, which during - # a pre_call guardrail can be the live ``auth`` phase span. Emit with the - # guardrail's actual execution window so a pre_call guardrail is placed - # before the LLM call rather than at emission time. One entry in, one span - # out — the module-level entry point routes each entry to this single - # registered logger so a guardrail is never emitted more than once. + # A guardrail is a sibling of the LLM call under the request's root span, so parent it to + # the explicit anchor, never the active span (which during a pre_call guardrail can be the + # live ``auth`` phase span). Emit with the guardrail's actual execution window so a + # pre_call guardrail is placed before the LLM call rather than at emission time. data = GuardrailSpanData.from_logging_entry(entry) self._emitter.emit( SpanRole.GUARDRAIL, @@ -768,21 +685,10 @@ def select_global_otel_v2_logger( ) -> "OpenTelemetryV2": """The single ``OpenTelemetryV2`` whose provider should become the OTel global. - The callback factory designates one logger as canonical the moment it builds - the first one (``_init_otel_logger_on_litellm_proxy`` sets - ``proxy_server.open_telemetry_logger``), and every other v2 entry point — - guardrail, identity seeding, phase spans — already routes through that same - ``registered`` owner. Reuse it here too so the global provider has one source - of truth instead of a second, independently-derived guess; this is the logger - a preset (arize, langfuse, …) folds the ``OTEL_*`` base exporter and its own - exporter into, so the FastAPI server span and the gen-ai spans share one - provider and one trace. - - Fall back to ``in_memory_loggers`` for the SDK path, where no proxy global is - set (selecting from there, not ``service_callback``, which a preset logger does - not always reach), and build a generic logger from ``OTEL_*`` only when none was - configured at all. Each fallback still avoids the second generic logger that - orphaned the gen-ai spans onto a different backend than the server span. + Prefer the canonical ``registered`` owner every other v2 entry point routes through, so the + server span and gen-ai spans share one provider and one trace. Fall back to a v2 logger in + ``in_memory_loggers`` (the SDK path), then build a generic one from ``OTEL_*`` only when none + was configured; each fallback avoids the second generic logger that orphaned the gen-ai spans. """ if registered is not None: return registered @@ -797,15 +703,9 @@ def publish_global_otel_v2_provider( ) -> "OpenTelemetryV2": """Select the single v2 logger and publish its provider as the OTel global. - The proxy calls this once at startup, after callbacks are initialized, so the - preset logger already exists; it passes ``registered`` (the canonical owner the - factory designated as ``proxy_server.open_telemetry_logger``) so the global - provider reuses the same logger the rest of the v2 code emits through (see - :func:`select_global_otel_v2_logger`). Both ``registered`` and - ``set_global_provider`` (the proxy passes - ``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. + Called once at startup; ``registered`` (the canonical owner) makes the global reuse the + logger the rest of the v2 code emits through, and both it and ``set_global_provider`` are + injected so the publish step is unit-testable without touching real global OTel state. """ logger = select_global_otel_v2_logger(in_memory_loggers, registered=registered) set_global_provider(logger._tracer_provider) @@ -824,13 +724,9 @@ def _registered_v2_logger() -> "OpenTelemetryV2 | None": def emit_guardrail_span(entry: "StandardLoggingGuardrailInformation") -> None: """Emit a guardrail span on the registered v2 OTel logger. - Called by the guardrail-recording code the moment a guardrail finishes, so a - span is produced regardless of whether a post-call hook later runs (it does - not on the pass-through allow path). Routes through the single canonical - logger — the same one every other v2 entry point uses — so a guardrail - recorded once yields exactly one span; fanning out across every reachable - ``OpenTelemetryV2`` instance double-emits the same entry. Best-effort: span - emission must never break guardrail evaluation. + Called when a guardrail finishes, so a span is produced even when no post-call hook runs (the + pass-through allow path). Routes through the single canonical logger so a guardrail yields + exactly one span. Best-effort: emission must never break guardrail evaluation. """ logger = _registered_v2_logger() if logger is None: diff --git a/litellm/integrations/otel/model/metadata.py b/litellm/integrations/otel/model/metadata.py index 4db37996301..4c34acf91a5 100644 --- a/litellm/integrations/otel/model/metadata.py +++ b/litellm/integrations/otel/model/metadata.py @@ -1,37 +1,18 @@ """The single translation layer between a request's metadata and the spans. -Every relevant field litellm exposes about a request — the user-facing model, -the model actually dispatched to the provider, the deployment, and the caller's -identity (team, key, end-user) — is parsed **once**, here, out of the -``StandardLoggingPayload`` (or a ``UserAPIKeyAuth`` at the auth boundary). Span -data, baggage promotion, and the mappers then read these typed fields instead of -each digging into the raw ``metadata`` / ``hidden_params`` dicts. +Every field is parsed **once**, here, out of the ``StandardLoggingPayload`` (or a +``UserAPIKeyAuth`` at the auth boundary), so span data, baggage, and the mappers +read typed fields instead of the raw ``metadata`` / ``hidden_params`` dicts. -Two models live here because a request's identity is known *before* its model -resolution is: +:class:`RequestIdentity` holds caller identity (team/key/end-user), seeded into +Baggage at the auth boundary before routing has picked a deployment; its +``provider_model`` is absent from that early seed and filled only from the payload +at close. :class:`RequestContext` is the fully-resolved view at close, wrapping it. -* :class:`RequestIdentity` — team / key / end-user, seeded into Baggage at the - auth boundary (``from_user_api_key_auth``), before routing has picked a - deployment. ``provider_model`` is therefore absent from that early seed and is - only filled in from the payload once the call closes. -* :class:`RequestContext` — the full picture available at close: the resolved - request vs. provider model split, plus the response model, model group, model - id, and api base, wrapping the :class:`RequestIdentity`. - -The request-vs-provider model split is the subtle part. On the proxy a caller -asks for a *model group* (e.g. ``gpt-4o``) that routes to a concrete deployment -(e.g. ``azure/my-deployment``); the two are distinct and both worth recording. -``StandardLoggingPayload`` exposes them as: - -* ``model_group`` — the user-facing name the caller requested. -* ``model`` — already reconstructed (see ``reconstruct_model_name``) to the name - litellm dispatched to the provider (the deployment, provider-prefixed). -* ``hidden_params.litellm_model_name`` — a secondary source for the dispatched - model (populated only on some call paths, e.g. files). - -So ``gen_ai.request.model`` is the *group* (falling back to the call model on the -SDK path, which has no group), and ``litellm.provider.model`` is the *dispatched* -model. They coincide on the SDK path, which is correct. +The request-vs-provider model split: a caller asks for a *model group* (``gpt-4o``) +that routes to a concrete deployment (``azure/my-deployment``). ``gen_ai.request.model`` +records the group and ``litellm.provider.model`` the dispatched model; they coincide +on the SDK path, which has no group. """ from __future__ import annotations @@ -55,33 +36,22 @@ class RequestIdentity: call_id: str | None = None team_id: str | None = None team_alias: str | None = None - # The team's free-form metadata, carried raw (empty/missing -> None) and - # filtered to an operator allowlist only at Baggage-promotion time, so an - # unconfigured deployment never promotes any of it. + # The team's free-form metadata, carried raw; filtered to an operator allowlist at Baggage-promotion time. team_metadata: Mapping[str, Any] | None = None key_hash: str | None = None end_user: str | None = None - # The model litellm dispatched to the provider. Only known once the call - # completes (routing has picked a deployment), so it's absent from the - # auth-time seed and filled only from the payload. + # The model litellm dispatched to the provider; known only at close, so absent from the auth-time seed. provider_model: str | None = None metadata: Mapping[str, str] = field(default_factory=dict) @classmethod def from_payload(cls, payload: "StandardLoggingPayload") -> "RequestIdentity": - """Parse caller identity out of a closed request's payload metadata. - - ``provider_model`` is resolved here too (see :func:`resolve_provider_model`) - so the identity carried into Baggage labels every span with the dispatched - model, not just the user-facing one. - """ + """Parse caller identity (incl. resolved ``provider_model``) from a closed request's payload.""" raw_meta = cast(Mapping[str, object], payload.get("metadata") or {}) metadata = {key: str(value) for key, value in raw_meta.items() if isinstance(value, (str, bool, int, float))} return cls( call_id=as_str(payload.get("litellm_call_id")) or as_str(payload.get("id")), - # StandardLoggingMetadata's canonical key is ``user_api_key_team_id``; - # the bare ``team_id`` is a legacy alias and is often empty, so prefer - # the canonical key and fall back to the alias. + # Prefer the canonical ``user_api_key_team_id``; the bare ``team_id`` is a legacy alias. team_id=as_str(raw_meta.get("user_api_key_team_id")) or as_str(raw_meta.get("team_id")), team_alias=as_str(raw_meta.get("user_api_key_team_alias")) or as_str(raw_meta.get("team_alias")), team_metadata=_team_metadata_dict(raw_meta.get("user_api_key_team_metadata")), @@ -93,14 +63,10 @@ class RequestIdentity: @classmethod def from_user_api_key_auth(cls, auth: object) -> "RequestIdentity": - """Identity from a ``UserAPIKeyAuth`` (duck-typed to keep this module - free of a proxy import). + """Identity from a ``UserAPIKeyAuth`` (duck-typed to avoid a proxy import). - Used in the pre-call hook to seed Baggage early — before any LLM, - guardrail, or service span is created — so the whole request's spans - inherit identity, not just the LLM-call span. Metadata sub-keys use the - ``user_api_key_*`` names that ``baggage.DEFAULT_BAGGAGE_METADATA_KEYS`` - promotes. + Seeds Baggage at the pre-call hook so every span inherits identity. Metadata + sub-keys use the ``user_api_key_*`` names Baggage promotion expects. """ get = lambda name: getattr(auth, name, None) # noqa: E731 metadata = { @@ -119,20 +85,14 @@ class RequestIdentity: team_metadata=_team_metadata_dict(get("team_metadata")), key_hash=as_str(get("api_key")), end_user=as_str(get("end_user_id")), - # ``provider_model`` is unknown at the auth boundary — routing hasn't - # picked a deployment yet — so it's only populated from the payload. + # ``provider_model`` is unknown at the auth boundary (routing hasn't picked a deployment yet). metadata=metadata, ) @dataclass(frozen=True) class RequestContext: - """The fully-resolved view of a closed request, parsed once from the payload. - - ``request_model`` is the user-facing requested model and ``provider_model`` - (on :attr:`identity`) is the model litellm dispatched to the provider; the two - differ on the proxy (group vs. deployment) and coincide on the SDK path. - """ + """The fully-resolved view of a closed request, parsed once from the payload.""" request_model: str response_model: str | None @@ -167,40 +127,23 @@ class RequestContext: # --- live-callback kwargs parsing ------------------------------------------- # -# -# The model and helpers below parse the *live* callback ``kwargs`` god object (and -# the raw pre/post-call ``data`` dicts) — the untyped request state that reaches a -# ``CustomLogger`` before, or instead of, a ``StandardLoggingPayload``. They live -# here, with the payload/auth parsers, so every read out of a request's raw dicts -# is in one place rather than scattered across the ``CustomLogger``. +# Parse the live callback ``kwargs`` god object and raw pre/post-call ``data`` dicts — +# the untyped request state reaching a ``CustomLogger`` before a ``StandardLoggingPayload``. @dataclass(frozen=True) class LLMCallEvent: - """The typed view of the live callback ``kwargs`` (``model_call_details``). + """The typed view of the live callback ``kwargs`` (``model_call_details``), parsed once.""" - litellm hands every callback an untyped ``kwargs`` god object. The fields the - OTel logger needs out of it are parsed **once**, here, so the ``CustomLogger`` - reads typed attributes instead of digging into the dict at each boundary. - """ - - # The ``litellm_call_id`` correlating ``pre_call`` with the close callback. - # Present in ``model_call_details`` at ``pre_call`` and in both the kwargs and - # the ``standard_logging_object`` at success/failure, so it's a stable key for - # the open-call carrier — no back-reference to the logging object required (the - # object isn't reachable from the callback kwargs at ``pre_call`` time). + # The ``litellm_call_id`` correlating ``pre_call`` with the close callback; the stable + # key for the open-call carrier. call_id: str | None - # The ``StandardLoggingPayload`` carried on a success/failure callback; ``None`` - # at ``pre_call``, or when the call closed before any payload materialized (so - # there is nothing to stamp on the span). + # The success/failure payload; ``None`` at ``pre_call`` or if the call closed with no payload. payload: "StandardLoggingPayload | None" otel_destinations: tuple[OtelDestination, ...] - # True for synthetic proxy-gate logs (auth / rate-limit rejections): they fire - # the ``pre_call`` hook but never made an upstream call, so they get no span. + # True for synthetic proxy-gate logs (auth/rate-limit rejections): no upstream call, so no span. is_no_upstream_call: bool - # A best-effort ``"{operation} {model}"`` name known at ``pre_call`` time. The - # span is renamed from the typed payload at close (``finish_span``); this only - # needs to be reasonable for a span that never gets closed (a leak). + # Best-effort ``"{operation} {model}"`` name at ``pre_call``; only matters for a leaked span (renamed at close). provisional_span_name: str time_to_first_chunk_seconds: float | None @@ -221,10 +164,8 @@ class LLMCallEvent: def time_to_first_chunk_seconds(kwargs: Mapping[str, Any]) -> float | None: - """Seconds from the upstream request being issued (``api_call_start_time``) - to the first streamed chunk (``completion_start_time``); ``None`` for - non-streaming calls, where ``completion_start_time`` is backfilled with the - end time and would not measure first-chunk latency.""" + """Seconds from upstream request (``api_call_start_time``) to first streamed chunk + (``completion_start_time``); ``None`` for non-streaming calls.""" optional_params = cast(Mapping[str, Any], kwargs.get("optional_params") or {}) if not optional_params.get("stream"): return None @@ -245,11 +186,7 @@ def _call_id(payload: "StandardLoggingPayload | None", kwargs: Mapping[str, Any] def model_from_request_data(data: object) -> str | None: - """The user-facing ``model`` from a pre-call ``data`` dict (``None`` if absent). - - Read at the auth boundary to label early Baggage before routing has resolved - a deployment; ``data`` is duck-typed since it arrives untyped from the proxy. - """ + """The user-facing ``model`` from a pre-call ``data`` dict (``None`` if absent).""" if isinstance(data, Mapping): return as_str(data.get("model")) return None @@ -258,16 +195,13 @@ def model_from_request_data(data: object) -> str | None: def resolve_provider_model(payload: "StandardLoggingPayload") -> str | None: """The model litellm dispatched to the provider, from the payload. - Prefers the explicit ``hidden_params.litellm_model_name`` (set on call paths - that know it, e.g. files), then the top-level ``model`` — which - ``reconstruct_model_name`` has already resolved to the deployment's - provider-prefixed name. Returns ``None`` only when neither is present. + Prefers ``metadata.deployment``, then ``hidden_params.litellm_model_name``, then the + top-level ``model`` (already resolved to the provider-prefixed deployment name). """ raw_meta = cast(Mapping[str, object], payload.get("metadata") or {}) hidden = cast(Mapping[str, object], payload.get("hidden_params") or {}) return ( - # ``deployment`` survives only on paths that don't strip it from metadata; - # harmless (and most precise) to prefer it when present. + # ``deployment`` (most precise) survives only on paths that don't strip it from metadata. as_str(raw_meta.get("deployment")) or as_str(hidden.get("litellm_model_name")) or as_str(payload.get("model")) ) @@ -280,13 +214,7 @@ def _model_info_id(model_info: object) -> str | None: def _team_metadata_dict(value: object) -> Mapping[str, Any] | None: - """The team's free-form metadata as a raw mapping, or ``None`` when missing - or empty. - - Carried raw on the identity and filtered to an operator allowlist only at - Baggage-promotion time (see ``baggage.promoted_baggage``), so an empty case - is dropped rather than carrying a useless ``{}``. - """ + """The team's free-form metadata as a raw mapping, or ``None`` when missing or empty.""" if isinstance(value, Mapping) and value: return dict(value) return None diff --git a/litellm/integrations/otel/plumbing/context.py b/litellm/integrations/otel/plumbing/context.py index d221503b4bb..8a50312c7d4 100644 --- a/litellm/integrations/otel/plumbing/context.py +++ b/litellm/integrations/otel/plumbing/context.py @@ -22,21 +22,11 @@ if TYPE_CHECKING: _PROPAGATOR = TraceContextTextMapPropagator() -# The request's root span — the FastAPI-owned SERVER span — captured ONCE when the -# proxy first resolves it, so request-level spans (the LLM call, guardrails) can -# parent to it EXPLICITLY instead of to whatever span happens to be active at the -# instant they are emitted. Ambient-only parenting (``get_current_span()``) is -# wrong at two boundaries: -# * inside the ``auth`` phase span the active span is the auth span, so an LLM / -# guardrail span emitted there would nest under auth instead of being its -# sibling; and -# * in a detached success task (pass-through logs success from a fire-and-forget -# ``asyncio.create_task``) the server span may not be active at all, orphaning -# the span into a brand-new trace. -# A ``ContextVar`` (not a request attribute) so it rides the request task's context -# 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. +# The request's root span (the FastAPI-owned SERVER span), captured once so request-level +# spans (LLM call, guardrails) parent to it explicitly rather than to whatever span is active — +# which under the ``auth`` phase span or a detached success task would misnest or orphan them. +# A ``ContextVar`` so it rides the request task and its ``create_task`` children (the async +# logging callbacks that close the span). _request_root_span: "ContextVar[Span | None]" = ContextVar("litellm_otel_request_root_span", default=None) _request_destinations: 'ContextVar[tuple["OtelDestination", ...]]' = ContextVar( @@ -59,9 +49,7 @@ def request_destinations() -> 'tuple["OtelDestination", ...]': def set_request_root_span(span: Span) -> None: """Anchor the request's root (server) span for explicit child parenting. - No-ops for a non-recordable span so a bad capture can never replace a good one - with a phantom parent. Idempotent — the proxy captures the same server span at - more than one entry point. + No-ops for a non-recordable span (a bad capture can't replace a good one); idempotent. """ if is_recordable_span(span): _request_root_span.set(span) @@ -73,11 +61,9 @@ def request_root_span() -> "Span | None": return span if is_recordable_span(span) else None -# The W3C trace-context carrier (``traceparent``/``tracestate``/``baggage``) the -# MCP client propagated in the current request's ``params._meta``. The MCP gateway -# sets it per message so the MCP span can parent to the client's span rather than -# to the transport. A ``ContextVar`` because, like the root-span anchor, it must -# ride the request task and be readable by the inline success-logging callback. +# The W3C trace-context carrier (``traceparent``/``tracestate``/``baggage``) the MCP client +# propagated in the current request's ``params._meta``; set per message. A ``ContextVar`` so it +# rides the request task and is readable by the inline success-logging callback. _mcp_message_trace_carrier: "ContextVar[Mapping[str, str] | None]" = ContextVar( "litellm_otel_mcp_message_trace_carrier", default=None ) @@ -99,17 +85,10 @@ def reset_mcp_message_trace_carrier(token: "Token[Mapping[str, str] | None]") -> # The transport span of the HTTP request carrying the CURRENT MCP message. -# -# ``_request_root_span`` above cannot be used for MCP: a *stateful* streamable-HTTP -# session runs every message on the single task spawned by that session's -# ``initialize`` POST, so the ContextVar the ASGI request task writes at auth time -# is frozen at ``initialize`` there and never sees the later ``tools/call`` POSTs. -# Reading it from the message handler parents every tool call in the session to the -# first request's server span and aims that call's ``error.*`` at it — a span that -# ended long ago, so the SDK drops the write and the failure reaches no request at -# all. The gateway instead resolves the current message's transport span on the -# request task and hands it over the same way it hands over per-request auth, and -# the handler publishes it here for the span emitter and the failure hook. +# ``_request_root_span`` can't serve MCP: a stateful streamable-HTTP session runs every message +# on the task of its ``initialize`` POST, so that anchor is frozen at ``initialize`` and never +# sees later ``tools/call`` POSTs. The gateway instead resolves each message's transport span on +# the request task and the handler publishes it here. _mcp_message_transport_span: "ContextVar[Span | None]" = ContextVar( "litellm_otel_mcp_message_transport_span", default=None ) @@ -118,21 +97,10 @@ _mcp_message_transport_span: "ContextVar[Span | None]" = ContextVar( def set_mcp_message_transport_span(span: object) -> "Token[Span | None]": """Publish the transport span of the request carrying the current MCP message. - Also re-anchors the request root, so everything else the message emits or stamps - — the identity attributes seeded onto the server span, a guardrail span, a - proxy-level failure — lands on this request instead of on the one that opened - the session. The MCP SDK dispatches each message on its own task, so the anchor - is scoped to this message; the handler re-publishes it for the next one either - way. Only a transport still open for writes is anchored: replacing the anchor - with a request that already answered would just move the dropped writes from one - finished span to another. - - Takes ``object`` because the gateway reads it back out of the ASGI scope, whose - values are untyped; anything that is not a usable span is stored as ``None`` - rather than trusted. - - Returns the reset token; the caller must reset it once the message is handled - so the transport never leaks to the next message on the same session task. + Also re-anchors the request root so everything the message emits lands on this request, not + the one that opened the session; only a transport still open for writes is anchored. Takes + ``object`` (untyped ASGI-scope value); a non-span is stored as ``None``. Returns the reset + token, which the caller must reset once the message is handled to avoid leaking it. """ transport = span if isinstance(span, Span) and is_recordable_span(span) else None if transport is not None and transport.is_recording(): @@ -145,15 +113,11 @@ def reset_mcp_message_transport_span(token: "Token[Span | None]") -> None: def mcp_message_transport_span() -> "Span | None": - """The published transport span, only while it is still open for writes. + """The published transport span, only while it is still recording (open for writes). - Recording — not merely valid — is the bar here because this span is the target - of ``error.*`` stamping from another task, and the publisher's validity check - cannot speak for a span that has since ended. A finished span keeps a valid - context forever, so it would otherwise be handed back for a write the SDK then - refuses. The POST carrying a ``tools/call`` stays open until the result is - written, so it is recording for the life of the call; a notification POST can - answer first, and this returns ``None`` for it rather than writing into the void. + Recording — not merely valid — is the bar because this span is the target of ``error.*`` + stamping from another task, and a finished span keeps a valid context forever but would + refuse the write. Returns ``None`` for a transport that has already answered. """ span = _mcp_message_transport_span.get() if span is None or not span.is_recording(): @@ -164,11 +128,9 @@ def mcp_message_transport_span() -> "Span | None": def _mcp_transport_span_context() -> "SpanContext | None": """The transport span an MCP message span should attach to. - Prefers the transport the gateway published for this specific message; falls - back to the ambient request anchor for paths that emit an MCP span on the - request task itself (the REST MCP endpoints, the SDK). Parenting and linking - only need the immutable context, and unlike ``mcp_message_transport_span`` they - stay correct against a transport that has already finished, so this does not + Prefers the transport the gateway published for this message; falls back to the ambient + request anchor for paths that emit an MCP span on the request task (REST MCP endpoints, SDK). + Only the immutable context is needed, so unlike ``mcp_message_transport_span`` this does not require the span to still be recording. """ published = _mcp_message_transport_span.get() @@ -199,16 +161,10 @@ def context_from_span(span: Span, context: Context | None = None) -> Context: def resolve_parent_context(threaded: Span | None = None) -> Context: """The context a child span should parent under. - Ambient-first: parent to the active OTel context (the server span, restored - by the logging worker or active in the request task), falling back to a span - passed explicitly (``threaded``) only when the ambient context has no - recordable span — e.g. a background service call with no request on the - stack. When neither is recordable the ambient context is returned unchanged, - so the span starts a new root trace. - - Only service/DB spans pass ``threaded`` (the ``parent_otel_span`` handed to - the service hook). Request-level spans — the LLM call and guardrails — are - created where the server span is genuinely ambient, so they never need it. + Ambient-first: parent to the active OTel context, falling back to an explicitly passed + ``threaded`` span only when the ambient context has no recordable span (a background service + call with no request on the stack); when neither is recordable the ambient context is returned + unchanged, so the span starts a new root trace. Only service/DB spans pass ``threaded``. """ ctx = get_current() if is_recordable_span(threaded) and not is_recordable_span(get_current_span(ctx)): @@ -219,15 +175,10 @@ def resolve_parent_context(threaded: Span | None = None) -> Context: def resolve_request_span_context() -> Context: """The parent context for a request-level span (the LLM call, a guardrail). - These are direct children of the request's root server span — siblings of the - ``auth`` phase span and of each other, never nested under whatever span is - momentarily active. So prefer the explicitly anchored root span; fall back to - ambient context only when there is no anchor (the SDK / no-proxy path), where - the span legitimately starts its own root trace. - - Unlike :func:`resolve_parent_context` (used by DB/service spans, which DO want - to nest under the active phase span, e.g. an auth DB lookup under ``auth``), - this never returns the active span when an anchor exists. + These are direct children of the request's root server span, never nested under whatever span + is momentarily active, so prefer the anchored root span; fall back to ambient only on the + SDK/no-proxy path with no anchor. Unlike :func:`resolve_parent_context`, this never returns + the active span when an anchor exists. """ root = request_root_span() if root is not None: @@ -240,33 +191,17 @@ def resolve_mcp_span_context( ) -> "tuple[Context, tuple[Link, ...]]": """Parent context + links for an MCP message span. - When the client propagates W3C trace context in the request's ``params._meta`` - (SEP-414), MCP and the underlying transport are independent lifecycles — one - streamable-HTTP session multiplexes many messages, and the client's own span is - the truthful parent. So, per the OTel GenAI MCP semconv: + Per the OTel GenAI MCP semconv: when the client propagates W3C trace context in + ``params._meta`` (SEP-414), parent to that remote context and link the transport span. + Almost no client implements SEP-414, so with no remote parent, parent to this message's + transport span (from :func:`_mcp_transport_span_context`, the current message's POST, so a + long-lived session doesn't glue every message under its first request) and add no link. With + neither, the span starts its own root trace. - * parent to the trace context the client propagated (a *remote* parent), and - * record the transport span as a *link*, never the parent. - - Almost no client implements SEP-414 yet, so in practice nothing is propagated. - Rooting the span there splits a single tool call into two disconnected traces - joined only by a link, which is how it surfaces in APM: the ``POST`` transaction - and the ``tools/call`` span share no trace. With no remote parent to honor, - parent to the transport span of the request carrying this message instead, so - the call stays in one trace; no link is added since the transport is now the - real parent. The transport comes from :func:`_mcp_transport_span_context`, which - is the *current message's* POST rather than whatever request happened to open - the session, so a long-lived session does not glue every message under its - first request. With neither a remote parent nor a transport the returned context - carries no span and the span legitimately starts its own root trace. - - Only trace context (``traceparent``/``tracestate``) is extracted, never the - client's W3C Baggage: ``params._meta`` is caller-controlled, and the otel - baggage processor stamps allowlisted baggage keys (``litellm.team.id``, - ``litellm.metadata.*``, ...) onto the span as attributes, so honoring remote - baggage would let a client spoof a span's identity attribution. The base context - for extraction is explicitly empty so an absent or malformed ``traceparent`` can - never fall through to the ambient (stale session) span. + Only ``traceparent``/``tracestate`` is extracted, never the client's Baggage: ``params._meta`` + is caller-controlled and honoring remote baggage would let a client spoof identity attribution. + The extraction base context is empty so a malformed ``traceparent`` can't fall through to the + ambient (stale session) span. """ source = carrier if carrier is not None else _mcp_message_trace_carrier.get() parent = _PROPAGATOR.extract(dict(source or {}), context=Context()) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index d32263c4afa..6bb067ed170 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -58,12 +58,9 @@ def to_otel_span_kind(kind: LiteLLMSpanKind) -> SpanKind: return _SPAN_KIND_BY_ROLE_KIND[kind] -# Custom exporter factories keyed by ``ExporterSpec.kind``. A preset registers -# one here when its destination needs construction logic the built-in kinds -# can't express — e.g. an exporter that fetches an auth token lazily on its -# first export (off the event loop) instead of blocking at config-build time. -# Keeping the registry here lets this module stay vendor-agnostic: the factory -# lives with the integration that needs it. +# Custom exporter factories keyed by ``ExporterSpec.kind``. A preset registers one when its +# destination needs construction logic the built-in kinds can't express (e.g. an exporter that +# fetches an auth token lazily on first export). Keeps this module vendor-agnostic. _EXPORTER_FACTORIES: dict[str, Callable[[ExporterSpec], SpanExporter]] = {} @@ -104,11 +101,8 @@ class LiteLLMBaggageSpanProcessor(SpanProcessor): def _otlp_traces_endpoint(endpoint: str | None) -> str | None: """Point an OTLP/HTTP base endpoint at the ``/v1/traces`` signal path. - ``OTEL_EXPORTER_OTLP_ENDPOINT`` is a base URL (e.g. ``http://host:4318``). - The OTLP/HTTP exporter only appends the ``/v1/traces`` path when it reads - that env var itself; when an endpoint is passed explicitly it is used - verbatim, so a base URL would POST to the root and the collector returns - 404. Append the signal path here (leaving an already-correct path intact). + An explicitly passed endpoint is used verbatim (unlike ``OTEL_EXPORTER_OTLP_ENDPOINT``), so a + base URL would POST to the root and 404; append the signal path here, leaving a correct path intact. """ if not endpoint: return endpoint @@ -130,20 +124,11 @@ def default_otlp_kind_for_backend(callback_name: "str | None") -> str: def destination_resource_attrs(destination: "OtelDestination") -> Mapping[str, str]: - """The backend-required Resource attributes a destination carries on every span. + """The destination's builder-declared Resource attributes (e.g. Arize's + ``model_id`` / ``arize.project.name``; empty for header-routed backends). - Backend-agnostic: each backend's destination builder (``presets.destinations``) - declares whatever Resource attributes its ingestion needs, and this just reads - them. Backends that route by auth header (langfuse, weave, generic OTLP) declare - none; Arize declares ``model_id`` / ``arize.project.name`` because it selects the - project from the Resource, not a header. New backends needing Resource-level - routing only have to populate ``resource_attributes`` in their builder. - - Shared by the two export paths that reach a per-tenant destination -- the - ``TenantFanOutSpanProcessor`` (proxy-internal spans) and the ``TenantTracerCache`` - clone provider (the gen-AI span) -- so the gen-AI span and its parents always - carry the SAME Resource and a backend like Arize renders one connected trace - instead of an orphaned subtree. + Both export paths to a destination -- the fan-out processor and the per-tenant + clone provider -- read these so the gen-AI span and its parents share one Resource. """ return dict(destination.resource_attributes) @@ -198,10 +183,8 @@ def build_span_exporter(config: OpenTelemetryV2Config) -> SpanExporter: def _otlp_metrics_endpoint(endpoint: str | None) -> str | None: """Point an OTLP/HTTP base endpoint at the ``/v1/metrics`` signal path. - The OTLP/HTTP exporter only appends ``/v1/metrics`` when it reads - ``OTEL_EXPORTER_OTLP_ENDPOINT`` itself; an explicitly passed endpoint is used - verbatim, so a base URL would POST to the root. Mirror ``_otlp_traces_endpoint`` - for the metrics signal (rewriting a sibling signal path when present). + Mirrors ``_otlp_traces_endpoint`` for the metrics signal (an explicitly passed endpoint is + used verbatim, so a base URL would POST to the root). """ if not endpoint: return endpoint @@ -217,9 +200,8 @@ def _otlp_metrics_endpoint(endpoint: str | None) -> str | None: def build_metric_reader(config: OpenTelemetryV2Config) -> "MetricReader": """Build a metric reader mirroring v1's exporter selection. - ``console`` (and any unrecognized kind) exports to the console; ``otlp_http`` - and ``otlp_grpc`` export over OTLP with the configured endpoint/headers. The - reader exports on a 5s period, matching v1. + ``console`` (and any unrecognized kind) exports to the console; ``otlp_http``/``otlp_grpc`` + export over OTLP with the configured endpoint/headers, on a 5s period. """ from opentelemetry.sdk.metrics.export import ( ConsoleMetricExporter, @@ -267,10 +249,8 @@ def build_metric_reader(config: OpenTelemetryV2Config) -> "MetricReader": def _otlp_logs_endpoint(endpoint: str | None) -> str | None: """Point an OTLP/HTTP base endpoint at the ``/v1/logs`` signal path. - The OTLP/HTTP exporter only appends ``/v1/logs`` when it reads - ``OTEL_EXPORTER_OTLP_ENDPOINT`` itself; an explicitly passed endpoint is used - verbatim, so a base URL would POST to the root. Mirror ``_otlp_traces_endpoint`` - for the logs signal (rewriting a sibling signal path when present). + Mirrors ``_otlp_traces_endpoint`` for the logs signal (an explicitly passed endpoint is used + verbatim, so a base URL would POST to the root). """ if not endpoint: return endpoint @@ -286,10 +266,8 @@ def _otlp_logs_endpoint(endpoint: str | None) -> str | None: def build_log_exporter(config: OpenTelemetryV2Config) -> LogExporter: """Build a log exporter mirroring the exporter selection of the other signals. - ``console`` (and any unrecognized kind) exports to the console; ``otlp_http`` - and ``otlp_grpc`` export over OTLP with the configured endpoint/headers; - ``in_memory`` buffers for tests. Like GenAI metrics, events ride the - single-destination shorthand fields, not the multi-exporter ``exporters`` list. + ``console`` (and any unrecognized kind) to the console; ``otlp_http``/``otlp_grpc`` over OTLP; + ``in_memory`` buffers for tests. Events ride the single-destination shorthand fields, not ``exporters``. """ kind = (config.exporter or "console").lower() if kind in ("in_memory", "inmemory", "memory"): @@ -324,11 +302,8 @@ def build_logger_provider( ) -> SDKLoggerProvider: """Build the :class:`LoggerProvider` GenAI events export through. - ``log_exporter`` is an explicit override (tests inject an - ``InMemoryLogExporter``); otherwise the exporter is selected from the config's - exporter kind via :func:`build_log_exporter`. Console and in-memory exporters - get a Simple processor (synchronous export, which tests rely on), everything - else a Batch processor — the same split as span processing. + ``log_exporter`` is an explicit override (tests); otherwise selected from the config via + :func:`build_log_exporter`. Console/in-memory get a Simple processor, everything else Batch. """ exporter = log_exporter if log_exporter is not None else build_log_exporter(config) provider = SDKLoggerProvider(resource=build_resource(config)) @@ -343,14 +318,12 @@ def resolve_logger_provider( config: OpenTelemetryV2Config, logger_provider: SDKLoggerProvider | None = None, ) -> SDKLoggerProvider | None: - """Resolve the :class:`LoggerProvider` GenAI events record through, or ``None`` - when the operator has opted out of the logs signal. + """Resolve the :class:`LoggerProvider` GenAI events record through, or ``None`` when the + operator opted out of the logs signal. - Same resolution order as :func:`resolve_meter_provider`: an injected provider - wins (DI/tests); an operator-configured SDK global is reused so events ride - their pipeline; an explicit ``NoOpLoggerProvider`` global is an opt-out and - yields ``None``, so no event is ever built. Only the default placeholder - global makes V2 build a provider from the config and publish it as the global. + An injected provider wins (DI/tests); an operator-configured SDK global is reused; an explicit + ``NoOpLoggerProvider`` is an opt-out (``None``). Only the default placeholder global makes V2 + build and publish one from the config. """ if logger_provider is not None: return logger_provider @@ -390,14 +363,10 @@ def resolve_meter_provider( ) -> MeterProvider: """Resolve the :class:`MeterProvider` GenAI metrics record through. - An injected provider wins (DI/tests). Otherwise reuse whatever the operator has - configured as the global, whether a real SDK provider or an explicit - ``NoOpMeterProvider``, so the GenAI histograms ride the operator's - readers/exporters and an explicit opt-out is honored. Only when the global is - still the default proxy placeholder does V2 build one from the config and - publish it as the global, mirroring how V2 owns trace export. The built - provider is the one returned, so its reader thread is always live, never - orphaned. + An injected provider wins (DI/tests); otherwise reuse the operator's configured global (a real + SDK provider or an explicit ``NoOpMeterProvider`` opt-out). Only the default proxy placeholder + makes V2 build and publish one from the config; the built provider is returned so its reader + thread stays live. """ if meter_provider is not None: return meter_provider @@ -431,24 +400,13 @@ def build_tracer_provider( tenant_fan_out_owner: str | None = None, attach_tenant_fan_out: bool = False, ) -> TracerProvider: - """Build the shared :class:`TracerProvider`. + """Build the shared :class:`TracerProvider`: the Baggage processor first, then one + ``SpanProcessor`` per ``config.exporters`` entry (``exporter`` overrides with a test exporter). - Attach the Baggage processor first (so identity attributes land on each - span before any export decision), then add one ``SpanProcessor`` per - ``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). - - ``attach_tenant_fan_out`` — attach a ``TenantFanOutSpanProcessor`` that forwards - each finished proxy-internal span (FastAPI server, ``auth`` phase, DB lookups, the - cost ledger) to the request's admin-resolved destinations. The MAIN v2 logger - provider always opts in, EVEN when no backend is named (the generic global logger - published for a destination-only deployment) -- otherwise the server span never - reaches the destination and its gen-AI child is orphaned. ``tenant_fan_out_owner`` - is the owning backend name when one exists; it is informational (the fan-out skips - the gen-AI span by attribute and forwards internal spans to every destination). - The per-tenant clone providers (built by ``TenantTracerCache``) pass neither, so - the LLM-call span exported through them is not also fanned out here. + ``attach_tenant_fan_out``/``tenant_fan_out_owner`` add a ``TenantFanOutSpanProcessor`` that + forwards proxy-internal spans to the request's destinations. The main v2 provider always opts + in (so the server span reaches the destination and its gen-AI child isn't orphaned); per-tenant + clone providers pass neither. """ provider = TracerProvider(resource=build_resource(config)) if baggage_processor is None: diff --git a/litellm/integrations/otel/plumbing/routing.py b/litellm/integrations/otel/plumbing/routing.py index 423dc17ff2c..29b6454994f 100644 --- a/litellm/integrations/otel/plumbing/routing.py +++ b/litellm/integrations/otel/plumbing/routing.py @@ -1,26 +1,16 @@ """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. +``TenantTracerCache`` routes the gen-AI LLM-call span, building per-tenant clone +``TracerProvider``s that export to the request's admin-owned destinations plus the +configured/global exporter. ``TenantFanOutSpanProcessor`` (at the bottom) forwards the +proxy-internal spans (server, auth, DB, cost) to every destination. Both read the request's +destinations from the same server-only contextvar, so a caller can neither redirect a trace +nor spawn providers. -A call's identity chain is assigned a set of admin-owned OTEL destinations -(``LLMCallEvent.otel_destinations``, resolved server-side from named credentials). -Its spans must export to ALL of them plus the configured/global exporter, so -``TenantTracerCache`` builds and caches ``TracerProvider``s that append one -``SpanProcessor`` per destination. The gen-AI span path (``tracers_for``) groups -destinations by their backend-required Resource attributes and builds one provider -per group, because a span carries exactly one Resource (a provider property) and a -backend like Arize selects its project FROM the Resource -- so two Arize projects -each get a correctly-tagged span instead of one last-wins merge. Header-routed -backends declare no Resource attributes, so their destinations stay in one group with -multiple exporters and route by per-exporter auth. With no destinations it hands back -the logger's default tracer (global only). Destinations are never request-derived, so -a caller can neither redirect a trace nor spawn providers. +The gen-AI path (``tracers_for``) groups destinations by their backend-required Resource +attributes and builds one provider per group, because a span carries exactly one Resource and +a backend like Arize selects its project FROM it, so two Arize projects each get a correctly +tagged span instead of a last-wins merge. Empty destinations -> the logger's default (global only). """ from collections import OrderedDict @@ -61,29 +51,18 @@ class TenantTracerCache: ) # mutable-ok: bounded LRU tracer-provider cache def _evict_if_full(self) -> None: - """Drop the least-recently-used provider when over capacity, without a - synchronous ``shutdown``. Mirrors ``TenantFanOutSpanProcessor``: the - evicted provider's worker drains on its own and is reclaimed at process - exit, and the cache stays bounded.""" + """Drop the least-recently-used provider when over capacity (no synchronous + ``shutdown``; the evicted worker drains on its own and is reclaimed at process exit).""" if len(self._providers) > _MAX_CACHED_PROVIDERS: self._providers.popitem(last=False) def tracers_for(self, default: Tracer, destinations: "tuple[OtelDestination, ...]") -> "tuple[Tracer, ...]": """The tracers for this request's gen-AI span, one per distinct Resource group. - A span carries exactly one Resource (it's a property of the ``TracerProvider``), - but a backend like Arize selects its project FROM the Resource - (``arize.project.name`` / ``model_id``), so two Arize destinations with different - projects need two differently-tagged spans. Group the backend's resolved - destinations by ``destination_resource_attrs`` and return one tracer per group; - the caller emits the span once per tracer (mirroring how the fan-out processor - re-wraps proxy-internal spans per destination). - - Header-routed backends (langfuse, weave) declare no Resource attributes, so all - their destinations collapse into one empty-Resource group with one exporter each - and keep routing by per-exporter auth -- unchanged from the single-group path. - The configured/global exporters ride the FIRST group only, so the global receives - the span once. Empty ``destinations`` -> the logger's default tracer (deny). + A backend like Arize selects its project from the Resource, so destinations are grouped + by ``destination_resource_attrs`` and the caller emits the span once per tracer. The + configured/global exporters ride the FIRST group only, so the global receives the span + once. Empty ``destinations`` -> the logger's default tracer (deny). """ if not destinations: return (default,) @@ -97,10 +76,8 @@ class TenantTracerCache: ) -> "tuple[tuple[tuple[tuple[str, str], ...], tuple[OtelDestination, ...]], ...]": """Destinations grouped by their backend-required Resource attributes. - The key is a stable sorted tuple of ``destination_resource_attrs`` items. - Groups are returned in a deterministic order (sorted by key), so the - empty-Resource group (header-routed backends) sorts first and the - configured/global exporters attach to it. + Groups sort deterministically by key, so the empty-Resource group (header-routed + backends) sorts first and the configured/global exporters attach to it. """ from litellm.integrations.otel.plumbing.providers import ( destination_resource_attrs, @@ -140,9 +117,7 @@ class TenantTracerCache: def tracer_for(self, default: Tracer, destinations: "tuple[OtelDestination, ...]") -> Tracer: """Single merged tracer for ``destinations`` (one provider, one Resource). - The single-group primitive: kept for the destination-set cache mechanics and as - the building block ``tracers_for`` composes per group. The gen-AI span path uses - ``tracers_for`` so multiple Resource groups aren't last-wins merged. + The single-group primitive ``tracers_for`` composes per group. """ if not destinations: return default @@ -157,13 +132,10 @@ class TenantTracerCache: return get_tracer(provider, self._tracer_name) def _owned_otlp_kind(self) -> str: - """The OTLP transport of this integration's own exporter (langfuse -> http, - arize -> grpc), used for the destinations appended below. + """The OTLP transport of this integration's own exporter (langfuse -> http, arize -> grpc). - Prefer the admin's configured exporter kind for this backend; fall back to - the backend's intrinsic default (shared with the fan-out processor via - ``default_otlp_kind_for_backend``) so a lazily-activated backend with no - owned spec still picks the right transport (e.g. arize -> grpc, not http). + Prefers the admin's configured exporter kind; falls back to the backend's intrinsic + default so a lazily-activated backend with no owned spec still picks the right transport. """ from litellm.integrations.otel.plumbing.providers import ( default_otlp_kind_for_backend, @@ -180,24 +152,13 @@ class TenantTracerCache: *, include_base_exporters: bool = True, ) -> OpenTelemetryV2Config: - """Clone the config and APPEND one exporter per resolved destination. The shared - ``TracerProvider`` attaches one ``SpanProcessor`` per spec, so a single span is - emitted once and exported to every appended destination. Each appended exporter's - endpoint is the resolved host (the cross-host fix) with its own auth headers - (per-destination isolation). + """Clone the config, appending one exporter per resolved destination (its resolved host + and own auth headers) so one span exports to every destination. - ``include_base_exporters`` keeps the configured/global exporters too (so the - global still receives). ``tracers_for`` sets it only on the first Resource group, - so when a backend splits into multiple groups the global gets the span once - rather than once per group. - - The clone's Resource folds in the destinations' backend-required Resource - attributes (Arize needs ``model_id`` / ``arize.project.name``), via the same - ``destination_resource_attrs`` the fan-out path uses on proxy-internal spans. - Callers group destinations by those attributes first, so within one call all - ``destinations`` share a Resource and the merge is not lossy -- without this the - gen-AI span would reach Arize with only ``service.name`` while its parents - (fan-out) carry ``model_id``, orphaning the subtree.""" + ``include_base_exporters`` keeps the configured/global exporters; ``tracers_for`` sets it + only on the first Resource group so the global gets the span once, not once per group. The + clone's Resource folds in the destinations' ``destination_resource_attrs`` (Arize needs + ``model_id`` / ``arize.project.name``); callers group by those first so the merge is lossless.""" from litellm.integrations.otel.plumbing.providers import ( destination_resource_attrs, ) @@ -240,9 +201,8 @@ def _is_genai_span(span: ReadableSpan) -> bool: 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.""" + """Return ``span`` with its Resource augmented by the destination's required attributes, + via a shallow wrapper that leaves the original span untouched.""" from litellm.integrations.otel.plumbing.providers import ( destination_resource_attrs, ) @@ -255,12 +215,8 @@ def _with_destination_resource(span: ReadableSpan, destination: OtelDestination) 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. - """ + """A ``ReadableSpan`` view whose ``resource`` is overridden (for backend-specific attributes + like Arize's ``model_id``) without mutating the underlying span.""" def __init__(self, inner: ReadableSpan, resource: Resource) -> None: super().__init__( @@ -282,9 +238,8 @@ class _ResourceWrappedReadableSpan(ReadableSpan): 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. + Destinations come from a request-scoped contextvar set during auth, so the processor is + stateless across requests and concurrent requests are isolated by contextvars. """ def __init__(self, owner_callback_name: str | None) -> None: