mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-14 23:21:35 +00:00
fix(otel v2): name Langfuse traces from the langfuse_trace_name header or metadata.trace_name (#40793)
* fix(otel v2): name Langfuse traces from the langfuse_trace_name header or metadata.trace_name Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(otel v2): type the named-request helper in the Langfuse logger tests Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: yucheng <yucheng@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
dab7f6a86a
commit
bf146e2cac
8 changed files with 160 additions and 8 deletions
|
|
@ -3,7 +3,12 @@ from typing import TYPE_CHECKING, Final
|
|||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.otel.logger import OpenTelemetryV2
|
||||
from litellm.integrations.otel.mappers.langfuse import LANGFUSE_OBSERVATION_INPUT, LANGFUSE_OBSERVATION_OUTPUT
|
||||
from litellm.integrations.otel.mappers.langfuse import (
|
||||
LANGFUSE_OBSERVATION_INPUT,
|
||||
LANGFUSE_OBSERVATION_OUTPUT,
|
||||
LANGFUSE_TRACE_NAME,
|
||||
)
|
||||
from litellm.integrations.otel.model.metadata import caller_trace_name
|
||||
from litellm.integrations.otel.model.request_io import request_input, response_output, stream_output
|
||||
from litellm.integrations.otel.plumbing.context import request_root_span
|
||||
|
||||
|
|
@ -13,6 +18,18 @@ if TYPE_CHECKING:
|
|||
|
||||
|
||||
class LangfuseOpenTelemetryV2(OpenTelemetryV2):
|
||||
"""Names the trace from the request. Langfuse reads ``langfuse.trace.name`` off the root observation,
|
||||
and the proxy's root span is still recording when the LLM call starts."""
|
||||
|
||||
def log_pre_api_call(self, model: str, messages: object, kwargs: Mapping[str, object]) -> None:
|
||||
root: Final = request_root_span()
|
||||
name: Final = caller_trace_name(kwargs)
|
||||
if root is not None and root.is_recording() and name is not None:
|
||||
root.set_attribute(LANGFUSE_TRACE_NAME, name)
|
||||
super().log_pre_api_call(model, messages, kwargs)
|
||||
|
||||
|
||||
class LangfuseContentOpenTelemetryV2(LangfuseOpenTelemetryV2):
|
||||
"""Stamps the request's input and output on the root observation while it is still recording.
|
||||
|
||||
Langfuse shows a trace's input and output from its root observation. The proxy's root span ends
|
||||
|
|
|
|||
|
|
@ -554,6 +554,7 @@ class OpenTelemetryV2(CustomLogger):
|
|||
capture_content=self.config.capture_span_content,
|
||||
time_to_first_chunk_seconds=call.time_to_first_chunk_seconds,
|
||||
request_route=request_root_http_route(),
|
||||
trace_name=call.trace_name,
|
||||
)
|
||||
end_time_ns: Final = to_ns(end_time)
|
||||
if carrier is not None and carrier.span is not None:
|
||||
|
|
@ -984,8 +985,8 @@ def build_otel_v2_logger(
|
|||
|
||||
|
||||
def _logger_class(config: OpenTelemetryV2Config) -> type[OpenTelemetryV2]:
|
||||
if "langfuse" not in config.mapper_names or not config.capture_span_content:
|
||||
if "langfuse" not in config.mapper_names:
|
||||
return OpenTelemetryV2
|
||||
from litellm.integrations.otel.langfuse_logger import LangfuseOpenTelemetryV2
|
||||
from litellm.integrations.otel.langfuse_logger import LangfuseContentOpenTelemetryV2, LangfuseOpenTelemetryV2
|
||||
|
||||
return LangfuseOpenTelemetryV2
|
||||
return LangfuseContentOpenTelemetryV2 if config.capture_span_content else LangfuseOpenTelemetryV2
|
||||
|
|
|
|||
|
|
@ -28,6 +28,7 @@ from litellm.integrations.otel.model.payloads import (
|
|||
|
||||
LANGFUSE_OBSERVATION_INPUT: Final = "langfuse.observation.input"
|
||||
LANGFUSE_OBSERVATION_OUTPUT: Final = "langfuse.observation.output"
|
||||
LANGFUSE_TRACE_NAME: Final = "langfuse.trace.name"
|
||||
|
||||
|
||||
class LangfuseMapper:
|
||||
|
|
@ -36,6 +37,7 @@ class LangfuseMapper:
|
|||
"langfuse.observation.model.name": lambda d: d.request_model or None,
|
||||
"langfuse.observation.metadata.provider": lambda d: d.provider or None,
|
||||
"langfuse.observation.id": lambda d: d.identity.call_id or None,
|
||||
LANGFUSE_TRACE_NAME: lambda d: d.trace_name or None,
|
||||
"langfuse.trace.metadata.team_id": lambda d: d.identity.team_id or None,
|
||||
"langfuse.trace.metadata.team_alias": lambda d: d.identity.team_alias or None,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -48,6 +48,8 @@ from litellm.integrations.otel.model.utils import as_str, to_seconds
|
|||
if TYPE_CHECKING:
|
||||
from litellm.types.utils import StandardLoggingPayload
|
||||
|
||||
LANGFUSE_TRACE_NAME_HEADER: Final = "langfuse_trace_name"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RequestIdentity:
|
||||
|
|
@ -215,6 +217,7 @@ class LLMCallEvent:
|
|||
# needs to be reasonable for a span that never gets closed (a leak).
|
||||
provisional_span_name: str
|
||||
time_to_first_chunk_seconds: float | None
|
||||
trace_name: str | None
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, kwargs: Mapping[str, Any]) -> LLMCallEvent:
|
||||
|
|
@ -231,9 +234,30 @@ class LLMCallEvent:
|
|||
upstream_started=kwargs.get("api_call_start_time") is not None,
|
||||
provisional_span_name=f"{operation.value} {model}".strip(),
|
||||
time_to_first_chunk_seconds=time_to_first_chunk_seconds(kwargs),
|
||||
trace_name=caller_trace_name(kwargs),
|
||||
)
|
||||
|
||||
|
||||
def caller_trace_name(kwargs: Mapping[str, object]) -> str | None:
|
||||
request: Final = _as_str_mapping(kwargs.get("litellm_params"))
|
||||
if request is None:
|
||||
return None
|
||||
proxy_request: Final = _as_str_mapping(request.get("proxy_server_request"))
|
||||
headers: Final = _as_str_mapping(proxy_request.get("headers")) if proxy_request is not None else None
|
||||
from_header: Final = as_str(headers.get(LANGFUSE_TRACE_NAME_HEADER)) if headers is not None else None
|
||||
if from_header:
|
||||
return from_header
|
||||
return next(
|
||||
(
|
||||
name
|
||||
for key in ("metadata", "litellm_metadata")
|
||||
if (metadata := _as_str_mapping(request.get(key))) is not None
|
||||
and (name := as_str(metadata.get("trace_name")))
|
||||
),
|
||||
None,
|
||||
)
|
||||
|
||||
|
||||
def time_to_first_chunk_seconds(kwargs: Mapping[str, Any]) -> float | None:
|
||||
"""Seconds from the upstream request being issued (``api_call_start_time``)
|
||||
to the first streamed chunk (``completion_start_time``); ``None`` for
|
||||
|
|
|
|||
|
|
@ -387,6 +387,7 @@ class LLMCallSpanData:
|
|||
output_type: GenAIOutputType | None = None
|
||||
call_type: str | None = None
|
||||
request_route: str | None = None
|
||||
trace_name: str | None = None
|
||||
|
||||
@classmethod
|
||||
def from_standard_logging_payload(
|
||||
|
|
@ -395,6 +396,7 @@ class LLMCallSpanData:
|
|||
capture_content: bool = False,
|
||||
time_to_first_chunk_seconds: float | None = None,
|
||||
request_route: str | None = None,
|
||||
trace_name: str | None = None,
|
||||
) -> LLMCallSpanData:
|
||||
params: Final = cast(Mapping[str, object], payload.get("model_parameters") or {})
|
||||
# The single parse of the request's metadata — the request-vs-provider
|
||||
|
|
@ -436,6 +438,7 @@ class LLMCallSpanData:
|
|||
output_type=resolve_output_type(call_type),
|
||||
call_type=call_type or None,
|
||||
request_route=request_route or context.identity.request_route,
|
||||
trace_name=trace_name,
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -1,9 +1,9 @@
|
|||
"""Tests for ``LangfuseOpenTelemetryV2``: the root observation's input and output are stamped from the
|
||||
request-task hooks, while the root span is still recording, so Langfuse can show them on the trace."""
|
||||
"""Tests for the Langfuse OTel v2 loggers: the trace name and the root observation's input and output are
|
||||
stamped from the request task while the root span is still recording, so Langfuse can show them on the trace."""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
from collections.abc import AsyncIterator, Sequence
|
||||
from collections.abc import AsyncIterator, Mapping, Sequence
|
||||
from typing import Final
|
||||
|
||||
import pytest
|
||||
|
|
@ -14,7 +14,7 @@ from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanE
|
|||
|
||||
import litellm # noqa: E402
|
||||
from litellm.caching.dual_cache import DualCache # noqa: E402
|
||||
from litellm.integrations.otel.logger import build_otel_v2_logger # noqa: E402
|
||||
from litellm.integrations.otel.logger import OpenTelemetryV2, build_otel_v2_logger # noqa: E402
|
||||
from litellm.integrations.otel.model.config import OpenTelemetryV2Config, is_otel_v2_enabled # noqa: E402
|
||||
from litellm.integrations.otel.model.spans import LITELLM_PROXY_REQUEST_SPAN_NAME, SpanRole # noqa: E402
|
||||
from litellm.integrations.otel.plumbing import context as otel_context # noqa: E402
|
||||
|
|
@ -41,6 +41,7 @@ from litellm.types.utils import ( # noqa: E402
|
|||
|
||||
INPUT_ATTR: Final = "langfuse.observation.input"
|
||||
OUTPUT_ATTR: Final = "langfuse.observation.output"
|
||||
TRACE_NAME_ATTR: Final = "langfuse.trace.name"
|
||||
CHAT_DATA: Final = {"model": "gpt-5.4-mini", "messages": [{"role": "user", "content": "ping"}]}
|
||||
|
||||
|
||||
|
|
@ -306,6 +307,73 @@ def test_unrenderable_output_never_raises_into_the_request():
|
|||
assert INPUT_ATTR not in attrs and OUTPUT_ATTR not in attrs
|
||||
|
||||
|
||||
def _run_named_request(
|
||||
logger: OpenTelemetryV2, exporter: InMemorySpanExporter, litellm_params: Mapping[str, object]
|
||||
) -> tuple[Mapping[str, object], Mapping[str, object]]:
|
||||
response: Final = ModelResponse(choices=[Choices(message=Message(role="assistant", content="pong"))])
|
||||
root: Final = _start_root(logger)
|
||||
logger.log_pre_api_call(
|
||||
model="gpt-5.4-mini", messages=[], kwargs={"litellm_call_id": "call_1", "litellm_params": litellm_params}
|
||||
)
|
||||
root.end()
|
||||
payload: Final = {
|
||||
"call_type": "acompletion",
|
||||
"custom_llm_provider": "openai",
|
||||
"model": "gpt-5.4-mini",
|
||||
"messages": CHAT_DATA["messages"],
|
||||
"response": response.model_dump(),
|
||||
"status": "success",
|
||||
"litellm_call_id": "call_1",
|
||||
"metadata": {},
|
||||
"hidden_params": {},
|
||||
}
|
||||
asyncio.run(
|
||||
logger.async_log_success_event(
|
||||
{"standard_logging_object": payload, "litellm_params": litellm_params}, response, None, None
|
||||
)
|
||||
)
|
||||
generation: Final = next(
|
||||
span for span in exporter.get_finished_spans() if span.name != LITELLM_PROXY_REQUEST_SPAN_NAME
|
||||
)
|
||||
return _root_attrs(exporter), dict(generation.attributes or {})
|
||||
|
||||
|
||||
@pytest.mark.parametrize("capture", ["span_only", "no_content"])
|
||||
def test_langfuse_trace_name_header_names_the_root_and_the_generation_over_body_metadata(capture):
|
||||
logger, exporter = _logger(capture=capture)
|
||||
|
||||
root_attrs, generation_attrs = _run_named_request(
|
||||
logger,
|
||||
exporter,
|
||||
{
|
||||
"metadata": {"trace_name": "from-body"},
|
||||
"proxy_server_request": {"headers": {"langfuse_trace_name": "from-header"}},
|
||||
},
|
||||
)
|
||||
|
||||
assert root_attrs[TRACE_NAME_ATTR] == "from-header"
|
||||
assert generation_attrs[TRACE_NAME_ATTR] == "from-header"
|
||||
|
||||
|
||||
def test_body_metadata_trace_name_names_the_root_and_the_generation():
|
||||
logger, exporter = _logger()
|
||||
|
||||
root_attrs, generation_attrs = _run_named_request(
|
||||
logger, exporter, {"metadata": {"trace_name": "from-body"}, "proxy_server_request": {"headers": {}}}
|
||||
)
|
||||
|
||||
assert root_attrs[TRACE_NAME_ATTR] == "from-body"
|
||||
assert generation_attrs[TRACE_NAME_ATTR] == "from-body"
|
||||
|
||||
|
||||
def test_unnamed_request_leaves_the_trace_name_off_both_spans():
|
||||
logger, exporter = _logger()
|
||||
|
||||
root_attrs, generation_attrs = _run_named_request(logger, exporter, {"proxy_server_request": {"headers": {}}})
|
||||
|
||||
assert TRACE_NAME_ATTR not in root_attrs and TRACE_NAME_ATTR not in generation_attrs
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("capture", "mappers"),
|
||||
[("no_content", ("genai", "langfuse")), ("span_only", ("genai",))],
|
||||
|
|
|
|||
|
|
@ -28,6 +28,7 @@ from litellm.integrations.otel import (
|
|||
)
|
||||
from litellm.integrations.otel.mappers.genai import GenAIMapper
|
||||
from litellm.integrations.otel.model import spans as spans_mod
|
||||
from litellm.integrations.otel.model.metadata import LLMCallEvent, caller_trace_name
|
||||
from litellm.integrations.otel.model.payloads import (
|
||||
LLMCallSpanData,
|
||||
RequestIdentity,
|
||||
|
|
@ -722,6 +723,37 @@ def test_request_identity_falls_back_to_legacy_team_keys():
|
|||
assert ident.team_alias == "legacy"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("request_data", "expected"),
|
||||
[
|
||||
({"proxy_server_request": {"headers": {"langfuse_trace_name": "from-header"}}}, "from-header"),
|
||||
({"metadata": {"trace_name": "from-body"}}, "from-body"),
|
||||
({"litellm_metadata": {"trace_name": "from-anthropic-body"}}, "from-anthropic-body"),
|
||||
(
|
||||
{
|
||||
"proxy_server_request": {"headers": {"langfuse_trace_name": "from-header"}},
|
||||
"metadata": {"trace_name": "from-body"},
|
||||
},
|
||||
"from-header",
|
||||
),
|
||||
({"proxy_server_request": {"headers": {"langfuse_trace_name": ""}}, "metadata": {"trace_name": "body"}}, "body"),
|
||||
({"proxy_server_request": {"headers": {}}, "metadata": {"user_api_key_team_id": "t1"}}, None),
|
||||
({}, None),
|
||||
],
|
||||
ids=["header", "body", "anthropic-body", "header-beats-body", "blank-header-falls-through", "neither", "empty"],
|
||||
)
|
||||
def test_caller_trace_name_prefers_the_langfuse_header_over_body_metadata(request_data, expected):
|
||||
assert caller_trace_name({"litellm_params": request_data}) == expected
|
||||
assert LLMCallEvent.from_dict({"litellm_params": request_data}).trace_name == expected
|
||||
|
||||
|
||||
def test_llm_span_data_carries_the_caller_trace_name():
|
||||
data: Final = LLMCallSpanData.from_standard_logging_payload(_sample_payload(), trace_name="nightly-eval")
|
||||
|
||||
assert data.trace_name == "nightly-eval"
|
||||
assert LLMCallSpanData.from_standard_logging_payload(_sample_payload()).trace_name is None
|
||||
|
||||
|
||||
def test_llm_span_carries_proxy_request_route():
|
||||
"""The LLM span records the proxy route the request arrived on, so it can be
|
||||
filtered by endpoint (``/v1/responses`` vs ``/v1/chat/completions``) without
|
||||
|
|
|
|||
|
|
@ -134,6 +134,11 @@ def test_langfuse_mapper_observation_attrs():
|
|||
assert attrs["langfuse.trace.metadata.team_id"] == "t1"
|
||||
|
||||
|
||||
def test_langfuse_mapper_names_the_trace_from_the_caller():
|
||||
assert LangfuseMapper().map(_llm_call(trace_name="nightly-eval"))["langfuse.trace.name"] == "nightly-eval"
|
||||
assert "langfuse.trace.name" not in LangfuseMapper().map(_llm_call(trace_name=None))
|
||||
|
||||
|
||||
def test_langfuse_mapper_skips_when_no_messages():
|
||||
data = _llm_call(messages_in=(), choices_out=())
|
||||
attrs = LangfuseMapper().map(data)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue