diff --git a/litellm/integrations/arize/_utils.py b/litellm/integrations/arize/_utils.py index 0271cf1e03c..0f5addeb39d 100644 --- a/litellm/integrations/arize/_utils.py +++ b/litellm/integrations/arize/_utils.py @@ -9,6 +9,7 @@ from litellm.integrations.opentelemetry_utils.base_otel_llm_obs_attributes impor BaseLLMObsOTELAttributes, safe_set_attribute, ) +from litellm.integrations.otel.model.utils import as_str_mapping from litellm.litellm_core_utils.redact_messages import ( should_redact_message_logging, ) @@ -410,7 +411,14 @@ def _set_tool_attributes(span: "Span", optional_tools: list | None, metadata_too ) -def set_attributes(span: "Span", kwargs, response_obj, attributes: type[BaseLLMObsOTELAttributes]): +def set_attributes( + span: "Span", + kwargs, + response_obj, + attributes: type[BaseLLMObsOTELAttributes], + *, + emit_session_and_user: bool = True, +): """ Populates span with OpenInference-compliant LLM attributes for Arize and Phoenix tracing. """ @@ -471,7 +479,9 @@ def set_attributes(span: "Span", kwargs, response_obj, attributes: type[BaseLLMO # blank the attributes set by the main try-block above. New attributes are # written under new keys; existing attributes are not overwritten. slp: Final = kwargs.get("standard_logging_object") - _safe_emit("session/user attrs", _set_session_and_user_attrs, span, kwargs, slp) + if emit_session_and_user: + _safe_emit("session/user attrs", _set_session_and_user_attrs, span, kwargs, slp) + _safe_emit("request context attrs", _set_request_context_attrs, span, slp) _safe_emit("response cost", _set_response_cost_attr, span, slp) _safe_emit( "passthrough normalization", @@ -834,14 +844,12 @@ def _emit_input_message_extras(span: "Span", prefix: str, message: dict) -> None def _set_session_and_user_attrs(span: "Span", kwargs: dict, standard_logging_payload) -> None: - """Emit `SESSION_ID` / `USER_ID` / team metadata when source data exists. + """Emit `SESSION_ID` / `USER_ID` when source data exists. `SESSION_ID` is emitted only when an explicit end-user identifier exists (`metadata.user_api_key_end_user_id`). We deliberately do NOT fall back to `trace_id`, because that would create a distinct "session" for every - single request and distort Arize's Session-grouping analytics. The - `trace_id` is still emitted under its own `litellm.trace_id` key so - spans remain filterable by trace. + single request and distort Arize's Session-grouping analytics. USER_ID is *only* emitted when no upstream path (model_params.user or optional_params.user) has already set it, to avoid overwriting an @@ -857,10 +865,6 @@ def _set_session_and_user_attrs(span: "Span", kwargs: dict, standard_logging_pay if session_id: safe_set_attribute(span, SpanAttributes.SESSION_ID, str(session_id)) - trace_id: Final = standard_logging_payload.get("trace_id") - if trace_id: - safe_set_attribute(span, "litellm.trace_id", str(trace_id)) - optional_params: Final = kwargs.get("optional_params") or {} model_params: Final = standard_logging_payload.get("model_parameters") or {} has_user_already: Final = bool( @@ -872,6 +876,19 @@ def _set_session_and_user_attrs(span: "Span", kwargs: dict, standard_logging_pay if user_id: safe_set_attribute(span, SpanAttributes.USER_ID, str(user_id)) + +def _set_request_context_attrs(span: "Span", standard_logging_payload: object) -> None: + payload: Final = as_str_mapping(standard_logging_payload) + if payload is None: + return + + trace_id: Final = payload.get("trace_id") + if trace_id: + safe_set_attribute(span, "litellm.trace_id", str(trace_id)) + + metadata: Final = as_str_mapping(payload.get("metadata")) + if metadata is None: + return team_id: Final = metadata.get("user_api_key_team_id") if team_id: safe_set_attribute(span, "litellm.team_id", str(team_id)) diff --git a/litellm/integrations/langfuse/langfuse_otel.py b/litellm/integrations/langfuse/langfuse_otel.py index a96fac32c2a..342b8ef5148 100644 --- a/litellm/integrations/langfuse/langfuse_otel.py +++ b/litellm/integrations/langfuse/langfuse_otel.py @@ -6,10 +6,13 @@ from typing import TYPE_CHECKING, Any, Final, Optional from litellm._logging import verbose_logger from litellm.integrations.arize import _utils +from litellm.integrations.arize._utils import safe_set_attribute from litellm.integrations.langfuse.langfuse_otel_attributes import ( LangfuseLLMObsOTELAttributes, ) from litellm.integrations.opentelemetry import OpenTelemetry, OpenTelemetryConfig +from litellm.integrations.otel.model.trace_controls import metadata_bodies +from litellm.integrations.otel.model.utils import as_str, as_str_mapping from litellm.litellm_core_utils.safe_json_loads import safe_json_loads from litellm.types.integrations.langfuse_otel import ( LangfuseSpanAttributes, @@ -39,19 +42,20 @@ class LangfuseOtelLogger(OpenTelemetry): super().__init__(config=config, *args, **kwargs) @staticmethod - def set_langfuse_otel_attributes(span: Span, kwargs, response_obj): + def set_langfuse_otel_attributes(span: Span, kwargs: dict[str, object], response_obj) -> None: """ Sets OpenTelemetry span attributes for Langfuse observability. Uses the same attribute setting logic as Arize Phoenix for consistency. """ - _utils.set_attributes(span, kwargs, response_obj, LangfuseLLMObsOTELAttributes) + _utils.set_attributes(span, kwargs, response_obj, LangfuseLLMObsOTELAttributes, emit_session_and_user=False) span.set_attribute("langfuse.observation.type", "generation") ######################################################### # Set Langfuse specific attributes ######################################################### LangfuseOtelLogger._set_langfuse_specific_attributes(span=span, kwargs=kwargs, response_obj=response_obj) + LangfuseOtelLogger._set_trace_user_attribute(span=span, kwargs=kwargs) @staticmethod def _extract_langfuse_metadata(kwargs: dict) -> dict: @@ -255,6 +259,31 @@ class LangfuseOtelLogger(OpenTelemetry): LangfuseOtelLogger._set_observation_output(span=span, response_obj=response_obj) + @staticmethod + def _set_trace_user_attribute(span: Span, kwargs: dict[str, object]) -> None: + slp: Final = as_str_mapping(kwargs.get("standard_logging_object")) + slp_metadata: Final = as_str_mapping(slp.get("metadata")) if slp is not None else None + if slp is None or slp_metadata is None: + return + metadata: Final = as_str_mapping( + LangfuseOtelLogger._extract_langfuse_metadata(kwargs) # pyright: ignore[reportUnknownArgumentType,reportUnknownMemberType] # helper returns a loosely typed dict + ) + litellm_params: Final = as_str_mapping(kwargs.get("litellm_params")) or {} + bodies: Final = metadata_bodies(litellm_params) + caller: Final = (as_str(metadata.get("trace_user_id")) if metadata is not None else None) or next( + (value for body in bodies if (value := as_str(body.get("trace_user_id")))), None + ) + if caller is not None: + safe_set_attribute(span, LangfuseSpanAttributes.TRACE_USER_ID.value, caller) + return + end_user: Final = ( + slp_metadata.get("user_api_key_end_user_id") + or slp.get("end_user") + or next((value for body in bodies if (value := as_str(body.get("user_api_key_end_user_id")))), None) + ) + if end_user: + safe_set_attribute(span, LangfuseSpanAttributes.TRACE_USER_ID.value, str(end_user)) + @staticmethod def _get_langfuse_otel_host() -> str | None: """ diff --git a/litellm/integrations/otel/langfuse_logger.py b/litellm/integrations/otel/langfuse_logger.py index d029b153c52..2578777f303 100644 --- a/litellm/integrations/otel/langfuse_logger.py +++ b/litellm/integrations/otel/langfuse_logger.py @@ -9,7 +9,7 @@ from litellm.integrations.otel.mappers.langfuse import ( LangfuseMapper, ) 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.trace_controls import langfuse_trace_controls from litellm.integrations.otel.plumbing.context import request_root_span if TYPE_CHECKING: @@ -24,7 +24,7 @@ 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(langfuse_trace_controls(kwargs))) super().log_pre_api_call(model, messages, kwargs) diff --git a/litellm/integrations/otel/model/trace_controls.py b/litellm/integrations/otel/model/trace_controls.py index eac7b5c897b..dc64d5d5743 100644 --- a/litellm/integrations/otel/model/trace_controls.py +++ b/litellm/integrations/otel/model/trace_controls.py @@ -3,7 +3,7 @@ from __future__ import annotations from collections.abc import Mapping -from dataclasses import dataclass +from dataclasses import dataclass, replace from typing import Final from pydantic import TypeAdapter, ValidationError @@ -33,11 +33,7 @@ def caller_trace_controls(kwargs: Mapping[str, object]) -> TraceControls: return TraceControls() proxy_request: Final = as_str_mapping(request.get("proxy_server_request")) headers: Final = as_str_mapping(proxy_request.get("headers")) if proxy_request is not None else None - bodies: Final = tuple( - metadata - for key in ("metadata", "litellm_metadata") - if (metadata := as_str_mapping(request.get(key))) is not None - ) + bodies: Final = metadata_bodies(request) def scalar(control: str) -> str | None: from_header: Final = as_str(headers.get(f"{LANGFUSE_HEADER_PREFIX}{control}")) if headers is not None else None @@ -53,6 +49,28 @@ def caller_trace_controls(kwargs: Mapping[str, object]) -> TraceControls: ) +def langfuse_trace_controls(kwargs: Mapping[str, object]) -> TraceControls: + controls: Final = caller_trace_controls(kwargs) + if controls.user_id: + return controls + request: Final = as_str_mapping(kwargs.get("litellm_params")) + if request is None: + return controls + end_user: Final = next( + (value for body in metadata_bodies(request) if (value := as_str(body.get("user_api_key_end_user_id")))), + None, + ) + return replace(controls, user_id=end_user) + + +def metadata_bodies(request: Mapping[str, object]) -> tuple[Mapping[str, object], ...]: + return tuple( + metadata + for key in ("metadata", "litellm_metadata") + if (metadata := as_str_mapping(request.get(key))) is not None + ) + + def _str_items(value: object) -> tuple[str, ...]: try: items: Final = _ITEMS.validate_python(value) diff --git a/tests/integration/observability/_langfuse_otel.py b/tests/integration/observability/_langfuse_otel.py new file mode 100644 index 00000000000..6022d43e587 --- /dev/null +++ b/tests/integration/observability/_langfuse_otel.py @@ -0,0 +1,307 @@ +import json +from collections.abc import Iterator, Mapping +from contextlib import contextmanager +from pathlib import Path +from typing import Final + +import yaml +from integration._support.client import Gateway +from integration._support.process import owned_proxy +from integration._support.wire import Reply, Request, Wire, wire_server +from opentelemetry.proto.collector.trace.v1 import trace_service_pb2 +from opentelemetry.proto.common.v1.common_pb2 import AnyValue + +MARKER_JSON_KEYS: Final = ("llm.response.id", "gen_ai.response.id") + + +def _attribute_value(value: AnyValue) -> object: + kind: Final = value.WhichOneof("value") + return getattr(value, kind) if kind is not None else None + + +def _spans(body: bytes) -> tuple[tuple[str, dict[str, object]], ...]: + export: Final = trace_service_pb2.ExportTraceServiceRequest() + export.ParseFromString(body) + return tuple( + ( + span.trace_id.hex(), + {attribute.key: _attribute_value(attribute.value) for attribute in span.attributes}, + ) + for resource in export.resource_spans + for scope in resource.scope_spans + for span in scope.spans + ) + + +def _sse_frame(value: Mapping[str, object]) -> bytes: + return b"data: " + json.dumps(value, ensure_ascii=False).encode() + b"\n\n" + + +def _upstream_reply(marker: str) -> Reply: + return Reply( + body=json.dumps( + { + "id": marker, + "object": "chat.completion", + "created": 1, + "model": "gpt-4o-mini", + "choices": [ + { + "index": 0, + "message": {"role": "assistant", "content": f"reply {marker}"}, + "finish_reason": "stop", + } + ], + "usage": {"prompt_tokens": 3, "completion_tokens": 2, "total_tokens": 5}, + } + ).encode() + ) + + +def _upstream_stream_reply(marker: str) -> Reply: + chunk: Final = { + "id": marker, + "object": "chat.completion.chunk", + "created": 1, + "model": "gpt-4o-mini", + "choices": [{"index": 0, "delta": {"role": "assistant", "content": f"reply {marker}"}, "finish_reason": None}], + "usage": {"prompt_tokens": 3, "completion_tokens": 2, "total_tokens": 5}, + } + final: Final = { + "id": marker, + "object": "chat.completion.chunk", + "created": 1, + "model": "gpt-4o-mini", + "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}], + } + return Reply( + content_type="text/event-stream", + chunks=(_sse_frame(chunk), _sse_frame(final), b"data: [DONE]\n\n"), + ) + + +def _responses_reply(marker: str) -> Reply: + return Reply( + body=json.dumps( + { + "id": marker, + "object": "response", + "created_at": 1, + "model": "gpt-4o-mini", + "status": "completed", + "output": [ + { + "type": "message", + "role": "assistant", + "content": [{"type": "output_text", "text": f"reply {marker}"}], + } + ], + "usage": {"input_tokens": 3, "output_tokens": 2, "total_tokens": 5}, + } + ).encode() + ) + + +def _responses_stream_reply(marker: str) -> Reply: + response: Final = { + "id": marker, + "object": "response", + "created_at": 1, + "model": "gpt-4o-mini", + "status": "in_progress", + "output": [], + "usage": None, + } + completed: Final = { + "id": marker, + "object": "response", + "created_at": 1, + "model": "gpt-4o-mini", + "status": "completed", + "output": [ + { + "type": "message", + "id": "item-1", + "role": "assistant", + "content": [{"type": "output_text", "text": f"reply {marker}"}], + } + ], + "usage": {"input_tokens": 3, "output_tokens": 2, "total_tokens": 5}, + } + events: Final = ( + {"type": "response.created", "response": response}, + {"type": "response.output_item.added", "output_index": 0, "item": completed["output"][0]}, + { + "type": "response.output_text.delta", + "item_id": "item-1", + "output_index": 0, + "content_index": 0, + "delta": f"reply {marker}", + }, + { + "type": "response.output_text.done", + "item_id": "item-1", + "output_index": 0, + "content_index": 0, + "text": f"reply {marker}", + }, + {"type": "response.output_item.done", "output_index": 0, "item": completed["output"][0]}, + {"type": "response.completed", "response": completed}, + ) + return Reply( + content_type="text/event-stream", + chunks=tuple(_sse_frame(event) for event in events) + (b"data: [DONE]\n\n",), + ) + + +def _upstream_reply_for(request: Request, marker: str) -> Reply: + body: Final = json.loads(request.body) if request.body else {} + if request.target.endswith("/responses"): + if body.get("stream") is True: + return _responses_stream_reply(marker) + return _responses_reply(marker) + if body.get("stream") is True: + return _upstream_stream_reply(marker) + return _upstream_reply(marker) + + +def _marker_from_body(request: Request) -> str: + body: Final = json.loads(request.body) if request.body else {} + if body.get("input") is not None: + value: Final = body["input"] + if isinstance(value, str): + return value + return str(value) + messages: Final = body.get("messages") + if not messages: + return "" + return str(messages[0]["content"]) + + +def _sink(_request: Request) -> Reply: + return Reply(body=b"", content_type="application/x-protobuf") + + +def _drained_spans(sink: Wire, batches: list[bytes]) -> tuple[tuple[str, dict[str, object]], ...]: + batches.extend(request.body for request in sink.drain()) + return tuple(span for body in batches for span in _spans(body)) + + +def _generation_span_attributes(sink: Wire, batches: list[bytes], marker: str) -> tuple[dict[str, object], ...]: + return tuple( + attributes + for _trace_id, attributes in _drained_spans(sink, batches) + if attributes.get("llm.response.id") == marker or attributes.get("gen_ai.response.id") == marker + ) + + +def _span_attributes_containing_marker(sink: Wire, batches: list[bytes], marker: str) -> tuple[dict[str, object], ...]: + return tuple( + attributes + for _trace_id, attributes in _drained_spans(sink, batches) + if any(isinstance(value, str) and marker in value for value in attributes.values()) + ) + + +def _generation_marker_span_attributes(sink: Wire, batches: list[bytes], marker: str) -> tuple[dict[str, object], ...]: + return tuple( + attributes + for _trace_id, attributes in _drained_spans(sink, batches) + if attributes.get("langfuse.observation.type") == "generation" + and any(isinstance(value, str) and marker in value for value in attributes.values()) + ) + + +def _dedupe_spans(spans: tuple[tuple[str, dict[str, object]], ...]) -> tuple[tuple[str, dict[str, object]], ...]: + seen: Final = set() + unique: Final = [] + for trace_id, attributes in spans: + fingerprint: Final = (trace_id, frozenset(attributes.items())) + if fingerprint not in seen: + seen.add(fingerprint) + unique.append((trace_id, attributes)) + return tuple(unique) + + +def _span_containing_marker( + spans: tuple[tuple[str, dict[str, object]], ...], marker: str +) -> tuple[tuple[str, dict[str, object]], ...]: + return tuple( + (trace_id, attributes) + for trace_id, attributes in spans + if any(isinstance(value, str) and marker in value for value in attributes.values()) + ) + + +def _arize_generation_span_attributes(sink: Wire, batches: list[bytes], marker: str) -> tuple[dict[str, object], ...]: + return tuple( + attributes + for _trace_id, attributes in _drained_spans(sink, batches) + if attributes.get("llm.response.id") == marker and "session.id" in attributes + ) + + +def _trace_user_span_attributes(sink: Wire, batches: list[bytes], marker: str) -> tuple[dict[str, object], ...]: + spans: Final = _drained_spans(sink, batches) + generation_trace: Final = next( + ( + trace_id + for trace_id, attributes in spans + if attributes.get("llm.response.id") == marker or attributes.get("gen_ai.response.id") == marker + ), + None, + ) + if generation_trace is None: + return () + return tuple( + attributes for trace_id, attributes in spans if trace_id == generation_trace and "user.id" in attributes + ) + + +def _user_id_session_id(attributes: Mapping[str, object]) -> dict[str, object]: + return {key: attributes.get(key) for key in ("user.id", "session.id")} + + +def _proxy_config(directory: Path, name: str, callbacks: tuple[str, ...]) -> Path: + config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text()) + config["litellm_settings"].update({"callbacks": list(callbacks)}) + path: Final = directory / name + path.write_text(yaml.safe_dump(config)) + return path + + +@contextmanager +def _observability_proxy( + gateway: Gateway, + directory: Path, + overrides: Mapping[str, str], + *, + callbacks: tuple[str, ...] = ("langfuse_otel",), + config_name: str = "langfuse_otel.yaml", +) -> Iterator[Gateway]: + path: Final = _proxy_config(directory, config_name, callbacks) + with owned_proxy( + gateway, + directory, + {"OTEL_BSP_SCHEDULE_DELAY": "100", **overrides}, + config=path, + workers=2, + ) as candidate: + yield candidate + + +@contextmanager +def _langfuse_proxy( + gateway: Gateway, directory: Path, collector_url: str, overrides: Mapping[str, str] | None = None +) -> Iterator[Gateway]: + with _observability_proxy( + gateway, + directory, + { + "LANGFUSE_PUBLIC_KEY": "pk-integration", + "LANGFUSE_SECRET_KEY": "sk-integration", + "LANGFUSE_HOST": collector_url, + **(overrides or {}), + }, + ) as candidate: + yield candidate diff --git a/tests/integration/observability/test_langfuse_otel_identity.py b/tests/integration/observability/test_langfuse_otel_identity.py new file mode 100644 index 00000000000..8ad5e5509e6 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_identity.py @@ -0,0 +1,295 @@ +import uuid +from pathlib import Path +from typing import Final + +from _langfuse_otel import ( + _generation_span_attributes, + _langfuse_proxy, + _sink, + _span_attributes_containing_marker, + _trace_user_span_attributes, + _upstream_reply, +) +from integration._support.client import Gateway, eventually +from integration._support.wire import Reply, Request, wire_server + + +def test_langfuse_otel_header_end_user_lands_in_user_id_not_session_id_for_a_key_owned_by_an_internal_user( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = uuid.uuid4().hex + upstream_bodies: Final[list[bytes]] = [] + + def upstream(request: Request) -> Reply: + upstream_bodies.append(request.body) + return _upstream_reply(marker) + + with ( + wire_server(upstream) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + assert {key: attributes.get(key) for key in ("user.id", "session.id")} == { + "user.id": f"end-user-{marker}", + "session.id": None, + }, attributes + + +def test_langfuse_otel_header_end_user_lands_in_user_id_for_a_service_account_key( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = uuid.uuid4().hex + upstream_bodies: Final[list[bytes]] = [] + + def upstream(request: Request) -> Reply: + upstream_bodies.append(request.body) + return _upstream_reply(marker) + + with ( + wire_server(upstream) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + key: Final = scenario.key() + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + assert {key: attributes.get(key) for key in ("user.id", "session.id")} == { + "user.id": f"end-user-{marker}", + "session.id": None, + }, attributes + + +def test_langfuse_otel_body_user_is_never_a_session(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + upstream_bodies: Final[list[bytes]] = [] + + def upstream(request: Request) -> Reply: + upstream_bodies.append(request.body) + return _upstream_reply(marker) + + with ( + wire_server(upstream) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": f"end-user-{marker}", + "cache": {"no-cache": True}, + }, + ) + assert response.status_code == 200, response.text + assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + assert {key: attributes.get(key) for key in ("user.id", "session.id")} == { + "user.id": f"end-user-{marker}", + "session.id": None, + }, attributes + + +def test_langfuse_otel_caller_trace_user_id_wins_over_the_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + upstream_bodies: Final[list[bytes]] = [] + + def upstream(request: Request) -> Reply: + upstream_bodies.append(request.body) + return _upstream_reply(marker) + + with ( + wire_server(upstream) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": f"end-user-{marker}", + "metadata": {"trace_user_id": f"caller-{marker}"}, + "cache": {"no-cache": True}, + }, + ) + assert response.status_code == 200, response.text + assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + assert {key: attributes.get(key) for key in ("user.id", "session.id")} == { + "user.id": f"caller-{marker}", + "session.id": None, + }, attributes + + +def test_langfuse_otel_caller_session_id_stays_the_session_beside_the_end_user( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = uuid.uuid4().hex + upstream_bodies: Final[list[bytes]] = [] + + def upstream(request: Request) -> Reply: + upstream_bodies.append(request.body) + return _upstream_reply(marker) + + with ( + wire_server(upstream) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": f"end-user-{marker}", + "metadata": {"session_id": f"sess-{marker}"}, + "cache": {"no-cache": True}, + }, + ) + assert response.status_code == 200, response.text + assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + assert {key: attributes.get(key) for key in ("user.id", "session.id")} == { + "user.id": f"end-user-{marker}", + "session.id": f"sess-{marker}", + }, attributes + + +def test_langfuse_otel_v2_header_end_user_lands_in_user_id(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + upstream_bodies: Final[list[bytes]] = [] + + def upstream(request: Request) -> Reply: + upstream_bodies.append(request.body) + return _upstream_reply(marker) + + with ( + wire_server(upstream) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url, {"LITELLM_OTEL_V2": "1"}) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies + batches: Final[list[bytes]] = [] + user_spans: Final = eventually( + lambda: _trace_user_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + ) + assert {key: user_spans[0].get(key) for key in ("user.id", "session.id")} == { + "user.id": f"end-user-{marker}", + "session.id": None, + }, user_spans[0] + + +def test_langfuse_otel_messages_caller_trace_user_id_under_litellm_metadata_wins_over_the_end_user( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = uuid.uuid4().hex + upstream_bodies: Final[list[bytes]] = [] + + def upstream(request: Request) -> Reply: + upstream_bodies.append(request.body) + return _upstream_reply(marker) + + with ( + wire_server(upstream) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/messages", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "max_tokens": 5, + "litellm_metadata": {"trace_user_id": f"caller-{marker}"}, + "cache": {"no-cache": True}, + }, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _span_attributes_containing_marker(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + assert {key: attributes.get(key) for key in ("user.id", "session.id")} == { + "user.id": f"caller-{marker}", + "session.id": None, + }, attributes diff --git a/tests/integration/observability/test_langfuse_otel_identity_chaos.py b/tests/integration/observability/test_langfuse_otel_identity_chaos.py new file mode 100644 index 00000000000..4a627ea7e73 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_identity_chaos.py @@ -0,0 +1,300 @@ +import concurrent.futures +import os +import signal +import threading +import time +import uuid +from pathlib import Path +from typing import Final + +import httpx +from _langfuse_otel import ( + _dedupe_spans, + _drained_spans, + _langfuse_proxy, + _marker_from_body, + _proxy_config, + _sink, + _span_containing_marker, + _upstream_reply_for, +) +from integration._support.client import Gateway, eventually +from integration._support.process import group_members, owned_proxy_process +from integration._support.wire import Reply, Request, Wire, wire_server + + +def _endpoints(size: int) -> tuple[str, ...]: + return tuple( + "chat" if index < 10 else "chat_stream" if index < 20 else "messages" if index < 25 else "responses" + for index in range(size) + ) + + +def _burst_body(endpoint: str, model: str, marker: str) -> dict[str, object]: + if endpoint == "responses": + return {"model": model, "input": marker, "cache": {"no-cache": True}} + body: Final = {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}} + if endpoint == "messages": + return {**body, "max_tokens": 5} + if endpoint == "chat_stream": + return {**body, "stream": True} + return body + + +def _burst_path(endpoint: str) -> str: + if endpoint == "responses": + return "/v1/responses" + if endpoint == "messages": + return "/v1/messages" + return "/v1/chat/completions" + + +def _send(candidate: Gateway, model: str, endpoint: str, marker: str) -> dict[str, object]: + headers: Final = { + "Authorization": f"Bearer {candidate.key}", + "x-litellm-end-user-id": f"end-user-{marker}", + } + try: + if endpoint == "chat_stream": + with candidate.client.stream( + "POST", "/v1/chat/completions", json=_burst_body(endpoint, model, marker), headers=headers + ) as response: + return {"marker": marker, "status": response.status_code, "body": response.read().decode()} + response: Final = candidate.client.request( + "POST", _burst_path(endpoint), json=_burst_body(endpoint, model, marker), headers=headers + ) + return {"marker": marker, "status": response.status_code, "body": response.text} + except httpx.HTTPError as error: + return {"marker": marker, "status": 0, "body": repr(error)} + + +def _send_to_live_gateway(holder: dict[str, Gateway], model: str, endpoint: str, marker: str) -> dict[str, object]: + deadline: Final = time.monotonic() + 60 + outcome: Final[dict[str, object]] = {"marker": marker, "status": 0, "body": "no live proxy"} + attempts: Final = {"count": 0} + while time.monotonic() < deadline: + candidate: Final = holder.get("gateway") + if candidate is None: + time.sleep(0.25) + continue + result: Final = _send(candidate, model, endpoint, marker) + attempts["count"] += 1 + if result["status"] != 0: + result["attempts"] = attempts["count"] + return result + outcome.update(result) + time.sleep(0.25) + outcome["attempts"] = attempts["count"] + return outcome + + +def _landed_counts(sink: Wire, batches: list[bytes], markers: tuple[str, ...]) -> dict[str, int]: + spans: Final = _dedupe_spans(_drained_spans(sink, batches)) + generations: Final = tuple(span for span in spans if span[1].get("langfuse.observation.type") == "generation") + return {marker: len(_span_containing_marker(generations, marker)) for marker in markers} + + +def test_langfuse_otel_sink_outage_mid_burst_never_duplicates(gateway: Gateway, tmp_path: Path) -> None: + outage: Final = threading.Event() + + def outage_sink(request: Request) -> Reply: + if outage.is_set(): + return Reply(status=503) + return _sink(request) + + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(outage_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + markers: Final = tuple(uuid.uuid4().hex for _ in range(30)) + endpoints: Final = _endpoints(30) + results: Final[dict[str, dict[str, object]]] = {} + window: Final[set[str]] = set() + with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor: + futures: Final = { + executor.submit(_send, candidate, model, endpoints[index], markers[index]): markers[index] + for index in range(30) + } + for future in concurrent.futures.as_completed(futures): + marker: Final = futures[future] + results[marker] = future.result() + if len(results) == 10: + outage.set() + if len(results) == 20: + outage.clear() + if outage.is_set(): + window.add(marker) + failures: Final = {m: r for m, r in results.items() if r["status"] != 200} + assert not failures, failures + batches: Final[list[bytes]] = [] + counts: Final = eventually( + lambda: _landed_counts(collector, batches, markers), + lambda values: all(value == 1 for value in values.values()), + seconds=60, + return_last_on_timeout=True, + ) + missing: Final = {m for m, c in counts.items() if c == 0} + duplicates: Final = {m: c for m, c in counts.items() if c > 1} + assert not duplicates, duplicates + assert missing <= window, (missing, window, counts) + + +def test_langfuse_otel_slow_sink_never_duplicates_or_loses(gateway: Gateway, tmp_path: Path) -> None: + def slow_sink(request: Request) -> Reply: + time.sleep(1.5) + return _sink(request) + + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(slow_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + markers: Final = tuple(uuid.uuid4().hex for _ in range(30)) + endpoints: Final = _endpoints(30) + results: Final[dict[str, dict[str, object]]] = {} + with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor: + futures: Final = { + executor.submit(_send, candidate, model, endpoints[index], markers[index]): markers[index] + for index in range(30) + } + for future in concurrent.futures.as_completed(futures): + results[futures[future]] = future.result() + failures: Final = {m: r for m, r in results.items() if r["status"] != 200} + assert not failures, failures + batches: Final[list[bytes]] = [] + counts: Final = eventually( + lambda: _landed_counts(collector, batches, markers), + lambda values: all(value == 1 for value in values.values()), + seconds=60, + return_last_on_timeout=True, + ) + assert counts == {marker: 1 for marker in markers}, counts + + +def test_langfuse_otel_worker_kill_mid_burst_keeps_every_marker(gateway: Gateway, tmp_path: Path) -> None: + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + owned_proxy_process( + gateway, + tmp_path, + { + "LANGFUSE_PUBLIC_KEY": "pk-integration", + "LANGFUSE_SECRET_KEY": "sk-integration", + "LANGFUSE_HOST": collector.url, + "OTEL_BSP_SCHEDULE_DELAY": "100", + }, + config=_proxy_config(tmp_path, "langfuse_otel.yaml", ("langfuse_otel",)), + workers=2, + ) as owned, + owned.gateway.scenario() as scenario, + ): + candidate: Final = owned.gateway + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + markers: Final = tuple(uuid.uuid4().hex for _ in range(30)) + endpoints: Final = _endpoints(30) + results: Final[dict[str, dict[str, object]]] = {} + killed: Final = {"done": False} + with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor: + futures: Final = { + executor.submit(_send, candidate, model, endpoints[index], markers[index]): markers[index] + for index in range(30) + } + for future in concurrent.futures.as_completed(futures): + results[futures[future]] = future.result() + if len(results) >= 10 and not killed["done"]: + children: Final = tuple( + process for process in group_members(owned.process.pid) if process.pid != owned.process.pid + ) + assert children, "no uvicorn worker child found" + os.kill(children[0].pid, signal.SIGKILL) + killed["done"] = True + assert killed["done"], "worker kill never fired" + accepted: Final = {m for m, r in results.items() if r["status"] == 200} + failures: Final = {m: r for m, r in results.items() if r["status"] != 200} + survivor: Final = _send(candidate, model, "chat", uuid.uuid4().hex) + assert survivor["status"] == 200, survivor + accepted_markers: Final = tuple(sorted(accepted)) + batches: Final[list[bytes]] = [] + counts: Final = eventually( + lambda: _landed_counts(collector, batches, accepted_markers), + lambda values: all(value == 1 for value in values.values()), + seconds=60, + return_last_on_timeout=True, + ) + assert counts == {marker: 1 for marker in accepted_markers}, (counts, failures) + + +def test_langfuse_otel_proxy_restart_mid_burst_keeps_every_marker(gateway: Gateway, tmp_path: Path) -> None: + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + ): + holder: Final[dict[str, Gateway]] = {} + markers: Final = tuple(uuid.uuid4().hex for _ in range(30)) + endpoints: Final = _endpoints(30) + results: Final[dict[str, dict[str, object]]] = {} + restarted: Final = {"done": False} + one: Final = tmp_path / "one" + two: Final = tmp_path / "two" + one.mkdir() + two.mkdir() + context_one: Final = _langfuse_proxy(gateway, one, collector.url) + context_two: Final = _langfuse_proxy(gateway, two, collector.url) + candidate_one: Final = context_one.__enter__() + candidate_two: Final = {"gateway": None} + model_name: Final = uuid.uuid4().hex + created: Final = candidate_one.post( + "/model/new", + { + "model_name": model_name, + "litellm_params": { + "model": "openai/gpt-4o-mini", + "api_key": "integration-provider-key", + "api_base": provider.url + "/v1", + }, + }, + ) + model_id: Final = str(created["model_info"]["id"]) + holder["gateway"] = candidate_one + try: + with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor: + futures: Final = { + executor.submit( + _send_to_live_gateway, holder, model_name, endpoints[index], markers[index] + ): markers[index] + for index in range(30) + } + for future in concurrent.futures.as_completed(futures): + results[futures[future]] = future.result() + if len(results) < 8 or restarted["done"]: + continue + restarted["done"] = True + holder.pop("gateway") + context_one.__exit__(None, None, None) + candidate_two["gateway"] = context_two.__enter__() + holder["gateway"] = candidate_two["gateway"] + finally: + if candidate_two["gateway"] is not None: + candidate_two["gateway"].post("/model/delete", {"id": model_id}) + context_two.__exit__(None, None, None) + context_one.__exit__(None, None, None) + assert restarted["done"], "restart never fired" + accepted: Final = {m for m, r in results.items() if r["status"] == 200} + failures: Final = {m: r for m, r in results.items() if r["status"] != 200} + accepted_markers: Final = tuple(sorted(accepted)) + batches: Final[list[bytes]] = [] + counts: Final = eventually( + lambda: _landed_counts(collector, batches, accepted_markers), + lambda values: all(value >= 1 for value in values.values()), + seconds=60, + return_last_on_timeout=True, + ) + missing: Final = {m for m, c in counts.items() if c == 0} + over_counted: Final = {m: c for m, c in counts.items() if c > int(results[m].get("attempts", 1))} + assert not missing and not over_counted, (missing, over_counted, failures) diff --git a/tests/integration/observability/test_langfuse_otel_identity_edges.py b/tests/integration/observability/test_langfuse_otel_identity_edges.py new file mode 100644 index 00000000000..793a153dda6 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_identity_edges.py @@ -0,0 +1,425 @@ +import threading +import uuid +from pathlib import Path +from typing import Final + +import httpx +from _langfuse_otel import ( + _drained_spans, + _generation_marker_span_attributes, + _generation_span_attributes, + _langfuse_proxy, + _marker_from_body, + _sink, + _span_containing_marker, + _upstream_reply_for, +) +from integration._support.client import Gateway, eventually +from integration._support.wire import Reply, Request, Wire, wire_server + + +def _identity(attributes: dict[str, object]) -> dict[str, object]: + return {key: attributes.get(key) for key in ("user.id", "session.id")} + + +def _await_generation(collector: Wire, batches: list[bytes], marker: str) -> dict[str, object]: + return eventually( + lambda: _generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + + +def _end_user_request( + candidate: Gateway, model: str, marker: str, key: str | None = None, **extra: object +) -> httpx.Response: + return candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + **extra, + }, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + + +def _subsequent_request_still_lands(candidate: Gateway, collector: Wire, batches: list[bytes], model: str) -> None: + next_marker: Final = uuid.uuid4().hex + response: Final = _end_user_request(candidate, model, next_marker) + assert response.status_code == 200, response.text + assert _await_generation(collector, batches, next_marker) is not None + + +def test_langfuse_otel_five_kb_end_user_header_lands_in_user_id(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + end_user: Final = "e" * 5120 + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + headers={"x-litellm-end-user-id": end_user}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": end_user, "session.id": None}, attributes + + +def test_langfuse_otel_empty_end_user_header_writes_no_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + headers={"x-litellm-end-user-id": ""}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": None, "session.id": None}, attributes + _subsequent_request_still_lands(candidate, collector, batches, model) + + +def test_langfuse_otel_integer_body_user_does_not_crash(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": 123, + "cache": {"no-cache": True}, + }, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert attributes.get("session.id") is None, attributes + _subsequent_request_still_lands(candidate, collector, batches, model) + + +def test_langfuse_otel_list_body_user_does_not_crash(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": ["a", "b"], + "cache": {"no-cache": True}, + }, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert attributes.get("session.id") is None, attributes + _subsequent_request_still_lands(candidate, collector, batches, model) + + +def test_langfuse_otel_integer_trace_user_id_does_not_crash(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = _end_user_request(candidate, model, marker, metadata={"trace_user_id": 123}) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert attributes.get("user.id") in ("123", f"end-user-{marker}"), attributes + assert attributes.get("session.id") is None, attributes + _subsequent_request_still_lands(candidate, collector, batches, model) + + +def test_langfuse_otel_null_metadata_still_maps_the_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = _end_user_request(candidate, model, marker, metadata=None) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_duplicate_end_user_header_uses_the_first(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.client.request( + "POST", + "/v1/chat/completions", + json={ + "model": model, + "messages": [{"role": "user", "content": f"first-call-{marker}"}], + "cache": {"no-cache": True}, + }, + headers={ + "Authorization": f"Bearer {candidate.key}", + "x-litellm-end-user-id": f"first-{marker}", + }, + ) + assert response.status_code == 200, response.text + duplicate: Final = candidate.client.request( + "POST", + "/v1/chat/completions", + json={ + "model": model, + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + headers=[ + ("Authorization", f"Bearer {candidate.key}"), + ("x-litellm-end-user-id", f"first-{marker}"), + ("x-litellm-end-user-id", f"second-{marker}"), + ], + ) + assert duplicate.status_code == 200, duplicate.text + batches: Final[list[bytes]] = [] + spans: Final = eventually( + lambda: _generation_marker_span_attributes(collector, batches, marker), + lambda found: len(found) == 2, + seconds=30, + ) + assert all(attributes.get("user.id") == f"first-{marker}" for attributes in spans), spans + + +def test_langfuse_otel_bad_key_emits_no_generation_span(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + key="sk-not-a-real-key", + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 401, response.text + batches: Final[list[bytes]] = [] + _subsequent_request_still_lands(candidate, collector, batches, model) + offending: Final = tuple( + attributes + for attributes in _generation_marker_span_attributes(collector, batches, marker) + if attributes.get("user.id") not in (None, f"end-user-{marker}") + ) + assert offending == (), offending + + +def test_langfuse_otel_upstream_failure_never_invents_an_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server( + lambda request: ( + Reply(status=401, body=b'{"error": "upstream denied ' + marker.encode() + b'"}') + if marker.encode() in (request.body or b"") + else _upstream_reply_for(request, _marker_from_body(request)) + ) + ) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = _end_user_request(candidate, model, marker) + assert response.status_code == 401, response.text + batches: Final[list[bytes]] = [] + _subsequent_request_still_lands(candidate, collector, batches, model) + offending: Final = tuple( + attributes + for _t, attributes in _span_containing_marker(_drained_spans(collector, batches), marker) + if attributes.get("user.id") not in (None, f"end-user-{marker}") + ) + assert offending == (), offending + + +def test_langfuse_otel_sink_rejections_do_not_drop_the_proxy(gateway: Gateway, tmp_path: Path) -> None: + calls: Final[dict[str, int]] = {"count": 0} + lock: Final = threading.Lock() + + def rejecting_sink(request: Request) -> Reply: + with lock: + calls["count"] += 1 + seen: Final = calls["count"] + if seen == 1: + return Reply(status=403) + if seen == 2: + return Reply(status=404) + return _sink(request) + + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(rejecting_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + markers: Final = tuple(uuid.uuid4().hex for _ in range(3)) + batches: Final[list[bytes]] = [] + for index, item in enumerate(markers): + response: Final = _end_user_request(candidate, model, item) + assert response.status_code == 200, response.text + if index < 2: + eventually( + lambda collector=collector, batches=batches: ( + batches.extend(request.body for request in collector.drain()) or len(batches) + ), + lambda seen, index=index: seen >= index + 1, + seconds=30, + ) + attributes: Final = _await_generation(collector, batches, markers[2]) + assert attributes.get("user.id") == f"end-user-{markers[2]}", attributes + _subsequent_request_still_lands(candidate, collector, batches, model) + + +def test_langfuse_otel_unknown_model_error_reaches_the_caller(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = _end_user_request(candidate, "no-such-model-" + marker[:12], marker) + assert response.status_code in (400, 404), response.text + batches: Final[list[bytes]] = [] + _subsequent_request_still_lands(candidate, collector, batches, model) + + +def test_langfuse_otel_empty_and_null_session_id_stay_absent(gateway: Gateway, tmp_path: Path) -> None: + markers: Final = tuple(uuid.uuid4().hex for _ in range(3)) + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + batches: Final[list[bytes]] = [] + sessions: Final = ({"session_id": ""}, {"session_id": None}, {}) + for marker, metadata in zip(markers, sessions): + response: Final = _end_user_request(candidate, model, marker, metadata=metadata) + assert response.status_code == 200, response.text + for marker, metadata in zip(markers, sessions): + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == { + "user.id": f"end-user-{marker}", + "session.id": metadata.get("session_id"), + }, attributes + + +def test_langfuse_otel_empty_trace_user_id_falls_back_to_the_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = _end_user_request(candidate, model, marker, metadata={"trace_user_id": ""}) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_three_identical_requests_each_land_the_end_user(gateway: Gateway, tmp_path: Path) -> None: + markers: Final = tuple(uuid.uuid4().hex for _ in range(3)) + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + batches: Final[list[bytes]] = [] + for marker in markers: + response: Final = _end_user_request(candidate, model, marker) + assert response.status_code == 200, response.text + for marker in markers: + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_litellm_metadata_session_id_on_responses(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/responses", + { + "model": model, + "input": marker, + "litellm_metadata": {"session_id": f"sess-{marker}"}, + "cache": {"no-cache": True}, + }, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _generation_marker_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": f"sess-{marker}"}, attributes diff --git a/tests/integration/observability/test_langfuse_otel_identity_surfaces.py b/tests/integration/observability/test_langfuse_otel_identity_surfaces.py new file mode 100644 index 00000000000..420906b0fa9 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_identity_surfaces.py @@ -0,0 +1,512 @@ +import uuid +from pathlib import Path +from typing import Final + +from _langfuse_otel import ( + _arize_generation_span_attributes, + _generation_marker_span_attributes, + _generation_span_attributes, + _langfuse_proxy, + _observability_proxy, + _sink, + _trace_user_span_attributes, + _upstream_reply_for, +) +from anthropic import Anthropic, AsyncAnthropic +from integration._support.client import Gateway, eventually +from integration._support.database import read_rows +from integration._support.wire import Wire, wire_server +from openai import AsyncOpenAI, OpenAI + + +def _openai_client(gateway: Gateway, key: str, marker: str) -> OpenAI: + return OpenAI( + api_key=key, + base_url=str(gateway.client.base_url).rstrip("/") + "/v1", + timeout=30, + max_retries=0, + default_headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + + +def _async_openai_client(gateway: Gateway, key: str, marker: str) -> AsyncOpenAI: + return AsyncOpenAI( + api_key=key, + base_url=str(gateway.client.base_url).rstrip("/") + "/v1", + timeout=30, + max_retries=0, + default_headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + + +def _anthropic_client(gateway: Gateway, key: str, marker: str) -> Anthropic: + return Anthropic( + api_key=key, + base_url=str(gateway.client.base_url).rstrip("/"), + timeout=30, + max_retries=0, + default_headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + + +def _async_anthropic_client(gateway: Gateway, key: str, marker: str) -> AsyncAnthropic: + return AsyncAnthropic( + api_key=key, + base_url=str(gateway.client.base_url).rstrip("/"), + timeout=30, + max_retries=0, + default_headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + + +def _await_generation(collector: Wire, batches: list[bytes], marker: str) -> dict[str, object]: + return eventually( + lambda: _generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + + +def _await_marker_span(collector: Wire, batches: list[bytes], marker: str) -> dict[str, object]: + return eventually( + lambda: _generation_marker_span_attributes(collector, batches, marker), + lambda spans: len(spans) == 1, + seconds=30, + )[0] + + +def _identity(attributes: dict[str, object]) -> dict[str, object]: + return {key: attributes.get(key) for key in ("user.id", "session.id")} + + +def test_langfuse_otel_customer_id_header_lands_in_user_id(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + headers={"x-litellm-customer-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_langfuse_trace_user_id_header_wins_over_the_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + headers={ + "x-litellm-end-user-id": f"end-user-{marker}", + "langfuse_trace_user_id": f"caller-{marker}", + }, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"caller-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_no_end_user_never_exposes_the_internal_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = 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 + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": None, "session.id": None}, attributes + + +def test_langfuse_otel_team_key_end_user_keeps_the_litellm_attributes(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + team: Final = scenario.team(team_alias=f"team-alias-{marker[:12]}") + key: Final = scenario.key(team_id=team, key_alias=f"key-alias-{marker[:12]}") + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + assert attributes.get("litellm.team_id") == team, attributes + assert attributes.get("litellm.team_alias") == f"team-alias-{marker[:12]}", attributes + assert attributes.get("litellm.key_alias") == f"key-alias-{marker[:12]}", attributes + + +def test_langfuse_otel_header_end_user_beats_the_body_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": f"body-user-{marker}", + "cache": {"no-cache": True}, + }, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + rows: Final = eventually( + lambda: read_rows( + 'SELECT end_user FROM "LiteLLM_SpendLogs" WHERE request_id = %s', (str(response.json()["id"]),) + ), + lambda values: len(values) == 1, + seconds=70, + ) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + assert attributes.get("user.id") == rows[0]["end_user"], (attributes, rows) + + +def test_langfuse_otel_openai_sdk_streaming_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + with _openai_client(candidate, key, marker) as client: + stream: Final = client.chat.completions.create( + model=model, + messages=[{"role": "user", "content": marker}], + stream=True, + extra_body={"cache": {"no-cache": True}}, + ) + with stream: + chunks: Final = tuple(stream) + assert {chunk.id for chunk in chunks} == {marker} + assert f"reply {marker}" in "".join( + choice.delta.content or "" for chunk in chunks for choice in chunk.choices + ) + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +async def test_langfuse_otel_openai_async_sdk_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + async with _async_openai_client(candidate, key, marker) as client: + completion: Final = await client.chat.completions.create( + model=model, + messages=[{"role": "user", "content": marker}], + extra_body={"cache": {"no-cache": True}}, + ) + assert completion.id == marker + batches: Final[list[bytes]] = [] + attributes: Final = _await_generation(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_responses_api_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + with _openai_client(candidate, key, marker) as client: + completion: Final = client.responses.create( + model=model, input=marker, extra_body={"cache": {"no-cache": True}} + ) + assert marker in completion.output_text, completion.output_text + batches: Final[list[bytes]] = [] + attributes: Final = _await_marker_span(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +async def test_langfuse_otel_responses_api_async_streaming_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + async with _async_openai_client(candidate, key, marker) as client: + stream: Final = await client.responses.create( + model=model, input=marker, stream=True, extra_body={"cache": {"no-cache": True}} + ) + events: Final = tuple([event async for event in stream]) + assert events, events + batches: Final[list[bytes]] = [] + attributes: Final = _await_marker_span(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_anthropic_sdk_messages_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + with _anthropic_client(candidate, key, marker) as client: + message: Final = client.messages.create( + model=model, max_tokens=5, messages=[{"role": "user", "content": marker}] + ) + assert marker in message.content[0].text, message.model_dump_json() + batches: Final[list[bytes]] = [] + attributes: Final = _await_marker_span(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +async def test_langfuse_otel_anthropic_async_streaming_messages_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + async with _async_anthropic_client(candidate, key, marker) as client: + stream: Final = await client.messages.create( + model=model, max_tokens=5, messages=[{"role": "user", "content": marker}], stream=True + ) + events: Final = tuple([event async for event in stream]) + assert events, events + batches: Final[list[bytes]] = [] + attributes: Final = _await_marker_span(collector, batches, marker) + assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes + + +def test_langfuse_otel_v2_caller_trace_user_id_wins_over_the_end_user(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _langfuse_proxy(gateway, tmp_path, collector.url, {"LITELLM_OTEL_V2": "1"}) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "metadata": {"trace_user_id": f"caller-{marker}"}, + "cache": {"no-cache": True}, + }, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + user_spans: Final = eventually( + lambda: _trace_user_span_attributes(collector, batches, marker), + lambda spans: len(spans) >= 1, + seconds=30, + ) + assert all( + _identity(attributes) == {"user.id": f"caller-{marker}", "session.id": None} for attributes in user_spans + ), user_spans + + +def test_arize_phoenix_header_end_user_keeps_session_mapping(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _observability_proxy( + gateway, + tmp_path, + { + "PHOENIX_COLLECTOR_ENDPOINT": collector.url + "/v1/traces", + "PHOENIX_API_KEY": "phoenix-integration", + }, + callbacks=("arize_phoenix",), + config_name="arize_phoenix.yaml", + ) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _arize_generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) >= 1, + seconds=30, + )[0] + assert _identity(attributes) == {"user.id": owner, "session.id": f"end-user-{marker}"}, attributes + assert attributes.get("litellm.trace_id") is not None, attributes + + +def test_arize_phoenix_caller_session_id_stays_session(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as collector, + _observability_proxy( + gateway, + tmp_path, + { + "PHOENIX_COLLECTOR_ENDPOINT": collector.url + "/v1/traces", + "PHOENIX_API_KEY": "phoenix-integration", + }, + callbacks=("arize_phoenix",), + config_name="arize_phoenix.yaml", + ) as candidate, + candidate.scenario() as scenario, + ): + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": f"end-user-{marker}", + "metadata": {"session_id": f"sess-{marker}"}, + "cache": {"no-cache": True}, + }, + ) + assert response.status_code == 200, response.text + batches: Final[list[bytes]] = [] + attributes: Final = eventually( + lambda: _arize_generation_span_attributes(collector, batches, marker), + lambda spans: len(spans) >= 1, + seconds=30, + )[0] + assert _identity(attributes) == { + "user.id": f"end-user-{marker}", + "session.id": f"end-user-{marker}", + }, attributes + assert attributes.get("litellm.trace_id") == f"sess-{marker}", attributes + + +def test_langfuse_otel_and_arize_phoenix_together_keep_each_mapping(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + with ( + wire_server(lambda request: _upstream_reply_for(request, marker)) as provider, + wire_server(_sink) as langfuse_collector, + wire_server(_sink) as phoenix_collector, + _observability_proxy( + gateway, + tmp_path, + { + "LANGFUSE_PUBLIC_KEY": "pk-integration", + "LANGFUSE_SECRET_KEY": "sk-integration", + "LANGFUSE_HOST": langfuse_collector.url, + "PHOENIX_COLLECTOR_ENDPOINT": phoenix_collector.url + "/v1/traces", + "PHOENIX_API_KEY": "phoenix-integration", + }, + callbacks=("langfuse_otel", "arize_phoenix"), + config_name="dual_callbacks.yaml", + ) as candidate, + candidate.scenario() as scenario, + ): + owner: Final = scenario.user() + key: Final = scenario.key(user_id=owner) + model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1") + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}, + key=key, + headers={"x-litellm-end-user-id": f"end-user-{marker}"}, + ) + assert response.status_code == 200, response.text + langfuse_batches: Final[list[bytes]] = [] + phoenix_batches: Final[list[bytes]] = [] + langfuse_attributes: Final = _await_generation(langfuse_collector, langfuse_batches, marker) + phoenix_attributes: Final = eventually( + lambda: _arize_generation_span_attributes(phoenix_collector, phoenix_batches, marker), + lambda spans: len(spans) >= 1, + seconds=30, + )[0] + assert _identity(langfuse_attributes) == { + "user.id": f"end-user-{marker}", + "session.id": None, + }, langfuse_attributes + assert _identity(phoenix_attributes) == { + "user.id": owner, + "session.id": f"end-user-{marker}", + }, phoenix_attributes diff --git a/tests/unit/integrations/arize/test_arize_utils.py b/tests/unit/integrations/arize/test_arize_utils.py index 167b083e147..5490e1503fb 100644 --- a/tests/unit/integrations/arize/test_arize_utils.py +++ b/tests/unit/integrations/arize/test_arize_utils.py @@ -1513,3 +1513,32 @@ def test_arize_mcp_emitter_is_inert_without_a_standard_logging_object(): written = {c.args[0]: c.args[1] for c in span.set_attribute.call_args_list} assert SpanAttributes.TOOL_NAME not in written + + +def test_arize_session_and_user_attrs_still_emit_from_key_metadata_by_default(): + from unittest.mock import MagicMock + + span = MagicMock() + kwargs = { + "model": "gpt-4o", + "messages": [{"role": "user", "content": "hello"}], + "standard_logging_object": { + "call_type": "acompletion", + "model_parameters": {}, + "metadata": { + "user_api_key_end_user_id": "end-1", + "user_api_key_user_id": "internal-1", + "user_api_key_team_id": "team-1", + }, + "trace_id": "trace-1", + }, + "optional_params": {}, + "litellm_params": {"custom_llm_provider": "openai"}, + } + + ArizeLogger.set_arize_attributes(span, kwargs, {"id": "chatcmpl-1", "choices": [], "usage": {}}) + + span.set_attribute.assert_any_call(SpanAttributes.SESSION_ID, "end-1") + span.set_attribute.assert_any_call(SpanAttributes.USER_ID, "internal-1") + span.set_attribute.assert_any_call("litellm.trace_id", "trace-1") + span.set_attribute.assert_any_call("litellm.team_id", "team-1") diff --git a/tests/unit/integrations/otel/test_langfuse_logger.py b/tests/unit/integrations/otel/test_langfuse_logger.py index aca9dcc8a5e..bf0afec6f54 100644 --- a/tests/unit/integrations/otel/test_langfuse_logger.py +++ b/tests/unit/integrations/otel/test_langfuse_logger.py @@ -419,6 +419,31 @@ def test_langfuse_user_and_session_headers_beat_body_metadata_on_both_spans(): assert attrs["session.id"] == "from-header-s" +def test_the_proxy_end_user_fills_user_id_when_the_caller_names_no_trace_user(): + logger, exporter = _logger() + + root_attrs, generation_attrs = _run_named_request( + logger, exporter, {"metadata": {"user_api_key_end_user_id": "end-1"}, "proxy_server_request": {"headers": {}}} + ) + + assert root_attrs["user.id"] == "end-1" + + +def test_a_callers_trace_user_id_still_wins_over_the_proxy_end_user(): + logger, exporter = _logger() + + root_attrs, generation_attrs = _run_named_request( + logger, + exporter, + { + "metadata": {"trace_user_id": "caller-1", "user_api_key_end_user_id": "end-1"}, + "proxy_server_request": {"headers": {}}, + }, + ) + + assert root_attrs["user.id"] == "caller-1" + + def test_caller_metadata_cannot_override_the_proxy_team_identity(): logger, exporter = _logger() response: Final = ModelResponse(choices=[Choices(message=Message(role="assistant", content="pong"))]) diff --git a/tests/unit/integrations/otel/test_otel_v2_sources_of_truth.py b/tests/unit/integrations/otel/test_otel_v2_sources_of_truth.py index 57f4557c6f7..1bdab85cf24 100644 --- a/tests/unit/integrations/otel/test_otel_v2_sources_of_truth.py +++ b/tests/unit/integrations/otel/test_otel_v2_sources_of_truth.py @@ -46,7 +46,11 @@ from litellm.integrations.otel.model.spans import ( root_roles, validate_registry, ) -from litellm.integrations.otel.model.trace_controls import TraceControls, caller_trace_controls +from litellm.integrations.otel.model.trace_controls import ( + TraceControls, + caller_trace_controls, + langfuse_trace_controls, +) @pytest.fixture(autouse=True) @@ -1342,6 +1346,43 @@ def test_caller_trace_controls_carry_user_session_and_tags(request_data, expecte assert LLMCallEvent.from_dict({"litellm_params": request_data}).trace == expected +@pytest.mark.parametrize( + ("request_data", "expected"), + [ + ({"metadata": {"user_api_key_end_user_id": "end-1"}}, TraceControls(user_id="end-1")), + ({"litellm_metadata": {"user_api_key_end_user_id": "end-2"}}, TraceControls(user_id="end-2")), + ( + {"metadata": {"trace_user_id": "caller-1", "user_api_key_end_user_id": "end-1"}}, + TraceControls(user_id="caller-1"), + ), + ( + { + "proxy_server_request": {"headers": {"langfuse_trace_user_id": "header-1"}}, + "metadata": {"user_api_key_end_user_id": "end-1"}, + }, + TraceControls(user_id="header-1"), + ), + ( + {"metadata": {"user_api_key_end_user_id": "end-1", "session_id": "s-1"}}, + TraceControls(user_id="end-1", session_id="s-1"), + ), + ({"metadata": {"user_api_key_user_id": "internal-1"}}, TraceControls()), + ({}, TraceControls()), + ], + ids=[ + "body-end-user", + "anthropic-end-user", + "body-caller-wins", + "header-caller-wins", + "session-kept", + "internal-user-ignored", + "empty", + ], +) +def test_langfuse_trace_controls_fall_back_to_the_proxy_end_user(request_data, expected): + assert langfuse_trace_controls({"litellm_params": request_data}) == expected + + def test_llm_span_data_carries_the_caller_trace_controls(): controls: Final = TraceControls(name="nightly-eval", user_id="u1", session_id="s1", tags=("a", "b")) data: Final = LLMCallSpanData.from_standard_logging_payload(_sample_payload(), trace=controls) diff --git a/tests/unit/integrations/test_langfuse_otel.py b/tests/unit/integrations/test_langfuse_otel.py index 0a9ce55fe16..48473523352 100644 --- a/tests/unit/integrations/test_langfuse_otel.py +++ b/tests/unit/integrations/test_langfuse_otel.py @@ -112,7 +112,7 @@ class TestLangfuseOtelIntegration: ) mock_set_attributes.assert_called_once_with( - mock_span, mock_kwargs, mock_response, LangfuseLLMObsOTELAttributes + mock_span, mock_kwargs, mock_response, LangfuseLLMObsOTELAttributes, emit_session_and_user=False ) mock_span.set_attribute.assert_any_call( "langfuse.observation.type", "generation" @@ -709,7 +709,7 @@ class TestLangfuseOtelResponsesAPI: # Verify that set_attributes was called for general attributes mock_set_attributes.assert_called_once_with( - mock_span, kwargs, mock_response, LangfuseLLMObsOTELAttributes + mock_span, kwargs, mock_response, LangfuseLLMObsOTELAttributes, emit_session_and_user=False ) # Verify that Langfuse-specific attributes were set @@ -1006,5 +1006,115 @@ class TestLangfuseOtelResponsesAPI: assert output_data[0]["arguments"] == {} +class TestLangfuseOtelTraceIdentity: + def _recording_span(self): + from opentelemetry.sdk.trace import TracerProvider + + return TracerProvider().get_tracer("test").start_span("generation") + + def _kwargs(self, slp_metadata=None, litellm_metadata=None, model_parameters=None, slp_extra=None): + return { + "model": "gpt-4o", + "messages": [{"role": "user", "content": "hello"}], + "optional_params": {}, + "litellm_params": {"metadata": litellm_metadata or {}, "custom_llm_provider": "openai"}, + "standard_logging_object": { + "call_type": "acompletion", + "model_parameters": model_parameters or {}, + "metadata": slp_metadata or {}, + **(slp_extra or {}), + }, + } + + def _response_obj(self): + return { + "id": "chatcmpl-1", + "model": "gpt-4o", + "choices": [{"index": 0, "message": {"role": "assistant", "content": "hi"}, "finish_reason": "stop"}], + "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}, + } + + def _identity(self, kwargs): + span = self._recording_span() + LangfuseOtelLogger.set_langfuse_otel_attributes(span, kwargs, self._response_obj()) + attributes = dict(span.attributes or {}) + return {key: attributes.get(key) for key in ("user.id", "session.id")}, attributes + + def test_header_end_user_beats_internal_key_owner_in_user_id(self): + identity, _ = self._identity( + self._kwargs( + slp_metadata={ + "user_api_key_end_user_id": "end-1", + "user_api_key_user_id": "internal-1", + } + ) + ) + assert identity == {"user.id": "end-1", "session.id": None} + + def test_header_end_user_lands_in_user_id_for_a_service_key(self): + identity, _ = self._identity(self._kwargs(slp_metadata={"user_api_key_end_user_id": "end-1"})) + assert identity == {"user.id": "end-1", "session.id": None} + + def test_body_user_is_never_a_session(self): + identity, _ = self._identity( + self._kwargs( + slp_metadata={"user_api_key_end_user_id": "body-user"}, + model_parameters={"user": "body-user"}, + ) + ) + assert identity == {"user.id": "body-user", "session.id": None} + + def test_caller_trace_user_id_wins_over_the_end_user(self): + identity, _ = self._identity( + self._kwargs( + slp_metadata={"user_api_key_end_user_id": "end-1"}, + litellm_metadata={"trace_user_id": "caller-1"}, + ) + ) + assert identity == {"user.id": "caller-1", "session.id": None} + + def test_caller_trace_user_id_under_litellm_metadata_wins_over_the_end_user(self): + kwargs = self._kwargs(slp_metadata={"user_api_key_end_user_id": "end-1"}) + kwargs["litellm_params"]["litellm_metadata"] = {"trace_user_id": "caller-1"} + identity, _ = self._identity(kwargs) + assert identity == {"user.id": "caller-1", "session.id": None} + + def test_end_user_only_under_litellm_metadata_lands_in_user_id(self): + kwargs = self._kwargs() + kwargs["litellm_params"]["litellm_metadata"] = {"user_api_key_end_user_id": "end-1"} + identity, _ = self._identity(kwargs) + assert identity == {"user.id": "end-1", "session.id": None} + + def test_caller_session_id_stays_the_session_beside_the_end_user(self): + identity, _ = self._identity( + self._kwargs( + slp_metadata={"user_api_key_end_user_id": "end-1"}, + litellm_metadata={"session_id": "sess-1"}, + ) + ) + assert identity == {"user.id": "end-1", "session.id": "sess-1"} + + def test_internal_user_without_an_end_user_never_lands_in_user_id(self): + identity, _ = self._identity(self._kwargs(slp_metadata={"user_api_key_user_id": "internal-1"})) + assert identity == {"user.id": None, "session.id": None} + + def test_request_context_attributes_still_emit(self): + _, attributes = self._identity( + self._kwargs( + slp_metadata={ + "user_api_key_end_user_id": "end-1", + "user_api_key_team_id": "team-1", + "user_api_key_team_alias": "team-alias", + "user_api_key_alias": "key-alias", + }, + slp_extra={"trace_id": "trace-1"}, + ) + ) + assert attributes["litellm.trace_id"] == "trace-1" + assert attributes["litellm.team_id"] == "team-1" + assert attributes["litellm.team_alias"] == "team-alias" + assert attributes["litellm.key_alias"] == "key-alias" + + if __name__ == "__main__": pytest.main([__file__])