mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
fix(langfuse_otel): derive trace-level name, input, output and tags on the generation span
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
40ec84caa2
commit
dfd956a788
9 changed files with 760 additions and 229 deletions
|
|
@ -95,7 +95,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 +125,41 @@ 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) -> 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: dict, metadata: dict) -> 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 = (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 = metadata.get("hidden_params", {}) or {}
|
||||
return f"cache_key:{hidden_params.get('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 = (
|
||||
(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 +188,26 @@ 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:
|
||||
if metadata.get("trace_name") is None and metadata.get("existing_trace_id") is None:
|
||||
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)
|
||||
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)
|
||||
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)
|
||||
safe_set_attribute(span, LangfuseSpanAttributes.TRACE_OUTPUT.value, output_json)
|
||||
|
||||
@staticmethod
|
||||
def _get_langfuse_otel_host() -> str | None:
|
||||
|
|
@ -445,3 +399,83 @@ class LangfuseOtelLogger(OpenTelemetry):
|
|||
"""
|
||||
Langfuse should not receive service failure logs.
|
||||
"""
|
||||
|
||||
|
||||
def _extract_choices_output(response_obj) -> str | None:
|
||||
from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
|
||||
|
||||
choices: Final = response_obj.get("choices", [])
|
||||
if not choices:
|
||||
return None
|
||||
message: Final = choices[0].get("message", {})
|
||||
tool_calls: Final = message.get("tool_calls")
|
||||
if tool_calls:
|
||||
transformed_tool_calls: Final = [
|
||||
{
|
||||
"id": response_obj.get("id", ""),
|
||||
"name": tool_call.get("function", {}).get("name", ""),
|
||||
"call_id": tool_call.get("id", ""),
|
||||
"type": "function_call",
|
||||
"arguments": _tool_call_arguments(tool_call.get("function", {}).get("arguments", "{}")),
|
||||
}
|
||||
for tool_call in tool_calls
|
||||
]
|
||||
return safe_dumps(transformed_tool_calls)
|
||||
output_data: Final = {
|
||||
key: value
|
||||
for key, value in (
|
||||
("role", message.get("role")),
|
||||
("content", message.get("content")),
|
||||
)
|
||||
if value is not None
|
||||
}
|
||||
return safe_dumps(output_data) if output_data else None
|
||||
|
||||
|
||||
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) -> str | None:
|
||||
from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
|
||||
|
||||
output: Final = response_obj.get("output", [])
|
||||
if not output:
|
||||
return None
|
||||
output_items: Final = tuple(_output_item(item) for item in output)
|
||||
rendered: Final = tuple(item for item in output_items if item is not None)
|
||||
return safe_dumps(list(rendered)) if rendered else None
|
||||
|
||||
|
||||
def _output_item(item) -> dict | None:
|
||||
if not hasattr(item, "type"):
|
||||
return None
|
||||
if item.type == "reasoning" and hasattr(item, "summary"):
|
||||
return next(
|
||||
(
|
||||
{"role": "reasoning_summary", "content": summary.text}
|
||||
for summary in item.summary
|
||||
if hasattr(summary, "text")
|
||||
),
|
||||
None,
|
||||
)
|
||||
if item.type == "message":
|
||||
return {
|
||||
"role": getattr(item, "role", "assistant"),
|
||||
"content": getattr(getattr(item, "content", [{}])[0], "text", ""),
|
||||
}
|
||||
if item.type == "function_call":
|
||||
arguments: Final = getattr(item, "arguments", "{}")
|
||||
return {
|
||||
"id": getattr(item, "id", ""),
|
||||
"name": getattr(item, "name", ""),
|
||||
"call_id": getattr(item, "call_id", ""),
|
||||
"type": "function_call",
|
||||
"arguments": safe_json_loads(arguments, default={}) if isinstance(arguments, str) else arguments,
|
||||
}
|
||||
return None
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
@ -67,10 +79,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
|
||||
|
|
@ -85,10 +97,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),
|
||||
|
|
@ -99,6 +111,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),
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -257,6 +257,24 @@
|
|||
"tests/integration/observability/test_otel_text_completion_choices.py::test_otel_weave_output_keeps_text_completion_provider_fields_beside_the_synthesized_message": [
|
||||
"other.observability.otel.text_completion_choices_keep_provider_fields"
|
||||
],
|
||||
"tests/integration/observability/test_langfuse_otel_trace_fields.py::test_chat_completion_derives_trace_name_input_and_output": [
|
||||
"other.observability.langfuse_otel.chat_derives_trace_fields"
|
||||
],
|
||||
"tests/integration/observability/test_langfuse_otel_trace_fields.py::test_parented_chat_completion_derives_trace_fields_inside_inbound_trace": [
|
||||
"other.observability.langfuse_otel.parented_chat_derives_trace_fields_and_keeps_parent"
|
||||
],
|
||||
"tests/integration/observability/test_langfuse_otel_trace_fields.py::test_streaming_chat_completion_derives_trace_fields": [
|
||||
"other.observability.langfuse_otel.streaming_chat_derives_trace_fields"
|
||||
],
|
||||
"tests/integration/observability/test_langfuse_otel_trace_fields.py::test_caller_supplied_trace_name_and_tags_are_emitted": [
|
||||
"other.observability.langfuse_otel.caller_trace_name_and_tags_win"
|
||||
],
|
||||
"tests/integration/observability/test_langfuse_otel_trace_fields.py::test_responses_call_derives_trace_fields": [
|
||||
"other.observability.langfuse_otel.responses_derives_trace_fields"
|
||||
],
|
||||
"tests/integration/observability/test_langfuse_otel_trace_fields.py::test_otel_v2_derives_trace_name_input_and_output": [
|
||||
"other.observability.langfuse_otel.v2_derives_trace_fields"
|
||||
],
|
||||
"tests/integration/observability/test_guardrail_effects.py::test_guardrail_rewrites_system_and_user_in_actual_anthropic_request": [
|
||||
"other.observability.guardrails.rewrite_reaches_correct_anthropic_positions"
|
||||
],
|
||||
|
|
|
|||
|
|
@ -0,0 +1,332 @@
|
|||
import json
|
||||
import uuid
|
||||
from collections.abc import Iterator
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
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"
|
||||
|
||||
|
||||
@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
|
||||
return Reply(body=b"")
|
||||
|
||||
|
||||
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:
|
||||
batches: list[Request] = [] # mutable-ok: accumulated across eventually() polls
|
||||
|
||||
def spans() -> tuple[SpanRecord, ...]:
|
||||
batches.extend(collector.drain())
|
||||
return tuple(
|
||||
record
|
||||
for request in batches
|
||||
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:
|
||||
config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())
|
||||
config["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
|
||||
|
||||
def upstream(request: Request) -> Reply:
|
||||
assert request.target.endswith("/chat/completions"), request.target
|
||||
return 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]
|
||||
|
||||
def upstream(request: Request) -> Reply:
|
||||
assert request.target.endswith("/chat/completions"), request.target
|
||||
return 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
|
||||
|
||||
def upstream(request: Request) -> Reply:
|
||||
assert request.target.endswith("/chat/completions"), request.target
|
||||
return Reply(content_type="text/event-stream", chunks=_chat_stream(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}],
|
||||
"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
|
||||
|
||||
def upstream(request: Request) -> Reply:
|
||||
assert request.target.endswith("/chat/completions"), request.target
|
||||
return 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
|
||||
|
||||
def upstream(request: Request) -> Reply:
|
||||
assert request.target.endswith("/responses"), request.target
|
||||
return 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
|
||||
|
||||
def upstream(request: Request) -> Reply:
|
||||
assert request.target.endswith("/chat/completions"), request.target
|
||||
return 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")
|
||||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -332,3 +332,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"]
|
||||
|
|
|
|||
|
|
@ -104,19 +104,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 +119,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 +136,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 +203,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
|
||||
|
|
@ -249,9 +225,7 @@ class TestLangfuseOtelIntegration:
|
|||
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 +235,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 +260,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 +301,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 +318,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 +345,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 +355,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 +374,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 +527,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 +588,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 +648,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 +662,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 +712,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 +754,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
|
||||
|
|
@ -852,21 +795,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 = [
|
||||
|
|
@ -885,23 +822,18 @@ class TestLangfuseOtelResponsesAPI:
|
|||
|
||||
# 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..."
|
||||
|
||||
# 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[1]["content"] == "The weather in San Francisco is sunny, 20°C."
|
||||
|
||||
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 +852,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 +915,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 +936,166 @@ class TestLangfuseOtelResponsesAPI:
|
|||
|
||||
if __name__ == "__main__":
|
||||
pytest.main([__file__])
|
||||
|
||||
|
||||
class _RecordingSpan:
|
||||
def __init__(self):
|
||||
self.attributes = {}
|
||||
|
||||
def set_attribute(self, key, value):
|
||||
self.attributes[key] = value
|
||||
|
||||
|
||||
def _emitted(kwargs, response_obj=None):
|
||||
span = _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_without_trace_name_emits_no_name(self):
|
||||
attributes = _emitted(
|
||||
{
|
||||
"call_type": "acompletion",
|
||||
"litellm_params": {"metadata": {"existing_trace_id": "abc123"}},
|
||||
}
|
||||
)
|
||||
assert "langfuse.trace.name" not in attributes
|
||||
assert attributes["langfuse.trace.existing_id"] == "abc123"
|
||||
|
||||
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_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_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"]
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue