mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-05 02:41:56 +00:00
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.
This commit is contained in:
parent
762b9eb1b4
commit
da4583a110
5 changed files with 230 additions and 558 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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())
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue