This commit is contained in:
devin-ai-integration[bot] 2026-10-03 16:29:53 -04:00 • committed by GitHub
commit 062056c786
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 985 additions and 235 deletions

View file

@ -1,8 +1,10 @@
import base64
import json
import os
from collections.abc import Iterable, Mapping, Sequence
from datetime import datetime
from typing import TYPE_CHECKING, Any, Final, Optional
from itertools import chain
from typing import TYPE_CHECKING, Any, Final, Optional, Protocol, runtime_checkable
from litellm._logging import verbose_logger
from litellm.integrations.arize import _utils
@ -30,6 +32,27 @@ LANGFUSE_INGESTION_VERSION_HEADER: Final = "x-langfuse-ingestion-version"
LANGFUSE_INGESTION_VERSION: Final = "4"
@runtime_checkable
class _Gettable(Protocol):
def get(self, key: str, default: object = None) -> object: ...
def _as_gettable(value: object) -> _Gettable | None:
return value if isinstance(value, _Gettable) else None
def _attr(item: object, name: str, default: object = None) -> object:
return getattr(item, name, default)
def _preset_cache_key(kwargs: Mapping[str, object]) -> object:
import litellm
if litellm.cache is None:
return None
return litellm.cache._get_preset_cache_key_from_kwargs(**kwargs)
class LangfuseOtelLogger(OpenTelemetry):
def __init__(self, config=None, *args, **kwargs):
# Prevent LangfuseOtelLogger from modifying global environment variables by constructing config manually
@ -95,7 +118,6 @@ class LangfuseOtelLogger(OpenTelemetry):
"mask_output": LangfuseSpanAttributes.MASK_OUTPUT,
"trace_user_id": LangfuseSpanAttributes.TRACE_USER_ID,
"session_id": LangfuseSpanAttributes.SESSION_ID,
"tags": LangfuseSpanAttributes.TAGS,
"trace_name": LangfuseSpanAttributes.TRACE_NAME,
"trace_id": LangfuseSpanAttributes.TRACE_ID,
"trace_metadata": LangfuseSpanAttributes.TRACE_METADATA,
@ -126,97 +148,52 @@ class LangfuseOtelLogger(OpenTelemetry):
safe_set_attribute(span, enum_attr.value, value)
@staticmethod
def _set_observation_output(span: Span, response_obj):
"""Helper to set observation output attributes."""
from litellm.integrations.arize._utils import safe_set_attribute
from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
def _observation_output(response_obj: _Gettable | None) -> str | None:
"""Serialized observation output, or None when the response yields nothing."""
if not response_obj or not hasattr(response_obj, "get"):
return
return None
return _extract_output_items(response_obj) or _extract_choices_output(response_obj)
choices: Final = response_obj.get("choices", [])
if choices:
first_choice: Final = choices[0]
message: Final = first_choice.get("message", {})
tool_calls: Final = message.get("tool_calls")
if tool_calls:
transformed_tool_calls: Final = []
for tool_call in tool_calls:
function = tool_call.get("function", {})
arguments_str = function.get("arguments", "{}")
try:
arguments_obj = json.loads(arguments_str) if isinstance(arguments_str, str) else arguments_str
except json.JSONDecodeError:
arguments_obj = {}
langfuse_tool_call = {
"id": response_obj.get("id", ""),
"name": function.get("name", ""),
"call_id": tool_call.get("id", ""),
"type": "function_call",
"arguments": arguments_obj,
}
transformed_tool_calls.append(langfuse_tool_call)
safe_set_attribute(
span,
LangfuseSpanAttributes.OBSERVATION_OUTPUT.value,
safe_dumps(transformed_tool_calls),
)
else:
output_data: Final = {}
if message.get("role"):
output_data["role"] = message.get("role")
if message.get("content") is not None:
output_data["content"] = message.get("content")
if output_data:
safe_set_attribute(
span,
LangfuseSpanAttributes.OBSERVATION_OUTPUT.value,
safe_dumps(output_data),
)
@staticmethod
def _trace_tags(kwargs: Mapping[str, object], metadata: Mapping[str, object]) -> tuple[str, ...]:
"""Order-preserving dedupe of caller tags, request tags and langfuse_default_tags expansions."""
import litellm
output: Final = response_obj.get("output", [])
if output:
output_items_data: Final[list[dict]] = []
for item in output:
if hasattr(item, "type"):
item_type = item.type
if item_type == "reasoning" and hasattr(item, "summary"):
for summary in item.summary:
if hasattr(summary, "text"):
output_items_data.append(
{
"role": "reasoning_summary",
"content": summary.text,
}
)
elif item_type == "message":
output_items_data.append(
{
"role": getattr(item, "role", "assistant"),
"content": getattr(getattr(item, "content", [{}])[0], "text", ""),
}
)
elif item_type == "function_call":
arguments_str = getattr(item, "arguments", "{}")
arguments_obj = (
safe_json_loads(arguments_str, default={})
if isinstance(arguments_str, str)
else arguments_str
)
langfuse_tool_call = {
"id": getattr(item, "id", ""),
"name": getattr(item, "name", ""),
"call_id": getattr(item, "call_id", ""),
"type": "function_call",
"arguments": arguments_obj,
}
output_items_data.append(langfuse_tool_call)
if output_items_data:
safe_set_attribute(
span,
LangfuseSpanAttributes.OBSERVATION_OUTPUT.value,
safe_dumps(output_items_data),
caller_tags: Final = metadata.get("tags")
request_tags: Final = (_as_gettable(kwargs.get("standard_logging_object")) or {}).get("request_tags")
default_tags: Final = litellm.langfuse_default_tags
def _default_tag(key: str) -> str | None:
if key == "cache_hit":
return f"cache_hit:{kwargs.get('cache_hit', False)}"
if key == "cache_key":
hidden_params: Final = _as_gettable(metadata.get("hidden_params", {})) or {}
cache_key: Final = (
hidden_params.get("cache_key")
if hidden_params.get("cache_key") is not None
else _preset_cache_key(kwargs)
)
return f"cache_key:{cache_key}"
if key == "proxy_base_url":
proxy_base_url: Final = os.environ.get("PROXY_BASE_URL")
return f"proxy_base_url:{proxy_base_url}" if proxy_base_url is not None else None
if key in metadata and metadata[key] is not None:
return f"{key}:{metadata[key]}"
return None
expanded: Final = tuple(_default_tag(key) for key in default_tags) if isinstance(default_tags, list) else ()
candidates: Final = (
(
(caller_tags,)
if isinstance(caller_tags, str)
else tuple(tag for tag in caller_tags if isinstance(tag, str))
if isinstance(caller_tags, list)
else ()
)
+ (tuple(tag for tag in request_tags if isinstance(tag, str)) if isinstance(request_tags, list) else ())
+ tuple(tag for tag in expanded if tag is not None)
)
return tuple(dict.fromkeys(candidates))
@staticmethod
def _set_langfuse_specific_attributes(span: Span, kwargs, response_obj):
@ -245,15 +222,30 @@ class LangfuseOtelLogger(OpenTelemetry):
metadata: Final = LangfuseOtelLogger._extract_langfuse_metadata(kwargs)
LangfuseOtelLogger._set_metadata_attributes(span=span, metadata=metadata)
messages: Final = kwargs.get("messages")
if messages:
joins_existing_trace: Final = metadata.get("existing_trace_id") is not None
if metadata.get("trace_name") is None and not joins_existing_trace:
safe_set_attribute(
span,
LangfuseSpanAttributes.OBSERVATION_INPUT.value,
safe_dumps(messages),
LangfuseSpanAttributes.TRACE_NAME.value,
f"litellm-{kwargs.get('call_type') or 'completion'}",
)
LangfuseOtelLogger._set_observation_output(span=span, response_obj=response_obj)
if not joins_existing_trace:
tags: Final = LangfuseOtelLogger._trace_tags(kwargs, metadata)
if tags:
safe_set_attribute(span, LangfuseSpanAttributes.TAGS.value, json.dumps(list(tags)))
input_json: Final = safe_dumps(kwargs.get("messages")) if kwargs.get("messages") else None
if input_json is not None:
safe_set_attribute(span, LangfuseSpanAttributes.OBSERVATION_INPUT.value, input_json)
if not joins_existing_trace:
safe_set_attribute(span, LangfuseSpanAttributes.TRACE_INPUT.value, input_json)
output_json: Final = LangfuseOtelLogger._observation_output(response_obj)
if output_json is not None:
safe_set_attribute(span, LangfuseSpanAttributes.OBSERVATION_OUTPUT.value, output_json)
if not joins_existing_trace:
safe_set_attribute(span, LangfuseSpanAttributes.TRACE_OUTPUT.value, output_json)
@staticmethod
def _get_langfuse_otel_host() -> str | None:
@ -445,3 +437,105 @@ class LangfuseOtelLogger(OpenTelemetry):
"""
Langfuse should not receive service failure logs.
"""
def _extract_choices_output(response_obj: _Gettable) -> str | None:
from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
choices: Final = response_obj.get("choices", [])
if not isinstance(choices, list) or not choices:
return None
first_choice: Final = _as_gettable(choices[0])
if first_choice is None:
return None
message: Final = _as_gettable(first_choice.get("message", {}))
if message is None:
return None
tool_calls: Final = message.get("tool_calls")
if isinstance(tool_calls, list) and tool_calls:
transformed_tool_calls: Final = [
entry for tool_call in tool_calls if (entry := _transformed_tool_call(response_obj, tool_call)) is not None
]
return safe_dumps(transformed_tool_calls)
output_data: Final = {
key: value
for key, value in (
("role", message.get("role") or None),
("content", message.get("content")),
)
if value is not None
}
return safe_dumps(output_data) if output_data else None
def _transformed_tool_call(response_obj: _Gettable, tool_call: object) -> dict[str, object] | None:
call: Final = _as_gettable(tool_call)
if call is None:
return None
function: Final = _as_gettable(call.get("function", {}))
if function is None:
return None
return {
"id": response_obj.get("id", ""),
"name": function.get("name", ""),
"call_id": call.get("id", ""),
"type": "function_call",
"arguments": _tool_call_arguments(function.get("arguments", "{}")),
}
def _tool_call_arguments(arguments: object) -> object:
if not isinstance(arguments, str):
return arguments
try:
return json.loads(arguments)
except json.JSONDecodeError:
return {}
def _extract_output_items(response_obj: _Gettable) -> str | None:
from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
output: Final = response_obj.get("output", [])
if not isinstance(output, list) or not output:
return None
rendered: Final = tuple(chain.from_iterable(map(_output_items, output)))
return safe_dumps(list(rendered)) if rendered else None
def _output_items(item: object) -> tuple[dict[str, object], ...]:
item_type: Final = _attr(item, "type")
if item_type == "reasoning":
summaries: Final = _attr(item, "summary")
if not isinstance(summaries, Iterable):
return ()
return tuple(
{"role": "reasoning_summary", "content": _attr(summary, "text")}
for summary in summaries
if hasattr(summary, "text")
)
if item_type == "message":
content_items: Final = _attr(item, "content", [{}])
first_content: Final = (
content_items[0]
if isinstance(content_items, Sequence) and not isinstance(content_items, (str, bytes)) and content_items
else {}
)
return (
{
"role": _attr(item, "role", "assistant"),
"content": _attr(first_content, "text", ""),
},
)
if item_type == "function_call":
arguments: Final = _attr(item, "arguments", "{}")
return (
{
"id": _attr(item, "id", ""),
"name": _attr(item, "name", ""),
"call_id": _attr(item, "call_id", ""),
"type": "function_call",
"arguments": safe_json_loads(arguments, default={}) if isinstance(arguments, str) else arguments,
},
)
return ()

View file

@ -10,6 +10,7 @@ from litellm.integrations.otel.mappers.langfuse import (
)
from litellm.integrations.otel.model.request_io import request_input, response_output, stream_output
from litellm.integrations.otel.model.trace_controls import caller_trace_controls
from litellm.integrations.otel.model.utils import as_str
from litellm.integrations.otel.plumbing.context import request_root_span
if TYPE_CHECKING:
@ -24,7 +25,9 @@ class LangfuseOpenTelemetryV2(OpenTelemetryV2):
def log_pre_api_call(self, model: str, messages: object, kwargs: Mapping[str, object]) -> None:
root: Final = request_root_span()
if root is not None and root.is_recording():
root.set_attributes(LangfuseMapper.trace_attributes(caller_trace_controls(kwargs)))
root.set_attributes(
LangfuseMapper.trace_attributes(caller_trace_controls(kwargs), as_str(kwargs.get("call_type")))
)
super().log_pre_api_call(model, messages, kwargs)

View file

@ -31,12 +31,24 @@ from litellm.integrations.otel.model.trace_controls import TraceControls
LANGFUSE_OBSERVATION_INPUT: Final = "langfuse.observation.input"
LANGFUSE_OBSERVATION_OUTPUT: Final = "langfuse.observation.output"
LANGFUSE_TRACE_INPUT: Final = "langfuse.trace.input"
LANGFUSE_TRACE_OUTPUT: Final = "langfuse.trace.output"
LANGFUSE_TRACE_NAME: Final = "langfuse.trace.name"
LANGFUSE_TRACE_USER_ID: Final = "user.id"
LANGFUSE_TRACE_SESSION_ID: Final = "session.id"
LANGFUSE_TRACE_TAGS: Final = "langfuse.trace.tags"
def _observation_input(data: LLMCallSpanData) -> AttrValue | None:
return serialize_messages(data.messages_in)
def _observation_output(data: LLMCallSpanData) -> AttrValue | None:
if data.embedding_output is not None:
return data.embedding_output.as_json()
return serialize_messages(output_messages(data))
class LangfuseMapper:
_LLM_CALL_ATTRS: dict[str, Callable[[LLMCallSpanData], AttrValue | None]] = {
"langfuse.observation.type": lambda _: "generation",
@ -71,10 +83,10 @@ class LangfuseMapper:
"langfuse.observation.model.parameters": lambda d: json_if(
collect(LangfuseMapper._MODEL_PARAMS, d.request_params)
),
LANGFUSE_OBSERVATION_INPUT: lambda d: serialize_messages(d.messages_in),
LANGFUSE_OBSERVATION_OUTPUT: lambda d: (
d.embedding_output.as_json() if d.embedding_output is not None else serialize_messages(output_messages(d))
),
LANGFUSE_OBSERVATION_INPUT: _observation_input,
LANGFUSE_OBSERVATION_OUTPUT: _observation_output,
LANGFUSE_TRACE_INPUT: _observation_input,
LANGFUSE_TRACE_OUTPUT: _observation_output,
"langfuse.observation.usage_details": lambda d: json_if(collect(LangfuseMapper._USAGE_FIELDS, d.usage)),
"langfuse.observation.cost_details": lambda d: (
json.dumps({"total": d.response_cost}) if d.response_cost is not None else None
@ -89,10 +101,10 @@ class LangfuseMapper:
return {}
@staticmethod
def trace_attributes(trace: TraceControls) -> AttributeMap:
def trace_attributes(trace: TraceControls, call_type: str | None = None) -> AttributeMap:
return drop_none_pairs(
(
(LANGFUSE_TRACE_NAME, trace.name or None),
(LANGFUSE_TRACE_NAME, trace.name or (f"litellm-{call_type}" if call_type else None)),
(LANGFUSE_TRACE_USER_ID, trace.user_id or None),
(LANGFUSE_TRACE_SESSION_ID, trace.session_id or None),
(LANGFUSE_TRACE_TAGS, trace.tags or None),
@ -103,6 +115,6 @@ class LangfuseMapper:
def _llm_call(cls, data: LLMCallSpanData) -> AttributeMap:
return {
**collect(cls._LLM_CALL_ATTRS, data),
**cls.trace_attributes(data.trace),
**cls.trace_attributes(data.trace, data.call_type),
**collect(cls._BLOB_ATTRS, data),
}

View file

@ -31,6 +31,8 @@ class LangfuseSpanAttributes(str, Enum):
OBSERVATION_OUTPUT = "langfuse.observation.output"
# ---- Trace-level metadata ----
TRACE_INPUT = "langfuse.trace.input"
TRACE_OUTPUT = "langfuse.trace.output"
TRACE_USER_ID = "user.id"
SESSION_ID = "session.id"
TAGS = "langfuse.trace.tags"

View file

@ -0,0 +1,411 @@
import json
import threading
import uuid
from collections.abc import Callable, Iterator
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass
from pathlib import Path
from queue import SimpleQueue
from typing import Final
import pytest
import yaml
from integration._support.client import Gateway, eventually, gateway_from_environment
from integration._support.process import owned_proxy
from integration._support.wire import Reply, Request, Wire, wire_server
from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceRequest
REPLY_TEXT: Final = "langfuse-otel scripted reply"
_SINK_OUTAGE: Final = threading.Event()
_ACCEPTED: Final[SimpleQueue[Request]] = SimpleQueue()
@dataclass(frozen=True, slots=True)
class SpanRecord:
trace_id: str
parent_span_id: str
attributes: dict[str, object]
def _attribute_value(value: object) -> object:
kind: Final = value.WhichOneof("value")
return getattr(value, kind) if kind is not None else None
def _span_records(body: bytes) -> tuple[SpanRecord, ...]:
export: Final = ExportTraceServiceRequest.FromString(body)
return tuple(
SpanRecord(
trace_id=span.trace_id.hex(),
parent_span_id=span.parent_span_id.hex(),
attributes={pair.key: _attribute_value(pair.value) for pair in span.attributes},
)
for resource in export.resource_spans
for scope in resource.scope_spans
for span in scope.spans
)
def _otlp_sink(request: Request) -> Reply:
if request.target.endswith("/v1/traces"):
assert request.headers.get("x-langfuse-ingestion-version") == "4", request.headers
if _SINK_OUTAGE.is_set():
return Reply(status=503, body=b'{"error": "scripted outage"}')
_ACCEPTED.put(request)
return Reply(body=b"")
def _scripted_upstream(suffix: str, reply: Reply) -> Callable[[Request], Reply]:
def respond(request: Request) -> Reply:
if request.target.endswith(suffix):
return reply
return Reply(body=b'{"object": "list", "data": []}')
return respond
def _chaos_upstream(request: Request) -> Reply:
if not request.target.endswith("/chat/completions"):
return Reply(body=b'{"object": "list", "data": []}')
body: Final = json.loads(request.body)
marker: Final = body["messages"][0]["content"]
if body.get("stream"):
return Reply(content_type="text/event-stream", chunks=_chat_stream(marker))
return Reply(body=_chat_completion(marker))
def _chat_completion(marker: str) -> bytes:
return json.dumps(
{
"id": marker,
"object": "chat.completion",
"created": 1,
"model": "gpt-4o-mini",
"choices": [
{
"index": 0,
"message": {"role": "assistant", "content": REPLY_TEXT},
"finish_reason": "stop",
}
],
"usage": {"prompt_tokens": 5, "completion_tokens": 4, "total_tokens": 9},
}
).encode()
def _responses_result(marker: str) -> bytes:
return json.dumps(
{
"id": marker,
"object": "response",
"created_at": 1,
"status": "completed",
"model": "gpt-4o-mini",
"output": [
{
"type": "message",
"id": "msg-lit8281",
"status": "completed",
"role": "assistant",
"content": [{"type": "output_text", "text": REPLY_TEXT, "annotations": []}],
}
],
"usage": {"input_tokens": 5, "output_tokens": 4, "total_tokens": 9},
}
).encode()
def _chat_stream(marker: str) -> tuple[bytes, ...]:
def frame(delta: dict[str, object], finish: str | None = None) -> bytes:
return (
b"data: "
+ json.dumps(
{
"id": marker,
"object": "chat.completion.chunk",
"created": 1,
"model": "gpt-4o-mini",
"choices": [{"index": 0, "delta": delta, "finish_reason": finish}],
}
).encode()
+ b"\n\n"
)
return (
frame({"role": "assistant", "content": "langfuse-otel "}),
frame({"content": "scripted reply"}),
frame({}, finish="stop"),
b"data: [DONE]\n\n",
)
def _generation_span(collector: Wire, marker: str) -> SpanRecord:
def spans() -> tuple[SpanRecord, ...]:
return tuple(
record
for request in collector.drain()
if request.target.endswith("/v1/traces")
for record in _span_records(request.body)
if record.attributes.get("langfuse.observation.type") == "generation"
and marker in str(tuple(record.attributes.values()))
)
found: Final = eventually(spans, lambda values: len(values) == 1, seconds=70)
return found[0]
def _assert_derived_trace_fields(span: SpanRecord, call_type: str) -> None:
attributes: Final = span.attributes
assert attributes.get("langfuse.observation.input") is not None, attributes
assert attributes.get("langfuse.observation.output") is not None, attributes
assert attributes.get("langfuse.trace.name") == f"litellm-{call_type}", attributes
assert attributes.get("langfuse.trace.input") == attributes.get("langfuse.observation.input"), attributes
assert attributes.get("langfuse.trace.output") == attributes.get("langfuse.observation.output"), attributes
def _langfuse_config(directory: Path, callbacks: tuple[str, ...]) -> Path:
base: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())
config: Final = {**base, "litellm_settings": {**base["litellm_settings"], "callbacks": list(callbacks)}}
path: Final = directory / "langfuse.yaml"
path.write_text(yaml.safe_dump(config))
return path
@pytest.fixture(scope="module")
def langfuse_sink() -> Iterator[Wire]:
with wire_server(_otlp_sink) as wire:
yield wire
@pytest.fixture(scope="module")
def langfuse_v1_proxy(langfuse_sink: Wire, tmp_path_factory: pytest.TempPathFactory) -> Iterator[Gateway]:
directory: Final = tmp_path_factory.mktemp("langfuse-v1")
with gateway_from_environment() as gateway:
with owned_proxy(
gateway,
directory,
{
"LANGFUSE_PUBLIC_KEY": "pk-lf-test",
"LANGFUSE_SECRET_KEY": "sk-lf-test",
"LANGFUSE_HOST": langfuse_sink.url,
},
config=_langfuse_config(directory, ("langfuse_otel",)),
) as proxy:
yield proxy
@pytest.fixture(scope="module")
def langfuse_v2_proxy(langfuse_sink: Wire, tmp_path_factory: pytest.TempPathFactory) -> Iterator[Gateway]:
directory: Final = tmp_path_factory.mktemp("langfuse-v2")
with gateway_from_environment() as gateway:
with owned_proxy(
gateway,
directory,
{
"LITELLM_OTEL_V2": "1",
"OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT": "span_only",
"LANGFUSE_PUBLIC_KEY": "pk-lf-test",
"LANGFUSE_SECRET_KEY": "sk-lf-test",
"LANGFUSE_HOST": langfuse_sink.url,
},
config=_langfuse_config(directory, ("langfuse_otel",)),
) as proxy:
yield proxy
@pytest.mark.covers("other.observability.langfuse_otel.chat_derives_trace_fields")
def test_chat_completion_derives_trace_name_input_and_output(langfuse_v1_proxy: Gateway, langfuse_sink: Wire) -> None:
marker: Final = "lit8281-chat-" + uuid.uuid4().hex
upstream: Final = _scripted_upstream("/chat/completions", Reply(body=_chat_completion(marker)))
with wire_server(upstream) as provider, langfuse_v1_proxy.scenario() as scenario:
model: Final = scenario.model(api_base=provider.url + "/v1")
response: Final = langfuse_v1_proxy.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}]},
)
assert response.status_code == 200, response.text
span: Final = _generation_span(langfuse_sink, marker)
assert span.attributes["llm.response.id"] == marker, span.attributes
assert span.attributes["langfuse.observation.type"] == "generation", span.attributes
_assert_derived_trace_fields(span, "acompletion")
assert marker in str(span.attributes["langfuse.observation.input"]), span.attributes
assert REPLY_TEXT in str(span.attributes["langfuse.observation.output"]), span.attributes
@pytest.mark.covers("other.observability.langfuse_otel.parented_chat_derives_trace_fields_and_keeps_parent")
def test_parented_chat_completion_derives_trace_fields_inside_inbound_trace(
langfuse_v1_proxy: Gateway, langfuse_sink: Wire
) -> None:
marker: Final = "lit8281-parent-" + uuid.uuid4().hex
trace_id: Final = uuid.uuid4().hex
parent_id: Final = uuid.uuid4().hex[:16]
upstream: Final = _scripted_upstream("/chat/completions", Reply(body=_chat_completion(marker)))
with wire_server(upstream) as provider, langfuse_v1_proxy.scenario() as scenario:
model: Final = scenario.model(api_base=provider.url + "/v1")
response: Final = langfuse_v1_proxy.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}]},
headers={"traceparent": f"00-{trace_id}-{parent_id}-01"},
)
assert response.status_code == 200, response.text
span: Final = _generation_span(langfuse_sink, marker)
assert span.trace_id == trace_id, (span.trace_id, trace_id)
assert span.parent_span_id == parent_id, (span.parent_span_id, parent_id)
_assert_derived_trace_fields(span, "acompletion")
@pytest.mark.covers("other.observability.langfuse_otel.streaming_chat_derives_trace_fields")
def test_streaming_chat_completion_derives_trace_fields(langfuse_v1_proxy: Gateway, langfuse_sink: Wire) -> None:
marker: Final = "lit8281-stream-" + uuid.uuid4().hex
with (
wire_server(
_scripted_upstream(
"/chat/completions", Reply(content_type="text/event-stream", chunks=_chat_stream(marker))
)
) as provider,
langfuse_v1_proxy.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1")
response: Final = langfuse_v1_proxy.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"stream": True,
},
)
assert response.status_code == 200, response.text
assert "scripted reply" in response.text, response.text
span: Final = _generation_span(langfuse_sink, marker)
_assert_derived_trace_fields(span, "acompletion")
assert REPLY_TEXT in str(span.attributes["langfuse.observation.output"]), span.attributes
@pytest.mark.covers("other.observability.langfuse_otel.caller_trace_name_and_tags_win")
def test_caller_supplied_trace_name_and_tags_are_emitted(langfuse_v1_proxy: Gateway, langfuse_sink: Wire) -> None:
marker: Final = "lit8281-caller-" + uuid.uuid4().hex
upstream: Final = _scripted_upstream("/chat/completions", Reply(body=_chat_completion(marker)))
with wire_server(upstream) as provider, langfuse_v1_proxy.scenario() as scenario:
model: Final = scenario.model(api_base=provider.url + "/v1")
response: Final = langfuse_v1_proxy.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"metadata": {"trace_name": "caller-trace", "tags": ["caller-tag"]},
},
)
assert response.status_code == 200, response.text
span: Final = _generation_span(langfuse_sink, marker)
assert span.attributes.get("langfuse.trace.name") == "caller-trace", span.attributes
tags: Final = json.loads(str(span.attributes["langfuse.trace.tags"]))
assert tags[0] == "caller-tag", span.attributes
@pytest.mark.covers("other.observability.langfuse_otel.responses_derives_trace_fields")
def test_responses_call_derives_trace_fields(langfuse_v1_proxy: Gateway, langfuse_sink: Wire) -> None:
marker: Final = "lit8281-resp-" + uuid.uuid4().hex
upstream: Final = _scripted_upstream("/responses", Reply(body=_responses_result(marker)))
with wire_server(upstream) as provider, langfuse_v1_proxy.scenario() as scenario:
model: Final = scenario.model(api_base=provider.url + "/v1")
response: Final = langfuse_v1_proxy.request(
"POST",
"/v1/responses",
{"model": model, "input": marker},
)
assert response.status_code == 200, response.text
span: Final = _generation_span(langfuse_sink, marker)
_assert_derived_trace_fields(span, "aresponses")
@pytest.mark.covers("other.observability.langfuse_otel.v2_derives_trace_fields")
def test_otel_v2_derives_trace_name_input_and_output(langfuse_v2_proxy: Gateway, langfuse_sink: Wire) -> None:
marker: Final = "lit8281-v2-" + uuid.uuid4().hex
upstream: Final = _scripted_upstream("/chat/completions", Reply(body=_chat_completion(marker)))
with wire_server(upstream) as provider, langfuse_v2_proxy.scenario() as scenario:
model: Final = scenario.model(api_base=provider.url + "/v1")
response: Final = langfuse_v2_proxy.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}]},
)
assert response.status_code == 200, response.text
span: Final = _generation_span(langfuse_sink, marker)
_assert_derived_trace_fields(span, "acompletion")
@pytest.mark.covers("other.observability.langfuse_otel.sink_outage_mid_burst_lands_every_span_exactly_once")
def test_sink_outage_mid_burst_lands_every_generation_span_exactly_once(
langfuse_v1_proxy: Gateway, langfuse_sink: Wire
) -> None:
markers: Final = tuple("lit8281-chaos-" + uuid.uuid4().hex for _ in range(24))
landed: list[Request] = [] # mutable-ok: accumulated across eventually() polls
try:
_SINK_OUTAGE.set()
with wire_server(_chaos_upstream) as provider, langfuse_v1_proxy.scenario() as scenario:
model: Final = scenario.model(api_base=provider.url + "/v1")
with ThreadPoolExecutor(max_workers=8) as pool:
futures: Final = tuple(
pool.submit(
langfuse_v1_proxy.request,
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
**({"stream": True} if index % 2 else {}),
},
)
for index, marker in enumerate(markers)
)
responses: Final = tuple(future.result() for future in futures)
for index, response in enumerate(responses):
assert response.status_code == 200, response.text
if index % 2:
assert "scripted reply" in response.text, response.text
refused_requests: Final = eventually(
lambda: tuple(r for r in langfuse_sink.drain() if r.target.endswith("/v1/traces")),
lambda rs: len(rs) >= 1,
seconds=70,
)
refused: Final = len(refused_requests)
_SINK_OUTAGE.clear()
def accepted() -> tuple[SpanRecord, ...]:
landed.extend(tuple(_ACCEPTED.get_nowait() for _ in range(_ACCEPTED.qsize())))
return tuple(
record
for request in landed
for record in _span_records(request.body)
if record.attributes.get("langfuse.observation.type") == "generation"
and record.attributes.get("llm.response.id") in markers
)
landed_spans: Final = eventually(accepted, lambda records: len(records) == 24, seconds=120)
finally:
_SINK_OUTAGE.clear()
counts: Final = {
marker: sum(1 for record in landed_spans if record.attributes["llm.response.id"] == marker)
for marker in markers
}
assert counts == {marker: 1 for marker in markers}, counts
assert refused >= 1, refused
for span in landed_spans:
_assert_derived_trace_fields(span, "acompletion")

