diff --git a/litellm/integrations/langfuse/langfuse_otel.py b/litellm/integrations/langfuse/langfuse_otel.py index a96fac32c2a..334130a16ef 100644 --- a/litellm/integrations/langfuse/langfuse_otel.py +++ b/litellm/integrations/langfuse/langfuse_otel.py @@ -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 () diff --git a/litellm/integrations/otel/langfuse_logger.py b/litellm/integrations/otel/langfuse_logger.py index d029b153c52..9ef4da60564 100644 --- a/litellm/integrations/otel/langfuse_logger.py +++ b/litellm/integrations/otel/langfuse_logger.py @@ -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) diff --git a/litellm/integrations/otel/mappers/langfuse.py b/litellm/integrations/otel/mappers/langfuse.py index 68860931b76..f589e79f129 100644 --- a/litellm/integrations/otel/mappers/langfuse.py +++ b/litellm/integrations/otel/mappers/langfuse.py @@ -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), } diff --git a/litellm/types/integrations/langfuse_otel.py b/litellm/types/integrations/langfuse_otel.py index c58dc567cda..23b79e19d4d 100644 --- a/litellm/types/integrations/langfuse_otel.py +++ b/litellm/types/integrations/langfuse_otel.py @@ -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" diff --git a/tests/integration/observability/test_langfuse_otel_trace_fields.py b/tests/integration/observability/test_langfuse_otel_trace_fields.py new file mode 100644 index 00000000000..2db434ce243 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_trace_fields.py @@ -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") diff --git a/tests/unit/integrations/otel/test_langfuse_logger.py b/tests/unit/integrations/otel/test_langfuse_logger.py index aca9dcc8a5e..0eaac41276a 100644 --- a/tests/unit/integrations/otel/test_langfuse_logger.py +++ b/tests/unit/integrations/otel/test_langfuse_logger.py @@ -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( diff --git a/tests/unit/integrations/otel/test_otel_v2_vendor_mappers.py b/tests/unit/integrations/otel/test_otel_v2_vendor_mappers.py index cdff9c960f3..6309ba898b4 100644 --- a/tests/unit/integrations/otel/test_otel_v2_vendor_mappers.py +++ b/tests/unit/integrations/otel/test_otel_v2_vendor_mappers.py @@ -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"] diff --git a/tests/unit/integrations/test_langfuse_otel.py b/tests/unit/integrations/test_langfuse_otel.py index 0a9ce55fe16..ff9acc9085e 100644 --- a/tests/unit/integrations/test_langfuse_otel.py +++ b/tests/unit/integrations/test_langfuse_otel.py @@ -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"]