mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-06 02:48:13 +00:00
refactor(otel/v2): merge fan-out processor into routing; describe abstractions in docstrings
Fold TenantFanOutSpanProcessor from plumbing/fan_out.py into plumbing/routing.py so one module owns both span-routing paths: TenantTracerCache for the gen-AI LLM-call span (per-tenant clone providers) and TenantFanOutSpanProcessor for proxy-internal spans on the main provider. Both read the request's destinations from the same server-only contextvar, so there is a single source of truth for routing. Rewrite the OtelDestination and dynamic-params-destination docstrings to describe what the abstraction is (a resolved endpoint plus auth headers the exporter sends) rather than narrating what it is not.
This commit is contained in:
parent
9733ce5809
commit
61c23a9430
6 changed files with 202 additions and 221 deletions
|
|
@ -1,11 +1,10 @@
|
|||
"""The resolved, admin-owned OTLP destination.
|
||||
"""The resolved OTLP destination a request's traces export to.
|
||||
|
||||
A trace destination is admin-owned infrastructure config, never request data.
|
||||
The proxy resolves a key/team's bound named credential into this typed,
|
||||
backend-agnostic target (an endpoint plus auth headers) server-side, and the v2
|
||||
logger exports through it. Every OTEL backend -- Langfuse, Arize, Weave, a
|
||||
self-hosted collector -- reduces to this shape; the per-backend field mapping
|
||||
lives in ``litellm.integrations.otel.destinations``.
|
||||
A destination is a backend-agnostic target: an endpoint plus the auth headers the
|
||||
exporter sends. The proxy builds it from the named logging credential bound to the
|
||||
request's identity chain, and the v2 logger exports through it. Every OTEL backend
|
||||
-- Langfuse, Arize, Weave, a self-hosted collector -- reduces to this shape; the
|
||||
per-backend field mapping lives in ``litellm.integrations.otel.presets.destinations``.
|
||||
"""
|
||||
|
||||
from pydantic import BaseModel, ConfigDict, Field
|
||||
|
|
|
|||
|
|
@ -51,13 +51,12 @@ if TYPE_CHECKING:
|
|||
|
||||
|
||||
def _otel_destinations(dynamic_params: Any) -> tuple[OtelDestination, ...]:
|
||||
"""The admin-resolved OTLP destinations carried on ``standard_callback_dynamic_params``.
|
||||
"""The admin-resolved OTLP destinations for this call, parsed off
|
||||
``standard_callback_dynamic_params``.
|
||||
|
||||
Server-set only (the proxy resolves the exporters assigned to the request's
|
||||
identity chain and strips any client value), so this is the sole source the v2
|
||||
router trusts -- request-supplied vendor credentials are never read here. A
|
||||
request fans out to every destination here; each logger keeps only the ones
|
||||
tagged with its own backend.
|
||||
The proxy resolves the destinations assigned to the request's identity chain and
|
||||
places them here server-side. A request fans out to every destination; each
|
||||
logger keeps only the ones tagged with its own backend.
|
||||
"""
|
||||
if not isinstance(dynamic_params, Mapping):
|
||||
return ()
|
||||
|
|
|
|||
|
|
@ -1,204 +0,0 @@
|
|||
"""Per-request span fan-out to admin-resolved destinations.
|
||||
|
||||
Attached to the main ``TracerProvider`` so every span on it (the FastAPI server
|
||||
span, the proxy's ``auth`` phase span, DB lookups, the post-call cost ledger)
|
||||
ships to every per-tenant destination the admin assigned for this request, on
|
||||
top of the provider's configured global exporters.
|
||||
|
||||
Spans emitted through the per-tenant ``TenantTracerCache`` clone providers (the
|
||||
gen-AI LLM-call span and its MCP-tool sibling) reach tenant backends through
|
||||
the clone's own exporters; the clone's provider has its own processor list and
|
||||
does NOT carry this fan-out processor, so a span is exported once per backend.
|
||||
|
||||
Each backend often requires backend-specific Resource attributes (Arize rejects
|
||||
spans missing ``model_id`` / ``arize.project.name``), so the fan-out wraps each
|
||||
forwarded span with the destination's expected Resource before handing it to
|
||||
the per-destination exporter.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections import OrderedDict
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from opentelemetry.context import Context
|
||||
from opentelemetry.sdk.resources import Resource
|
||||
from opentelemetry.sdk.trace import ReadableSpan, Span, SpanProcessor
|
||||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.otel.model.config import ExporterSpec
|
||||
from litellm.integrations.otel.plumbing.context import request_destinations
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from litellm.integrations.otel.model.destination import OtelDestination
|
||||
|
||||
# Bound on cached per-destination processors. One processor per
|
||||
# ``(endpoint, sorted(headers))`` pair, so the working set is one entry per
|
||||
# admin-resolved tenant credential -- a real-world deployment with hundreds of
|
||||
# tenants stays well under this. Evicted entries are dropped (not shut down; see
|
||||
# the eviction site) and reclaimed at process exit.
|
||||
_MAX_CACHED_PROCESSORS = 256
|
||||
|
||||
|
||||
def _processor_key(destination: OtelDestination) -> tuple:
|
||||
return (destination.endpoint, tuple(sorted(destination.headers.items())))
|
||||
|
||||
|
||||
class TenantFanOutSpanProcessor(SpanProcessor):
|
||||
"""Forward each finished span to every admin-resolved per-tenant destination.
|
||||
|
||||
The destinations are looked up from a request-scoped contextvar set during
|
||||
auth, so the processor is stateless across requests and isolation across
|
||||
concurrent requests is guaranteed by Python's contextvars.
|
||||
"""
|
||||
|
||||
def __init__(self, owner_callback_name: str | None) -> None:
|
||||
self._owner = owner_callback_name
|
||||
# Built lazily so we avoid importing the providers module at class
|
||||
# definition time (which would create a circular import with routing).
|
||||
self._processors: OrderedDict[tuple, SpanProcessor] = OrderedDict()
|
||||
|
||||
def on_start(self, span: Span, parent_context: Context | None = None) -> None:
|
||||
return None
|
||||
|
||||
def on_end(self, span: ReadableSpan) -> None:
|
||||
destinations = request_destinations()
|
||||
if not destinations:
|
||||
return
|
||||
# The gen-AI LLM-call span (and the MCP tool-call sibling) is already
|
||||
# routed to per-tenant destinations by the per-backend v2 logger via
|
||||
# ``TenantTracerCache`` -- the logger picks the right attribute mapper
|
||||
# (OpenInference for arize, GenAI semconv for langfuse_otel) and ships
|
||||
# through the clone provider's appended exporter. Forwarding it here too
|
||||
# would deliver a SECOND copy with the wrong vocabulary and a fresh
|
||||
# span_id, surfacing in the destination as an orphaned duplicate. Skip.
|
||||
if _is_genai_span(span):
|
||||
return
|
||||
# Proxy-internal spans (FastAPI server, ``auth`` phase, postgres lookups,
|
||||
# post-call cost ledger) are generic OTel semantic-convention spans with
|
||||
# no backend-specific vocabulary, so they ship to EVERY admin-resolved
|
||||
# destination this request fans out to, regardless of the destination's
|
||||
# ``callback_name``. The owner discriminator only matters for the gen-AI
|
||||
# span (handled by the skip above).
|
||||
for destination in destinations:
|
||||
processor = self._processor_for(destination)
|
||||
if processor is None:
|
||||
continue
|
||||
try:
|
||||
processor.on_end(_with_destination_resource(span, destination))
|
||||
except Exception as exc: # noqa: BLE001
|
||||
verbose_logger.debug(
|
||||
"OTel V2 fan-out: forwarding span to %s failed: %s",
|
||||
destination.endpoint,
|
||||
exc,
|
||||
)
|
||||
|
||||
def shutdown(self) -> None:
|
||||
for processor in self._processors.values():
|
||||
try:
|
||||
processor.shutdown()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
verbose_logger.debug("OTel V2 fan-out: processor shutdown failed: %s", exc)
|
||||
self._processors.clear()
|
||||
|
||||
def force_flush(self, timeout_millis: int = 30000) -> bool:
|
||||
all_ok = True
|
||||
for processor in self._processors.values():
|
||||
try:
|
||||
if not processor.force_flush(timeout_millis):
|
||||
all_ok = False
|
||||
except Exception: # noqa: BLE001
|
||||
all_ok = False
|
||||
return all_ok
|
||||
|
||||
def _processor_for(self, destination: OtelDestination) -> SpanProcessor | None:
|
||||
key = _processor_key(destination)
|
||||
cached = self._processors.get(key)
|
||||
if cached is not None:
|
||||
self._processors.move_to_end(key)
|
||||
return cached
|
||||
from litellm.integrations.otel.plumbing.providers import (
|
||||
_exporter_from_spec,
|
||||
_processor_for as _build_processor,
|
||||
default_otlp_kind_for_backend,
|
||||
)
|
||||
|
||||
try:
|
||||
spec = ExporterSpec(
|
||||
kind=default_otlp_kind_for_backend(destination.callback_name),
|
||||
endpoint=destination.endpoint,
|
||||
headers=destination.header_string(),
|
||||
owner=None,
|
||||
)
|
||||
exporter = _exporter_from_spec(spec)
|
||||
processor = _build_processor(exporter, use_simple=False)
|
||||
except Exception as exc:
|
||||
verbose_logger.debug(
|
||||
"OTel V2 fan-out: failed to build processor for %s: %s",
|
||||
destination.endpoint,
|
||||
exc,
|
||||
)
|
||||
return None
|
||||
self._processors[key] = processor
|
||||
if len(self._processors) > _MAX_CACHED_PROCESSORS:
|
||||
# Evict the LRU entry but do NOT shut it down here: a
|
||||
# ``BatchSpanProcessor`` may still hold spans queued on its exporter
|
||||
# thread, and calling ``shutdown`` synchronously can drop or raise on
|
||||
# those in-flight spans. Dropping the reference lets the worker drain
|
||||
# naturally and be reclaimed at process exit. The cache is bounded at
|
||||
# ``_MAX_CACHED_PROCESSORS``, so the un-shut-down working set stays
|
||||
# bounded.
|
||||
self._processors.popitem(last=False)
|
||||
return processor
|
||||
|
||||
|
||||
# Attribute set on every gen-AI LLM-call span by the v2 emitter. Used as the
|
||||
# unambiguous skip signal: only the LLM-call span carries this, and the
|
||||
# per-backend v2 logger already routes it to per-tenant destinations through
|
||||
# the TenantTracerCache clone provider's appended exporter.
|
||||
_GENAI_SPAN_ATTR = "gen_ai.operation.name"
|
||||
|
||||
|
||||
def _is_genai_span(span: ReadableSpan) -> bool:
|
||||
attributes = span.attributes or {}
|
||||
return _GENAI_SPAN_ATTR in attributes
|
||||
|
||||
|
||||
def _with_destination_resource(span: ReadableSpan, destination: OtelDestination) -> ReadableSpan:
|
||||
"""Return ``span`` with its Resource augmented by the destination's required
|
||||
attributes. The original span object is left untouched; a shallow wrapper
|
||||
reuses every other field and only swaps the ``resource`` property."""
|
||||
from litellm.integrations.otel.plumbing.providers import (
|
||||
destination_resource_attrs,
|
||||
)
|
||||
|
||||
extra = destination_resource_attrs(destination)
|
||||
if not extra:
|
||||
return span
|
||||
merged = Resource.create({**dict(span.resource.attributes), **extra})
|
||||
return _ResourceWrappedReadableSpan(span, merged)
|
||||
|
||||
|
||||
class _ResourceWrappedReadableSpan(ReadableSpan):
|
||||
"""A ``ReadableSpan`` view whose ``resource`` is overridden.
|
||||
|
||||
The OTLP exporter reads each span's ``resource`` when serializing; for
|
||||
backend-specific attributes (Arize's ``model_id``) we substitute the
|
||||
destination's expected Resource without mutating the underlying span.
|
||||
"""
|
||||
|
||||
def __init__(self, inner: ReadableSpan, resource: Resource) -> None:
|
||||
super().__init__(
|
||||
name=inner.name,
|
||||
context=inner.context,
|
||||
parent=inner.parent,
|
||||
resource=resource,
|
||||
attributes=inner.attributes,
|
||||
events=inner.events,
|
||||
links=inner.links,
|
||||
kind=inner.kind,
|
||||
status=inner.status,
|
||||
start_time=inner.start_time,
|
||||
end_time=inner.end_time,
|
||||
instrumentation_scope=inner.instrumentation_scope,
|
||||
)
|
||||
|
|
@ -346,7 +346,7 @@ def build_tracer_provider(
|
|||
provider.add_span_processor(baggage_processor)
|
||||
|
||||
if attach_tenant_fan_out or tenant_fan_out_owner is not None:
|
||||
from litellm.integrations.otel.plumbing.fan_out import (
|
||||
from litellm.integrations.otel.plumbing.routing import (
|
||||
TenantFanOutSpanProcessor,
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -1,4 +1,12 @@
|
|||
"""Per-request multi-tenant tracer routing with fan-out.
|
||||
"""Per-request multi-tenant tracer routing and span fan-out.
|
||||
|
||||
This module owns both halves of routing a request's spans to its admin-owned
|
||||
destinations. ``TenantTracerCache`` handles the gen-AI LLM-call span by building
|
||||
per-tenant clone providers (grouped by backend Resource attributes).
|
||||
``TenantFanOutSpanProcessor`` (at the bottom) handles the proxy-internal spans on
|
||||
the main provider (server, auth, DB, cost ledger) by forwarding each to every
|
||||
destination. Both read the request's destinations from the same server-only
|
||||
contextvar, so there is one source of truth for where a request's traces go.
|
||||
|
||||
A call's identity chain is assigned a set of admin-owned OTEL destinations
|
||||
(``LLMCallEvent.otel_destinations``, resolved server-side from named credentials).
|
||||
|
|
@ -17,12 +25,15 @@ a caller can neither redirect a trace nor spawn providers.
|
|||
|
||||
from collections import OrderedDict
|
||||
|
||||
from opentelemetry.sdk.trace import TracerProvider
|
||||
from opentelemetry.context import Context
|
||||
from opentelemetry.sdk.resources import Resource
|
||||
from opentelemetry.sdk.trace import ReadableSpan, Span, SpanProcessor, TracerProvider
|
||||
from opentelemetry.trace import Tracer
|
||||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.otel.model.config import ExporterSpec, OpenTelemetryV2Config
|
||||
from litellm.integrations.otel.model.destination import OtelDestination
|
||||
from litellm.integrations.otel.plumbing.context import request_destinations
|
||||
from litellm.integrations.otel.plumbing.providers import (
|
||||
build_tracer_provider,
|
||||
get_tracer,
|
||||
|
|
@ -224,3 +235,179 @@ class TenantTracerCache:
|
|||
"resource_attributes": merged_resource_attrs,
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
# --- Proxy-internal span fan-out ------------------------------------------- #
|
||||
#
|
||||
# ``TenantTracerCache`` above routes the gen-AI LLM-call span (and its MCP-tool
|
||||
# sibling) to per-tenant destinations through clone providers. The processor below
|
||||
# handles the OTHER span class: the proxy-internal spans (FastAPI server span, the
|
||||
# ``auth`` phase, DB lookups, the post-call cost ledger) emitted on the MAIN
|
||||
# provider. It forwards each to every admin-resolved destination for the request,
|
||||
# reading them from the same server-only contextvar the cache's callers set, so both
|
||||
# routing paths share one source of truth for where a request's traces go.
|
||||
|
||||
# Bound on cached per-destination processors. One processor per
|
||||
# ``(endpoint, sorted(headers))`` pair, so the working set is one entry per
|
||||
# admin-resolved tenant credential -- a real-world deployment with hundreds of
|
||||
# tenants stays well under this. Evicted entries are dropped (not shut down; see
|
||||
# the eviction site) and reclaimed at process exit.
|
||||
_MAX_CACHED_PROCESSORS = 256
|
||||
|
||||
# Attribute set on every gen-AI LLM-call span by the v2 emitter. Used as the
|
||||
# unambiguous skip signal: only the LLM-call span carries this, and the per-backend
|
||||
# v2 logger already routes it to per-tenant destinations through the
|
||||
# TenantTracerCache clone provider's appended exporter.
|
||||
_GENAI_SPAN_ATTR = "gen_ai.operation.name"
|
||||
|
||||
|
||||
def _processor_key(destination: OtelDestination) -> tuple:
|
||||
return (destination.endpoint, tuple(sorted(destination.headers.items())))
|
||||
|
||||
|
||||
def _is_genai_span(span: ReadableSpan) -> bool:
|
||||
attributes = span.attributes or {}
|
||||
return _GENAI_SPAN_ATTR in attributes
|
||||
|
||||
|
||||
def _with_destination_resource(span: ReadableSpan, destination: OtelDestination) -> ReadableSpan:
|
||||
"""Return ``span`` with its Resource augmented by the destination's required
|
||||
attributes. The original span object is left untouched; a shallow wrapper reuses
|
||||
every other field and only swaps the ``resource`` property."""
|
||||
from litellm.integrations.otel.plumbing.providers import (
|
||||
destination_resource_attrs,
|
||||
)
|
||||
|
||||
extra = destination_resource_attrs(destination)
|
||||
if not extra:
|
||||
return span
|
||||
merged = Resource.create({**dict(span.resource.attributes), **extra})
|
||||
return _ResourceWrappedReadableSpan(span, merged)
|
||||
|
||||
|
||||
class _ResourceWrappedReadableSpan(ReadableSpan):
|
||||
"""A ``ReadableSpan`` view whose ``resource`` is overridden.
|
||||
|
||||
The OTLP exporter reads each span's ``resource`` when serializing; for
|
||||
backend-specific attributes (Arize's ``model_id``) we substitute the
|
||||
destination's expected Resource without mutating the underlying span.
|
||||
"""
|
||||
|
||||
def __init__(self, inner: ReadableSpan, resource: Resource) -> None:
|
||||
super().__init__(
|
||||
name=inner.name,
|
||||
context=inner.context,
|
||||
parent=inner.parent,
|
||||
resource=resource,
|
||||
attributes=inner.attributes,
|
||||
events=inner.events,
|
||||
links=inner.links,
|
||||
kind=inner.kind,
|
||||
status=inner.status,
|
||||
start_time=inner.start_time,
|
||||
end_time=inner.end_time,
|
||||
instrumentation_scope=inner.instrumentation_scope,
|
||||
)
|
||||
|
||||
|
||||
class TenantFanOutSpanProcessor(SpanProcessor):
|
||||
"""Forward each finished proxy-internal span to every admin-resolved destination.
|
||||
|
||||
The destinations are looked up from a request-scoped contextvar set during auth,
|
||||
so the processor is stateless across requests and isolation across concurrent
|
||||
requests is guaranteed by Python's contextvars.
|
||||
"""
|
||||
|
||||
def __init__(self, owner_callback_name: str | None) -> None:
|
||||
self._owner = owner_callback_name
|
||||
self._processors: OrderedDict[tuple, SpanProcessor] = OrderedDict()
|
||||
|
||||
def on_start(self, span: Span, parent_context: Context | None = None) -> None:
|
||||
return None
|
||||
|
||||
def on_end(self, span: ReadableSpan) -> None:
|
||||
destinations = request_destinations()
|
||||
if not destinations:
|
||||
return
|
||||
# The gen-AI LLM-call span (and the MCP tool-call sibling) is already routed
|
||||
# to per-tenant destinations by the per-backend v2 logger via
|
||||
# ``TenantTracerCache`` -- the logger picks the right attribute mapper
|
||||
# (OpenInference for arize, GenAI semconv for langfuse_otel) and ships through
|
||||
# the clone provider's appended exporter. Forwarding it here too would deliver
|
||||
# a SECOND copy with the wrong vocabulary and a fresh span_id, surfacing in the
|
||||
# destination as an orphaned duplicate. Skip.
|
||||
if _is_genai_span(span):
|
||||
return
|
||||
# Proxy-internal spans (FastAPI server, ``auth`` phase, postgres lookups,
|
||||
# post-call cost ledger) are generic OTel semantic-convention spans with no
|
||||
# backend-specific vocabulary, so they ship to EVERY admin-resolved destination
|
||||
# this request fans out to, regardless of the destination's ``callback_name``.
|
||||
for destination in destinations:
|
||||
processor = self._processor_for(destination)
|
||||
if processor is None:
|
||||
continue
|
||||
try:
|
||||
processor.on_end(_with_destination_resource(span, destination))
|
||||
except Exception as exc: # noqa: BLE001
|
||||
verbose_logger.debug(
|
||||
"OTel V2 fan-out: forwarding span to %s failed: %s",
|
||||
destination.endpoint,
|
||||
exc,
|
||||
)
|
||||
|
||||
def shutdown(self) -> None:
|
||||
for processor in self._processors.values():
|
||||
try:
|
||||
processor.shutdown()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
verbose_logger.debug("OTel V2 fan-out: processor shutdown failed: %s", exc)
|
||||
self._processors.clear()
|
||||
|
||||
def force_flush(self, timeout_millis: int = 30000) -> bool:
|
||||
all_ok = True
|
||||
for processor in self._processors.values():
|
||||
try:
|
||||
if not processor.force_flush(timeout_millis):
|
||||
all_ok = False
|
||||
except Exception: # noqa: BLE001
|
||||
all_ok = False
|
||||
return all_ok
|
||||
|
||||
def _processor_for(self, destination: OtelDestination) -> SpanProcessor | None:
|
||||
key = _processor_key(destination)
|
||||
cached = self._processors.get(key)
|
||||
if cached is not None:
|
||||
self._processors.move_to_end(key)
|
||||
return cached
|
||||
from litellm.integrations.otel.plumbing.providers import (
|
||||
_exporter_from_spec,
|
||||
_processor_for as _build_processor,
|
||||
default_otlp_kind_for_backend,
|
||||
)
|
||||
|
||||
try:
|
||||
spec = ExporterSpec(
|
||||
kind=default_otlp_kind_for_backend(destination.callback_name),
|
||||
endpoint=destination.endpoint,
|
||||
headers=destination.header_string(),
|
||||
owner=None,
|
||||
)
|
||||
exporter = _exporter_from_spec(spec)
|
||||
processor = _build_processor(exporter, use_simple=False)
|
||||
except Exception as exc:
|
||||
verbose_logger.debug(
|
||||
"OTel V2 fan-out: failed to build processor for %s: %s",
|
||||
destination.endpoint,
|
||||
exc,
|
||||
)
|
||||
return None
|
||||
self._processors[key] = processor
|
||||
if len(self._processors) > _MAX_CACHED_PROCESSORS:
|
||||
# Evict the LRU entry but do NOT shut it down here: a
|
||||
# ``BatchSpanProcessor`` may still hold spans queued on its exporter
|
||||
# thread, and calling ``shutdown`` synchronously can drop or raise on those
|
||||
# in-flight spans. Dropping the reference lets the worker drain naturally
|
||||
# and be reclaimed at process exit. The cache is bounded, so the
|
||||
# un-shut-down working set stays bounded.
|
||||
self._processors.popitem(last=False)
|
||||
return processor
|
||||
|
|
|
|||
|
|
@ -33,7 +33,7 @@ from litellm.integrations.otel.plumbing.context import ( # noqa: E402
|
|||
set_request_destinations,
|
||||
_request_destinations,
|
||||
)
|
||||
from litellm.integrations.otel.plumbing.fan_out import ( # noqa: E402
|
||||
from litellm.integrations.otel.plumbing.routing import ( # noqa: E402
|
||||
TenantFanOutSpanProcessor,
|
||||
_processor_key,
|
||||
)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue