fix(otel v2): preserve operator spans on destination failure

This commit is contained in:
Yucheng He 2026-09-05 01:45:26 -07:00
parent 594340cdc1
commit 29765dc8c7
5 changed files with 362 additions and 23 deletions

View file

@ -198,6 +198,11 @@ class OpenTelemetryV2(CustomLogger):
self._open_llm_calls: OrderedDict[str, _LLMCallSpan] = OrderedDict()
self._init_otel_logger_on_litellm_proxy()
@property
def tracer_provider(self) -> TracerProvider:
"""The provider this logger emits through, read-only to its callers."""
return self._tracer_provider
def _init_metrics(self, meter_provider: "MeterProvider | None") -> "GenAIMetricRecorder | None":
"""Create the six GenAI histograms when metrics are enabled, else ``None``.
@ -872,8 +877,8 @@ def publish_global_otel_v2_provider(
through; see :func:`attach_tenant_fan_out`.
"""
logger: Final = select_global_otel_v2_logger(in_memory_loggers, registered=registered)
attach_tenant_fan_out(logger._tracer_provider, logger.config)
set_global_provider(logger._tracer_provider)
attach_tenant_fan_out(logger.tracer_provider, logger.config)
set_global_provider(logger.tracer_provider)
return logger

View file

@ -8,7 +8,7 @@ from collections.abc import Callable, Iterable, Mapping
from types import MappingProxyType
from typing import TYPE_CHECKING, Any, Final, Literal
from opentelemetry import _logs, baggage, metrics
from opentelemetry import _logs, baggage, metrics, trace
from opentelemetry._events import EventLogger
from opentelemetry._logs import LoggerProvider, NoOpLoggerProvider
from opentelemetry.context import Context
@ -259,17 +259,28 @@ class _DrainPool:
worker.start()
def submit(self, processor: SpanProcessor) -> None:
"""Queue ``processor`` for closing, or close it here once the pool is retired.
"""Queue ``processor`` for closing, or hand it off once the pool is retired.
The check and the put share one lock. Reading a closed flag on its own leaves
room for :meth:`close` to run in between, and the processor would land behind
the sentinels every worker has already exited on.
Past close there is no worker left to take it, and the caller is whichever
thread just ended a span, so closing it inline would park that thread on a
network flush the shutdown deadline has already stopped waiting for. The extra
thread is bounded by the same close: the fan-out stops handing processors out
at that point, so only the ones already exporting when it happened arrive here.
"""
with self._lock:
if not self._closed:
self._pending.put(processor)
return
_shutdown_quietly(processor)
threading.Thread(
target=_shutdown_quietly,
args=(processor,),
daemon=True,
name="litellm-otel-destination-drain-straggler",
).start()
def close(self, timeout: float | None = None) -> None:
"""Retire the workers once they have closed everything already queued.
@ -440,6 +451,29 @@ class TenantFanOutSpanProcessor(SpanProcessor):
except Exception: # noqa: BLE001 # one exporter's flush failure must not fail the whole flush
return False
def deliverable(self, destinations: Iterable["OtelDestination"]) -> tuple["OtelDestination", ...]:
"""The subset of ``destinations`` this fan-out can actually export to.
A destination whose exporter will not build (a protocol whose package is not
installed, a malformed endpoint) has to be dropped before the request anchors
it, not when its first span ends. By then the operator's own exporter has been
told to hold that backend's spans back for this request, so dropping there
loses the span outright instead of leaving it where it would have gone with no
override at all.
"""
return tuple(destination for destination in destinations if self._buildable(destination))
def _buildable(self, destination: "OtelDestination") -> bool:
"""Whether a processor for ``destination`` exists or can be built right now."""
with self._lock:
if self._closed:
return False
built: Final = self._cached_or_built_locked(destination)
drained: Final = self._drainable_locked()
for shed in drained:
self._drain.submit(shed)
return built is not None
def _acquire(self, destination: "OtelDestination") -> SpanProcessor | None:
"""The processor for ``destination``, marked busy until ``_release``.
@ -448,13 +482,10 @@ class TenantFanOutSpanProcessor(SpanProcessor):
thread with all but the winner shed. Building an exporter opens no connection,
so the cost of holding the lock is a constructor, once per destination.
"""
key: Final = destination.cache_key()
with self._lock:
if self._closed:
return None
if (cached := self._processors.get(key)) is not None:
self._processors.move_to_end(key)
processor: Final = cached if cached is not None else self._build_locked(destination, key)
processor: Final = self._cached_or_built_locked(destination)
if processor is None:
return None
self._exporting[id(processor)] = self._exporting.get(id(processor), 0) + 1
@ -463,6 +494,13 @@ class TenantFanOutSpanProcessor(SpanProcessor):
self._drain.submit(shed)
return processor
def _cached_or_built_locked(self, destination: "OtelDestination") -> SpanProcessor | None:
key: Final = destination.cache_key()
if (cached := self._processors.get(key)) is not None:
self._processors.move_to_end(key)
return cached
return self._build_locked(destination, key)
def _build_locked(self, destination: "OtelDestination", key: object) -> SpanProcessor | None:
built: Final = self._build(destination)
if built is None:
@ -853,19 +891,50 @@ def attach_tenant_fan_out(provider: TracerProvider, config: OpenTelemetryV2Confi
provider.add_span_processor(TenantFanOutSpanProcessor(operator_sinks=operator_sink_keys(config)))
def deliverable_destinations(
destinations: Iterable["OtelDestination"],
provider: trace.TracerProvider | None = None,
) -> tuple["OtelDestination", ...]:
"""The destinations a request can anchor, given what is published to carry them.
Anchoring a destination is what tells the operator's own exporter to stand down
for that backend, so one nothing can deliver has to be dropped here: with no
fan-out attached, or with an exporter that will not build, the request keeps
exactly the routing it would have had without any override.
"""
fan_out: Final = next(
(
processor
for processor in _attached_processors(provider if provider is not None else trace.get_tracer_provider())
if isinstance(processor, TenantFanOutSpanProcessor)
),
None,
)
return fan_out.deliverable(destinations) if fan_out is not None else ()
def operator_sink_keys(config: OpenTelemetryV2Config | None) -> frozenset[_SinkKey]:
"""The accounts the operator's own exporters write to, in destination terms.
An exporter with no endpoint of its own resolves one from the environment at
export time, so it has no comparable identity and is left out.
export time, so it has no comparable identity and is left out, and so is one
that never reaches the wire: a console kind ignores the endpoint, and a
header-gated spec with no credentials is skipped when the provider is built.
"""
if config is None:
return frozenset()
return frozenset(
key for spec in config.exporters if (key := _sink_key(spec.endpoint, parse_headers(spec.headers))) is not None
key
for spec in config.exporters
if _exports_to_the_wire(spec) and (key := _sink_key(spec.endpoint, parse_headers(spec.headers))) is not None
)
def _exports_to_the_wire(spec: ExporterSpec) -> bool:
"""Whether ``build_tracer_provider`` gives ``spec`` an exporter that sends OTLP."""
return exporter_transport(spec.kind) != "headerless" and not (spec.requires_headers and not spec.headers)
def _sink_key(endpoint: str | None, headers: Mapping[str, str]) -> "_SinkKey | None":
"""The account an exporter writes to, or ``None`` when it has no fixed one.
@ -886,7 +955,7 @@ def _credential_name(header: str) -> str:
return _CREDENTIAL_ALIASES.get(normalized, normalized)
def _attached_processors(provider: TracerProvider) -> "tuple[SpanProcessor, ...]":
def _attached_processors(provider: trace.TracerProvider) -> "tuple[SpanProcessor, ...]":
"""The processors already on ``provider``, or empty when the SDK hides them."""
multi: Final = getattr(provider, "_active_span_processor", None)
return tuple(getattr(multi, "_span_processors", ()))

View file

@ -200,6 +200,7 @@ if TYPE_CHECKING:
from mcp.types import EmbeddedResource, ImageContent, TextContent
from litellm.integrations.otel.logger import OpenTelemetryV2
from litellm.integrations.otel.model.config import OpenTelemetryV2Config
from litellm.llms.base_llm.passthrough.transformation import BasePassthroughConfig
try:
from litellm_enterprise.enterprise_callbacks.callback_controls import (
@ -4800,27 +4801,41 @@ def _maybe_construct_otel_v2(callback_name: str, _in_memory_loggers: list[Custom
Returns ``None`` when V2 is off OR when there's no preset registered for
``callback_name`` — callers should then fall through to the legacy path.
A preset that needs operator credentials it cannot find is allowed to build
anyway, exporting nowhere, only while this request has a key/team destination
for that backend: the exporter-less logger exists to let the fan-out carry those
spans without a second detached copy. With no such destination the preset raises
as it always did and the caller falls through to the legacy path, so the proxy
never publishes a provider that exports nowhere for a backend the operator
configured and no tenant can use.
"""
from litellm.integrations.otel.model.config import is_otel_v2_enabled
if not is_otel_v2_enabled():
return None
from litellm.integrations.otel.logger import OpenTelemetryV2, build_otel_v2_logger
from litellm.integrations.otel.plumbing.context import destination_backends
from litellm.integrations.otel.presets import PRESET_BY_CALLBACK
preset_fn: Final = PRESET_BY_CALLBACK.get(callback_name)
if preset_fn is None:
return None
serves_a_destination: Final = callback_name in destination_backends()
for callback in _in_memory_loggers:
if isinstance(callback, OpenTelemetryV2) and getattr(callback, "callback_name", None) == callback_name:
if (
isinstance(callback, OpenTelemetryV2)
and getattr(callback, "callback_name", None) == callback_name
and (serves_a_destination or not _exports_nowhere(callback.config))
):
return callback
try:
config: Final = preset_fn(allow_missing_credentials=True)
config: Final = preset_fn(allow_missing_credentials=serves_a_destination)
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):
if _exports_nowhere(config):
verbose_logger.warning(
"OTel V2: no operator credentials for '%s'; only key/team destinations will receive its traces",
callback_name,
@ -4830,6 +4845,11 @@ def _maybe_construct_otel_v2(callback_name: str, _in_memory_loggers: list[Custom
return v2_logger
def _exports_nowhere(config: "OpenTelemetryV2Config") -> bool:
"""Whether every exporter in ``config`` is waiting on credentials it never got."""
return all(spec.requires_headers and not spec.headers for spec in config.exporters)
def _maybe_auto_initialize_arize_phoenix(_in_memory_loggers: list[CustomLogger]) -> None:
"""
Auto-initialize ArizePhoenixLogger when Phoenix env vars are detected.

View file

@ -2852,17 +2852,23 @@ def _seed_request_destinations(user_api_key_dict: UserAPIKeyAuth) -> None:
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.
Only destinations the published fan-out can build are anchored. Anchoring one is
what tells the operator's exporter to hold that backend's spans back under
``override``, so an unbuildable one would leave the span with nowhere to go.
The ``postgres`` spans under ``auth`` close before this runs, because they are the
reads that resolve the identity being read here. They never reach the tenant's
account, and they are never withheld from the operator's backend, whichever mode
is set.
"""
try:
from litellm.integrations.otel.plumbing.context import set_request_destinations
from litellm.integrations.otel.plumbing.providers import deliverable_destinations
from litellm.proxy.litellm_pre_call_utils import (
resolve_tenant_otel_destinations,
)
set_request_destinations(resolve_tenant_otel_destinations(user_api_key_dict))
set_request_destinations(deliverable_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)

View file

@ -32,6 +32,7 @@ from litellm.integrations.otel.plumbing.providers import (
_OverriddenBackendFilter,
_sink_key,
build_tracer_provider,
deliverable_destinations,
operator_sink_keys,
)
from litellm.integrations.otel.plumbing.routing import TenantTracerCache, get_tracer
@ -322,6 +323,48 @@ class TestRoutingMode:
assert operator_sink_keys(config) == frozenset({self.OPERATOR_SINK})
def test_operator_sink_keys_skips_exporters_that_never_reach_the_wire(self):
"""A console kind ignores the endpoint and a header-gated spec with no
credentials is dropped when the provider is built, so treating either as an
account the operator writes to would silently withhold a team's own spans
under additive."""
config = OpenTelemetryV2Config(
exporters=(
ExporterSpec(kind="otlp_http", endpoint=self.OPERATOR_SINK[0], headers="authorization=Basic op"),
ExporterSpec(kind="console", endpoint="http://team.local/v1/traces"),
ExporterSpec(kind="otlp_http", endpoint="http://gated.local/v1/traces", requires_headers=True),
)
)
assert operator_sink_keys(config) == frozenset({self.OPERATOR_SINK})
def test_a_team_pointing_at_a_credential_less_operator_exporter_still_gets_its_spans(self, monkeypatch):
"""Under additive the fan-out skips a destination the operator already writes
to. An exporter the provider never built writes nothing, so skipping it would
cost the team every span."""
monkeypatch.setenv("LITELLM_OTEL_TENANT_DESTINATION_MODE", "additive")
gated_endpoint = "http://gated.local/v1/traces"
destination = OtelDestination(endpoint=gated_endpoint, callback_name="newrelic")
config = OpenTelemetryV2Config(
exporters=(ExporterSpec(kind="otlp_http", endpoint=gated_endpoint, requires_headers=True),)
)
dest_exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(
TenantFanOutSpanProcessor(
processor_factory=lambda _d: SimpleSpanProcessor(dest_exporter),
operator_sinks=operator_sink_keys(config),
)
)
def run():
set_request_destinations((destination,))
emit(provider)
in_fresh_context(run)
assert [s.name for s in dest_exporter.get_finished_spans()] == ["chat gpt-4"]
def test_the_operators_own_langfuse_and_a_team_naming_it_are_one_account(self, monkeypatch):
"""The two sides are built by different code that writes the endpoint and the
header names differently, so comparing them raw silently never matches."""
@ -473,6 +516,82 @@ class TestFanOut:
assert attempts == [LANGFUSE_DEST.endpoint]
assert reached_the_end == [True]
def test_an_unbuildable_destination_leaves_the_span_with_the_operator(self):
"""Anchoring the destination is what makes the operator's exporter stand down
for the backend, so a destination nothing can deliver to must never be anchored,
or the span reaches neither account."""
global_exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(_OverriddenBackendFilter(SimpleSpanProcessor(global_exporter), "langfuse_otel"))
provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=lambda _d: None))
def run():
set_request_destinations(deliverable_destinations((LANGFUSE_DEST,), provider))
emit(provider)
return request_destinations()
anchored = in_fresh_context(run)
assert anchored == ()
assert [s.name for s in global_exporter.get_finished_spans()] == ["chat gpt-4"]
def test_a_buildable_destination_is_still_anchored_and_still_overrides(self):
global_exporter, dest_exporter = InMemorySpanExporter(), InMemorySpanExporter()
provider = wired_provider(dest_exporter, global_exporter)
def run():
set_request_destinations(deliverable_destinations((LANGFUSE_DEST,), provider))
emit(provider)
return request_destinations()
anchored = in_fresh_context(run)
assert anchored == (LANGFUSE_DEST,)
assert global_exporter.get_finished_spans() == ()
assert [s.name for s in dest_exporter.get_finished_spans()] == ["chat gpt-4"]
def test_only_the_unbuildable_destination_is_dropped_from_a_mixed_set(self):
dest_exporter = InMemorySpanExporter()
other = LANGFUSE_DEST.model_copy(update={"endpoint": "http://broken.local/otel"})
fan_out = TenantFanOutSpanProcessor(
processor_factory=lambda d: None if d.endpoint == other.endpoint else SimpleSpanProcessor(dest_exporter)
)
assert fan_out.deliverable((other, LANGFUSE_DEST)) == (LANGFUSE_DEST,)
def test_no_fan_out_means_nothing_is_anchored(self):
"""With nothing to carry the spans to the tenant, anchoring would only stop the
operator's exporter from writing them."""
provider = TracerProvider()
assert deliverable_destinations((LANGFUSE_DEST,), provider) == ()
def test_a_closed_fan_out_anchors_nothing(self):
fan_out = TenantFanOutSpanProcessor(processor_factory=lambda _d: SimpleSpanProcessor(InMemorySpanExporter()))
provider = TracerProvider()
provider.add_span_processor(fan_out)
fan_out.shutdown()
assert deliverable_destinations((LANGFUSE_DEST,), provider) == ()
def test_the_processor_built_to_check_deliverability_is_the_one_that_exports(self):
built = []
def factory(_destination):
built.append(SimpleSpanProcessor(InMemorySpanExporter()))
return built[-1]
provider = TracerProvider()
provider.add_span_processor(TenantFanOutSpanProcessor(processor_factory=factory))
def run():
set_request_destinations(deliverable_destinations((LANGFUSE_DEST,), provider))
emit(provider)
in_fresh_context(run)
assert len(built) == 1
def test_one_processor_is_reused_across_spans_of_the_same_destination(self):
built = []
@ -750,18 +869,91 @@ class TestPresetDegradation:
with pytest.raises(ValueError, match="LANGFUSE_PUBLIC_KEY"):
langfuse_preset()
def test_a_credential_less_proxy_still_builds_the_v2_logger(self, monkeypatch):
def test_a_credential_less_proxy_builds_the_v2_logger_for_a_team_destination(self, monkeypatch):
from litellm.litellm_core_utils.litellm_logging import _maybe_construct_otel_v2
credential_less_proxy(monkeypatch)
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
def run():
set_request_destinations((LANGFUSE_DEST,))
return _maybe_construct_otel_v2("langfuse_otel", [])
is_otel_v2_enabled.cache_clear()
logger = in_fresh_context(run)
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)
def test_a_credential_less_proxy_with_no_destinations_falls_back_to_the_legacy_path(self, monkeypatch):
"""Nothing can use a credential-less langfuse here, so the operator has to get
the same story as before v2: the legacy integration, not a global provider
that exports nowhere."""
from litellm.litellm_core_utils.litellm_logging import _maybe_construct_otel_v2
credential_less_proxy(monkeypatch)
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
is_otel_v2_enabled.cache_clear()
logger = _maybe_construct_otel_v2("langfuse_otel", [])
logger = in_fresh_context(_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)
assert logger is None
def test_a_destination_for_one_backend_does_not_degrade_another(self, monkeypatch):
from litellm.litellm_core_utils.litellm_logging import _maybe_construct_otel_v2
credential_less_proxy(monkeypatch)
monkeypatch.delenv("WANDB_API_KEY", raising=False)
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
def run():
set_request_destinations((LANGFUSE_DEST,))
return _maybe_construct_otel_v2("weave_otel", [])
is_otel_v2_enabled.cache_clear()
logger = in_fresh_context(run)
is_otel_v2_enabled.cache_clear()
assert logger is None
def test_the_exporter_less_logger_is_not_reused_by_a_request_without_destinations(self, monkeypatch):
"""Reusing it would let one team's destination decide how every later request
without one is logged, long after the degrade was justified."""
from litellm.litellm_core_utils.litellm_logging import _maybe_construct_otel_v2
credential_less_proxy(monkeypatch)
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
loggers = []
def with_destination():
set_request_destinations((LANGFUSE_DEST,))
return _maybe_construct_otel_v2("langfuse_otel", loggers)
is_otel_v2_enabled.cache_clear()
degraded = in_fresh_context(with_destination)
plain = in_fresh_context(_maybe_construct_otel_v2, "langfuse_otel", loggers)
is_otel_v2_enabled.cache_clear()
assert degraded is not None
assert plain is None
def test_a_credentialed_logger_is_still_reused_across_requests(self, monkeypatch):
from litellm.litellm_core_utils.litellm_logging import _maybe_construct_otel_v2
monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-lf-1")
monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-lf-1")
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
loggers = []
is_otel_v2_enabled.cache_clear()
first = in_fresh_context(_maybe_construct_otel_v2, "langfuse_otel", loggers)
second = in_fresh_context(_maybe_construct_otel_v2, "langfuse_otel", loggers)
is_otel_v2_enabled.cache_clear()
assert first is not None
assert second is first
class TestContextIsolation:
@ -991,6 +1183,21 @@ class TestEvictionSafety:
assert held.shutdown_calls == 1
def test_a_recently_used_destination_is_not_the_one_evicted(self):
"""Without the refresh the cache sheds by insertion order, so the busiest
destination is the one whose exporter is rebuilt on every overflow."""
from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS
fan_out, built = self._fan_out()
for index in range(_MAX_CACHED_DESTINATION_PROCESSORS):
fan_out._release(fan_out._acquire(self._dest(index)))
fan_out._release(fan_out._acquire(self._dest(0)))
fan_out._release(fan_out._acquire(self._dest(_MAX_CACHED_DESTINATION_PROCESSORS)))
self._settle(fan_out, built[1])
assert built[1].shutdown_calls == 1
assert built[0].shutdown_calls == 0, "the destination used most recently was the one shed"
def test_an_idle_evicted_processor_is_closed_off_the_export_path(self):
from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS
@ -1179,8 +1386,40 @@ class TestEvictionSafety:
fan_out.shutdown()
fan_out._drain.submit(stray)
for _ in range(500):
if stray.shutdown_calls:
break
time.sleep(0.02)
assert stray.shutdown_calls == 1
def test_releasing_a_straggler_after_shutdown_does_not_block_the_span_thread(self):
"""The teardown deadline has already expired by then, so closing the straggler
inline would park whichever thread just ended a span on the very flush the
deadline gave up waiting for."""
import threading
never = threading.Event()
class Stuck(self.Recording):
def shutdown(self):
never.wait()
def factory(_destination):
return Stuck()
fan_out = TenantFanOutSpanProcessor(processor_factory=factory, shutdown_drain_seconds=0.05)
held = fan_out._acquire(self._dest(0))
fan_out.shutdown()
released = threading.Event()
caller = threading.Thread(target=lambda: (fan_out._release(held), released.set()), daemon=True)
caller.start()
came_back = released.wait(timeout=5)
never.set()
assert came_back, "the thread that ended the span was left holding a stuck teardown"
def test_shutdown_waits_out_an_export_that_lands_inside_the_bound(self):
"""Without the wait the closing is left to a daemon thread, which the
interpreter can retire before it runs, so the last spans never reach the