feat(otel v2): send a key's or team's whole trace to its own destination

A key or team that configures its own Langfuse, Arize, Weave or New Relic
credentials used to get a single detached span in its account while the rest of
the request trace stayed on the operator's backend, so neither side held a
complete trace. Resolve the destination during auth, forward every span of the
request to it, and hold the same request back from the operator's exporter for
that backend, so the tenant gets the tree the operator would have seen and the
operator gets nothing for that request.

Also let a credential-mandatory preset build without the operator's own env
credentials. Without that, a proxy whose teams each bring their own account fell
back to the legacy integration and never ran a line of the v2 path.
This commit is contained in:
Yucheng He 2026-09-03 14:42:07 -07:00
parent a9f8a8d794
commit bcba86e426
20 changed files with 952 additions and 21 deletions

View file

@ -180,7 +180,9 @@ class OpenTelemetryV2(CustomLogger):
self.config: OpenTelemetryV2Config = config or OpenTelemetryV2Config(**kwargs)
self.callback_name = callback_name
self._tracer_provider: TracerProvider = (
tracer_provider if tracer_provider is not None else build_tracer_provider(self.config)
tracer_provider
if tracer_provider is not None
else build_tracer_provider(self.config, tenant_overrides=True)
)
self.tracer: Tracer = get_tracer(self._tracer_provider, LITELLM_TRACER_NAME)
self._metrics_recorder = self._init_metrics(meter_provider)

View file

@ -0,0 +1,55 @@
"""The resolved OTLP destination a request's traces export to.
A destination is a backend-agnostic target: an endpoint plus the auth headers the
exporter sends. The proxy builds one per backend from the key or team logging
config resolved at auth, and the fan-out span processor exports the request's
spans through it. Every OTEL backend reduces to this shape; the per-backend field
mapping lives in ``litellm.integrations.otel.presets.destinations``.
"""
from collections.abc import Mapping
from typing import Final
from urllib.parse import quote
from pydantic import BaseModel, ConfigDict, Field
class OtelDestination(BaseModel):
model_config = ConfigDict(frozen=True)
endpoint: str
headers: Mapping[str, str] = Field(default_factory=dict)
resource_attributes: Mapping[str, str] = Field(default_factory=dict)
callback_name: str | None = None
protocol: str | None = Field(
default=None,
description=(
"OTLP transport for this endpoint (``otlp_http`` / ``otlp_grpc``). The "
"backend's intrinsic default is used when unset. A backend whose own cloud "
"endpoint is gRPC can still be pointed at an HTTP collector, which the "
"scheme alone cannot express: Arize's own ``https://otlp.arize.com/v1`` is gRPC."
),
)
def header_string(self) -> str:
"""Render headers as the ``k=v,k2=v2`` form an ``ExporterSpec`` expects.
Values are percent-encoded because ``providers.parse_headers`` decodes them
with the SDK's W3C-Baggage parser: a value carrying a ``,`` or ``=`` (a
Langfuse project name, a base64 Authorization payload ending in ``==``)
would otherwise be split into bogus pairs on the way back out.
"""
return ",".join(f"{key}={quote(value, safe='')}" for key, value in self.headers.items())
def cache_key(self) -> tuple[str, tuple[tuple[str, str], ...], tuple[tuple[str, str], ...], str | None]:
"""Identity for processor reuse: two requests naming the same destination
must share one exporter rather than minting a connection pool each."""
return (
self.endpoint,
tuple(sorted(self.headers.items())),
tuple(sorted(self.resource_attributes.items())),
self.protocol,
)
NO_DESTINATIONS: Final[tuple[OtelDestination, ...]] = ()

View file

@ -2,7 +2,7 @@
from collections.abc import Mapping
from contextvars import ContextVar, Token
from typing import Final
from typing import TYPE_CHECKING, Final
from opentelemetry import baggage
from opentelemetry.context import Context, get_current
@ -21,6 +21,9 @@ from opentelemetry.trace.propagation.tracecontext import (
from litellm.integrations.otel.model.semconv import HTTP
if TYPE_CHECKING:
from litellm.integrations.otel.model.destination import OtelDestination
_PROPAGATOR: Final = TraceContextTextMapPropagator()
# The request's root span — the FastAPI-owned SERVER span — captured ONCE when the
@ -304,3 +307,34 @@ def extract_traceparent(headers: Mapping[str, str]) -> Context | None:
return None
carrier: Final = {str(key).lower(): value for key, value in headers.items()}
return _PROPAGATOR.extract(carrier)
# The OTLP destinations this request's key or team pointed its traces at, resolved
# once during auth. A ``ContextVar`` for the same reason the root span above is one:
# it rides the request task's context into the ``asyncio.create_task`` children that
# close the LLM span, and it is visible to every ``SpanProcessor.on_end`` that fires
# on the request task. Never reset -- it dies with the task.
_request_destinations: Final['ContextVar[tuple["OtelDestination", ...]]'] = ContextVar(
"litellm_otel_request_destinations", default=()
)
def set_request_destinations(destinations: 'tuple["OtelDestination", ...]') -> None:
"""Anchor the destinations this request exports to."""
_request_destinations.set(destinations)
def request_destinations() -> 'tuple["OtelDestination", ...]':
"""The destinations resolved for this request, empty outside a proxy request."""
return _request_destinations.get()
def overridden_backends() -> frozenset[str]:
"""Backends whose global exporters this request must NOT reach.
A team destination is an override, not an addition: once the request resolved a
destination for a backend, that backend's operator-level exporters are suppressed
for every span of the request, so the tenant's traffic reaches the tenant's
account and nowhere else.
"""
return frozenset(d.callback_name for d in _request_destinations.get() if d.callback_name)

View file

@ -1,5 +1,7 @@
"""Provider / exporter factory + the Baggage span processor."""
import threading
from collections import OrderedDict
from collections.abc import Callable, Iterable
from typing import TYPE_CHECKING, Any, Final, Literal
@ -20,6 +22,7 @@ from opentelemetry.sdk._logs.export import (
from opentelemetry.sdk.metrics import MeterProvider as SDKMeterProvider
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import ReadableSpan, SpanProcessor, TracerProvider
from opentelemetry.sdk.trace import Span as SDKSpan
from opentelemetry.sdk.trace.export import (
BatchSpanProcessor,
ConsoleSpanExporter,
@ -32,15 +35,22 @@ from opentelemetry.sdk.trace.export.in_memory_span_exporter import (
from opentelemetry.trace import Span, SpanKind, Tracer
from opentelemetry.util.re import parse_env_headers
from litellm._logging import verbose_logger
from litellm._version import version as litellm_version
from litellm.integrations.otel.model.config import ExporterSpec, OpenTelemetryV2Config
from litellm.integrations.otel.model.semconv import LiteLLM
from litellm.integrations.otel.model.spans import LiteLLMSpanKind
from litellm.integrations.otel.plumbing.context import (
overridden_backends,
request_destinations,
)
if TYPE_CHECKING:
from opentelemetry.metrics import Meter
from opentelemetry.sdk.metrics.export import MetricReader
from litellm.integrations.otel.model.destination import OtelDestination
_SPAN_KIND_BY_ROLE_KIND: Final[dict[LiteLLMSpanKind, SpanKind]] = {
LiteLLMSpanKind.SERVER: SpanKind.SERVER,
LiteLLMSpanKind.CLIENT: SpanKind.CLIENT,
@ -194,6 +204,178 @@ def _processor_for(exporter: SpanExporter, use_simple: bool | None) -> SpanProce
return SimpleSpanProcessor(exporter) if use_simple else BatchSpanProcessor(exporter)
#: Distinct tenant destinations whose exporters stay alive. Each holds a connection
#: pool and a batch thread, so the cache is bounded and evicts least-recently-used.
_MAX_CACHED_DESTINATION_PROCESSORS: Final = 32
class _ResourceWrappedReadableSpan(ReadableSpan):
"""A ``ReadableSpan`` view with an overridden Resource, leaving the original alone."""
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,
)
def _with_destination_resource(span: ReadableSpan, destination: "OtelDestination") -> ReadableSpan:
extra: Final = destination.resource_attributes
if not extra:
return span
merged: Final = Resource.create(
{**dict(span.resource.attributes), **dict(extra)} # mutable-ok: the OTel SDK takes a concrete attribute mapping
)
return _ResourceWrappedReadableSpan(span, merged)
class TenantFanOutSpanProcessor(SpanProcessor):
"""Export every finished span to each destination this request resolved.
Destinations ride a request-scoped ``ContextVar`` set during auth, so the
processor keeps no per-request state and concurrent requests stay isolated.
Every span is forwarded, the gen-AI span included: the tenant's account gets the
tree the operator's would have received, still parented, because the forwarded
view keeps the original span's trace and parent ids.
"""
def __init__(
self,
processor_factory: 'Callable[["OtelDestination"], SpanProcessor | None] | None' = None,
) -> None:
self._lock: Final = threading.Lock()
self._build: Final = processor_factory if processor_factory is not None else _destination_processor
self._processors: OrderedDict[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded LRU
def on_start(self, span: SDKSpan, parent_context: Context | None = None) -> None:
return None
def on_end(self, span: ReadableSpan) -> None:
for destination in request_destinations():
processor = self._processor_for(destination) # rebind-ok: loop variable; pyright forbids Final in a loop
if processor is None:
continue
try:
processor.on_end(_with_destination_resource(span, destination))
except Exception as exc: # noqa: BLE001 # one destination's failure must not cost the others their span
verbose_logger.debug("OTel V2 fan-out: forwarding to %s failed: %s", destination.endpoint, exc)
def shutdown(self) -> None:
# Snapshot first: ``on_end`` mutates the cache on whichever thread ends a span
# and can run concurrently with this SDK-driven shutdown, so iterating the live
# mapping risks a "mutated during iteration" the per-item except cannot catch.
for processor in self._snapshot():
try:
processor.shutdown()
except Exception as exc: # noqa: BLE001 # one processor's shutdown must not abort the rest
verbose_logger.debug("OTel V2 fan-out: processor shutdown failed: %s", exc)
with self._lock:
self._processors.clear()
def force_flush(self, timeout_millis: int = 30000) -> bool:
results: Final = tuple(self._flush_one(processor, timeout_millis) for processor in self._snapshot())
return all(results)
def _snapshot(self) -> tuple[SpanProcessor, ...]:
with self._lock:
return tuple(self._processors.values())
@staticmethod
def _flush_one(processor: SpanProcessor, timeout_millis: int) -> bool:
try:
return processor.force_flush(timeout_millis)
except Exception: # noqa: BLE001 # one exporter's flush failure must not fail the whole flush
return False
def _processor_for(self, destination: "OtelDestination") -> SpanProcessor | None:
key: Final = destination.cache_key()
with self._lock:
cached: Final = self._processors.get(key)
if cached is not None:
self._processors.move_to_end(key)
return cached
built: Final = self._build(destination)
if built is None:
return None
with self._lock:
existing: Final = self._processors.get(key)
if existing is not None:
# Another thread won the race; drop ours rather than leak its thread.
_shutdown_quietly(built)
return existing
self._processors[key] = built
evicted: Final = (
self._processors.popitem(last=False)[1]
if len(self._processors) > _MAX_CACHED_DESTINATION_PROCESSORS
else None
)
if evicted is not None:
_shutdown_quietly(evicted)
return built
def _destination_processor(destination: "OtelDestination") -> SpanProcessor | None:
"""A batching OTLP processor aimed at ``destination``, or ``None`` if unbuildable."""
try:
spec: Final = ExporterSpec(
kind=destination.protocol or "otlp_http",
endpoint=destination.endpoint,
headers=destination.header_string(),
owner=None,
)
return _processor_for(_exporter_from_spec(spec), use_simple=False)
except Exception as exc: # noqa: BLE001 # a malformed destination must not break the request or the other destinations
verbose_logger.debug("OTel V2 fan-out: no processor for %s: %s", destination.endpoint, exc)
return None
def _shutdown_quietly(processor: SpanProcessor) -> None:
try:
processor.shutdown()
except Exception as exc: # noqa: BLE001 # defensive: shedding a spare processor must not raise
verbose_logger.debug("OTel V2 fan-out: discarding processor failed: %s", exc)
class _OverriddenBackendFilter(SpanProcessor):
"""Hold a span back from ``owner``'s operator-level exporter when the request
pointed ``owner`` at a tenant's own account.
A team destination is an override rather than an addition, and a ``SpanProcessor``
cannot veto its siblings (``SynchronousMultiSpanProcessor.on_end`` ignores return
values), so suppression has to wrap the exporter's own processor. Dropping here
also keeps the span out of ``BatchSpanProcessor``'s bounded queue instead of
filling it with spans that will never ship.
"""
def __init__(self, inner: SpanProcessor, owner: str) -> None:
self._inner: Final = inner
self._owner: Final = owner
def on_start(self, span: SDKSpan, parent_context: Context | None = None) -> None:
self._inner.on_start(span, parent_context)
def on_end(self, span: ReadableSpan) -> None:
if self._owner in overridden_backends():
return
self._inner.on_end(span)
def shutdown(self) -> None:
self._inner.shutdown()
def force_flush(self, timeout_millis: int = 30000) -> bool:
return self._inner.force_flush(timeout_millis)
def build_span_exporter(config: OpenTelemetryV2Config) -> SpanExporter:
"""Build a single exporter from the top-level config fields.
@ -437,6 +619,7 @@ def build_tracer_provider(
exporter: SpanExporter | None = None,
baggage_processor: SpanProcessor | None = None,
use_simple_processor: bool | None = None,
tenant_overrides: bool = False,
) -> TracerProvider:
"""Build the shared :class:`TracerProvider`.
@ -445,6 +628,12 @@ def build_tracer_provider(
``config.exporters`` entry — this is what fans spans out to multiple
backends. ``exporter`` and ``use_simple_processor`` are explicit overrides:
pass a single exporter to attach exactly that one (used by tests).
``tenant_overrides`` belongs to the operator-level provider alone: it wraps each
owned exporter so a request that pointed that backend at a key's or team's own
account skips it, and adds the fan-out processor that delivers to that account
instead. The per-tenant providers this same function builds must leave it off,
or they would filter out the very spans they exist to carry.
"""
provider: Final = TracerProvider(resource=build_resource(config))
if baggage_processor is None:
@ -461,12 +650,16 @@ def build_tracer_provider(
if spec.requires_headers and not spec.headers:
continue
exp = _exporter_from_spec(spec)
provider.add_span_processor(
_processor_for(
exp,
(spec.use_simple_processor if spec.use_simple_processor is not None else use_simple_processor),
)
processor = _processor_for(
exp,
(spec.use_simple_processor if spec.use_simple_processor is not None else use_simple_processor),
)
owner = spec.owner.value if spec.owner is not None else None
provider.add_span_processor(
_OverriddenBackendFilter(processor, owner) if tenant_overrides and owner is not None else processor
)
if tenant_overrides:
provider.add_span_processor(TenantFanOutSpanProcessor())
return provider

View file

@ -25,6 +25,7 @@ from opentelemetry.trace import Tracer
from litellm._logging import verbose_logger
from litellm.constants import OTEL_SERVICE_NAME_METADATA_KEYS
from litellm.integrations.otel.model.config import ExporterSpec, OpenTelemetryV2Config
from litellm.integrations.otel.plumbing.context import overridden_backends
from litellm.integrations.otel.plumbing.providers import (
build_tracer_provider,
exporter_transport,
@ -231,7 +232,14 @@ class TenantTracerCache:
concurrent overflow eviction can't shut it down between selection and
the caller's span start. The caller must ``release`` it exactly once.
"""
credential_headers: Final = self._credential_headers(dynamic_params)
# An overridden backend is delivered by the fan-out processor, which carries the
# whole trace. Routing here too would detach this span onto a second provider,
# so the tenant would get the request tree plus a stray one-span trace.
credential_headers: Final = (
_NO_HEADERS
if self._callback_name is not None and self._callback_name in overridden_backends()
else self._credential_headers(dynamic_params)
)
project_headers: Final = self._project_headers(auth_metadata)
service_name: Final = tenant_service_name(auth_metadata)
if not credential_headers and not project_headers and service_name is None:

View file

@ -39,6 +39,7 @@ class _AgentOpsSettings(BaseSettings):
def agentops_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
"""Build the AgentOps config without any network I/O.

View file

@ -11,7 +11,10 @@ from litellm.integrations.otel.model.config import (
ExporterSpec,
OpenTelemetryV2Config,
)
from litellm.integrations.otel.presets.utils import ensure_mappers
from litellm.integrations.otel.presets.utils import (
credential_gated_exporters,
ensure_mappers,
)
from litellm.types.utils import StandardCallbackDynamicParams
@ -26,10 +29,22 @@ class _ArizeSettings(BaseSettings):
def arize_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
arize_cfg: Final = _V1ArizeLogger.get_arize_config()
headers: Final = _arize_headers(arize_cfg)
base: Final = config_overrides or OpenTelemetryV2Config()
mappers: Final = ensure_mappers(base.mapper_names, "openinference")
try:
arize_cfg: Final = _V1ArizeLogger.get_arize_config()
except Exception:
if not allow_missing_credentials:
raise
return base.model_copy(
update={ # mutable-ok: pydantic model_copy takes a plain update mapping
"exporters": credential_gated_exporters(base.exporters, ExporterOwner.ARIZE_AX),
"mapper_names": mappers,
}
)
headers: Final = _arize_headers(arize_cfg)
return base.model_copy(
update={
"exporters": [
@ -41,7 +56,7 @@ def arize_preset(
owner=ExporterOwner.ARIZE_AX,
),
],
"mapper_names": ensure_mappers(base.mapper_names, "openinference"),
"mapper_names": mappers,
"resource_attributes": {
**base.resource_attributes,
**({"model_id": arize_cfg.project_name} if arize_cfg.project_name else {}),

View file

@ -18,6 +18,18 @@ class Preset(Protocol):
``config_overrides`` lets one preset layer onto another's config (or onto
test-supplied defaults); the factory calls presets with no arguments.
``allow_missing_credentials`` lets a credential-mandatory backend (langfuse /
arize / weave) degrade to an exporter-less, mapper-only config instead of
raising when the operator set no env credentials of their own. That is a real
deployment: every team brings its own account and the operator keeps none, and
without it the whole V2 path silently falls back to the legacy integration, so
no team destination is ever reached. Credential-optional backends ignore it.
"""
def __call__(self, *, config_overrides: OpenTelemetryV2Config | None = None) -> OpenTelemetryV2Config: ...
def __call__(
self,
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config: ...

View file

@ -0,0 +1,118 @@
"""Map a key's or team's callback vars to the OTLP destination its traces export to.
The auth path resolves one destination per backend the caller configured, and the
fan-out span processor exports the whole request through it. Header building is
delegated to each preset's existing ``*_dynamic_headers`` builder, so a destination
authenticates exactly the way the per-request tracer route already did; only the
endpoint needs a per-backend rule, because a backend's host is either fixed, taken
from a region table, or named by the tenant alongside its own key pair.
"""
import os
from collections.abc import Callable, Mapping
from types import MappingProxyType
from typing import Final
from litellm.integrations.otel.model.destination import OtelDestination
from litellm.types.utils import StandardCallbackDynamicParams
#: gRPC is Arize's own transport; an explicitly named HTTP collector overrides it.
_ARIZE_GRPC_ENDPOINT: Final = "https://otlp.arize.com/v1"
def _langfuse_endpoint(params: StandardCallbackDynamicParams) -> str | None:
"""The tenant's own Langfuse host, else the operator's, else Langfuse US cloud.
Falling back to the operator's host is safe and is what V1 does: the tenant's
own key pair still selects its own project, and a self-hosted deployment where
every team lives on one Langfuse server is the common shape.
"""
from litellm.integrations.langfuse.langfuse_otel import (
LANGFUSE_CLOUD_US_ENDPOINT,
LangfuseOtelLogger,
)
host: Final = params.get("langfuse_host") or LangfuseOtelLogger._get_langfuse_otel_host() # pyright: ignore[reportPrivateUsage] # reuse the backend's own env host resolver rather than duplicating it
if not host:
return LANGFUSE_CLOUD_US_ENDPOINT
normalized: Final = host if host.startswith("http") else f"https://{host}"
return f"{normalized.rstrip('/')}/api/public/otel"
def _arize_endpoint(params: StandardCallbackDynamicParams) -> str | None:
return os.environ.get("ARIZE_ENDPOINT") or _ARIZE_GRPC_ENDPOINT
def _arize_protocol(params: StandardCallbackDynamicParams) -> str | None:
return "otlp_http" if os.environ.get("ARIZE_HTTP_ENDPOINT") and not os.environ.get("ARIZE_ENDPOINT") else None
def _weave_endpoint(params: StandardCallbackDynamicParams) -> str | None:
from litellm.integrations.weave.weave_otel import WEAVE_BASE_URL, WEAVE_OTEL_ENDPOINT
return WEAVE_BASE_URL + WEAVE_OTEL_ENDPOINT
def _newrelic_endpoint(params: StandardCallbackDynamicParams) -> str | None:
from litellm.integrations.otel.presets.newrelic import newrelic_dynamic_endpoint
return newrelic_dynamic_endpoint(params)
#: Callback name -> endpoint resolver. A backend is destination-capable exactly
#: when it appears here AND in ``DYNAMIC_HEADERS_BY_CALLBACK``: without a header
#: builder the destination would carry no tenant credentials, and the exporter
#: would post the tenant's traffic to the operator's account.
_ENDPOINT_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynamicParams], str | None]]] = MappingProxyType(
{
"langfuse_otel": _langfuse_endpoint,
"arize": _arize_endpoint,
"weave_otel": _weave_endpoint,
"newrelic": _newrelic_endpoint,
}
)
_PROTOCOL_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynamicParams], str | None]]] = MappingProxyType(
{
"arize": _arize_protocol,
}
)
_NO_ATTRS: Final[Mapping[str, str]] = MappingProxyType({})
def destination_capable_backends() -> frozenset[str]:
"""Backends a key or team can point at its own account."""
from litellm.integrations.otel.presets import DYNAMIC_HEADERS_BY_CALLBACK
return frozenset(_ENDPOINT_BY_CALLBACK) & frozenset(DYNAMIC_HEADERS_BY_CALLBACK)
def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -> OtelDestination | None:
"""The destination ``params`` names for ``callback_name``, or ``None``.
``None`` means the caller configured nothing usable for this backend, so the
request keeps the operator's global exporters. A partial config (a host with
no key pair) resolves to ``None`` rather than to the operator's endpoint with
the tenant's host, which would post the operator's credentials elsewhere.
"""
from litellm.integrations.otel.presets import DYNAMIC_HEADERS_BY_CALLBACK
header_builder: Final = DYNAMIC_HEADERS_BY_CALLBACK.get(callback_name)
endpoint_builder: Final = _ENDPOINT_BY_CALLBACK.get(callback_name)
if header_builder is None or endpoint_builder is None:
return None
headers: Final = header_builder(params)
if not headers:
return None
endpoint: Final = endpoint_builder(params)
if not endpoint:
return None
protocol_builder: Final = _PROTOCOL_BY_CALLBACK.get(callback_name)
return OtelDestination(
endpoint=endpoint,
headers=MappingProxyType(dict(headers)),
resource_attributes=_NO_ATTRS,
callback_name=callback_name,
protocol=protocol_builder(params) if protocol_builder is not None else None,
)

View file

@ -10,17 +10,32 @@ from litellm.integrations.otel.model.config import (
ExporterSpec,
OpenTelemetryV2Config,
)
from litellm.integrations.otel.presets.utils import ensure_mappers
from litellm.integrations.otel.presets.utils import (
credential_gated_exporters,
ensure_mappers,
)
from litellm.types.utils import StandardCallbackDynamicParams
def langfuse_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
cfg: Final = _V1Langfuse.get_langfuse_otel_config()
kind: Final = cfg.exporter if isinstance(cfg.exporter, str) else "otlp_http"
base: Final = config_overrides or OpenTelemetryV2Config()
mappers: Final = ensure_mappers(base.mapper_names, "langfuse")
try:
cfg: Final = _V1Langfuse.get_langfuse_otel_config()
except Exception:
if not allow_missing_credentials:
raise
return base.model_copy(
update={ # mutable-ok: pydantic model_copy takes a plain update mapping
"exporters": credential_gated_exporters(base.exporters, ExporterOwner.LANGFUSE_OTEL),
"mapper_names": mappers,
}
)
kind: Final = cfg.exporter if isinstance(cfg.exporter, str) else "otlp_http"
return base.model_copy(
update={
"exporters": [
@ -32,7 +47,7 @@ def langfuse_preset(
owner=ExporterOwner.LANGFUSE_OTEL,
),
],
"mapper_names": ensure_mappers(base.mapper_names, "langfuse"),
"mapper_names": mappers,
}
)

View file

@ -9,6 +9,7 @@ from litellm.integrations.otel.presets.utils import ensure_mappers
def langtrace_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
"""Compose the Langtrace mapper on top of the customer's OTLP destination.

View file

@ -13,6 +13,7 @@ from litellm.integrations.otel.model.config import (
def levo_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
cfg: Final = _V1Levo.get_levo_config()
base: Final = config_overrides or OpenTelemetryV2Config()

View file

@ -44,6 +44,7 @@ class _NewRelicSettings(BaseSettings):
def newrelic_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
settings: Final = _NewRelicSettings()
base: Final = config_overrides or OpenTelemetryV2Config()

View file

@ -60,6 +60,7 @@ def phoenix_project_headers(auth_metadata: Mapping[str, str] | None) -> Mapping[
def phoenix_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
cfg: Final = _V1Phoenix.get_arize_phoenix_config()
headers: Final = cfg.otlp_auth_headers if hasattr(cfg, "otlp_auth_headers") else None

View file

@ -3,6 +3,11 @@
from collections.abc import Iterable
from typing import Final
from litellm.integrations.otel.model.config import ExporterOwner, ExporterSpec
#: What ``OpenTelemetryV2Config._normalize`` folds in when no destination is configured.
_DEFAULT_SHORTHAND_EXPORTER: Final = ExporterSpec()
def ensure_mappers(mapper_names: Iterable[str], *names: str) -> list[str]:
"""Return ``mapper_names`` with each of ``names`` appended if not already present.
@ -15,3 +20,21 @@ def ensure_mappers(mapper_names: Iterable[str], *names: str) -> list[str]:
if name not in result:
result.append(name)
return result
def credential_gated_exporters(
exporters: "Iterable[ExporterSpec]", owner: "ExporterOwner"
) -> "tuple[ExporterSpec, ...]":
"""``exporters`` with the operator's destination replaced by a header-gated one.
Used when a credential-mandatory backend is asked to build without the operator's
own credentials, so only key/team destinations receive spans. Two things have to
happen for that to mean "export nowhere": the placeholder console spec that
``OpenTelemetryV2Config`` folds in for an empty exporter list is dropped, or every
span would be printed to stdout, and the gated spec keeps the owner so the
override filter still recognises which backend this provider speaks for.
"""
return (
*(spec for spec in exporters if spec != _DEFAULT_SHORTHAND_EXPORTER),
ExporterSpec(owner=owner, requires_headers=True),
)

View file

@ -7,7 +7,10 @@ from litellm.integrations.otel.model.config import (
ExporterSpec,
OpenTelemetryV2Config,
)
from litellm.integrations.otel.presets.utils import ensure_mappers
from litellm.integrations.otel.presets.utils import (
credential_gated_exporters,
ensure_mappers,
)
from litellm.integrations.weave.weave_otel import (
_get_weave_authorization_header,
get_weave_otel_config,
@ -18,9 +21,21 @@ from litellm.types.utils import StandardCallbackDynamicParams
def weave_preset(
*,
config_overrides: OpenTelemetryV2Config | None = None,
allow_missing_credentials: bool = False,
) -> OpenTelemetryV2Config:
weave_cfg: Final = get_weave_otel_config()
base: Final = config_overrides or OpenTelemetryV2Config()
mappers: Final = ensure_mappers(base.mapper_names, "openinference", "weave")
try:
weave_cfg: Final = get_weave_otel_config()
except Exception:
if not allow_missing_credentials:
raise
return base.model_copy(
update={ # mutable-ok: pydantic model_copy takes a plain update mapping
"exporters": credential_gated_exporters(base.exporters, ExporterOwner.WEAVE_OTEL),
"mapper_names": mappers,
}
)
return base.model_copy(
update={
"exporters": [
@ -33,7 +48,7 @@ def weave_preset(
),
],
# Weave consumes OpenInference + a small Weave-specific overlay.
"mapper_names": ensure_mappers(base.mapper_names, "openinference", "weave"),
"mapper_names": mappers,
}
)

View file

@ -4815,11 +4815,16 @@ def _maybe_construct_otel_v2(callback_name: str, _in_memory_loggers: list[Custom
if isinstance(callback, OpenTelemetryV2) and getattr(callback, "callback_name", None) == callback_name:
return callback
try:
config: Final = preset_fn()
config: Final = preset_fn(allow_missing_credentials=True)
except Exception:
# If env vars are missing or the preset raises, defer to the legacy path
# so customers get the same error story they had before V2 landed.
return None
if all(spec.requires_headers and not spec.headers for spec in config.exporters):
verbose_logger.warning(
"OTel V2: no operator credentials for '%s'; only key/team destinations will receive its traces",
callback_name,
)
v2_logger: Final = build_otel_v2_logger(config=config, callback_name=callback_name)
_in_memory_loggers.append(v2_logger)
return v2_logger

View file

@ -2845,6 +2845,28 @@ async def _authorize_authenticated_request(
@tracer.wrap()
def _seed_request_destinations(user_api_key_dict: UserAPIKeyAuth) -> None:
"""Anchor the OTLP destinations this key or team overrides its traces to.
Called inside the ``auth`` phase span so that span reaches the tenant's account
as well, and on the request task so the ``ContextVar`` is inherited by the logging
tasks that close the LLM span. Best-effort: trace routing must never fail auth.
The two ``postgres`` spans under ``auth`` close before this runs, because they are
the reads that resolve the identity being read here, so they keep going to the
operator's backend alone.
"""
try:
from litellm.integrations.otel.plumbing.context import set_request_destinations
from litellm.proxy.litellm_pre_call_utils import (
resolve_tenant_otel_destinations,
)
set_request_destinations(resolve_tenant_otel_destinations(user_api_key_dict))
except Exception as exc: # noqa: BLE001 # telemetry routing is best-effort and must never break authentication
verbose_proxy_logger.debug("OTel V2: tenant destination resolution failed: %s", exc)
async def user_api_key_auth(
request: Request,
api_key: str = fastapi.Security(api_key_header),
@ -2891,6 +2913,7 @@ async def user_api_key_auth(
raise body_parse_exception
raise
user_api_key_auth_obj.budget_reservation = None
_seed_request_destinations(user_api_key_auth_obj)
# A body that never parsed is authenticated (so the trace carries identity
# and this ``auth`` span) but not authorized: there is no model to check it

View file

@ -10,6 +10,7 @@ from types import MappingProxyType
from typing import TYPE_CHECKING, Any, Final, cast
from fastapi import HTTPException, Request
from pydantic import TypeAdapter
from pydantic import ValidationError as PydanticValidationError
from starlette.datastructures import Headers
@ -157,6 +158,7 @@ from litellm.types.utils import (
CustomPricingLiteLLMParams,
LlmProviders,
ProviderSpecificHeader,
StandardCallbackDynamicParams,
StandardLoggingUserAPIKeyMetadata,
SupportedCacheControls,
)
@ -170,6 +172,7 @@ _ENABLE_TEAM_STALE_ALIAS_BYPASS: bool | None = None
if TYPE_CHECKING:
from litellm.integrations.otel.model.destination import OtelDestination
from litellm.proxy.proxy_server import ProxyConfig as _ProxyConfig
from litellm.types.proxy.policy_engine import Policy, PolicyMatchContext
@ -974,6 +977,51 @@ def _get_dynamic_logging_metadata(
return callback_settings_obj
_TENANT_OTEL_PARAMS: Final = TypeAdapter(StandardCallbackDynamicParams)
def _tenant_otel_params(callback_vars: Mapping[str, str]) -> StandardCallbackDynamicParams:
try:
return _TENANT_OTEL_PARAMS.validate_python(callback_vars)
except PydanticValidationError:
return StandardCallbackDynamicParams()
def resolve_tenant_otel_destinations(
user_api_key_dict: UserAPIKeyAuth,
) -> "tuple[OtelDestination, ...]":
"""The OTLP destinations this request's key or team config overrides its traces to.
Key settings win over team settings outright, the same precedence
``_get_dynamic_logging_metadata`` applies, so one caller never exports the same
backend to two accounts. Returns empty when OTEL V2 is off, when neither level
named a destination-capable backend, or when the config is incomplete, and the
request then keeps the operator's own exporters.
"""
from litellm.integrations.otel.model.config import is_otel_v2_enabled
from litellm.integrations.otel.presets.destinations import destination_for
if not is_otel_v2_enabled():
return ()
entries: Final = KeyAndTeamLoggingSettings.get_key_dynamic_logging_settings(
user_api_key_dict
) or KeyAndTeamLoggingSettings.get_team_dynamic_logging_settings(user_api_key_dict)
if not entries:
return ()
resolved: Final = tuple(
destination
for item in entries
if (callback := _get_validated_callback_metadata(item=item, source="otel-destination")) is not None
if (destination := destination_for(callback.callback_name, _tenant_otel_params(callback.callback_vars)))
is not None
)
return tuple(
destination
for index, destination in enumerate(resolved)
if destination.callback_name not in tuple(earlier.callback_name for earlier in resolved[:index])
)
def clean_headers(
headers: Headers,
litellm_key_header_name: str | None = None,

View file

@ -0,0 +1,360 @@
"""Key/team OTLP destinations override the operator's exporters for that backend."""
import contextvars
from collections.abc import Mapping
import pytest
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
from litellm.integrations.otel.model.config import (
ExporterOwner,
ExporterSpec,
OpenTelemetryV2Config,
is_otel_v2_enabled,
)
from litellm.integrations.otel.model.destination import OtelDestination
from litellm.integrations.otel.plumbing.context import (
overridden_backends,
request_destinations,
set_request_destinations,
)
from litellm.integrations.otel.plumbing.providers import (
TenantFanOutSpanProcessor,
_OverriddenBackendFilter,
build_tracer_provider,
)
from litellm.integrations.otel.plumbing.routing import TenantTracerCache, get_tracer
from litellm.integrations.otel.presets.destinations import (
destination_capable_backends,
destination_for,
)
from litellm.integrations.otel.presets.langfuse import langfuse_preset
from litellm.proxy._types import UserAPIKeyAuth
from litellm.proxy.litellm_pre_call_utils import resolve_tenant_otel_destinations
LANGFUSE_DEST = OtelDestination(
endpoint="http://tenant.local/api/public/otel",
headers={"Authorization": "Basic dGVuYW50"},
callback_name="langfuse_otel",
)
def in_fresh_context(fn, *args):
"""Run ``fn`` in its own context so one test's destinations never leak."""
return contextvars.copy_context().run(fn, *args)
def emit(provider: TracerProvider, name: str = "chat gpt-4") -> None:
with get_tracer(provider, "litellm").start_as_current_span(name):
pass
def wired_provider(dest_exporter: InMemorySpanExporter, global_exporter: InMemorySpanExporter) -> TracerProvider:
"""The operator's provider: one owned exporter plus the tenant fan-out."""
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(global_exporter), "langfuse_otel"))
provider.add_span_processor(
TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
)
return provider
class TestOverrideSuppression:
def test_operator_exporter_keeps_the_span_when_no_destination_is_resolved(self):
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
in_fresh_context(emit, provider)
assert [s.name for s in global_exporter.get_finished_spans()] == ["chat gpt-4"]
assert dest_exporter.get_finished_spans() == ()
def test_operator_exporter_is_skipped_once_the_backend_is_overridden(self):
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
def run():
set_request_destinations((LANGFUSE_DEST,))
emit(provider)
in_fresh_context(run)
assert global_exporter.get_finished_spans() == ()
assert [s.name for s in dest_exporter.get_finished_spans()] == ["chat gpt-4"]
def test_a_backend_the_request_did_not_override_still_exports(self):
arize_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(arize_exporter), "arize"))
provider.add_span_processor(
TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
)
def run():
set_request_destinations((LANGFUSE_DEST,))
emit(provider)
in_fresh_context(run)
assert [s.name for s in arize_exporter.get_finished_spans()] == ["chat gpt-4"]
class TestFanOut:
def test_every_span_of_the_request_reaches_the_destination_in_one_trace(self):
"""The whole tree, gen-AI span included, parented as the operator would see it."""
dest_exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(
TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter))
)
tracer = get_tracer(provider, "litellm")
def run():
set_request_destinations((LANGFUSE_DEST,))
with tracer.start_as_current_span("POST /v1/chat/completions"):
with tracer.start_as_current_span("auth /v1/chat/completions"):
pass
with tracer.start_as_current_span("chat gpt-4"):
pass
in_fresh_context(run)
spans = dest_exporter.get_finished_spans()
by_name = {s.name: s for s in spans}
assert set(by_name) == {"POST /v1/chat/completions", "auth /v1/chat/completions", "chat gpt-4"}
root = by_name["POST /v1/chat/completions"]
assert len({s.context.trace_id for s in spans}) == 1, "the tenant must receive one connected trace"
for child in ("auth /v1/chat/completions", "chat gpt-4"):
assert by_name[child].parent.span_id == root.context.span_id
def test_two_destinations_each_receive_their_own_copy(self):
first, second = InMemorySpanExporter(), InMemorySpanExporter()
by_endpoint = {"http://a.local": first, "http://b.local": second}
provider = TracerProvider()
provider.add_span_processor(
TenantFanOutSpanProcessor(
processor_factory=lambda d: SimpleSpanProcessor(by_endpoint[d.endpoint]),
)
)
def run():
set_request_destinations(
(
OtelDestination(endpoint="http://a.local", callback_name="langfuse_otel"),
OtelDestination(endpoint="http://b.local", callback_name="arize"),
)
)
emit(provider)
in_fresh_context(run)
assert [s.name for s in first.get_finished_spans()] == ["chat gpt-4"]
assert [s.name for s in second.get_finished_spans()] == ["chat gpt-4"]
def test_a_destination_that_cannot_build_a_processor_is_skipped_quietly(self):
"""An unbuildable destination must not cost the caller its request."""
attempts = []
reached_the_end = []
def factory(destination):
attempts.append(destination.endpoint)
provider = TracerProvider()
provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory))
def run():
set_request_destinations((LANGFUSE_DEST,))
emit(provider)
reached_the_end.append(True)
in_fresh_context(run)
assert attempts == [LANGFUSE_DEST.endpoint]
assert reached_the_end == [True]
def test_one_processor_is_reused_across_spans_of_the_same_destination(self):
built = []
def factory(_destination):
processor = SimpleSpanProcessor(InMemorySpanExporter())
built.append(processor)
return processor
provider = TracerProvider()
provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory))
def run():
set_request_destinations((LANGFUSE_DEST,))
emit(provider, "one")
emit(provider, "two")
in_fresh_context(run)
assert len(built) == 1
class TestProviderWiring:
def test_build_tracer_provider_only_filters_when_asked(self):
config = OpenTelemetryV2Config(exporters=[ExporterSpec(kind="in_memory", owner=ExporterOwner.LANGFUSE_OTEL)])
operator = build_tracer_provider(config, tenant_overrides=True)
tenant = build_tracer_provider(config)
def kinds(provider):
return [type(p).__name__ for p in provider._active_span_processor._span_processors]
assert "_OverriddenBackendFilter" in kinds(operator)
assert "TenantFanOutSpanProcessor" in kinds(operator)
assert "_OverriddenBackendFilter" not in kinds(tenant), "a per-tenant provider must not filter itself out"
assert "TenantFanOutSpanProcessor" not in kinds(tenant)
class TestRouting:
def test_an_overridden_backend_is_not_detached_onto_a_second_provider(self):
config = OpenTelemetryV2Config(
exporters=[ExporterSpec(kind="otlp_http", endpoint="http://op.local", owner=ExporterOwner.LANGFUSE_OTEL)]
)
cache = TenantTracerCache(config, "langfuse_otel", "litellm")
default = get_tracer(TracerProvider(), "litellm")
params = {"langfuse_public_key": "pk", "langfuse_secret_key": "sk"}
assert cache.route_for(default, params).detached is True
def run():
set_request_destinations((LANGFUSE_DEST,))
return cache.route_for(default, params)
route = in_fresh_context(run)
assert route.detached is False
assert route.tracer is default
assert route.provider is None
class TestDestinationResolution:
def test_a_langfuse_key_pair_and_host_become_a_destination(self, monkeypatch):
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
is_otel_v2_enabled.cache_clear()
auth = UserAPIKeyAuth(
team_metadata={
"logging": [
{
"callback_name": "langfuse_otel",
"callback_type": "success",
"callback_vars": {
"langfuse_public_key": "pk-team",
"langfuse_secret_key": "sk-team",
"langfuse_host": "http://team.local",
},
}
]
}
)
destinations = resolve_tenant_otel_destinations(auth)
assert [d.endpoint for d in destinations] == ["http://team.local/api/public/otel"]
assert destinations[0].callback_name == "langfuse_otel"
def test_the_key_wins_over_the_team_for_the_same_backend(self, monkeypatch):
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
is_otel_v2_enabled.cache_clear()
def entry(host: str) -> Mapping[str, object]:
return {
"callback_name": "langfuse_otel",
"callback_type": "success",
"callback_vars": {
"langfuse_public_key": "pk",
"langfuse_secret_key": "sk",
"langfuse_host": host,
},
}
auth = UserAPIKeyAuth(
metadata={"logging": [entry("http://key.local")]},
team_metadata={"logging": [entry("http://team.local")]},
)
assert [d.endpoint for d in resolve_tenant_otel_destinations(auth)] == ["http://key.local/api/public/otel"]
def test_nothing_resolves_while_otel_v2_is_off(self, monkeypatch):
monkeypatch.delenv("LITELLM_OTEL_V2", raising=False)
is_otel_v2_enabled.cache_clear()
auth = UserAPIKeyAuth(
team_metadata={
"logging": [
{
"callback_name": "langfuse_otel",
"callback_type": "success",
"callback_vars": {"langfuse_public_key": "pk", "langfuse_secret_key": "sk"},
}
]
}
)
assert resolve_tenant_otel_destinations(auth) == ()
def test_a_host_without_its_key_pair_resolves_to_nothing(self):
assert destination_for("langfuse_otel", {"langfuse_host": "http://team.local"}) is None
def test_a_backend_with_no_dynamic_credentials_has_no_destination(self):
assert "arize_phoenix" not in destination_capable_backends()
assert destination_for("arize_phoenix", {"arize_api_key": "k"}) is None
def test_the_destination_header_string_survives_the_exporter_round_trip(self):
from litellm.integrations.otel.plumbing.providers import parse_headers
destination = destination_for(
"langfuse_otel",
{"langfuse_public_key": "pk", "langfuse_secret_key": "sk", "langfuse_host": "http://x"},
)
assert parse_headers(destination.header_string())["authorization"] == destination.headers["Authorization"]
class TestPresetDegradation:
def test_a_credential_less_langfuse_exports_nowhere_instead_of_to_the_console(self, monkeypatch, capsys):
"""``_normalize`` folds a console exporter in for an empty list, which would
print every span on a proxy whose teams bring their own credentials."""
monkeypatch.delenv("LANGFUSE_PUBLIC_KEY", raising=False)
monkeypatch.delenv("LANGFUSE_SECRET_KEY", raising=False)
config = langfuse_preset(allow_missing_credentials=True)
provider = build_tracer_provider(config, tenant_overrides=True)
capsys.readouterr()
in_fresh_context(emit, provider)
provider.force_flush()
assert capsys.readouterr().out == ""
assert "langfuse" in config.mapper_names
def test_langfuse_still_raises_for_a_global_callback_with_no_credentials(self, monkeypatch):
monkeypatch.delenv("LANGFUSE_PUBLIC_KEY", raising=False)
monkeypatch.delenv("LANGFUSE_SECRET_KEY", raising=False)
with pytest.raises(ValueError, match="LANGFUSE_PUBLIC_KEY"):
langfuse_preset()
def test_a_credential_less_proxy_still_builds_the_v2_logger(self, monkeypatch):
from litellm.litellm_core_utils.litellm_logging import _maybe_construct_otel_v2
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
monkeypatch.delenv("LANGFUSE_PUBLIC_KEY", raising=False)
monkeypatch.delenv("LANGFUSE_SECRET_KEY", raising=False)
is_otel_v2_enabled.cache_clear()
logger = _maybe_construct_otel_v2("langfuse_otel", [])
is_otel_v2_enabled.cache_clear()
assert logger is not None, "team-only deployments must not fall back to the legacy integration"
assert all(spec.requires_headers and not spec.headers for spec in logger.config.exporters)
class TestContextIsolation:
def test_destinations_do_not_leak_between_requests(self):
def first():
set_request_destinations((LANGFUSE_DEST,))
return overridden_backends()
assert in_fresh_context(first) == frozenset({"langfuse_otel"})
assert in_fresh_context(request_destinations) == ()