View file

@ -314,7 +314,9 @@ def _run_named_request(
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}
model="gpt-5.4-mini",
messages=[],
kwargs={"litellm_call_id": "call_1", "litellm_params": litellm_params, "call_type": "acompletion"},
)
root.end()
payload: Final = {
@ -367,12 +369,13 @@ def test_body_metadata_trace_name_names_the_root_and_the_generation():
assert generation_attrs[TRACE_NAME_ATTR] == "from-body"
def test_unnamed_request_leaves_the_trace_name_off_both_spans():
def test_unnamed_request_derives_the_trace_name_on_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
assert root_attrs[TRACE_NAME_ATTR] == "litellm-acompletion"
assert generation_attrs[TRACE_NAME_ATTR] == "litellm-acompletion"
@pytest.mark.parametrize("capture", ["span_only", "no_content"])
@ -397,7 +400,7 @@ def test_body_metadata_user_session_and_tags_land_on_the_root_and_the_generation
assert attrs["user.id"] == "user-42"
assert attrs["session.id"] == "session-7"
assert tuple(attrs["langfuse.trace.tags"]) == ("prod", "eval", "nightly")
assert TRACE_NAME_ATTR not in attrs
assert attrs[TRACE_NAME_ATTR] == "litellm-acompletion"
def test_langfuse_user_and_session_headers_beat_body_metadata_on_both_spans():
@ -461,11 +464,16 @@ def test_a_request_without_trace_controls_stamps_none_of_them():
logger, exporter = _logger()
root_attrs, generation_attrs = _run_named_request(
logger, exporter, {"metadata": {"user_api_key_team_id": "t1", "tags": []}, "proxy_server_request": {"headers": {}}}
logger,
exporter,
{"metadata": {"user_api_key_team_id": "t1", "tags": []}, "proxy_server_request": {"headers": {}}},
)
assert set(TRACE_CONTROL_ATTRS).isdisjoint(root_attrs)
assert set(TRACE_CONTROL_ATTRS).isdisjoint(generation_attrs)
caller_only: Final = tuple(attr for attr in TRACE_CONTROL_ATTRS if attr != TRACE_NAME_ATTR)
assert set(caller_only).isdisjoint(root_attrs)
assert set(caller_only).isdisjoint(generation_attrs)
assert root_attrs[TRACE_NAME_ATTR] == "litellm-acompletion"
assert generation_attrs[TRACE_NAME_ATTR] == "litellm-acompletion"
@pytest.mark.parametrize(

View file

@ -400,3 +400,34 @@ def test_resolve_mappers_composition_layers_vocabularies():
def test_resolve_mappers_rejects_unknown_name():
with pytest.raises(ValueError, match="unknown mapper name 'nope'"):
resolve_mappers(["genai", "nope"])
def test_langfuse_mapper_derives_trace_name_input_and_output():
attrs = LangfuseMapper().map(_llm_call(call_type="acompletion"))
assert attrs["langfuse.trace.name"] == "litellm-acompletion"
assert attrs["langfuse.trace.input"] == attrs["langfuse.observation.input"]
assert attrs["langfuse.trace.output"] == attrs["langfuse.observation.output"]
def test_langfuse_mapper_derived_name_yields_to_the_caller_name():
attrs = LangfuseMapper().map(_llm_call(call_type="acompletion", trace=TraceControls(name="caller-trace")))
assert attrs["langfuse.trace.name"] == "caller-trace"
def test_langfuse_mapper_without_call_type_or_name_emits_no_trace_name():
attrs = LangfuseMapper().map(_llm_call(call_type=None, trace=TraceControls()))
assert "langfuse.trace.name" not in attrs
def test_langfuse_root_trace_attributes_match_the_generation_pairs():
data = _llm_call(call_type="aresponses", trace=TraceControls())
generation = LangfuseMapper().map(data)
root = LangfuseMapper.trace_attributes(TraceControls(), "aresponses")
assert root == {"langfuse.trace.name": "litellm-aresponses"}
assert generation["langfuse.trace.name"] == root["langfuse.trace.name"]
assert generation["langfuse.trace.input"] == generation["langfuse.observation.input"]
assert generation["langfuse.trace.output"] == generation["langfuse.observation.output"]

View file

@ -1,5 +1,7 @@
import json
import os
from collections.abc import Mapping
from typing import Final
from unittest.mock import MagicMock, patch
import pytest
@ -104,19 +106,13 @@ class TestLangfuseOtelIntegration:
mock_kwargs = {"test": "kwargs"}
mock_response = {"test": "response"}
with patch(
"litellm.integrations.arize._utils.set_attributes"
) as mock_set_attributes:
LangfuseOtelLogger.set_langfuse_otel_attributes(
mock_span, mock_kwargs, mock_response
)
with patch("litellm.integrations.arize._utils.set_attributes") as mock_set_attributes:
LangfuseOtelLogger.set_langfuse_otel_attributes(mock_span, mock_kwargs, mock_response)
mock_set_attributes.assert_called_once_with(
mock_span, mock_kwargs, mock_response, LangfuseLLMObsOTELAttributes
)
mock_span.set_attribute.assert_any_call(
"langfuse.observation.type", "generation"
)
mock_span.set_attribute.assert_any_call("langfuse.observation.type", "generation")
def test_set_langfuse_environment_attribute(self):
"""Test that Langfuse environment is set correctly when environment variable is present."""
@ -125,17 +121,11 @@ class TestLangfuseOtelIntegration:
test_env = "staging"
with patch.dict(os.environ, {"LANGFUSE_TRACING_ENVIRONMENT": test_env}):
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(
mock_span, mock_kwargs, {}
)
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(mock_span, mock_kwargs, {})
# safe_set_attribute(span, key, value) → positional args
mock_safe_set_attribute.assert_called_once_with(
mock_span, "langfuse.environment", test_env
)
mock_safe_set_attribute.assert_any_call(mock_span, "langfuse.environment", test_env)
def test_set_langfuse_environment_attribute_prefers_dynamic_param(self):
"""Per-key/team langfuse_environment beats the deployment env var."""
@ -148,18 +138,10 @@ class TestLangfuseOtelIntegration:
self.attributes[key] = value
span = _RecordingSpan()
mock_kwargs = {
"standard_callback_dynamic_params": {
"langfuse_environment": "team-a-env"
}
}
mock_kwargs = {"standard_callback_dynamic_params": {"langfuse_environment": "team-a-env"}}
with patch.dict(
os.environ, {"LANGFUSE_TRACING_ENVIRONMENT": "deployment-wide"}
):
LangfuseOtelLogger._set_langfuse_specific_attributes(
span, mock_kwargs, {}
)
with patch.dict(os.environ, {"LANGFUSE_TRACING_ENVIRONMENT": "deployment-wide"}):
LangfuseOtelLogger._set_langfuse_specific_attributes(span, mock_kwargs, {})
assert span.attributes["langfuse.environment"] == "team-a-env"
@ -223,12 +205,8 @@ class TestLangfuseOtelIntegration:
kwargs = {"litellm_params": {"metadata": metadata}}
# Capture calls to safe_set_attribute
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(
MagicMock(), kwargs, None
)
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None)
# Build expected calls manually for clarity
from litellm.types.integrations.langfuse_otel import LangfuseSpanAttributes
@ -243,15 +221,12 @@ class TestLangfuseOtelIntegration:
LangfuseSpanAttributes.TRACE_USER_ID.value: "user-123",
LangfuseSpanAttributes.SESSION_ID.value: "sess-456",
# Lists / dicts should be JSON strings
LangfuseSpanAttributes.TAGS.value: json.dumps(["tagA", "tagB"]),
LangfuseSpanAttributes.TRACE_NAME.value: "trace-name",
LangfuseSpanAttributes.TRACE_ID.value: "traceid", # stripped dashes
LangfuseSpanAttributes.TRACE_METADATA.value: json.dumps({"k": "v"}),
LangfuseSpanAttributes.RELEASE.value: "rel-1",
LangfuseSpanAttributes.EXISTING_TRACE_ID.value: "existing-id",
LangfuseSpanAttributes.UPDATE_TRACE_KEYS.value: json.dumps(
["key1", "key2"]
),
LangfuseSpanAttributes.UPDATE_TRACE_KEYS.value: json.dumps(["key1", "key2"]),
LangfuseSpanAttributes.DEBUG_LANGFUSE.value: True,
}
@ -261,9 +236,7 @@ class TestLangfuseOtelIntegration:
for call in mock_safe_set_attribute.call_args_list
}
assert (
actual == expected
), "Mismatch between expected and actual OTEL attribute mapping."
assert actual == expected, "Mismatch between expected and actual OTEL attribute mapping."
@pytest.mark.parametrize(
"metadata, expected_version",
@ -288,16 +261,10 @@ class TestLangfuseOtelIntegration:
def test_version_emitted_on_langfuse_v4_key(self, metadata, expected_version):
kwargs = {"litellm_params": {"metadata": {"trace_release": "rel-9", **metadata}}}
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(
MagicMock(), kwargs, None
)
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None)
emitted = {
call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list
}
emitted = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list}
if expected_version is None:
assert "langfuse.version" not in emitted
@ -335,12 +302,8 @@ class TestLangfuseOtelIntegration:
"messages": [{"role": "user", "content": "What's the weather in Tokyo?"}],
}
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(
MagicMock(), kwargs, response_obj
)
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, response_obj)
expect_output = {
LangfuseSpanAttributes.OBSERVATION_INPUT.value: [
@ -356,11 +319,10 @@ class TestLangfuseOtelIntegration:
actual = {
call.args[1]: json.loads(call.args[2])
for call in mock_safe_set_attribute.call_args_list
if call.args[1] in expect_output
}
assert (
actual == expect_output
), "Mismatch in observation input/output OTEL attributes."
assert actual == expect_output, "Mismatch in observation input/output OTEL attributes."
def test_set_langfuse_specific_attributes_with_tool_calls(self):
"""Test that _set_langfuse_specific_attributes correctly sets observation.output with tool calls in Langfuse format."""
@ -384,9 +346,7 @@ class TestLangfuseOtelIntegration:
"content": None,
"tool_calls": [
ChatCompletionMessageToolCall(
function=Function(
arguments='{"location":"Tokyo"}', name="get_weather"
),
function=Function(arguments='{"location":"Tokyo"}', name="get_weather"),
id="call_123",
type="function",
)
@ -396,12 +356,8 @@ class TestLangfuseOtelIntegration:
],
)
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(
MagicMock(), {}, response_obj
)
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), {}, response_obj)
expected = {
LangfuseSpanAttributes.OBSERVATION_OUTPUT.value: [
@ -419,10 +375,9 @@ class TestLangfuseOtelIntegration:
actual = {
call.args[1]: json.loads(call.args[2])
for call in mock_safe_set_attribute.call_args_list
if call.args[1] in expected
}
assert (
actual == expected
), "Mismatch in observation output OTEL attribute for tool calls."
assert actual == expected, "Mismatch in observation output OTEL attribute for tool calls."
def test_construct_dynamic_otel_headers_with_langfuse_keys(self):
"""Test that construct_dynamic_otel_headers creates proper auth headers when langfuse keys are provided."""
@ -573,9 +528,7 @@ class TestLangfuseOtelKeyDynamicConfig:
logger = LangfuseOtelLogger()
assert logger.OTEL_EXPORTER == "console"
tracer = logger.get_tracer_to_use_for_request(
{"standard_callback_dynamic_params": self._dynamic_params()}
)
tracer = logger.get_tracer_to_use_for_request({"standard_callback_dynamic_params": self._dynamic_params()})
assert tracer is not logger.tracer
assert len(logger._tracer_provider_cache) == 1
@ -636,9 +589,7 @@ class TestLangfuseOtelKeyDynamicConfig:
with self._clean_env():
logger = LangfuseOtelLogger()
with patch.object(otel_module.verbose_logger, "debug", side_effect=_spy):
logger.get_tracer_to_use_for_request(
{"standard_callback_dynamic_params": self._dynamic_params()}
)
logger.get_tracer_to_use_for_request({"standard_callback_dynamic_params": self._dynamic_params()})
logged = "\n".join(recorded_arguments)
assert "initializing span processor" in logged
@ -698,12 +649,8 @@ class TestLangfuseOtelResponsesAPI:
LangfuseLLMObsOTELAttributes,
)
with patch(
"litellm.integrations.arize._utils.set_attributes"
) as mock_set_attributes:
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
with patch("litellm.integrations.arize._utils.set_attributes") as mock_set_attributes:
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
logger = LangfuseOtelLogger()
logger.set_langfuse_otel_attributes(mock_span, kwargs, mock_response)
@ -716,9 +663,7 @@ class TestLangfuseOtelResponsesAPI:
mock_safe_set_attribute.assert_any_call(
mock_span, "langfuse.generation.name", "responses_test_generation"
)
mock_safe_set_attribute.assert_any_call(
mock_span, "langfuse.trace.name", "responses_api_trace"
)
mock_safe_set_attribute.assert_any_call(mock_span, "langfuse.trace.name", "responses_api_trace")
def test_responses_api_metadata_extraction(self):
"""Test that metadata is correctly extracted from ResponsesAPI kwargs."""
@ -768,9 +713,7 @@ class TestLangfuseOtelResponsesAPI:
mock_span = MagicMock()
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(mock_span, kwargs, {})
# Verify specific attributes were set
@ -812,11 +755,12 @@ class TestLangfuseOtelResponsesAPI:
def test_responses_api_with_output(self):
"""Test Langfuse OTEL logger with Responses API output (reasoning + message)."""
from openai.types.responses import (
ResponseReasoningItem,
ResponseOutputMessage,
ResponseOutputText,
ResponseReasoningItem,
)
from openai.types.responses.response_reasoning_item import Summary
from litellm.types.integrations.langfuse_otel import LangfuseSpanAttributes
# Create Responses API response with reasoning and message
@ -831,7 +775,11 @@ class TestLangfuseOtelResponsesAPI:
Summary(
text="Let me analyze this problem step by step...",
type="summary_text",
)
),
Summary(
text="Now checking the forecast data.",
type="summary_text",
),
],
),
ResponseOutputMessage(
@ -852,21 +800,15 @@ class TestLangfuseOtelResponsesAPI:
kwargs = {
"call_type": "responses",
"messages": [
{"role": "user", "content": "What's the weather in San Francisco?"}
],
"messages": [{"role": "user", "content": "What's the weather in San Francisco?"}],
"model": "gpt-4o",
"optional_params": {},
}
mock_span = MagicMock()
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(
mock_span, kwargs, response_obj
)
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(mock_span, kwargs, response_obj)
# Verify observation output was set
output_calls = [
@ -879,29 +821,33 @@ class TestLangfuseOtelResponsesAPI:
output_json = output_calls[0].args[2]
output_data = json.loads(output_json)
# Verify output contains reasoning and message
# Verify output contains both reasoning summaries and the message
assert isinstance(output_data, list)
assert len(output_data) == 2
assert len(output_data) == 3
# Verify reasoning summary
assert output_data[0]["role"] == "reasoning_summary"
assert (
output_data[0]["content"]
== "Let me analyze this problem step by step..."
)
assert output_data[0]["content"] == "Let me analyze this problem step by step..."
assert output_data[1]["role"] == "reasoning_summary"
assert output_data[1]["content"] == "Now checking the forecast data."
# Verify message
assert output_data[1]["role"] == "assistant"
assert (
output_data[1]["content"]
== "The weather in San Francisco is sunny, 20°C."
)
assert output_data[2]["role"] == "assistant"
assert output_data[2]["content"] == "The weather in San Francisco is sunny, 20°C."
trace_output_calls = [
call
for call in mock_safe_set_attribute.call_args_list
if call.args[1] == LangfuseSpanAttributes.TRACE_OUTPUT.value
]
assert len(trace_output_calls) > 0, "trace.output should be set"
assert trace_output_calls[0].args[2] == output_json
def test_responses_api_with_function_calls(self):
"""Test Langfuse OTEL logger with Responses API function_call output."""
from litellm.types.integrations.langfuse_otel import LangfuseSpanAttributes
from openai.types.responses import ResponseFunctionToolCall
from litellm.types.integrations.langfuse_otel import LangfuseSpanAttributes
# Create Responses API response with function call
response_obj = ResponsesAPIResponse(
id="response-789",
@ -920,21 +866,15 @@ class TestLangfuseOtelResponsesAPI:
kwargs = {
"call_type": "responses",
"messages": [
{"role": "user", "content": "What's the weather in San Francisco?"}
],
"messages": [{"role": "user", "content": "What's the weather in San Francisco?"}],
"model": "gpt-4o",
"optional_params": {},
}
mock_span = MagicMock()
with patch(
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(
mock_span, kwargs, response_obj
)
with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute:
LangfuseOtelLogger._set_langfuse_specific_attributes(mock_span, kwargs, response_obj)
# Verify observation output was set
output_calls = [
@ -989,9 +929,11 @@ class TestLangfuseOtelResponsesAPI:
mock_span = MagicMock()
with patch( # test-quality-ok: the span attribute sink is the observable boundary; sibling tests in this class stub the same seam
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute:
with (
patch( # test-quality-ok: the span attribute sink is the observable boundary; sibling tests in this class stub the same seam
"litellm.integrations.arize._utils.safe_set_attribute"
) as mock_safe_set_attribute
):
LangfuseOtelLogger._set_langfuse_specific_attributes(mock_span, kwargs, response_obj)
output_calls = [
@ -1008,3 +950,250 @@ class TestLangfuseOtelResponsesAPI:
if __name__ == "__main__":
pytest.main([__file__])
class _RecordingSpan:
def __init__(self) -> None:
self.attributes: dict[str, object] = {}
def set_attribute(self, key: str, value: object) -> None:
self.attributes[key] = value
def _emitted(kwargs: Mapping[str, object], response_obj: object = None) -> Mapping[str, object]:
span: Final = _RecordingSpan()
LangfuseOtelLogger._set_langfuse_specific_attributes(span, kwargs, response_obj)
return span.attributes
class TestDerivedTraceFields:
def test_chat_call_derives_trace_name_input_and_output(self):
from litellm.types.utils import Choices, ModelResponse
response_obj = ModelResponse(
id="chatcmpl-derived",
model="gpt-4o",
choices=[
Choices(
finish_reason="stop",
message={"role": "assistant", "content": "Sunny."},
)
],
)
attributes = _emitted(
{
"call_type": "acompletion",
"messages": [{"role": "user", "content": "weather?"}],
"litellm_params": {"metadata": {}},
},
response_obj,
)
assert attributes["langfuse.trace.name"] == "litellm-acompletion"
assert attributes["langfuse.trace.input"] == attributes["langfuse.observation.input"]
assert attributes["langfuse.trace.output"] == attributes["langfuse.observation.output"]
assert json.loads(attributes["langfuse.trace.input"]) == [{"role": "user", "content": "weather?"}]
assert json.loads(attributes["langfuse.trace.output"]) == {"role": "assistant", "content": "Sunny."}
def test_missing_call_type_defaults_to_completion(self):
attributes = _emitted(
{
"messages": [{"role": "user", "content": "hi"}],
"litellm_params": {"metadata": {}},
}
)
assert attributes["langfuse.trace.name"] == "litellm-completion"
def test_caller_trace_name_wins_and_is_the_only_name(self):
attributes = _emitted(
{
"call_type": "acompletion",
"litellm_params": {"metadata": {"trace_name": "caller-trace"}},
}
)
assert attributes["langfuse.trace.name"] == "caller-trace"
def test_existing_trace_id_leaves_derived_trace_fields_off(self, monkeypatch):
import litellm
from litellm.types.utils import Choices, ModelResponse
monkeypatch.setattr(litellm, "langfuse_default_tags", ["cache_hit"])
response_obj = ModelResponse(
id="chatcmpl-joined",
model="gpt-4o",
choices=[
Choices(
finish_reason="stop",
message={"role": "assistant", "content": "hi"},
)
],
)
attributes = _emitted(
{
"call_type": "acompletion",
"messages": [{"role": "user", "content": "hi"}],
"cache_hit": True,
"litellm_params": {
"metadata": {"existing_trace_id": "abc123", "tags": ["experiment"]},
},
"standard_logging_object": {"request_tags": ["injected"]},
},
response_obj,
)
assert "langfuse.trace.name" not in attributes
assert attributes["langfuse.trace.existing_id"] == "abc123"
assert "langfuse.trace.input" not in attributes
assert "langfuse.trace.output" not in attributes
assert "langfuse.trace.tags" not in attributes
assert "langfuse.observation.input" in attributes
assert "langfuse.observation.output" in attributes
def test_responses_api_output_mirrors_observation_output_items(self):
from openai.types.responses import ResponseFunctionToolCall
response_obj = ResponsesAPIResponse(
id="response-derived",
created_at=1625247700,
output=[
ResponseFunctionToolCall(
id="fc-1",
type="function_call",
name="get_weather",
call_id="call-1",
arguments='{"city":"sf"}',
status="completed",
)
],
)
attributes = _emitted(
{
"call_type": "aresponses",
"litellm_params": {"metadata": {}},
},
response_obj,
)
assert attributes["langfuse.trace.name"] == "litellm-aresponses"
assert attributes["langfuse.trace.output"] == attributes["langfuse.observation.output"]
assert json.loads(attributes["langfuse.trace.output"]) == [
{
"id": "fc-1",
"name": "get_weather",
"call_id": "call-1",
"type": "function_call",
"arguments": {"city": "sf"},
}
]
def test_tags_merge_caller_request_and_default_tags_in_order(self, monkeypatch):
import litellm
monkeypatch.setattr(litellm, "langfuse_default_tags", ["cache_hit", "user_api_key_alias"])
attributes = _emitted(
{
"call_type": "acompletion",
"cache_hit": True,
"litellm_params": {
"metadata": {"tags": ["a", "b"], "user_api_key_alias": "k1"},
},
"standard_logging_object": {"request_tags": ["b", "c"]},
}
)
assert json.loads(attributes["langfuse.trace.tags"]) == [
"a",
"b",
"c",
"cache_hit:True",
"user_api_key_alias:k1",
]
def test_cache_key_default_tag_falls_back_to_preset_cache_key(self, monkeypatch):
import litellm
class _PresetCache:
@staticmethod
def _get_preset_cache_key_from_kwargs(**kwargs) -> str:
return "preset-abc"
monkeypatch.setattr(litellm, "langfuse_default_tags", ["cache_key"])
monkeypatch.setattr(litellm, "cache", _PresetCache())
attributes = _emitted(
{
"call_type": "acompletion",
"litellm_params": {"metadata": {}},
}
)
assert json.loads(attributes["langfuse.trace.tags"]) == ["cache_key:preset-abc"]
def test_cache_key_default_tag_prefers_hidden_params_over_preset(self, monkeypatch):
import litellm
class _PresetCache:
@staticmethod
def _get_preset_cache_key_from_kwargs(**kwargs) -> str:
return "preset-abc"
monkeypatch.setattr(litellm, "langfuse_default_tags", ["cache_key"])
monkeypatch.setattr(litellm, "cache", _PresetCache())
attributes = _emitted(
{
"call_type": "acompletion",
"litellm_params": {"metadata": {"hidden_params": {"cache_key": "explicit-key"}}},
}
)
assert json.loads(attributes["langfuse.trace.tags"]) == ["cache_key:explicit-key"]
def test_caller_tags_string_is_a_single_tag(self):
attributes = _emitted(
{
"call_type": "acompletion",
"litellm_params": {"metadata": {"tags": "single"}},
}
)
assert json.loads(attributes["langfuse.trace.tags"]) == ["single"]
def test_no_tags_anywhere_leaves_trace_tags_unset(self):
attributes = _emitted(
{
"call_type": "acompletion",
"litellm_params": {"metadata": {}},
}
)
assert "langfuse.trace.tags" not in attributes
def test_empty_string_role_is_omitted_from_output(self):
response_obj = {
"id": "chatcmpl-norole",
"choices": [{"message": {"role": "", "content": "hi"}}],
}
attributes = _emitted(
{
"call_type": "acompletion",
"litellm_params": {"metadata": {}},
},
response_obj,
)
assert json.loads(attributes["langfuse.observation.output"]) == {"content": "hi"}
def test_no_messages_leaves_trace_input_unset_but_emits_output(self):
from litellm.types.utils import Choices, ModelResponse
response_obj = ModelResponse(
id="chatcmpl-nomsg",
model="gpt-4o",
choices=[
Choices(
finish_reason="stop",
message={"role": "assistant", "content": "hi"},
)
],
)
attributes = _emitted(
{
"call_type": "acompletion",
"litellm_params": {"metadata": {}},
},
response_obj,
)
assert "langfuse.trace.input" not in attributes
assert "langfuse.observation.input" not in attributes
assert attributes["langfuse.trace.output"] == attributes["langfuse.observation.output"]