fix(otel): populate Langfuse root observation IO

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Devin AI 2026-08-27 06:49:32 +00:00
parent 807ee7f232
commit 88e44717a4
3 changed files with 70 additions and 9 deletions

View file

@ -24,6 +24,7 @@ from litellm._logging import verbose_logger
from litellm.integrations.custom_logger import CustomLogger
from litellm.integrations.otel.emitter import SpanEmitter, stamp_error
from litellm.integrations.otel.mappers import resolve_mappers
from litellm.integrations.otel.mappers.langfuse import LangfuseMapper
from litellm.integrations.otel.model.baggage import promoted_baggage
from litellm.integrations.otel.model.config import OpenTelemetryV2Config
from litellm.integrations.otel.model.metadata import (
@ -550,6 +551,7 @@ class OpenTelemetryV2(CustomLogger):
# logging object) cannot re-emit through the deferred branch.
self._emitter.mark_emitted(call_id, SpanRole.LLM_CALL)
self._emitter.finish_span(SpanRole.LLM_CALL, carrier.span, data, end_time_ns=end_time_ns)
self._stamp_langfuse_root_observation(data, carrier.span)
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
@ -561,7 +563,7 @@ class OpenTelemetryV2(CustomLogger):
parent_ctx: Final = self._seed_identity_baggage(
data.identity, data.request_model, resolve_request_span_context()
)
return self._emitter.emit(
span: Final = self._emitter.emit(
SpanRole.LLM_CALL,
data,
parent_context=(set_span_in_context(INVALID_SPAN, parent_ctx) if route.detached else parent_ctx),
@ -570,9 +572,22 @@ class OpenTelemetryV2(CustomLogger):
tracer=route.tracer,
links=_request_trace_links(parent_ctx) if route.detached else None,
)
self._stamp_langfuse_root_observation(data, span)
return span
finally:
self._tenant_tracers.release(route.provider)
def _stamp_langfuse_root_observation(self, data: LLMCallSpanData, llm_span: Span | None) -> None:
if "langfuse" not in self.config.mapper_names or llm_span is None:
return
root: Final = request_root_span()
if root is None or not root.is_recording():
return
if root.get_span_context().trace_id != llm_span.get_span_context().trace_id:
return
for key, value in LangfuseMapper.root_observation_io(data).items():
root.set_attribute(key, value)
# ====================================================================== #
# Service hooks
# ====================================================================== #

View file

@ -77,3 +77,11 @@ class LangfuseMapper:
**collect(cls._LLM_CALL_ATTRS, data),
**collect(cls._BLOB_ATTRS, data),
}
@classmethod
def root_observation_io(cls, data: LLMCallSpanData) -> AttributeMap:
return {
key: value
for key, value in collect(cls._BLOB_ATTRS, data).items()
if key in ("langfuse.observation.input", "langfuse.observation.output")
}

View file

@ -9,8 +9,8 @@ hooks, proxy SERVER span lifecycle (start + setters), parent-context resolution
import asyncio
import contextlib
import os
from unittest.mock import patch
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
import pytest
@ -28,6 +28,16 @@ from litellm.integrations.otel import ( # noqa: E402
LiteLLM,
OpenTelemetryV2Config,
)
from litellm.integrations.otel.logger import OpenTelemetryV2 # noqa: E402
from litellm.integrations.otel.model.config import ( # noqa: E402
CaptureMessageContent,
ExporterSpec,
)
from litellm.integrations.otel.model.spans import ( # noqa: E402
LITELLM_PROXY_REQUEST_SPAN_NAME,
SpanRole,
)
from litellm.integrations.otel.model.utils import to_ns, to_seconds # noqa: E402
from litellm.integrations.otel.plumbing import providers # noqa: E402
from litellm.integrations.otel.plumbing.context import ( # noqa: E402
reset_mcp_message_trace_carrier,
@ -36,13 +46,6 @@ from litellm.integrations.otel.plumbing.context import ( # noqa: E402
set_mcp_message_transport_span,
set_request_root_span,
)
from litellm.integrations.otel.logger import OpenTelemetryV2 # noqa: E402
from litellm.integrations.otel.model.config import ExporterSpec # noqa: E402
from litellm.integrations.otel.model.spans import ( # noqa: E402
LITELLM_PROXY_REQUEST_SPAN_NAME,
SpanRole,
)
from litellm.integrations.otel.model.utils import to_ns, to_seconds # noqa: E402
# --------------------------------------------------------------------------- #
# Fixtures
@ -959,6 +962,41 @@ def test_llm_span_parents_to_ambient_server_span():
assert llm_span.parent.span_id == server.get_span_context().span_id
def test_langfuse_mapper_stamps_input_output_on_root_observation():
cfg = OpenTelemetryV2Config(
exporter="in_memory",
legacy_compat=False,
mapper_names=["langfuse"],
capture_message_content=CaptureMessageContent.SPAN_ONLY,
)
exporter = InMemorySpanExporter()
tracer_provider = providers.build_tracer_provider(cfg, exporter=exporter)
logger = OpenTelemetryV2(config=cfg, tracer_provider=tracer_provider)
server = logger._emitter.start_span(SpanRole.PROXY_REQUEST, LITELLM_PROXY_REQUEST_SPAN_NAME)
set_request_root_span(server)
payload = _payload(
messages=[{"role": "user", "content": "Say hi"}],
response={
"id": "resp_1",
"model": "gpt-4o",
"choices": [
{
"finish_reason": "stop",
"message": {"role": "assistant", "content": "Hi"},
}
],
},
)
_emit_llm(logger, _kwargs(payload=payload))
server.end()
by_name = {span.name: span for span in exporter.get_finished_spans()}
root_attrs = by_name[LITELLM_PROXY_REQUEST_SPAN_NAME].attributes
assert root_attrs["langfuse.observation.input"] == ('[{"role": "user", "content": "Say hi"}]')
assert root_attrs["langfuse.observation.output"] == ('[{"role": "assistant", "content": "Hi"}]')
def test_llm_span_is_root_without_ambient_server_span():
"""No server span at ``pre_call`` → creation is deferred and the span is a
root of its own trace (the SDK / no-proxy path)."""