diff --git a/litellm/integrations/langfuse/langfuse_otel.py b/litellm/integrations/langfuse/langfuse_otel.py index a96fac32c2a..8a36b45707e 100644 --- a/litellm/integrations/langfuse/langfuse_otel.py +++ b/litellm/integrations/langfuse/langfuse_otel.py @@ -29,6 +29,14 @@ LANGFUSE_CLOUD_US_ENDPOINT: Final = "https://us.cloud.langfuse.com/api/public/ot LANGFUSE_INGESTION_VERSION_HEADER: Final = "x-langfuse-ingestion-version" LANGFUSE_INGESTION_VERSION: Final = "4" +_TRACE_IDENTITY_FIELDS: Final = ( + "user_api_key_alias", + "user_api_key_user_id", + "user_api_key_end_user_id", + "user_api_key_team_id", + "user_api_key_team_alias", +) + class LangfuseOtelLogger(OpenTelemetry): def __init__(self, config=None, *args, **kwargs): @@ -125,6 +133,30 @@ class LangfuseOtelLogger(OpenTelemetry): value = str(value) safe_set_attribute(span, enum_attr.value, value) + @staticmethod + def _set_request_metadata_attributes(span: Span, kwargs: dict[str, object]) -> None: + from litellm.integrations.arize._utils import safe_set_attribute + from litellm.integrations.langfuse.langfuse import log_requester_metadata + from litellm.litellm_core_utils.redact_messages import redact_user_api_key_info + from litellm.litellm_core_utils.safe_json_dumps import safe_dumps + + standard_logging_object: Final = kwargs.get("standard_logging_object") + request_metadata: Final = ( + standard_logging_object.get("metadata") if isinstance(standard_logging_object, dict) else None + ) + if not isinstance(request_metadata, dict): + return + observation_metadata: Final = { # mutable-ok: stays a real dict for safe_dumps and .get reads below + key: value + for key, value in log_requester_metadata(redact_user_api_key_info(metadata=request_metadata)).items() + if value is not None + } + safe_set_attribute(span, LangfuseSpanAttributes.OBSERVATION_METADATA.value, safe_dumps(observation_metadata)) + trace_prefix: Final = LangfuseSpanAttributes.TRACE_METADATA.value + for field in _TRACE_IDENTITY_FIELDS: + if value := observation_metadata.get(field): + safe_set_attribute(span, f"{trace_prefix}.{field}", value) + @staticmethod def _set_observation_output(span: Span, response_obj): """Helper to set observation output attributes.""" @@ -244,6 +276,7 @@ class LangfuseOtelLogger(OpenTelemetry): metadata: Final = LangfuseOtelLogger._extract_langfuse_metadata(kwargs) LangfuseOtelLogger._set_metadata_attributes(span=span, metadata=metadata) + LangfuseOtelLogger._set_request_metadata_attributes(span=span, kwargs=kwargs) messages: Final = kwargs.get("messages") if messages: diff --git a/litellm/integrations/otel/mappers/langfuse.py b/litellm/integrations/otel/mappers/langfuse.py index 68860931b76..68895d996ce 100644 --- a/litellm/integrations/otel/mappers/langfuse.py +++ b/litellm/integrations/otel/mappers/langfuse.py @@ -12,6 +12,7 @@ the JSON-serialized payloads. ``trace_attributes`` maps the caller's trace contr import json from collections.abc import Callable +from types import MappingProxyType from typing import Final from litellm.integrations.otel.mappers.base import AttributeMap, AttrValue, SpanData @@ -28,6 +29,8 @@ from litellm.integrations.otel.model.payloads import ( LLMUsage, ) from litellm.integrations.otel.model.trace_controls import TraceControls +from litellm.integrations.otel.model.utils import as_str +from litellm.litellm_core_utils.safe_json_dumps import safe_dumps LANGFUSE_OBSERVATION_INPUT: Final = "langfuse.observation.input" LANGFUSE_OBSERVATION_OUTPUT: Final = "langfuse.observation.output" @@ -35,6 +38,19 @@ 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" +LANGFUSE_OBSERVATION_METADATA: Final = "langfuse.observation.metadata" +LANGFUSE_TRACE_METADATA_PREFIX: Final = "langfuse.trace.metadata." +TRACE_IDENTITY_FIELDS: Final = ( + "user_api_key_alias", + "user_api_key_user_id", + "user_api_key_end_user_id", + "user_api_key_team_id", + "user_api_key_team_alias", +) + + +def _identity_field(name: str) -> Callable[[LLMCallSpanData], AttrValue | None]: + return lambda d: as_str(d.request_metadata.get(name)) or None class LangfuseMapper: @@ -45,6 +61,9 @@ class LangfuseMapper: "langfuse.observation.id": lambda d: d.identity.call_id or None, "langfuse.trace.metadata.team_id": lambda d: d.identity.team_id or None, "langfuse.trace.metadata.team_alias": lambda d: d.identity.team_alias or None, + **MappingProxyType( + {f"{LANGFUSE_TRACE_METADATA_PREFIX}{name}": _identity_field(name) for name in TRACE_IDENTITY_FIELDS} + ), } # Sub-tables folded into their respective JSON blobs. @@ -68,6 +87,11 @@ class LangfuseMapper: # JSON-payload attributes: each builder returns the serialized blob or None. _BLOB_ATTRS: dict[str, Callable[[LLMCallSpanData], AttrValue | None]] = { + LANGFUSE_OBSERVATION_METADATA: lambda d: ( + safe_dumps(dict(d.request_metadata)) # mutable-ok: safe_dumps only serializes real dicts + if d.request_metadata + else None + ), "langfuse.observation.model.parameters": lambda d: json_if( collect(LangfuseMapper._MODEL_PARAMS, d.request_params) ), diff --git a/litellm/integrations/otel/model/metadata.py b/litellm/integrations/otel/model/metadata.py index 7cb64debfe0..e716fe88103 100644 --- a/litellm/integrations/otel/model/metadata.py +++ b/litellm/integrations/otel/model/metadata.py @@ -42,9 +42,11 @@ from types import MappingProxyType from typing import TYPE_CHECKING, Any, Final, cast from litellm.constants import LITELLM_LOGGING_NO_UPSTREAM_LLM_CALL, SESSION_ID_GENERATED_METADATA_KEY +from litellm.integrations.langfuse.langfuse import log_requester_metadata from litellm.integrations.otel.model.semconv import resolve_operation from litellm.integrations.otel.model.trace_controls import TraceControls, caller_trace_controls from litellm.integrations.otel.model.utils import as_str, as_str_mapping, to_seconds +from litellm.litellm_core_utils.redact_messages import redact_user_api_key_info if TYPE_CHECKING: from litellm.types.utils import StandardLoggingPayload @@ -417,6 +419,23 @@ def _model_info_id(model_info: object) -> str | None: return None +def exported_request_metadata(payload: StandardLoggingPayload) -> Mapping[str, object]: + """The request metadata as the logging callbacks export it: ``user_api_key_*`` + dropped when ``litellm.redact_user_api_key_info`` is on, header keys nested + under ``requester_metadata``, ``None`` values dropped.""" + raw_meta: Final = cast( # cast-ok: StandardLoggingPayload.get returns Any | None + Mapping[str, object], payload.get("metadata") or MappingProxyType({}) + ) + redacted: Final = cast( # cast-ok: redact_user_api_key_info is untyped + Mapping[str, object], + redact_user_api_key_info(metadata=dict(raw_meta)), # mutable-ok: the redactor requires a real dict + ) + exported: Final = cast( # cast-ok: log_requester_metadata is untyped + Mapping[str, object], log_requester_metadata(redacted) + ) + return MappingProxyType({key: value for key, value in exported.items() if value is not None}) + + def _team_metadata_dict(value: object) -> Mapping[str, object] | None: """The team's free-form metadata as a raw mapping, or ``None`` when missing or empty. diff --git a/litellm/integrations/otel/model/payloads.py b/litellm/integrations/otel/model/payloads.py index e06cc1d0407..42174c4d350 100644 --- a/litellm/integrations/otel/model/payloads.py +++ b/litellm/integrations/otel/model/payloads.py @@ -12,7 +12,7 @@ from urllib.parse import urlsplit from typing_extensions import ReadOnly, TypedDict -from litellm.integrations.otel.model.metadata import RequestContext, RequestIdentity +from litellm.integrations.otel.model.metadata import RequestContext, RequestIdentity, exported_request_metadata from litellm.integrations.otel.model.semconv import ( GenAIOperation, GenAIOutputType, @@ -409,6 +409,7 @@ class LLMCallSpanData: response_cost: float | None server: ServerInfo | None identity: RequestIdentity + request_metadata: Mapping[str, object] = field(default_factory=dict) is_streaming: bool | None = None cost: LLMCost = field(default_factory=LLMCost) tools: tuple[ToolDefinition, ...] = () @@ -476,6 +477,7 @@ class LLMCallSpanData: cost=LLMCost.from_breakdown(cast("Mapping[str, object] | None", payload.get("cost_breakdown"))), server=ServerInfo.from_api_base(context.api_base), identity=context.identity, + request_metadata=exported_request_metadata(payload), is_streaming=as_bool(payload.get("stream")), tools=_extract_tools(params), messages_in=_dicts(payload.get("messages")) if capture_content else (), diff --git a/litellm/types/integrations/langfuse_otel.py b/litellm/types/integrations/langfuse_otel.py index c58dc567cda..11f30a2bfa7 100644 --- a/litellm/types/integrations/langfuse_otel.py +++ b/litellm/types/integrations/langfuse_otel.py @@ -29,6 +29,7 @@ class LangfuseSpanAttributes(str, Enum): # ---- Observation input/output ---- OBSERVATION_INPUT = "langfuse.observation.input" OBSERVATION_OUTPUT = "langfuse.observation.output" + OBSERVATION_METADATA = "langfuse.observation.metadata" # ---- Trace-level metadata ---- TRACE_USER_ID = "user.id" diff --git a/tests/integration/observability/test_langfuse_otel_metadata.py b/tests/integration/observability/test_langfuse_otel_metadata.py new file mode 100644 index 00000000000..bf18bc967f7 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_metadata.py @@ -0,0 +1,1129 @@ +import base64 +import binascii +import json +import threading +import uuid +from collections.abc import Callable, Iterator, Mapping +from concurrent.futures import ThreadPoolExecutor, wait +from contextlib import contextmanager +from dataclasses import dataclass +from pathlib import Path +from typing import Final + +import httpx +import psutil +import yaml +from anthropic import Anthropic, AsyncAnthropic +from integration._support.client import Gateway, Scenario, eventually +from integration._support.process import OwnedProxy, owned_proxy_process +from integration._support.wire import Reply, Request, Wire, wire_server +from openai import AsyncOpenAI, OpenAI +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceRequest + +_IDENTITIES: Final = ( + "user_api_key_alias", + "user_api_key_user_id", + "user_api_key_end_user_id", + "user_api_key_team_id", + "user_api_key_team_alias", +) + + +@dataclass(frozen=True, slots=True) +class _Rig: + candidate: Gateway + provider: Wire + collector: Wire + proxy: OwnedProxy + + +def _span_attributes(body: bytes) -> tuple[dict[str, object], ...]: + request: Final = ExportTraceServiceRequest.FromString(body) + return tuple( + {attribute.key: getattr(attribute.value, attribute.value.WhichOneof("value")) for attribute in span.attributes} + for resource in request.resource_spans + for scope in resource.scope_spans + for span in scope.spans + ) + + +def _sse_chat(response_id: str) -> tuple[bytes, ...]: + chunk: Final = { + "id": response_id, + "object": "chat.completion.chunk", + "created": 1, + "model": "gpt-4o-mini", + } + return ( + f"data: {json.dumps({**chunk, 'choices': [{'index': 0, 'delta': {'role': 'assistant', 'content': 'langfuse-echo'}, 'finish_reason': None}]})}\n\n".encode(), + f"data: {json.dumps({**chunk, 'choices': [{'index': 0, 'delta': {}, 'finish_reason': 'stop'}]})}\n\n".encode(), + f"data: {json.dumps({**chunk, 'choices': [], 'usage': {'prompt_tokens': 3, 'completion_tokens': 2, 'total_tokens': 5}})}\n\n".encode(), + b"data: [DONE]\n\n", + ) + + +def _sse_messages(response_id: str) -> tuple[bytes, ...]: + def event(name: str, payload: Mapping[str, object]) -> bytes: + return f"event: {name}\ndata: {json.dumps(payload)}\n\n".encode() + + message: Final = { + "id": response_id, + "type": "message", + "role": "assistant", + "model": "claude-sonnet-4-6", + "content": [], + "stop_reason": None, + "usage": {"input_tokens": 3, "output_tokens": 1}, + } + return ( + event("message_start", {"type": "message_start", "message": message}), + event( + "content_block_start", + {"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}}, + ), + event( + "content_block_delta", + {"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": "langfuse-echo"}}, + ), + event("content_block_stop", {"type": "content_block_stop", "index": 0}), + event( + "message_delta", + {"type": "message_delta", "delta": {"stop_reason": "end_turn"}, "usage": {"output_tokens": 2}}, + ), + event("message_stop", {"type": "message_stop"}), + ) + + +def _sse_responses(response_id: str) -> tuple[bytes, ...]: + def event(name: str, payload: Mapping[str, object]) -> bytes: + return f"event: {name}\ndata: {json.dumps(payload)}\n\n".encode() + + response: Final = { + "id": response_id, + "object": "response", + "model": "gpt-4o-mini", + "status": "completed", + "output": [ + { + "type": "message", + "id": f"item-{response_id}", + "status": "completed", + "role": "assistant", + "content": [{"type": "output_text", "text": "langfuse-echo"}], + } + ], + "usage": {"input_tokens": 3, "output_tokens": 2, "total_tokens": 5}, + } + return ( + event( + "response.created", + {"type": "response.created", "response": {**response, "status": "in_progress", "output": []}}, + ), + event( + "response.output_text.delta", + { + "type": "response.output_text.delta", + "response_id": response_id, + "item_id": f"item-{response_id}", + "output_index": 0, + "content_index": 0, + "delta": "langfuse-echo", + }, + ), + event("response.completed", {"type": "response.completed", "response": response}), + ) + + +def _chat_body(response_id: str) -> dict[str, object]: + return { + "id": response_id, + "object": "chat.completion", + "created": 1, + "model": "gpt-4o-mini", + "choices": [ + {"index": 0, "message": {"role": "assistant", "content": "langfuse-echo"}, "finish_reason": "stop"} + ], + "usage": {"prompt_tokens": 3, "completion_tokens": 2, "total_tokens": 5}, + } + + +def _message_body(response_id: str) -> dict[str, object]: + return { + "id": response_id, + "type": "message", + "role": "assistant", + "model": "claude-sonnet-4-6", + "content": [{"type": "text", "text": "langfuse-echo"}], + "stop_reason": "end_turn", + "usage": {"input_tokens": 3, "output_tokens": 2}, + } + + +_RESPONSE_FIELDS: Final = { + "created_at": 1, + "error": None, + "incomplete_details": None, + "instructions": None, + "metadata": {}, + "parallel_tool_calls": True, + "temperature": 1.0, + "tool_choice": "auto", + "tools": [], + "top_p": 1.0, + "max_output_tokens": None, + "previous_response_id": None, + "store": True, + "truncation": "disabled", + "user": None, +} + + +def _response_body(response_id: str) -> dict[str, object]: + return { + **_RESPONSE_FIELDS, + "id": response_id, + "object": "response", + "model": "gpt-4o-mini", + "status": "completed", + "output": [ + { + "type": "message", + "id": f"item-{response_id}", + "status": "completed", + "role": "assistant", + "content": [{"type": "output_text", "text": "langfuse-echo"}], + } + ], + "usage": {"input_tokens": 3, "output_tokens": 2, "total_tokens": 5}, + } + + +@contextmanager +def _langfuse_rig( + gateway: Gateway, + tmp_path: Path, + settings: Mapping[str, object], + marker: str, + *, + sink_reply: Callable[[Request], Reply] | None = None, + overrides: Mapping[str, str] | None = None, +) -> Iterator[_Rig]: + def upstream(request: Request) -> Reply: + if request.method == "GET": + return Reply(body=b'{"data":[]}') + body: Final = json.loads(request.body or b"{}") + suffix: Final = uuid.uuid4().hex[:8] + if "-fail-" in json.dumps(body): + return Reply( + status=500, body=b'{"error": {"message": "scripted provider failure", "type": "server_error"}}' + ) + if request.target.endswith("/chat/completions"): + chat_id: Final = f"chatcmpl-{marker}-{suffix}" + if body.get("stream"): + return Reply(content_type="text/event-stream", chunks=_sse_chat(chat_id)) + return Reply(body=json.dumps(_chat_body(chat_id)).encode()) + if request.target.endswith("/messages"): + message_id: Final = f"msg-{marker}-{suffix}" + if body.get("stream"): + return Reply(content_type="text/event-stream", chunks=_sse_messages(message_id)) + return Reply(body=json.dumps(_message_body(message_id)).encode()) + if request.target.endswith("/responses"): + resp_id: Final = f"resp-{marker}-{suffix}" + if body.get("stream"): + return Reply(content_type="text/event-stream", chunks=_sse_responses(resp_id)) + return Reply(body=json.dumps(_response_body(resp_id)).encode()) + raise AssertionError(request.target) + + def sink(request: Request) -> Reply: + return sink_reply(request) if sink_reply is not None else Reply() + + with wire_server(upstream) as provider, wire_server(sink) as collector: + raw_config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text()) + config: Final = { + **raw_config, + "litellm_settings": {**raw_config["litellm_settings"], "callbacks": ["langfuse_otel"], **settings}, + } + path: Final = tmp_path / "langfuse_otel.yaml" + path.write_text(yaml.safe_dump(config)) + with owned_proxy_process( + gateway, + tmp_path, + { + "LANGFUSE_PUBLIC_KEY": "pk-lf-integration", + "LANGFUSE_SECRET_KEY": "sk-lf-integration", + "LANGFUSE_HOST": collector.url, + **{ + name: (collector.url if value == "__collector__" else value) + for name, value in dict(overrides or {}).items() + }, + }, + config=path, + workers=2, + ) as owned: + yield _Rig(owned.gateway, provider, collector, owned) + + +def _response_id(value: object) -> str | None: + if not isinstance(value, str): + return None + if value.startswith("resp_"): + try: + decoded: Final = base64.b64decode(value[5:] + "===").decode() + except (binascii.Error, UnicodeDecodeError): + return value + _, _, suffix = decoded.partition("response_id:") + if suffix: + return suffix.split(";", 1)[0] + return value + + +def _span_response_id(span: Mapping[str, object]) -> str | None: + return _response_id(span.get("llm.response.id")) + + +def _watched_spans(collector: Wire) -> Callable[[], tuple[dict[str, object], ...]]: + drained: list[ + dict[str, object] + ] = [] # mutable-ok: drain() consumes the wire queue, so observed spans must accumulate across eventually polls + + def snapshot() -> tuple[dict[str, object], ...]: + drained.extend(attributes for batch in collector.drain() for attributes in _span_attributes(batch.body)) + return tuple(drained) + + return snapshot + + +def _observation_span(collector: Wire, response_id: str) -> dict[str, object]: + def spans() -> tuple[dict[str, object], ...]: + return tuple( + attributes + for batch in collector.drain() + for attributes in _span_attributes(batch.body) + if _span_response_id(attributes) == _response_id(response_id) + ) + + observed: Final = eventually(spans, lambda values: len(values) == 1, seconds=25) + return observed[0] + + +def _marked_span(collector: Wire, marker: str, needle: str) -> dict[str, object]: + def spans() -> tuple[dict[str, object], ...]: + return tuple( + attributes + for batch in collector.drain() + for attributes in _span_attributes(batch.body) + if needle in json.dumps(attributes, default=str) and marker in json.dumps(attributes, default=str) + ) + + observed: Final = eventually(spans, lambda values: len(values) == 1, seconds=25) + return observed[0] + + +def _scenario_identity( + scenario: Scenario, + provider: Wire, + *, + spend_logs_metadata: Mapping[str, object] | None = None, + team: bool = True, +) -> tuple[str, dict[str, str | None]]: + team_alias: Final = ("lf-team-" + uuid.uuid4().hex) if team else None + team_id: Final = scenario.team(team_alias=team_alias) if team else None + key_alias: Final = "lf-key-" + uuid.uuid4().hex + key: Final = scenario.key( + **({"team_id": team_id} if team_id else {}), + key_alias=key_alias, + metadata={"spend_logs_metadata": spend_logs_metadata or {"ticket": "LIT-8283"}}, + ) + model: Final = scenario.model(api_base=provider.url + "/v1") + end_user: Final = "lf-end-user-" + uuid.uuid4().hex + return key, { + "key_alias": key_alias, + "team_id": team_id, + "team_alias": team_alias, + "end_user": end_user, + "model": model, + } + + +def _assert_identity( + attrs: Mapping[str, object], + expected: Mapping[str, object], + *, + spend_logs_metadata: Mapping[str, object] | None = None, +) -> None: + assert "metadata" in attrs, sorted(attrs) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + assert {key: observation.get(key) for key in expected} == dict(expected), observation + assert not [key for key, value in observation.items() if value is None], observation + if spend_logs_metadata is not None: + assert observation["spend_logs_metadata"] == dict(spend_logs_metadata), observation + flattened: Final = { + attribute: attrs.get(attribute) + for attribute in (f"langfuse.trace.metadata.{field}" for field in _IDENTITIES) + if attribute in attrs + } + assert flattened == {f"langfuse.trace.metadata.{field}": value for field, value in expected.items()}, attrs + + +def _chat_identity_fields(identity: Mapping[str, str | None]) -> dict[str, object]: + return { + "user_api_key_alias": identity["key_alias"], + **({"user_api_key_team_id": identity["team_id"]} if identity["team_id"] else {}), + **({"user_api_key_team_alias": identity["team_alias"]} if identity["team_alias"] else {}), + **({"user_api_key_end_user_id": identity["end_user"]} if identity["end_user"] else {}), + } + + +def test_langfuse_otel_emits_request_metadata_under_langfuse_observation_and_trace_keys( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "lf-meta-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "user": identity["end_user"], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + assert response.json()["choices"][0]["message"]["content"] == "langfuse-echo" + + attrs: Final = _observation_span(rig.collector, response.json()["id"]) + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +def test_langfuse_otel_metadata_redacts_user_api_key_fields_like_vanilla_langfuse( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "lf-redact-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {"redact_user_api_key_info": True}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key: Final = scenario.key( + key_alias="lf-redact-" + uuid.uuid4().hex, + metadata={"spend_logs_metadata": {"ticket": "LIT-8283"}}, + ) + model: Final = scenario.model(api_base=rig.provider.url + "/v1") + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + + attrs: Final = _observation_span(rig.collector, response.json()["id"]) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + assert not [name for name in observation if name.startswith("user_api_key")], observation + assert not [key_ for key_, value in observation.items() if value is None], observation + assert observation["spend_logs_metadata"] == {"ticket": "LIT-8283"}, observation + assert not [name for name in attrs if name.startswith("langfuse.trace.metadata.")], sorted(attrs) + + +def test_langfuse_otel_metadata_reaches_langfuse_for_openai_sdk_stream(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-sdkstream-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + with OpenAI( + base_url=str(rig.candidate.client.base_url).rstrip("/") + "/v1", + api_key=key, + http_client=httpx.Client(trust_env=False, timeout=30), + ) as client: + chunks: Final = tuple( + client.chat.completions.create( + model=identity["model"], + messages=[{"role": "user", "content": marker}], + stream=True, + user=identity["end_user"], + extra_body={"cache": {"no-cache": True}}, + ) + ) + assert chunks, "stream produced no chunks" + response_id: Final = chunks[-1].id + attrs: Final = _observation_span(rig.collector, response_id) + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +async def test_langfuse_otel_metadata_reaches_langfuse_for_openai_sdk_async(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-sdkasync-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + async with AsyncOpenAI( + base_url=str(rig.candidate.client.base_url).rstrip("/") + "/v1", + api_key=key, + http_client=httpx.AsyncClient(trust_env=False, timeout=30), + ) as client: + response: Final = await client.chat.completions.create( + model=identity["model"], + messages=[{"role": "user", "content": marker}], + user=identity["end_user"], + extra_body={"cache": {"no-cache": True}}, + ) + attrs: Final = _observation_span(rig.collector, response.id) + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +def test_langfuse_otel_metadata_reaches_langfuse_for_anthropic_sdk_messages(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-anthropic-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + model: Final = scenario.model(model="anthropic/claude-sonnet-4-6", api_base=rig.provider.url) + with Anthropic( + base_url=str(rig.candidate.client.base_url).rstrip("/"), + api_key=key, + http_client=httpx.Client(trust_env=False, timeout=30), + ) as client: + message: Final = client.messages.create( + model=model, + max_tokens=64, + messages=[{"role": "user", "content": marker}], + extra_body={"litellm_metadata": {"tags": ["lf8283"]}, "metadata": {"user_id": identity["end_user"]}}, + ) + attrs: Final = _observation_span(rig.collector, message.id) + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +async def test_langfuse_otel_metadata_reaches_langfuse_for_anthropic_sdk_stream( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "lf-anthstream-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + model: Final = scenario.model(model="anthropic/claude-sonnet-4-6", api_base=rig.provider.url) + async with ( + AsyncAnthropic( + base_url=str(rig.candidate.client.base_url).rstrip("/"), + api_key=key, + http_client=httpx.AsyncClient(trust_env=False, timeout=30), + ) as client, + client.messages.stream( + model=model, + max_tokens=64, + messages=[{"role": "user", "content": marker}], + extra_body={"metadata": {"user_id": identity["end_user"]}}, + ) as stream, + ): + final: Final = await stream.get_final_message() + attrs: Final = _observation_span(rig.collector, final.id) + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +def test_langfuse_otel_metadata_reaches_langfuse_for_responses_api(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-responses-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + client: Final = OpenAI( + base_url=str(rig.candidate.client.base_url).rstrip("/") + "/v1", + api_key=key, + http_client=httpx.Client(trust_env=False, timeout=30), + ) + response: Final = client.responses.create( + model=identity["model"], + input=marker, + user=identity["end_user"], + extra_body={"cache": {"no-cache": True}}, + ) + attrs: Final = _observation_span(rig.collector, response.id) + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +def test_langfuse_otel_metadata_reaches_langfuse_for_responses_sse_stream(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-respstream-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + with rig.candidate.client.stream( + "POST", + "/v1/responses", + json={ + "model": identity["model"], + "input": marker, + "user": identity["end_user"], + "stream": True, + "cache": {"no-cache": True}, + }, + headers={"Authorization": f"Bearer {key}"}, + ) as stream: + events: Final = tuple( + json.loads(line.removeprefix("data: ")) + for line in stream.iter_lines() + if line.startswith("data: ") and line != "data: [DONE]" + ) + response_id: Final = next( + event["response_id"] for event in events if event.get("type") == "response.output_text.delta" + ) + attrs: Final = _observation_span(rig.collector, response_id) + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +def test_langfuse_otel_caller_trace_metadata_survives_flattened_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-deploy-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "user": identity["end_user"], + "metadata": {"trace_metadata": {"deploy": "blue"}, "tags": ["lf8283"]}, + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + attrs: Final = _observation_span(rig.collector, response.json()["id"]) + assert json.loads(str(attrs["langfuse.trace.metadata"])) == {"deploy": "blue"}, attrs + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +def test_langfuse_otel_failure_span_carries_request_metadata(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-fail-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": f"{marker} scripted -fail- upstream"}], + "user": identity["end_user"], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code >= 500, response.text + attrs: Final = _marked_span(rig.collector, marker, "langfuse.observation.") + _assert_identity(attrs, _chat_identity_fields(identity), spend_logs_metadata={"ticket": "LIT-8283"}) + + +def test_langfuse_otel_teamless_key_emits_no_team_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-noteam-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider, team=False) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + attrs: Final = _observation_span(rig.collector, response.json()["id"]) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + for field in ("user_api_key_team_id", "user_api_key_team_alias", "user_api_key_end_user_id"): + assert field not in observation, observation + assert not [key_ for key_, value in observation.items() if value is None], observation + assert { + attribute: attrs.get(attribute) + for attribute in (f"langfuse.trace.metadata.{field}" for field in _IDENTITIES) + if attribute in attrs + } == {"langfuse.trace.metadata.user_api_key_alias": identity["key_alias"]}, attrs + + +def test_langfuse_otel_teamless_key_keeps_caller_trace_metadata(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-noteam-meta-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider, team=False) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "metadata": {"trace_metadata": {"deploy": "blue"}}, + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + attrs: Final = _observation_span(rig.collector, response.json()["id"]) + assert json.loads(str(attrs["langfuse.trace.metadata"])) == {"deploy": "blue"}, attrs + assert "langfuse.trace.metadata.user_api_key_team_id" not in attrs, attrs + assert "langfuse.trace.metadata.user_api_key_team_alias" not in attrs, attrs + assert attrs.get("langfuse.trace.metadata.user_api_key_alias") == identity["key_alias"], attrs + + +def test_langfuse_otel_spend_logs_metadata_roundtrips_oversized_and_nonstring_values( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "lf-bigmeta-" + uuid.uuid4().hex + spend_logs_metadata: Final = { + "blob": "x" * 5000, + "count": 7, + "steps": ["a", "b"], + "nested": {"inner": {"leaf": True}}, + } + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider, spend_logs_metadata=spend_logs_metadata) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "user": identity["end_user"], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + attrs: Final = _observation_span(rig.collector, response.json()["id"]) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + assert observation["spend_logs_metadata"] == dict(spend_logs_metadata), observation + _assert_identity(attrs, _chat_identity_fields(identity)) + + +def test_langfuse_otel_unauthenticated_request_leaks_no_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-unauth-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + denied: Final = rig.candidate.client.post( + "/v1/chat/completions", + json={ + "model": identity["model"], + "messages": [{"role": "user", "content": f"{marker}-denied"}], + }, + headers={}, + ) + assert denied.status_code == 401, denied.text + watch: Final = _watched_spans(rig.collector) + first: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert first.status_code == 200, first.text + first_id: Final = first.json()["id"] + eventually(watch, lambda spans: sum(_span_response_id(s) == first_id for s in spans) == 1, seconds=25) + second: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker + "-settle"}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert second.status_code == 200, second.text + second_id: Final = second.json()["id"] + settled: Final = eventually( + watch, lambda spans: sum(_span_response_id(s) == second_id for s in spans) == 1, seconds=25 + ) + denied_spans: Final = tuple(span for span in settled if f"{marker}-denied" in json.dumps(span, default=str)) + assert len(denied_spans) == 1, [s.get("llm.response.id") for s in settled] + denied_span: Final = denied_spans[0] + assert not [name for name in denied_span if name.startswith("langfuse.trace.metadata.user_api_key")], ( + denied_span + ) + leaked: Final = json.loads(str(denied_span.get("langfuse.observation.metadata", "{}"))) + for field in ( + "user_api_key_alias", + "user_api_key_team_id", + "user_api_key_team_alias", + "user_api_key_end_user_id", + ): + assert field not in leaked, leaked + + +def test_langfuse_otel_repeated_identical_requests_emit_one_span_each(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-repeat-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + watch: Final = _watched_spans(rig.collector) + responses: Final = tuple( + rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + for _ in range(3) + ) + assert all(response.status_code == 200 for response in responses), [r.text for r in responses] + response_ids: Final = tuple(response.json()["id"] for response in responses) + assert len(set(response_ids)) == 3, response_ids + settled: Final = eventually( + watch, + lambda spans: all(sum(_span_response_id(s) == rid for s in spans) == 1 for rid in response_ids), + seconds=40, + ) + for rid in response_ids: + assert sum(_span_response_id(s) == rid for s in settled) == 1 + + +def test_langfuse_otel_sink_outage_mid_burst_lands_every_response_id_once(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-chaos-outage-" + uuid.uuid4().hex + outage: Final = threading.Event() + rejected: Final = [] # mutable-ok: the wire server thread records each 503 it serves + + def sink_reply(request: Request) -> Reply: + if outage.is_set(): + rejected.append(1) + return Reply(status=503, body=b'{"error": "sink outage"}') + return Reply() + + with ( + _langfuse_rig(gateway, tmp_path, {}, marker, sink_reply=sink_reply) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + anthropic_model: Final = scenario.model(model="anthropic/claude-sonnet-4-6", api_base=rig.provider.url) + watch: Final = _watched_spans(rig.collector) + + def call(index: int) -> str: + path: Final = ("/v1/chat/completions", "/v1/messages", "/v1/responses")[index % 3] + streaming: Final = index % 2 == 0 + body: Final = { + "/v1/chat/completions": { + "model": identity["model"], + "messages": [{"role": "user", "content": f"{marker}-{index}"}], + "stream": streaming, + "cache": {"no-cache": True}, + }, + "/v1/messages": { + "model": anthropic_model, + "messages": [{"role": "user", "content": f"{marker}-{index}"}], + "max_tokens": 64, + "stream": streaming, + }, + "/v1/responses": {"model": identity["model"], "input": f"{marker}-{index}", "stream": streaming}, + }[path] + if streaming: + with rig.candidate.client.stream( + "POST", path, json=body, headers={"Authorization": f"Bearer {key}"} + ) as stream: + assert stream.status_code == 200, f"burst {index}: {stream.status_code}" + events: Final = tuple( + line.removeprefix("data: ") + for line in stream.iter_lines() + if line.startswith("data: ") and not line.endswith("[DONE]") + ) + payloads: Final = tuple(json.loads(line) for line in events) + found: Final = next( + payload.get("response_id") or (payload.get("message") or {}).get("id") or payload.get("id") + for payload in payloads + if payload.get("response_id") or payload.get("message") or payload.get("id") + ) + return str(found) + response: Final = rig.candidate.client.post(path, json=body, headers={"Authorization": f"Bearer {key}"}) + assert response.status_code == 200, f"burst {index}: {response.text}" + return response.json()["id"] + + with ThreadPoolExecutor(max_workers=30) as pool: + futures: Final = tuple(pool.submit(call, index) for index in range(30)) + while sum(f.done() for f in futures) < 10: + wait(futures, timeout=0.05) + outage.set() + eventually( + lambda: (sum(f.done() for f in futures), len(rejected)), + lambda progress: progress[0] >= 25 and progress[1] >= 1, + seconds=20, + ) + outage.clear() + response_ids: Final = tuple(f.result() for f in futures) + assert len(set(response_ids)) == 30, response_ids + settled: Final = eventually( + watch, + lambda spans: all(sum(_span_response_id(s) == rid for s in spans) == 1 for rid in response_ids), + seconds=80, + ) + assert sorted(_span_response_id(span) for span in settled if _span_response_id(span) in response_ids) == sorted( + response_ids + ) + assert len(rejected) >= 1, "sink outage never served a 503" + + +def test_langfuse_otel_worker_kill_mid_burst_keeps_serving_and_exports(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-chaos-kill-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + watch: Final = _watched_spans(rig.collector) + + def call(index: int) -> tuple[int, str]: + try: + response: Final = rig.candidate.client.post( + "/v1/chat/completions", + json={ + "model": identity["model"], + "messages": [{"role": "user", "content": f"{marker}-{index}"}], + "cache": {"no-cache": True}, + }, + headers={"Authorization": f"Bearer {key}"}, + timeout=30, + ) + except httpx.HTTPError as error: + return 0, repr(error) + return response.status_code, response.json()["id"] if response.status_code == 200 else response.text + + workers: Final = tuple( + child + for child in psutil.Process(rig.proxy.process.pid).children(recursive=True) + if "spawn_main" in " ".join(child.cmdline()) + ) + assert len(workers) >= 2, f"expected multiple uvicorn workers, got {workers!r}" + victim: Final = workers[0] + with ThreadPoolExecutor(max_workers=20) as pool: + futures: Final = tuple(pool.submit(call, index) for index in range(20)) + while sum(f.done() for f in futures) < 5: + wait(futures, timeout=0.05) + victim.kill() + outcomes: Final = tuple(f.result() for f in futures) + assert not victim.is_running() or victim.status() == psutil.STATUS_ZOMBIE + succeeded: Final = tuple(identifier for status, identifier in outcomes if status == 200) + assert succeeded, outcomes + post_kill: Final = tuple(call(20 + index) for index in range(5)) + post_kill_ids: Final = tuple(identifier for _, identifier in post_kill) + assert all(status == 200 for status, _ in post_kill), post_kill + settled: Final = eventually( + watch, + lambda spans: all(sum(_span_response_id(s) == rid for s in spans) == 1 for rid in post_kill_ids), + seconds=60, + ) + for rid in succeeded: + assert sum(_span_response_id(s) == rid for s in settled) <= 1, rid + + +def test_langfuse_otel_slow_sink_does_not_deadlock_exports(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-chaos-slow-" + uuid.uuid4().hex + + def sink_reply(request: Request) -> Reply: + gate: Final = threading.Event() + threading.Timer(2, gate.set).start() + return Reply(chunks=(b"{}",), gate_after_first=gate) + + with ( + _langfuse_rig(gateway, tmp_path, {}, marker, sink_reply=sink_reply) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + watch: Final = _watched_spans(rig.collector) + responses: Final = tuple( + rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": f"{marker}-{index}"}], + "cache": {"no-cache": True}, + }, + key=key, + ) + for index in range(10) + ) + assert all(response.status_code == 200 for response in responses), [r.text for r in responses] + response_ids: Final = tuple(response.json()["id"] for response in responses) + settled: Final = eventually( + watch, + lambda spans: all(sum(_span_response_id(s) == rid for s in spans) == 1 for rid in response_ids), + seconds=60, + ) + assert sorted(_span_response_id(span) for span in settled if _span_response_id(span) in response_ids) == sorted( + response_ids + ) + + +def test_langfuse_otel_v2_mapper_emits_request_metadata_and_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-v2-" + uuid.uuid4().hex + with ( + _langfuse_rig( + gateway, + tmp_path, + {}, + marker, + overrides={"LITELLM_OTEL_V2": "1"}, + ) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "user": identity["end_user"], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + watch: Final = _watched_spans(rig.collector) + observed: Final = eventually( + watch, + lambda spans: sum(s.get("gen_ai.response.id") == response.json()["id"] for s in spans) == 1, + seconds=25, + ) + attrs: Final = next(s for s in observed if s.get("gen_ai.response.id") == response.json()["id"]) + assert attrs.get("langfuse.trace.metadata.team_id") == identity["team_id"], attrs + assert attrs.get("langfuse.trace.metadata.team_alias") == identity["team_alias"], attrs + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + expected: Final = _chat_identity_fields(identity) + assert {key: observation.get(key) for key in expected} == dict(expected), observation + assert observation["spend_logs_metadata"] == {"ticket": "LIT-8283"}, observation + assert not [key_ for key_, value in observation.items() if value is None], observation + flattened: Final = { + attribute: attrs.get(attribute) + for attribute in (f"langfuse.trace.metadata.{field}" for field in _IDENTITIES) + if attribute in attrs + } + assert flattened == {f"langfuse.trace.metadata.{field}": value for field, value in expected.items()}, attrs + + +def test_langfuse_otel_v2_mapper_redacts_user_api_key_fields(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-v2-redact-" + uuid.uuid4().hex + with ( + _langfuse_rig( + gateway, + tmp_path, + {"redact_user_api_key_info": True}, + marker, + overrides={"LITELLM_OTEL_V2": "1"}, + ) as rig, + rig.candidate.scenario() as scenario, + ): + key: Final = scenario.key( + key_alias="lf-v2-redact-" + uuid.uuid4().hex, + metadata={"spend_logs_metadata": {"ticket": "LIT-8283"}}, + ) + model: Final = scenario.model(api_base=rig.provider.url + "/v1") + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + watch: Final = _watched_spans(rig.collector) + observed: Final = eventually( + watch, + lambda spans: sum(s.get("gen_ai.response.id") == response.json()["id"] for s in spans) == 1, + seconds=25, + ) + attrs: Final = next(s for s in observed if s.get("gen_ai.response.id") == response.json()["id"]) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + assert not [name for name in observation if name.startswith("user_api_key")], observation + assert not [key_ for key_, value in observation.items() if value is None], observation + assert observation["spend_logs_metadata"] == {"ticket": "LIT-8283"}, observation + assert not [name for name in attrs if name.startswith("langfuse.trace.metadata.user_api_key_")], sorted(attrs) + + +def test_langfuse_otel_v2_mapper_teamless_key_emits_no_team_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-v2-noteam-" + uuid.uuid4().hex + with ( + _langfuse_rig( + gateway, + tmp_path, + {}, + marker, + overrides={"LITELLM_OTEL_V2": "1"}, + ) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider, team=False) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + watch: Final = _watched_spans(rig.collector) + observed: Final = eventually( + watch, + lambda spans: sum(s.get("gen_ai.response.id") == response.json()["id"] for s in spans) == 1, + seconds=25, + ) + attrs: Final = next(s for s in observed if s.get("gen_ai.response.id") == response.json()["id"]) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + for field in ("user_api_key_team_id", "user_api_key_team_alias", "user_api_key_end_user_id"): + assert field not in observation, observation + assert not [key_ for key_, value in observation.items() if value is None], observation + assert { + attribute: attrs.get(attribute) + for attribute in (f"langfuse.trace.metadata.{field}" for field in _IDENTITIES) + if attribute in attrs + } == {"langfuse.trace.metadata.user_api_key_alias": identity["key_alias"]}, attrs 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..cd54465fc66 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,101 @@ 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_emits_request_metadata_and_identity(): + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": { + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-1", + "user_api_key_team_id": "team-9", + "user_api_key_team_alias": "team nine", + "user_api_key_max_budget": None, + "spend_logs_metadata": {"ticket": "LIT-8283"}, + "requester_metadata": {"headers": {"x-tenant": "acme"}}, + }, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + + assert json.loads(attrs["langfuse.observation.metadata"]) == { + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-1", + "user_api_key_team_id": "team-9", + "user_api_key_team_alias": "team nine", + "spend_logs_metadata": {"ticket": "LIT-8283"}, + "requester_metadata": {"headers": {"x-tenant": "acme"}}, + } + assert attrs["langfuse.trace.metadata.user_api_key_alias"] == "prod-key" + assert attrs["langfuse.trace.metadata.user_api_key_user_id"] == "user-1" + assert attrs["langfuse.trace.metadata.user_api_key_end_user_id"] == "end-1" + assert attrs["langfuse.trace.metadata.user_api_key_team_id"] == "team-9" + assert attrs["langfuse.trace.metadata.user_api_key_team_alias"] == "team nine" + assert attrs["langfuse.trace.metadata.team_id"] == "team-9" + assert attrs["langfuse.trace.metadata.team_alias"] == "team nine" + assert attrs["langfuse.observation.metadata.provider"] == "openai" + + +def test_langfuse_mapper_skips_empty_identity_values(): + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": { + "user_api_key_alias": None, + "user_api_key_user_id": "", + "user_api_key_end_user_id": "end-1", + }, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + + assert "langfuse.trace.metadata.user_api_key_alias" not in attrs + assert "langfuse.trace.metadata.user_api_key_user_id" not in attrs + assert attrs["langfuse.trace.metadata.user_api_key_end_user_id"] == "end-1" + + +def test_langfuse_mapper_redacts_user_api_key_fields(): + import litellm + + saved: Final = litellm.redact_user_api_key_info + try: + litellm.redact_user_api_key_info = True + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": { + "user_api_key_alias": "prod-key", + "user_api_key_team_id": "team-9", + "spend_logs_metadata": {"ticket": "LIT-8283"}, + }, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + finally: + litellm.redact_user_api_key_info = saved + + exported: Final = json.loads(attrs["langfuse.observation.metadata"]) + assert not [key for key in exported if key.startswith("user_api_key")] + assert exported["spend_logs_metadata"] == {"ticket": "LIT-8283"} + assert not [key for key in attrs if key.startswith("langfuse.trace.metadata.user_api_key_")] + assert attrs["langfuse.trace.metadata.team_id"] == "team-9" + + +def test_langfuse_mapper_omits_observation_metadata_without_payload_metadata(): + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": {}, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + + assert json.loads(attrs["langfuse.observation.metadata"]) == {"requester_metadata": {}} diff --git a/tests/unit/integrations/test_langfuse_otel.py b/tests/unit/integrations/test_langfuse_otel.py index 0a9ce55fe16..98fd7ebf41d 100644 --- a/tests/unit/integrations/test_langfuse_otel.py +++ b/tests/unit/integrations/test_langfuse_otel.py @@ -1,9 +1,11 @@ import json import os +from typing import Final from unittest.mock import MagicMock, patch import pytest +import litellm from litellm.integrations.langfuse.langfuse_otel import LangfuseOtelLogger from litellm.integrations.opentelemetry import OpenTelemetryConfig from litellm.types.llms.openai import ResponsesAPIResponse @@ -482,6 +484,115 @@ class TestLangfuseOtelIntegration: assert isinstance(config, OpenTelemetryConfig) # Endpoint assertion removed as side effect is gone + def test_request_metadata_is_emitted_under_langfuse_observation_and_trace_metadata_keys(self) -> None: + request_metadata: Final = { + "user_api_key_hash": "hash123", + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-user-1", + "user_api_key_team_id": "team-1", + "user_api_key_team_alias": "team-a", + "spend_logs_metadata": {"env": "prod"}, + "requester_ip_address": "10.0.0.1", + "requester_metadata": {"ticket": "LIT-8283"}, + } + kwargs: Final = { + "litellm_params": {"metadata": {}}, + "standard_logging_object": {"metadata": request_metadata}, + } + + with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute: + LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None) + + actual: Final = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list} + + assert json.loads(actual["langfuse.observation.metadata"]) == {**request_metadata} + assert {key: value for key, value in actual.items() if key.startswith("langfuse.trace.metadata.")} == { + "langfuse.trace.metadata.user_api_key_alias": "prod-key", + "langfuse.trace.metadata.user_api_key_user_id": "user-1", + "langfuse.trace.metadata.user_api_key_end_user_id": "end-user-1", + "langfuse.trace.metadata.user_api_key_team_id": "team-1", + "langfuse.trace.metadata.user_api_key_team_alias": "team-a", + } + + def test_request_metadata_redaction_matches_vanilla_langfuse(self) -> None: + previous_flag: Final = litellm.redact_user_api_key_info + litellm.redact_user_api_key_info = True + request_metadata: Final = { + "user_api_key_hash": "hash123", + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-user-1", + "user_api_key_team_id": "team-1", + "user_api_key_team_alias": "team-a", + "spend_logs_metadata": {"env": "prod"}, + "requester_ip_address": "10.0.0.1", + "requester_metadata": {"ticket": "LIT-8283"}, + } + kwargs: Final = { + "litellm_params": {"metadata": {}}, + "standard_logging_object": {"metadata": request_metadata}, + } + + try: + with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute: + LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None) + + actual: Final = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list} + + assert json.loads(actual["langfuse.observation.metadata"]) == { + "spend_logs_metadata": {"env": "prod"}, + "requester_ip_address": "10.0.0.1", + "requester_metadata": {"ticket": "LIT-8283"}, + } + assert not [key for key in actual if key.startswith("langfuse.trace.metadata.")] + finally: + litellm.redact_user_api_key_info = previous_flag + + def test_request_metadata_drops_null_fields_and_skips_empty_trace_identities(self) -> None: + request_metadata: Final = { + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": None, + "user_api_key_team_id": "", + "user_api_key_team_alias": None, + "team_id": None, + "team_alias": None, + "spend_logs_metadata": {"env": "prod"}, + } + kwargs: Final = { + "litellm_params": {"metadata": {}}, + "standard_logging_object": {"metadata": request_metadata}, + } + + with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute: + LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None) + + actual: Final = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list} + + assert json.loads(actual["langfuse.observation.metadata"]) == { + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_team_id": "", + "spend_logs_metadata": {"env": "prod"}, + "requester_metadata": {}, + } + assert {key: value for key, value in actual.items() if key.startswith("langfuse.trace.metadata.")} == { + "langfuse.trace.metadata.user_api_key_alias": "prod-key", + "langfuse.trace.metadata.user_api_key_user_id": "user-1", + } + + def test_request_metadata_keys_are_absent_without_standard_logging_object(self) -> None: + kwargs: Final = {"litellm_params": {"metadata": {"trace_metadata": {"k": "v"}}}} + + with patch("litellm.integrations.arize._utils.safe_set_attribute") as mock_safe_set_attribute: + LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None) + + actual: Final = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list} + + assert [key for key in actual if key.startswith("langfuse.trace.metadata")] == ["langfuse.trace.metadata"] + assert "langfuse.observation.metadata" not in actual + class TestLangfuseOtelKeyDynamicConfig: """Key/team-scoped Langfuse credentials must define the full export target