From c134e5bf1fd24e1713775ce6d672e72a5e018a1b Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 09:46:48 +0000 Subject: [PATCH] test(langfuse_otel): audit cells for end user identity mapping Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- tests/integration/contracts.json | 102 ++++ .../observability/_langfuse_otel.py | 307 ++++++++++ .../test_langfuse_otel_identity.py | 122 +--- .../test_langfuse_otel_identity_chaos.py | 305 ++++++++++ .../test_langfuse_otel_identity_edges.py | 442 +++++++++++++++ .../test_langfuse_otel_identity_surfaces.py | 528 ++++++++++++++++++ 6 files changed, 1693 insertions(+), 113 deletions(-) create mode 100644 tests/integration/observability/_langfuse_otel.py create mode 100644 tests/integration/observability/test_langfuse_otel_identity_chaos.py create mode 100644 tests/integration/observability/test_langfuse_otel_identity_edges.py create mode 100644 tests/integration/observability/test_langfuse_otel_identity_surfaces.py diff --git a/tests/integration/contracts.json b/tests/integration/contracts.json index 50ff6072ec2..18d968316f2 100644 --- a/tests/integration/contracts.json +++ b/tests/integration/contracts.json @@ -2048,6 +2048,108 @@ "mgmt.key.reset_spend.returns_promptly_with_wedged_coordination_redis", "mgmt.auth_cache_invalidation.publish_parked_by_short_redis_wedge_lands_after_recovery", "mgmt.auth_cache_invalidation.burst_with_worker_kill_keeps_serving_while_redis_wedged" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_customer_id_header_lands_in_user_id": [ + "other.observability.langfuse_otel.customer_id_header_end_user_in_user_id" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_langfuse_trace_user_id_header_wins_over_the_end_user": [ + "other.observability.langfuse_otel.langfuse_trace_user_id_header_wins" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_no_end_user_never_exposes_the_internal_user": [ + "other.observability.langfuse_otel.no_end_user_no_user_id" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_team_key_end_user_keeps_the_litellm_attributes": [ + "other.observability.langfuse_otel.team_key_end_user_with_team_attributes" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_header_end_user_beats_the_body_user": [ + "other.observability.langfuse_otel.header_end_user_beats_body_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_openai_sdk_streaming_end_user": [ + "other.observability.langfuse_otel.openai_sdk_streaming_end_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_openai_async_sdk_end_user": [ + "other.observability.langfuse_otel.openai_async_sdk_end_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_responses_api_end_user": [ + "other.observability.langfuse_otel.responses_api_end_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_responses_api_async_streaming_end_user": [ + "other.observability.langfuse_otel.responses_api_async_streaming_end_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_anthropic_sdk_messages_end_user": [ + "other.observability.langfuse_otel.anthropic_sdk_messages_end_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_anthropic_async_streaming_messages_end_user": [ + "other.observability.langfuse_otel.anthropic_async_streaming_messages_end_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_v2_caller_trace_user_id_wins_over_the_end_user": [ + "other.observability.langfuse_otel.v2_caller_trace_user_id_wins" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_arize_phoenix_header_end_user_keeps_session_mapping": [ + "other.observability.langfuse_otel.arize_phoenix_unchanged_header_end_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_arize_phoenix_caller_session_id_stays_session": [ + "other.observability.langfuse_otel.arize_phoenix_unchanged_caller_session" + ], + "tests/integration/observability/test_langfuse_otel_identity_surfaces.py::test_langfuse_otel_and_arize_phoenix_together_keep_each_mapping": [ + "other.observability.langfuse_otel.dual_callbacks_arize_unchanged" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_five_kb_end_user_header_lands_in_user_id": [ + "other.observability.langfuse_otel.end_user_5kb_header" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_empty_end_user_header_writes_no_identity": [ + "other.observability.langfuse_otel.empty_end_user_header" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_integer_body_user_does_not_crash": [ + "other.observability.langfuse_otel.integer_body_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_list_body_user_does_not_crash": [ + "other.observability.langfuse_otel.list_body_user" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_integer_trace_user_id_does_not_crash": [ + "other.observability.langfuse_otel.integer_trace_user_id" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_null_metadata_still_maps_the_end_user": [ + "other.observability.langfuse_otel.null_metadata" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_duplicate_end_user_header_uses_the_first": [ + "other.observability.langfuse_otel.duplicate_end_user_header" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_bad_key_emits_no_generation_span": [ + "other.observability.langfuse_otel.unauthenticated_no_span" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_upstream_failure_never_invents_an_identity": [ + "other.observability.langfuse_otel.upstream_failure_identity" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_sink_rejections_do_not_drop_the_proxy": [ + "other.observability.langfuse_otel.sink_rejections_survived" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_unknown_model_error_reaches_the_caller": [ + "other.observability.langfuse_otel.unknown_model" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_empty_and_null_session_id_stay_absent": [ + "other.observability.langfuse_otel.session_id_empty_null_missing" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_empty_trace_user_id_falls_back_to_the_end_user": [ + "other.observability.langfuse_otel.empty_trace_user_id" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_three_identical_requests_each_land_the_end_user": [ + "other.observability.langfuse_otel.three_identical_requests" + ], + "tests/integration/observability/test_langfuse_otel_identity_edges.py::test_langfuse_otel_litellm_metadata_session_id_on_responses": [ + "other.observability.langfuse_otel.litellm_metadata_session_on_responses" + ], + "tests/integration/observability/test_langfuse_otel_identity_chaos.py::test_langfuse_otel_sink_outage_mid_burst_never_duplicates": [ + "other.observability.langfuse_otel.burst_sink_outage" + ], + "tests/integration/observability/test_langfuse_otel_identity_chaos.py::test_langfuse_otel_slow_sink_never_duplicates_or_loses": [ + "other.observability.langfuse_otel.burst_slow_sink" + ], + "tests/integration/observability/test_langfuse_otel_identity_chaos.py::test_langfuse_otel_worker_kill_mid_burst_keeps_every_marker": [ + "other.observability.langfuse_otel.worker_kill_mid_burst" + ], + "tests/integration/observability/test_langfuse_otel_identity_chaos.py::test_langfuse_otel_proxy_restart_mid_burst_keeps_every_marker": [ + "other.observability.langfuse_otel.proxy_restart_mid_burst" ] }, "browser": { 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 index 2c9ed62fde0..cfa440f59f0 100644 --- a/tests/integration/observability/test_langfuse_otel_identity.py +++ b/tests/integration/observability/test_langfuse_otel_identity.py @@ -1,122 +1,18 @@ -import json import uuid -from collections.abc import Iterator, Mapping -from contextlib import contextmanager from pathlib import Path from typing import Final import pytest -import yaml +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.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 - - -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 _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 _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 _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 - ) - - -@contextmanager -def _langfuse_proxy( - gateway: Gateway, directory: Path, collector_url: str, overrides: Mapping[str, str] | None = None -) -> Iterator[Gateway]: - config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text()) - config["litellm_settings"].update({"callbacks": ["langfuse_otel"]}) - path: Final = directory / "langfuse_otel.yaml" - path.write_text(yaml.safe_dump(config)) - with owned_proxy( - gateway, - directory, - { - "LANGFUSE_PUBLIC_KEY": "pk-integration", - "LANGFUSE_SECRET_KEY": "sk-integration", - "LANGFUSE_HOST": collector_url, - "OTEL_BSP_SCHEDULE_DELAY": "100", - **(overrides or {}), - }, - config=path, - ) as candidate: - yield candidate +from integration._support.wire import Reply, Request, wire_server @pytest.mark.covers("other.observability.langfuse_otel.header_end_user_in_user_id_over_internal_user") 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..a5c6346292d --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_identity_chaos.py @@ -0,0 +1,305 @@ +import concurrent.futures +import os +import signal +import threading +import time +import uuid +from pathlib import Path +from typing import Final + +import httpx +import pytest +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} + + +@pytest.mark.covers("other.observability.langfuse_otel.burst_sink_outage") +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) + + +@pytest.mark.covers("other.observability.langfuse_otel.burst_slow_sink") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.worker_kill_mid_burst") +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) + + +@pytest.mark.covers("other.observability.langfuse_otel.proxy_restart_mid_burst") +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..eba412a65e5 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_identity_edges.py @@ -0,0 +1,442 @@ +import threading +import uuid +from pathlib import Path +from typing import Final + +import httpx +import pytest +from _langfuse_otel import ( + _drained_spans, + _generation_marker_span_attributes, + _generation_span_attributes, + _langfuse_proxy, + _sink, + _span_containing_marker, + _marker_from_body, + _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 + + +@pytest.mark.covers("other.observability.langfuse_otel.end_user_5kb_header") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.empty_end_user_header") +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) + + +@pytest.mark.covers("other.observability.langfuse_otel.integer_body_user") +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) + + +@pytest.mark.covers("other.observability.langfuse_otel.list_body_user") +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) + + +@pytest.mark.covers("other.observability.langfuse_otel.integer_trace_user_id") +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) + + +@pytest.mark.covers("other.observability.langfuse_otel.null_metadata") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.duplicate_end_user_header") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.unauthenticated_no_span") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.upstream_failure_identity") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.sink_rejections_survived") +def test_langfuse_otel_sink_rejections_do_not_drop_the_proxy(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = uuid.uuid4().hex + 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) + + +@pytest.mark.covers("other.observability.langfuse_otel.unknown_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) + + +@pytest.mark.covers("other.observability.langfuse_otel.session_id_empty_null_missing") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.empty_trace_user_id") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.three_identical_requests") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.litellm_metadata_session_on_responses") +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..6c51da46bb4 --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_identity_surfaces.py @@ -0,0 +1,528 @@ +import uuid +from pathlib import Path +from typing import Final + +import pytest +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 Reply, Request, 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")} + + +@pytest.mark.covers("other.observability.langfuse_otel.customer_id_header_end_user_in_user_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 + + +@pytest.mark.covers("other.observability.langfuse_otel.langfuse_trace_user_id_header_wins") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.no_end_user_no_user_id") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.team_key_end_user_with_team_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 + + +@pytest.mark.covers("other.observability.langfuse_otel.header_end_user_beats_body_user") +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) + + +@pytest.mark.covers("other.observability.langfuse_otel.openai_sdk_streaming_end_user") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.openai_async_sdk_end_user") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.responses_api_end_user") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.responses_api_async_streaming_end_user") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.anthropic_sdk_messages_end_user") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.anthropic_async_streaming_messages_end_user") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.v2_caller_trace_user_id_wins") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.arize_phoenix_unchanged_header_end_user") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.arize_phoenix_unchanged_caller_session") +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 + + +@pytest.mark.covers("other.observability.langfuse_otel.dual_callbacks_arize_unchanged") +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