mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
fix(otel/v2): emit deferred llm-call span when v2 logger is lazily activated
When a v2 callback is instantiated inside the success path (because the destination resolver appended its backend to success_callback after pre_call already iterated the callback list), no carrier exists at close, and the gen-ai span was silently dropped. The pre-existing early-return covered the auth-gate case (no payload, no upstream call) but conflated it with the lazy-activation case (real payload, real destinations naming this backend). Close now falls through to the existing deferred-emit path when the payload is present and the admin-resolved destinations name this backend, so a per-tenant exporter receives the span. The auth-gate semantics are preserved because rejection has no payload. Adds three regression tests around _close_llm_call to pin both halves of the new guard.
This commit is contained in:
parent
ffcea908a6
commit
83a94f78c4
2 changed files with 146 additions and 65 deletions
|
|
@ -52,6 +52,7 @@ from litellm.integrations.otel.model.spans import SpanRole, span_role_for_servic
|
|||
from litellm.integrations.otel.model.utils import to_ns
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from litellm.integrations.otel.model.destination import OtelDestination
|
||||
from litellm.types.utils import (
|
||||
StandardLoggingGuardrailInformation,
|
||||
StandardLoggingPayload,
|
||||
|
|
@ -109,19 +110,13 @@ 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)
|
||||
)
|
||||
self.tracer: Tracer = get_tracer(self._tracer_provider, LITELLM_TRACER_NAME)
|
||||
self._metrics_recorder = self._init_metrics(meter_provider)
|
||||
self._metric_filter_error_logged = False
|
||||
self._emitter = SpanEmitter(
|
||||
self.tracer, self.config, mappers=resolve_mappers(self.config.mapper_names)
|
||||
)
|
||||
self._tenant_tracers = TenantTracerCache(
|
||||
self.config, callback_name, LITELLM_TRACER_NAME
|
||||
)
|
||||
self._emitter = SpanEmitter(self.tracer, self.config, mappers=resolve_mappers(self.config.mapper_names))
|
||||
self._tenant_tracers = TenantTracerCache(self.config, callback_name, LITELLM_TRACER_NAME)
|
||||
self._open_llm_calls: "OrderedDict[str, _LLMCallSpan]" = OrderedDict()
|
||||
self._init_otel_logger_on_litellm_proxy()
|
||||
|
||||
|
|
@ -145,9 +140,7 @@ class OpenTelemetryV2(CustomLogger):
|
|||
|
||||
def _register_in_callback_list(self, callbacks: list) -> None:
|
||||
already_otel = any(
|
||||
cb.__class__.__module__.startswith(_OTEL_MODULES)
|
||||
for cb in callbacks
|
||||
if hasattr(cb, "__class__")
|
||||
cb.__class__.__module__.startswith(_OTEL_MODULES) for cb in callbacks if hasattr(cb, "__class__")
|
||||
)
|
||||
if not already_otel:
|
||||
callbacks.append(self)
|
||||
|
|
@ -174,9 +167,7 @@ class OpenTelemetryV2(CustomLogger):
|
|||
each logger exports only the destinations tagged with its own callback_name,
|
||||
so each backend's span keeps its own attribute vocabulary.
|
||||
"""
|
||||
return tuple(
|
||||
d for d in call.otel_destinations if d.callback_name == self.callback_name
|
||||
)
|
||||
return tuple(d for d in call.otel_destinations if d.callback_name == self.callback_name)
|
||||
|
||||
# ====================================================================== #
|
||||
# LLM-call callbacks — the span is opened at the ``pre_call`` boundary and
|
||||
|
|
@ -225,13 +216,9 @@ class OpenTelemetryV2(CustomLogger):
|
|||
call.provisional_span_name,
|
||||
parent_context=parent_context,
|
||||
start_time_ns=start_time_ns,
|
||||
tracer=self._tenant_tracers.tracer_for(
|
||||
self.tracer, self._destinations_for_backend(call)
|
||||
),
|
||||
tracer=self._tenant_tracers.tracer_for(self.tracer, self._destinations_for_backend(call)),
|
||||
)
|
||||
self._open_llm_calls[call_id] = _LLMCallSpan(
|
||||
span=span, start_time_ns=start_time_ns
|
||||
)
|
||||
self._open_llm_calls[call_id] = _LLMCallSpan(span=span, start_time_ns=start_time_ns)
|
||||
# Evict the oldest open call if the map is over budget. A call that opens
|
||||
# but never closes (a stream that only fires stream events) would linger
|
||||
# otherwise; the evicted span is simply dropped (never exported).
|
||||
|
|
@ -283,9 +270,7 @@ class OpenTelemetryV2(CustomLogger):
|
|||
no boundary to open it at), deduped on the call id by the emitter.
|
||||
"""
|
||||
raw_payload = kwargs.get("standard_logging_object")
|
||||
if not raw_payload or not is_mcp_tool_call(
|
||||
cast(Mapping[str, object], raw_payload)
|
||||
):
|
||||
if not raw_payload or not is_mcp_tool_call(cast(Mapping[str, object], raw_payload)):
|
||||
return False
|
||||
payload = cast("StandardLoggingPayload", raw_payload)
|
||||
data = MCPToolCallSpanData.from_standard_logging_payload(
|
||||
|
|
@ -313,42 +298,61 @@ class OpenTelemetryV2(CustomLogger):
|
|||
) -> Span | None:
|
||||
"""Finish the LLM-call span opened at ``pre_call`` (or create it deferred).
|
||||
|
||||
No carrier for this call id means ``pre_call`` never ran — the request was
|
||||
rejected at the gate or blocked by a pre-call guardrail before any upstream
|
||||
call — so there is nothing to record and no phantom span.
|
||||
Missing carrier has two shapes. ``pre_call`` genuinely never ran -- the
|
||||
request was rejected at the gate or blocked by a pre-call guardrail before
|
||||
any upstream call, so no payload exists and dropping is correct. OR this
|
||||
v2 instance was lazily activated AFTER ``pre_call`` iterated the callback
|
||||
list (the destination-resolver path: a credential resolved a backend the
|
||||
YAML didn't pre-list), so the upstream call DID happen, the payload IS
|
||||
set, and the admin-resolved destinations name this backend -- emit a
|
||||
deferred span with the success event's start time so the per-tenant
|
||||
exporter ships it.
|
||||
"""
|
||||
call = LLMCallEvent.from_dict(kwargs)
|
||||
call_id = call.call_id
|
||||
# ``pop`` is the dedup: this method runs from both the success and failure
|
||||
# paths, and whichever fires first removes the carrier and closes the span.
|
||||
carrier = self._open_llm_calls.pop(call_id, None) if call_id else None
|
||||
if carrier is None:
|
||||
return None
|
||||
payload = call.payload
|
||||
|
||||
if carrier is None:
|
||||
destinations = self._destinations_for_backend(call)
|
||||
if payload is None or not destinations:
|
||||
return None
|
||||
return self._emit_deferred_llm_call(payload, destinations, to_ns(start_time), to_ns(end_time))
|
||||
|
||||
end_time_ns = to_ns(end_time)
|
||||
if payload is None:
|
||||
if carrier.span is not None:
|
||||
# Opened at the boundary but the payload never materialized — end
|
||||
# it (named provisionally) so it isn't leaked as an open span.
|
||||
carrier.span.end(end_time=to_ns(end_time))
|
||||
carrier.span.end(end_time=end_time_ns)
|
||||
return None
|
||||
data = LLMCallSpanData.from_standard_logging_payload(
|
||||
payload, capture_content=self.config.capture_span_content
|
||||
)
|
||||
end_time_ns = to_ns(end_time)
|
||||
|
||||
data = LLMCallSpanData.from_standard_logging_payload(payload, capture_content=self.config.capture_span_content)
|
||||
if carrier.span is not None:
|
||||
# Born at the boundary: stamp attributes from the typed payload, set
|
||||
# status, and end it. Its parent (the server span) was captured at
|
||||
# creation from real ambient context.
|
||||
self._emitter.finish_span(
|
||||
SpanRole.LLM_CALL, carrier.span, data, end_time_ns=end_time_ns
|
||||
)
|
||||
self._emitter.finish_span(SpanRole.LLM_CALL, carrier.span, data, end_time_ns=end_time_ns)
|
||||
return carrier.span
|
||||
# Deferred: ``pre_call`` saw no recordable parent, so create the span now.
|
||||
# The worker copied the request task's context, which carries the anchored
|
||||
# root span — parent to it (ambient fallback on the SDK path). Seed identity
|
||||
# Baggage so the span — and the SDK path, which has none — is labeled
|
||||
# consistently.
|
||||
parent_ctx = resolve_request_span_context()
|
||||
return self._emit_deferred_llm_call(
|
||||
payload,
|
||||
self._destinations_for_backend(call),
|
||||
carrier.start_time_ns,
|
||||
end_time_ns,
|
||||
)
|
||||
|
||||
def _emit_deferred_llm_call(
|
||||
self,
|
||||
payload: "StandardLoggingPayload",
|
||||
destinations: "tuple[OtelDestination, ...]",
|
||||
start_time_ns: int,
|
||||
end_time_ns: int,
|
||||
) -> Span | None:
|
||||
"""Emit an LLM-call span outside the ``pre_call`` boundary.
|
||||
|
||||
Two callers: the SDK thread-pool path (carrier existed but ``pre_call``
|
||||
saw no recordable parent) and the destination-resolver path (this v2
|
||||
instance was born after ``pre_call`` ran, so no carrier was ever opened).
|
||||
Both anchor to the request's root span via the worker-copied context and
|
||||
seed identity Baggage so the span is labeled consistently.
|
||||
"""
|
||||
data = LLMCallSpanData.from_standard_logging_payload(payload, capture_content=self.config.capture_span_content)
|
||||
base_ctx = resolve_request_span_context()
|
||||
bag = promoted_baggage(
|
||||
data.identity,
|
||||
data.request_model,
|
||||
|
|
@ -356,17 +360,14 @@ class OpenTelemetryV2(CustomLogger):
|
|||
metadata_keys=tuple(self.config.baggage_metadata_keys),
|
||||
team_metadata_keys=tuple(self.config.baggage_team_metadata_keys),
|
||||
)
|
||||
if bag:
|
||||
parent_ctx = set_request_baggage(bag, context=parent_ctx)
|
||||
parent_ctx = set_request_baggage(bag, context=base_ctx) if bag else base_ctx
|
||||
return self._emitter.emit(
|
||||
SpanRole.LLM_CALL,
|
||||
data,
|
||||
parent_context=parent_ctx,
|
||||
start_time_ns=carrier.start_time_ns,
|
||||
start_time_ns=start_time_ns,
|
||||
end_time_ns=end_time_ns,
|
||||
tracer=self._tenant_tracers.tracer_for(
|
||||
self.tracer, self._destinations_for_backend(call)
|
||||
),
|
||||
tracer=self._tenant_tracers.tracer_for(self.tracer, destinations),
|
||||
)
|
||||
|
||||
# ====================================================================== #
|
||||
|
|
@ -432,12 +433,7 @@ class OpenTelemetryV2(CustomLogger):
|
|||
# zero-duration root with no context, so skip it. Real background work
|
||||
# (budget/reset jobs, spend flush) passes start/end times and still emits
|
||||
# as a root; anything with a parent emits regardless.
|
||||
if (
|
||||
error_override is None
|
||||
and start_time is None
|
||||
and end_time is None
|
||||
and parent_otel_span is None
|
||||
):
|
||||
if error_override is None and start_time is None and end_time is None and parent_otel_span is None:
|
||||
return None
|
||||
if error_override is not None and data.error is None:
|
||||
data = ServiceSpanData(
|
||||
|
|
@ -583,9 +579,7 @@ def select_global_otel_v2_logger(
|
|||
"""
|
||||
if registered is not None:
|
||||
return registered
|
||||
existing = next(
|
||||
(cb for cb in in_memory_loggers if isinstance(cb, OpenTelemetryV2)), None
|
||||
)
|
||||
existing = next((cb for cb in in_memory_loggers if isinstance(cb, OpenTelemetryV2)), None)
|
||||
return existing if existing is not None else OpenTelemetryV2()
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -538,6 +538,93 @@ def test_synthetic_error_log_produces_no_llm_span():
|
|||
assert "auth /chat/completions" in names # auth span itself still recorded
|
||||
|
||||
|
||||
def test_lazy_activation_emits_llm_span_when_destination_resolves(monkeypatch):
|
||||
"""LIT-3850 lazy-activation seam: a v2 instance born inside the success path
|
||||
(because the destination resolver appended its backend to ``success_callback``
|
||||
on this request) was not in the callback list when ``pre_call`` iterated, so
|
||||
no carrier was opened. The close must still emit the gen-ai span when the
|
||||
payload is present and the admin-resolved destinations name this backend, so
|
||||
the per-tenant exporter ships it. Without the fallthrough, this test sees
|
||||
zero spans."""
|
||||
logger, exporter = _logger()
|
||||
monkeypatch.setattr(logger, "callback_name", "in_memory")
|
||||
server = logger._emitter.start_span(
|
||||
SpanRole.PROXY_REQUEST, LITELLM_PROXY_REQUEST_SPAN_NAME
|
||||
)
|
||||
set_request_root_span(server)
|
||||
# Make ``tracer_for`` return the logger's default tracer regardless of the
|
||||
# destinations passed: the test asserts the close-path emitted the span, not
|
||||
# that the per-destination provider clone wired up an OTLP exporter (the
|
||||
# routing cache's job, covered separately). The default tracer is bound to
|
||||
# the in-memory exporter so the test can read the result.
|
||||
tracer_for_calls: list[tuple] = []
|
||||
|
||||
def _fake_tracer_for(default, destinations):
|
||||
tracer_for_calls.append(destinations)
|
||||
return default
|
||||
|
||||
monkeypatch.setattr(logger._tenant_tracers, "tracer_for", _fake_tracer_for)
|
||||
kwargs = _kwargs()
|
||||
kwargs["standard_callback_dynamic_params"] = {
|
||||
"otel_destinations": [
|
||||
{
|
||||
"callback_name": "in_memory",
|
||||
"endpoint": "https://otlp.example.com/v1",
|
||||
"headers": {"api_key": "k"},
|
||||
}
|
||||
]
|
||||
}
|
||||
assert "call_1" not in logger._open_llm_calls # no carrier opened
|
||||
asyncio.run(logger.async_log_success_event(kwargs, None, None, None))
|
||||
server.end()
|
||||
names = [s.name for s in exporter.get_finished_spans()]
|
||||
assert "chat gpt-4o" in names
|
||||
# ``tracer_for`` was invoked with exactly the resolved destination, proving
|
||||
# the deferred path used per-tenant routing rather than the default tracer
|
||||
# blindly.
|
||||
assert len(tracer_for_calls) == 1
|
||||
(dests,) = tracer_for_calls
|
||||
assert len(dests) == 1 and dests[0].endpoint == "https://otlp.example.com/v1"
|
||||
|
||||
|
||||
def test_close_without_carrier_and_without_destination_drops_silently():
|
||||
"""The pre-existing early-return semantics (auth gate / pre-call guardrail
|
||||
rejection with no destination resolving to this backend) must be preserved:
|
||||
no phantom span. The fix only widens emit-on-close when the admin-resolved
|
||||
destinations name this backend AND the payload exists."""
|
||||
logger, exporter = _logger()
|
||||
server = logger._emitter.start_span(
|
||||
SpanRole.PROXY_REQUEST, LITELLM_PROXY_REQUEST_SPAN_NAME
|
||||
)
|
||||
set_request_root_span(server)
|
||||
asyncio.run(logger.async_log_success_event(_kwargs(), None, None, None))
|
||||
server.end()
|
||||
names = [s.name for s in exporter.get_finished_spans()]
|
||||
assert "chat gpt-4o" not in names
|
||||
|
||||
|
||||
def test_close_without_carrier_drops_when_payload_missing(monkeypatch):
|
||||
"""No carrier + no payload = the auth-gate rejection case (no upstream call
|
||||
happened). Must drop even when destinations resolve, so a phantom span is
|
||||
never emitted for a request the gate refused."""
|
||||
logger, exporter = _logger()
|
||||
monkeypatch.setattr(logger, "callback_name", "in_memory")
|
||||
kwargs = {
|
||||
"litellm_params": {"metadata": {}},
|
||||
"standard_callback_dynamic_params": {
|
||||
"otel_destinations": [
|
||||
{
|
||||
"callback_name": "in_memory",
|
||||
"endpoint": "https://otlp.example.com/v1",
|
||||
"headers": {},
|
||||
}
|
||||
]
|
||||
},
|
||||
}
|
||||
asyncio.run(logger.async_log_success_event(kwargs, None, None, None))
|
||||
assert exporter.get_finished_spans() == ()
|
||||
|
||||
|
||||
def test_create_request_started_span_captures_anchor():
|
||||
"""``create_litellm_proxy_request_started_span`` doubles as the anchor capture
|
||||
point: the active server span becomes the request root for later spans."""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue