Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 18 additions & 1 deletion litellm/integrations/otel/langfuse_logger.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,12 @@

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

Expand All @@ -13,6 +18,18 @@


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)
Comment thread
greptile-apps[bot] marked this conversation as resolved.
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
Expand Down
7 changes: 4 additions & 3 deletions litellm/integrations/otel/logger.py
Original file line number Diff line number Diff line change
Expand Up @@ -554,6 +554,7 @@ def _finish_carrier(
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:
Expand Down Expand Up @@ -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
2 changes: 2 additions & 0 deletions litellm/integrations/otel/mappers/langfuse.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@

LANGFUSE_OBSERVATION_INPUT: Final = "langfuse.observation.input"
LANGFUSE_OBSERVATION_OUTPUT: Final = "langfuse.observation.output"
LANGFUSE_TRACE_NAME: Final = "langfuse.trace.name"


class LangfuseMapper:
Expand All @@ -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,
}
Expand Down
24 changes: 24 additions & 0 deletions litellm/integrations/otel/model/metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,8 @@
if TYPE_CHECKING:
from litellm.types.utils import StandardLoggingPayload

LANGFUSE_TRACE_NAME_HEADER: Final = "langfuse_trace_name"


@dataclass(frozen=True)
class RequestIdentity:
Expand Down Expand Up @@ -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:
Expand All @@ -231,9 +234,30 @@ def from_dict(cls, kwargs: Mapping[str, Any]) -> 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
Expand Down
3 changes: 3 additions & 0 deletions litellm/integrations/otel/model/payloads.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -395,6 +396,7 @@ def from_standard_logging_payload(
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
Expand Down Expand Up @@ -436,6 +438,7 @@ def from_standard_logging_payload(
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,
)


Expand Down
76 changes: 72 additions & 4 deletions tests/test_litellm/integrations/otel/test_langfuse_logger.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -14,7 +14,7 @@

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
Expand All @@ -41,6 +41,7 @@

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"}]}


Expand Down Expand Up @@ -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",))],
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
)
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,
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading