test(langfuse_otel): audit cells for end user identity mapping

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-23 09:46:48 +00:00
parent c9050a8fbe
commit c134e5bf1f
6 changed files with 1693 additions and 113 deletions

View file

@ -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": {

View file

@ -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

View file

@ -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")

View file

@ -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)

View file

@ -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

View file

@ -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