This commit is contained in:
devin-ai-integration[bot] 2026-10-03 16:29:53 -04:00 • committed by GitHub
commit 2208695b5d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
13 changed files with 2131 additions and 23 deletions

View file

@ -9,6 +9,7 @@ from litellm.integrations.opentelemetry_utils.base_otel_llm_obs_attributes impor
BaseLLMObsOTELAttributes,
safe_set_attribute,
)
from litellm.integrations.otel.model.utils import as_str_mapping
from litellm.litellm_core_utils.redact_messages import (
should_redact_message_logging,
)
@ -410,7 +411,14 @@ def _set_tool_attributes(span: "Span", optional_tools: list | None, metadata_too
)
def set_attributes(span: "Span", kwargs, response_obj, attributes: type[BaseLLMObsOTELAttributes]):
def set_attributes(
span: "Span",
kwargs,
response_obj,
attributes: type[BaseLLMObsOTELAttributes],
*,
emit_session_and_user: bool = True,
):
"""
Populates span with OpenInference-compliant LLM attributes for Arize and Phoenix tracing.
"""
@ -471,7 +479,9 @@ def set_attributes(span: "Span", kwargs, response_obj, attributes: type[BaseLLMO
# blank the attributes set by the main try-block above. New attributes are
# written under new keys; existing attributes are not overwritten.
slp: Final = kwargs.get("standard_logging_object")
_safe_emit("session/user attrs", _set_session_and_user_attrs, span, kwargs, slp)
if emit_session_and_user:
_safe_emit("session/user attrs", _set_session_and_user_attrs, span, kwargs, slp)
_safe_emit("request context attrs", _set_request_context_attrs, span, slp)
_safe_emit("response cost", _set_response_cost_attr, span, slp)
_safe_emit(
"passthrough normalization",
@ -834,14 +844,12 @@ def _emit_input_message_extras(span: "Span", prefix: str, message: dict) -> None
def _set_session_and_user_attrs(span: "Span", kwargs: dict, standard_logging_payload) -> None:
"""Emit `SESSION_ID` / `USER_ID` / team metadata when source data exists.
"""Emit `SESSION_ID` / `USER_ID` when source data exists.
`SESSION_ID` is emitted only when an explicit end-user identifier exists
(`metadata.user_api_key_end_user_id`). We deliberately do NOT fall back
to `trace_id`, because that would create a distinct "session" for every
single request and distort Arize's Session-grouping analytics. The
`trace_id` is still emitted under its own `litellm.trace_id` key so
spans remain filterable by trace.
single request and distort Arize's Session-grouping analytics.
USER_ID is *only* emitted when no upstream path (model_params.user or
optional_params.user) has already set it, to avoid overwriting an
@ -857,10 +865,6 @@ def _set_session_and_user_attrs(span: "Span", kwargs: dict, standard_logging_pay
if session_id:
safe_set_attribute(span, SpanAttributes.SESSION_ID, str(session_id))
trace_id: Final = standard_logging_payload.get("trace_id")
if trace_id:
safe_set_attribute(span, "litellm.trace_id", str(trace_id))
optional_params: Final = kwargs.get("optional_params") or {}
model_params: Final = standard_logging_payload.get("model_parameters") or {}
has_user_already: Final = bool(
@ -872,6 +876,19 @@ def _set_session_and_user_attrs(span: "Span", kwargs: dict, standard_logging_pay
if user_id:
safe_set_attribute(span, SpanAttributes.USER_ID, str(user_id))
def _set_request_context_attrs(span: "Span", standard_logging_payload: object) -> None:
payload: Final = as_str_mapping(standard_logging_payload)
if payload is None:
return
trace_id: Final = payload.get("trace_id")
if trace_id:
safe_set_attribute(span, "litellm.trace_id", str(trace_id))
metadata: Final = as_str_mapping(payload.get("metadata"))
if metadata is None:
return
team_id: Final = metadata.get("user_api_key_team_id")
if team_id:
safe_set_attribute(span, "litellm.team_id", str(team_id))

View file

@ -6,10 +6,13 @@ from typing import TYPE_CHECKING, Any, Final, Optional
from litellm._logging import verbose_logger
from litellm.integrations.arize import _utils
from litellm.integrations.arize._utils import safe_set_attribute
from litellm.integrations.langfuse.langfuse_otel_attributes import (
LangfuseLLMObsOTELAttributes,
)
from litellm.integrations.opentelemetry import OpenTelemetry, OpenTelemetryConfig
from litellm.integrations.otel.model.trace_controls import metadata_bodies
from litellm.integrations.otel.model.utils import as_str, as_str_mapping
from litellm.litellm_core_utils.safe_json_loads import safe_json_loads
from litellm.types.integrations.langfuse_otel import (
LangfuseSpanAttributes,
@ -39,19 +42,20 @@ class LangfuseOtelLogger(OpenTelemetry):
super().__init__(config=config, *args, **kwargs)
@staticmethod
def set_langfuse_otel_attributes(span: Span, kwargs, response_obj):
def set_langfuse_otel_attributes(span: Span, kwargs: dict[str, object], response_obj) -> None:
"""
Sets OpenTelemetry span attributes for Langfuse observability.
Uses the same attribute setting logic as Arize Phoenix for consistency.
"""
_utils.set_attributes(span, kwargs, response_obj, LangfuseLLMObsOTELAttributes)
_utils.set_attributes(span, kwargs, response_obj, LangfuseLLMObsOTELAttributes, emit_session_and_user=False)
span.set_attribute("langfuse.observation.type", "generation")
#########################################################
# Set Langfuse specific attributes
#########################################################
LangfuseOtelLogger._set_langfuse_specific_attributes(span=span, kwargs=kwargs, response_obj=response_obj)
LangfuseOtelLogger._set_trace_user_attribute(span=span, kwargs=kwargs)
@staticmethod
def _extract_langfuse_metadata(kwargs: dict) -> dict:
@ -255,6 +259,31 @@ class LangfuseOtelLogger(OpenTelemetry):
LangfuseOtelLogger._set_observation_output(span=span, response_obj=response_obj)
@staticmethod
def _set_trace_user_attribute(span: Span, kwargs: dict[str, object]) -> None:
slp: Final = as_str_mapping(kwargs.get("standard_logging_object"))
slp_metadata: Final = as_str_mapping(slp.get("metadata")) if slp is not None else None
if slp is None or slp_metadata is None:
return
metadata: Final = as_str_mapping(
LangfuseOtelLogger._extract_langfuse_metadata(kwargs) # pyright: ignore[reportUnknownArgumentType,reportUnknownMemberType] # helper returns a loosely typed dict
)
litellm_params: Final = as_str_mapping(kwargs.get("litellm_params")) or {}
bodies: Final = metadata_bodies(litellm_params)
caller: Final = (as_str(metadata.get("trace_user_id")) if metadata is not None else None) or next(
(value for body in bodies if (value := as_str(body.get("trace_user_id")))), None
)
if caller is not None:
safe_set_attribute(span, LangfuseSpanAttributes.TRACE_USER_ID.value, caller)
return
end_user: Final = (
slp_metadata.get("user_api_key_end_user_id")
or slp.get("end_user")
or next((value for body in bodies if (value := as_str(body.get("user_api_key_end_user_id")))), None)
)
if end_user:
safe_set_attribute(span, LangfuseSpanAttributes.TRACE_USER_ID.value, str(end_user))
@staticmethod
def _get_langfuse_otel_host() -> str | None:
"""

View file

@ -9,7 +9,7 @@ from litellm.integrations.otel.mappers.langfuse import (
LangfuseMapper,
)
from litellm.integrations.otel.model.request_io import request_input, response_output, stream_output
from litellm.integrations.otel.model.trace_controls import caller_trace_controls
from litellm.integrations.otel.model.trace_controls import langfuse_trace_controls
from litellm.integrations.otel.plumbing.context import request_root_span
if TYPE_CHECKING:
@ -24,7 +24,7 @@ class LangfuseOpenTelemetryV2(OpenTelemetryV2):
def log_pre_api_call(self, model: str, messages: object, kwargs: Mapping[str, object]) -> None:
root: Final = request_root_span()
if root is not None and root.is_recording():
root.set_attributes(LangfuseMapper.trace_attributes(caller_trace_controls(kwargs)))
root.set_attributes(LangfuseMapper.trace_attributes(langfuse_trace_controls(kwargs)))
super().log_pre_api_call(model, messages, kwargs)

View file

@ -3,7 +3,7 @@
from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass
from dataclasses import dataclass, replace
from typing import Final
from pydantic import TypeAdapter, ValidationError
@ -33,11 +33,7 @@ def caller_trace_controls(kwargs: Mapping[str, object]) -> TraceControls:
return TraceControls()
proxy_request: Final = as_str_mapping(request.get("proxy_server_request"))
headers: Final = as_str_mapping(proxy_request.get("headers")) if proxy_request is not None else None
bodies: Final = tuple(
metadata
for key in ("metadata", "litellm_metadata")
if (metadata := as_str_mapping(request.get(key))) is not None
)
bodies: Final = metadata_bodies(request)
def scalar(control: str) -> str | None:
from_header: Final = as_str(headers.get(f"{LANGFUSE_HEADER_PREFIX}{control}")) if headers is not None else None
@ -53,6 +49,28 @@ def caller_trace_controls(kwargs: Mapping[str, object]) -> TraceControls:
)
def langfuse_trace_controls(kwargs: Mapping[str, object]) -> TraceControls:
controls: Final = caller_trace_controls(kwargs)
if controls.user_id:
return controls
request: Final = as_str_mapping(kwargs.get("litellm_params"))
if request is None:
return controls
end_user: Final = next(
(value for body in metadata_bodies(request) if (value := as_str(body.get("user_api_key_end_user_id")))),
None,
)
return replace(controls, user_id=end_user)
def metadata_bodies(request: Mapping[str, object]) -> tuple[Mapping[str, object], ...]:
return tuple(
metadata
for key in ("metadata", "litellm_metadata")
if (metadata := as_str_mapping(request.get(key))) is not None
)
def _str_items(value: object) -> tuple[str, ...]:
try:
items: Final = _ITEMS.validate_python(value)

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

@ -0,0 +1,295 @@
import uuid
from pathlib import Path
from typing import Final
from _langfuse_otel import (
_generation_span_attributes,
_langfuse_proxy,
_sink,
_span_attributes_containing_marker,
_trace_user_span_attributes,
_upstream_reply,
)
from integration._support.client import Gateway, eventually
from integration._support.wire import Reply, Request, wire_server
def test_langfuse_otel_header_end_user_lands_in_user_id_not_session_id_for_a_key_owned_by_an_internal_user(
gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = uuid.uuid4().hex
upstream_bodies: Final[list[bytes]] = []
def upstream(request: Request) -> Reply:
upstream_bodies.append(request.body)
return _upstream_reply(marker)
with (
wire_server(upstream) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
assert {key: attributes.get(key) for key in ("user.id", "session.id")} == {
"user.id": f"end-user-{marker}",
"session.id": None,
}, attributes
def test_langfuse_otel_header_end_user_lands_in_user_id_for_a_service_account_key(
gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = uuid.uuid4().hex
upstream_bodies: Final[list[bytes]] = []
def upstream(request: Request) -> Reply:
upstream_bodies.append(request.body)
return _upstream_reply(marker)
with (
wire_server(upstream) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
key: Final = scenario.key()
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
assert {key: attributes.get(key) for key in ("user.id", "session.id")} == {
"user.id": f"end-user-{marker}",
"session.id": None,
}, attributes
def test_langfuse_otel_body_user_is_never_a_session(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
upstream_bodies: Final[list[bytes]] = []
def upstream(request: Request) -> Reply:
upstream_bodies.append(request.body)
return _upstream_reply(marker)
with (
wire_server(upstream) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"user": f"end-user-{marker}",
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
assert {key: attributes.get(key) for key in ("user.id", "session.id")} == {
"user.id": f"end-user-{marker}",
"session.id": None,
}, attributes
def test_langfuse_otel_caller_trace_user_id_wins_over_the_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
upstream_bodies: Final[list[bytes]] = []
def upstream(request: Request) -> Reply:
upstream_bodies.append(request.body)
return _upstream_reply(marker)
with (
wire_server(upstream) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"user": f"end-user-{marker}",
"metadata": {"trace_user_id": f"caller-{marker}"},
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
assert {key: attributes.get(key) for key in ("user.id", "session.id")} == {
"user.id": f"caller-{marker}",
"session.id": None,
}, attributes
def test_langfuse_otel_caller_session_id_stays_the_session_beside_the_end_user(
gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = uuid.uuid4().hex
upstream_bodies: Final[list[bytes]] = []
def upstream(request: Request) -> Reply:
upstream_bodies.append(request.body)
return _upstream_reply(marker)
with (
wire_server(upstream) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"user": f"end-user-{marker}",
"metadata": {"session_id": f"sess-{marker}"},
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
assert {key: attributes.get(key) for key in ("user.id", "session.id")} == {
"user.id": f"end-user-{marker}",
"session.id": f"sess-{marker}",
}, attributes
def test_langfuse_otel_v2_header_end_user_lands_in_user_id(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
upstream_bodies: Final[list[bytes]] = []
def upstream(request: Request) -> Reply:
upstream_bodies.append(request.body)
return _upstream_reply(marker)
with (
wire_server(upstream) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url, {"LITELLM_OTEL_V2": "1"}) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies
batches: Final[list[bytes]] = []
user_spans: Final = eventually(
lambda: _trace_user_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)
assert {key: user_spans[0].get(key) for key in ("user.id", "session.id")} == {
"user.id": f"end-user-{marker}",
"session.id": None,
}, user_spans[0]
def test_langfuse_otel_messages_caller_trace_user_id_under_litellm_metadata_wins_over_the_end_user(
gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = uuid.uuid4().hex
upstream_bodies: Final[list[bytes]] = []
def upstream(request: Request) -> Reply:
upstream_bodies.append(request.body)
return _upstream_reply(marker)
with (
wire_server(upstream) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/messages",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"max_tokens": 5,
"litellm_metadata": {"trace_user_id": f"caller-{marker}"},
"cache": {"no-cache": True},
},
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
assert any(marker.encode() in body for body in upstream_bodies), upstream_bodies
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _span_attributes_containing_marker(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
assert {key: attributes.get(key) for key in ("user.id", "session.id")} == {
"user.id": f"caller-{marker}",
"session.id": None,
}, attributes

View file

@ -0,0 +1,300 @@
import concurrent.futures
import os
import signal
import threading
import time
import uuid
from pathlib import Path
from typing import Final
import httpx
from _langfuse_otel import (
_dedupe_spans,
_drained_spans,
_langfuse_proxy,
_marker_from_body,
_proxy_config,
_sink,
_span_containing_marker,
_upstream_reply_for,
)
from integration._support.client import Gateway, eventually
from integration._support.process import group_members, owned_proxy_process
from integration._support.wire import Reply, Request, Wire, wire_server
def _endpoints(size: int) -> tuple[str, ...]:
return tuple(
"chat" if index < 10 else "chat_stream" if index < 20 else "messages" if index < 25 else "responses"
for index in range(size)
)
def _burst_body(endpoint: str, model: str, marker: str) -> dict[str, object]:
if endpoint == "responses":
return {"model": model, "input": marker, "cache": {"no-cache": True}}
body: Final = {"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}}
if endpoint == "messages":
return {**body, "max_tokens": 5}
if endpoint == "chat_stream":
return {**body, "stream": True}
return body
def _burst_path(endpoint: str) -> str:
if endpoint == "responses":
return "/v1/responses"
if endpoint == "messages":
return "/v1/messages"
return "/v1/chat/completions"
def _send(candidate: Gateway, model: str, endpoint: str, marker: str) -> dict[str, object]:
headers: Final = {
"Authorization": f"Bearer {candidate.key}",
"x-litellm-end-user-id": f"end-user-{marker}",
}
try:
if endpoint == "chat_stream":
with candidate.client.stream(
"POST", "/v1/chat/completions", json=_burst_body(endpoint, model, marker), headers=headers
) as response:
return {"marker": marker, "status": response.status_code, "body": response.read().decode()}
response: Final = candidate.client.request(
"POST", _burst_path(endpoint), json=_burst_body(endpoint, model, marker), headers=headers
)
return {"marker": marker, "status": response.status_code, "body": response.text}
except httpx.HTTPError as error:
return {"marker": marker, "status": 0, "body": repr(error)}
def _send_to_live_gateway(holder: dict[str, Gateway], model: str, endpoint: str, marker: str) -> dict[str, object]:
deadline: Final = time.monotonic() + 60
outcome: Final[dict[str, object]] = {"marker": marker, "status": 0, "body": "no live proxy"}
attempts: Final = {"count": 0}
while time.monotonic() < deadline:
candidate: Final = holder.get("gateway")
if candidate is None:
time.sleep(0.25)
continue
result: Final = _send(candidate, model, endpoint, marker)
attempts["count"] += 1
if result["status"] != 0:
result["attempts"] = attempts["count"]
return result
outcome.update(result)
time.sleep(0.25)
outcome["attempts"] = attempts["count"]
return outcome
def _landed_counts(sink: Wire, batches: list[bytes], markers: tuple[str, ...]) -> dict[str, int]:
spans: Final = _dedupe_spans(_drained_spans(sink, batches))
generations: Final = tuple(span for span in spans if span[1].get("langfuse.observation.type") == "generation")
return {marker: len(_span_containing_marker(generations, marker)) for marker in markers}
def test_langfuse_otel_sink_outage_mid_burst_never_duplicates(gateway: Gateway, tmp_path: Path) -> None:
outage: Final = threading.Event()
def outage_sink(request: Request) -> Reply:
if outage.is_set():
return Reply(status=503)
return _sink(request)
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(outage_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
markers: Final = tuple(uuid.uuid4().hex for _ in range(30))
endpoints: Final = _endpoints(30)
results: Final[dict[str, dict[str, object]]] = {}
window: Final[set[str]] = set()
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
futures: Final = {
executor.submit(_send, candidate, model, endpoints[index], markers[index]): markers[index]
for index in range(30)
}
for future in concurrent.futures.as_completed(futures):
marker: Final = futures[future]
results[marker] = future.result()
if len(results) == 10:
outage.set()
if len(results) == 20:
outage.clear()
if outage.is_set():
window.add(marker)
failures: Final = {m: r for m, r in results.items() if r["status"] != 200}
assert not failures, failures
batches: Final[list[bytes]] = []
counts: Final = eventually(
lambda: _landed_counts(collector, batches, markers),
lambda values: all(value == 1 for value in values.values()),
seconds=60,
return_last_on_timeout=True,
)
missing: Final = {m for m, c in counts.items() if c == 0}
duplicates: Final = {m: c for m, c in counts.items() if c > 1}
assert not duplicates, duplicates
assert missing <= window, (missing, window, counts)
def test_langfuse_otel_slow_sink_never_duplicates_or_loses(gateway: Gateway, tmp_path: Path) -> None:
def slow_sink(request: Request) -> Reply:
time.sleep(1.5)
return _sink(request)
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(slow_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
markers: Final = tuple(uuid.uuid4().hex for _ in range(30))
endpoints: Final = _endpoints(30)
results: Final[dict[str, dict[str, object]]] = {}
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
futures: Final = {
executor.submit(_send, candidate, model, endpoints[index], markers[index]): markers[index]
for index in range(30)
}
for future in concurrent.futures.as_completed(futures):
results[futures[future]] = future.result()
failures: Final = {m: r for m, r in results.items() if r["status"] != 200}
assert not failures, failures
batches: Final[list[bytes]] = []
counts: Final = eventually(
lambda: _landed_counts(collector, batches, markers),
lambda values: all(value == 1 for value in values.values()),
seconds=60,
return_last_on_timeout=True,
)
assert counts == {marker: 1 for marker in markers}, counts
def test_langfuse_otel_worker_kill_mid_burst_keeps_every_marker(gateway: Gateway, tmp_path: Path) -> None:
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
owned_proxy_process(
gateway,
tmp_path,
{
"LANGFUSE_PUBLIC_KEY": "pk-integration",
"LANGFUSE_SECRET_KEY": "sk-integration",
"LANGFUSE_HOST": collector.url,
"OTEL_BSP_SCHEDULE_DELAY": "100",
},
config=_proxy_config(tmp_path, "langfuse_otel.yaml", ("langfuse_otel",)),
workers=2,
) as owned,
owned.gateway.scenario() as scenario,
):
candidate: Final = owned.gateway
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
markers: Final = tuple(uuid.uuid4().hex for _ in range(30))
endpoints: Final = _endpoints(30)
results: Final[dict[str, dict[str, object]]] = {}
killed: Final = {"done": False}
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
futures: Final = {
executor.submit(_send, candidate, model, endpoints[index], markers[index]): markers[index]
for index in range(30)
}
for future in concurrent.futures.as_completed(futures):
results[futures[future]] = future.result()
if len(results) >= 10 and not killed["done"]:
children: Final = tuple(
process for process in group_members(owned.process.pid) if process.pid != owned.process.pid
)
assert children, "no uvicorn worker child found"
os.kill(children[0].pid, signal.SIGKILL)
killed["done"] = True
assert killed["done"], "worker kill never fired"
accepted: Final = {m for m, r in results.items() if r["status"] == 200}
failures: Final = {m: r for m, r in results.items() if r["status"] != 200}
survivor: Final = _send(candidate, model, "chat", uuid.uuid4().hex)
assert survivor["status"] == 200, survivor
accepted_markers: Final = tuple(sorted(accepted))
batches: Final[list[bytes]] = []
counts: Final = eventually(
lambda: _landed_counts(collector, batches, accepted_markers),
lambda values: all(value == 1 for value in values.values()),
seconds=60,
return_last_on_timeout=True,
)
assert counts == {marker: 1 for marker in accepted_markers}, (counts, failures)
def test_langfuse_otel_proxy_restart_mid_burst_keeps_every_marker(gateway: Gateway, tmp_path: Path) -> None:
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
):
holder: Final[dict[str, Gateway]] = {}
markers: Final = tuple(uuid.uuid4().hex for _ in range(30))
endpoints: Final = _endpoints(30)
results: Final[dict[str, dict[str, object]]] = {}
restarted: Final = {"done": False}
one: Final = tmp_path / "one"
two: Final = tmp_path / "two"
one.mkdir()
two.mkdir()
context_one: Final = _langfuse_proxy(gateway, one, collector.url)
context_two: Final = _langfuse_proxy(gateway, two, collector.url)
candidate_one: Final = context_one.__enter__()
candidate_two: Final = {"gateway": None}
model_name: Final = uuid.uuid4().hex
created: Final = candidate_one.post(
"/model/new",
{
"model_name": model_name,
"litellm_params": {
"model": "openai/gpt-4o-mini",
"api_key": "integration-provider-key",
"api_base": provider.url + "/v1",
},
},
)
model_id: Final = str(created["model_info"]["id"])
holder["gateway"] = candidate_one
try:
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
futures: Final = {
executor.submit(
_send_to_live_gateway, holder, model_name, endpoints[index], markers[index]
): markers[index]
for index in range(30)
}
for future in concurrent.futures.as_completed(futures):
results[futures[future]] = future.result()
if len(results) < 8 or restarted["done"]:
continue
restarted["done"] = True
holder.pop("gateway")
context_one.__exit__(None, None, None)
candidate_two["gateway"] = context_two.__enter__()
holder["gateway"] = candidate_two["gateway"]
finally:
if candidate_two["gateway"] is not None:
candidate_two["gateway"].post("/model/delete", {"id": model_id})
context_two.__exit__(None, None, None)
context_one.__exit__(None, None, None)
assert restarted["done"], "restart never fired"
accepted: Final = {m for m, r in results.items() if r["status"] == 200}
failures: Final = {m: r for m, r in results.items() if r["status"] != 200}
accepted_markers: Final = tuple(sorted(accepted))
batches: Final[list[bytes]] = []
counts: Final = eventually(
lambda: _landed_counts(collector, batches, accepted_markers),
lambda values: all(value >= 1 for value in values.values()),
seconds=60,
return_last_on_timeout=True,
)
missing: Final = {m for m, c in counts.items() if c == 0}
over_counted: Final = {m: c for m, c in counts.items() if c > int(results[m].get("attempts", 1))}
assert not missing and not over_counted, (missing, over_counted, failures)

View file

@ -0,0 +1,425 @@
import threading
import uuid
from pathlib import Path
from typing import Final
import httpx
from _langfuse_otel import (
_drained_spans,
_generation_marker_span_attributes,
_generation_span_attributes,
_langfuse_proxy,
_marker_from_body,
_sink,
_span_containing_marker,
_upstream_reply_for,
)
from integration._support.client import Gateway, eventually
from integration._support.wire import Reply, Request, Wire, wire_server
def _identity(attributes: dict[str, object]) -> dict[str, object]:
return {key: attributes.get(key) for key in ("user.id", "session.id")}
def _await_generation(collector: Wire, batches: list[bytes], marker: str) -> dict[str, object]:
return eventually(
lambda: _generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
def _end_user_request(
candidate: Gateway, model: str, marker: str, key: str | None = None, **extra: object
) -> httpx.Response:
return candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"cache": {"no-cache": True},
**extra,
},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
def _subsequent_request_still_lands(candidate: Gateway, collector: Wire, batches: list[bytes], model: str) -> None:
next_marker: Final = uuid.uuid4().hex
response: Final = _end_user_request(candidate, model, next_marker)
assert response.status_code == 200, response.text
assert _await_generation(collector, batches, next_marker) is not None
def test_langfuse_otel_five_kb_end_user_header_lands_in_user_id(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
end_user: Final = "e" * 5120
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
headers={"x-litellm-end-user-id": end_user},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": end_user, "session.id": None}, attributes
def test_langfuse_otel_empty_end_user_header_writes_no_identity(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
headers={"x-litellm-end-user-id": ""},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": None, "session.id": None}, attributes
_subsequent_request_still_lands(candidate, collector, batches, model)
def test_langfuse_otel_integer_body_user_does_not_crash(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"user": 123,
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert attributes.get("session.id") is None, attributes
_subsequent_request_still_lands(candidate, collector, batches, model)
def test_langfuse_otel_list_body_user_does_not_crash(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"user": ["a", "b"],
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert attributes.get("session.id") is None, attributes
_subsequent_request_still_lands(candidate, collector, batches, model)
def test_langfuse_otel_integer_trace_user_id_does_not_crash(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = _end_user_request(candidate, model, marker, metadata={"trace_user_id": 123})
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert attributes.get("user.id") in ("123", f"end-user-{marker}"), attributes
assert attributes.get("session.id") is None, attributes
_subsequent_request_still_lands(candidate, collector, batches, model)
def test_langfuse_otel_null_metadata_still_maps_the_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = _end_user_request(candidate, model, marker, metadata=None)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
def test_langfuse_otel_duplicate_end_user_header_uses_the_first(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.client.request(
"POST",
"/v1/chat/completions",
json={
"model": model,
"messages": [{"role": "user", "content": f"first-call-{marker}"}],
"cache": {"no-cache": True},
},
headers={
"Authorization": f"Bearer {candidate.key}",
"x-litellm-end-user-id": f"first-{marker}",
},
)
assert response.status_code == 200, response.text
duplicate: Final = candidate.client.request(
"POST",
"/v1/chat/completions",
json={
"model": model,
"messages": [{"role": "user", "content": marker}],
"cache": {"no-cache": True},
},
headers=[
("Authorization", f"Bearer {candidate.key}"),
("x-litellm-end-user-id", f"first-{marker}"),
("x-litellm-end-user-id", f"second-{marker}"),
],
)
assert duplicate.status_code == 200, duplicate.text
batches: Final[list[bytes]] = []
spans: Final = eventually(
lambda: _generation_marker_span_attributes(collector, batches, marker),
lambda found: len(found) == 2,
seconds=30,
)
assert all(attributes.get("user.id") == f"first-{marker}" for attributes in spans), spans
def test_langfuse_otel_bad_key_emits_no_generation_span(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key="sk-not-a-real-key",
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 401, response.text
batches: Final[list[bytes]] = []
_subsequent_request_still_lands(candidate, collector, batches, model)
offending: Final = tuple(
attributes
for attributes in _generation_marker_span_attributes(collector, batches, marker)
if attributes.get("user.id") not in (None, f"end-user-{marker}")
)
assert offending == (), offending
def test_langfuse_otel_upstream_failure_never_invents_an_identity(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(
lambda request: (
Reply(status=401, body=b'{"error": "upstream denied ' + marker.encode() + b'"}')
if marker.encode() in (request.body or b"")
else _upstream_reply_for(request, _marker_from_body(request))
)
) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = _end_user_request(candidate, model, marker)
assert response.status_code == 401, response.text
batches: Final[list[bytes]] = []
_subsequent_request_still_lands(candidate, collector, batches, model)
offending: Final = tuple(
attributes
for _t, attributes in _span_containing_marker(_drained_spans(collector, batches), marker)
if attributes.get("user.id") not in (None, f"end-user-{marker}")
)
assert offending == (), offending
def test_langfuse_otel_sink_rejections_do_not_drop_the_proxy(gateway: Gateway, tmp_path: Path) -> None:
calls: Final[dict[str, int]] = {"count": 0}
lock: Final = threading.Lock()
def rejecting_sink(request: Request) -> Reply:
with lock:
calls["count"] += 1
seen: Final = calls["count"]
if seen == 1:
return Reply(status=403)
if seen == 2:
return Reply(status=404)
return _sink(request)
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(rejecting_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
markers: Final = tuple(uuid.uuid4().hex for _ in range(3))
batches: Final[list[bytes]] = []
for index, item in enumerate(markers):
response: Final = _end_user_request(candidate, model, item)
assert response.status_code == 200, response.text
if index < 2:
eventually(
lambda collector=collector, batches=batches: (
batches.extend(request.body for request in collector.drain()) or len(batches)
),
lambda seen, index=index: seen >= index + 1,
seconds=30,
)
attributes: Final = _await_generation(collector, batches, markers[2])
assert attributes.get("user.id") == f"end-user-{markers[2]}", attributes
_subsequent_request_still_lands(candidate, collector, batches, model)
def test_langfuse_otel_unknown_model_error_reaches_the_caller(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = _end_user_request(candidate, "no-such-model-" + marker[:12], marker)
assert response.status_code in (400, 404), response.text
batches: Final[list[bytes]] = []
_subsequent_request_still_lands(candidate, collector, batches, model)
def test_langfuse_otel_empty_and_null_session_id_stay_absent(gateway: Gateway, tmp_path: Path) -> None:
markers: Final = tuple(uuid.uuid4().hex for _ in range(3))
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
batches: Final[list[bytes]] = []
sessions: Final = ({"session_id": ""}, {"session_id": None}, {})
for marker, metadata in zip(markers, sessions):
response: Final = _end_user_request(candidate, model, marker, metadata=metadata)
assert response.status_code == 200, response.text
for marker, metadata in zip(markers, sessions):
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {
"user.id": f"end-user-{marker}",
"session.id": metadata.get("session_id"),
}, attributes
def test_langfuse_otel_empty_trace_user_id_falls_back_to_the_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = _end_user_request(candidate, model, marker, metadata={"trace_user_id": ""})
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
def test_langfuse_otel_three_identical_requests_each_land_the_end_user(gateway: Gateway, tmp_path: Path) -> None:
markers: Final = tuple(uuid.uuid4().hex for _ in range(3))
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
batches: Final[list[bytes]] = []
for marker in markers:
response: Final = _end_user_request(candidate, model, marker)
assert response.status_code == 200, response.text
for marker in markers:
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
def test_langfuse_otel_litellm_metadata_session_id_on_responses(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, _marker_from_body(request))) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/responses",
{
"model": model,
"input": marker,
"litellm_metadata": {"session_id": f"sess-{marker}"},
"cache": {"no-cache": True},
},
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _generation_marker_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": f"sess-{marker}"}, attributes

View file

@ -0,0 +1,512 @@
import uuid
from pathlib import Path
from typing import Final
from _langfuse_otel import (
_arize_generation_span_attributes,
_generation_marker_span_attributes,
_generation_span_attributes,
_langfuse_proxy,
_observability_proxy,
_sink,
_trace_user_span_attributes,
_upstream_reply_for,
)
from anthropic import Anthropic, AsyncAnthropic
from integration._support.client import Gateway, eventually
from integration._support.database import read_rows
from integration._support.wire import Wire, wire_server
from openai import AsyncOpenAI, OpenAI
def _openai_client(gateway: Gateway, key: str, marker: str) -> OpenAI:
return OpenAI(
api_key=key,
base_url=str(gateway.client.base_url).rstrip("/") + "/v1",
timeout=30,
max_retries=0,
default_headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
def _async_openai_client(gateway: Gateway, key: str, marker: str) -> AsyncOpenAI:
return AsyncOpenAI(
api_key=key,
base_url=str(gateway.client.base_url).rstrip("/") + "/v1",
timeout=30,
max_retries=0,
default_headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
def _anthropic_client(gateway: Gateway, key: str, marker: str) -> Anthropic:
return Anthropic(
api_key=key,
base_url=str(gateway.client.base_url).rstrip("/"),
timeout=30,
max_retries=0,
default_headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
def _async_anthropic_client(gateway: Gateway, key: str, marker: str) -> AsyncAnthropic:
return AsyncAnthropic(
api_key=key,
base_url=str(gateway.client.base_url).rstrip("/"),
timeout=30,
max_retries=0,
default_headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
def _await_generation(collector: Wire, batches: list[bytes], marker: str) -> dict[str, object]:
return eventually(
lambda: _generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
def _await_marker_span(collector: Wire, batches: list[bytes], marker: str) -> dict[str, object]:
return eventually(
lambda: _generation_marker_span_attributes(collector, batches, marker),
lambda spans: len(spans) == 1,
seconds=30,
)[0]
def _identity(attributes: dict[str, object]) -> dict[str, object]:
return {key: attributes.get(key) for key in ("user.id", "session.id")}
def test_langfuse_otel_customer_id_header_lands_in_user_id(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
headers={"x-litellm-customer-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
def test_langfuse_otel_langfuse_trace_user_id_header_wins_over_the_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
headers={
"x-litellm-end-user-id": f"end-user-{marker}",
"langfuse_trace_user_id": f"caller-{marker}",
},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"caller-{marker}", "session.id": None}, attributes
def test_langfuse_otel_no_end_user_never_exposes_the_internal_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": None, "session.id": None}, attributes
def test_langfuse_otel_team_key_end_user_keeps_the_litellm_attributes(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
team: Final = scenario.team(team_alias=f"team-alias-{marker[:12]}")
key: Final = scenario.key(team_id=team, key_alias=f"key-alias-{marker[:12]}")
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
assert attributes.get("litellm.team_id") == team, attributes
assert attributes.get("litellm.team_alias") == f"team-alias-{marker[:12]}", attributes
assert attributes.get("litellm.key_alias") == f"key-alias-{marker[:12]}", attributes
def test_langfuse_otel_header_end_user_beats_the_body_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"user": f"body-user-{marker}",
"cache": {"no-cache": True},
},
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
rows: Final = eventually(
lambda: read_rows(
'SELECT end_user FROM "LiteLLM_SpendLogs" WHERE request_id = %s', (str(response.json()["id"]),)
),
lambda values: len(values) == 1,
seconds=70,
)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
assert attributes.get("user.id") == rows[0]["end_user"], (attributes, rows)
def test_langfuse_otel_openai_sdk_streaming_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
with _openai_client(candidate, key, marker) as client:
stream: Final = client.chat.completions.create(
model=model,
messages=[{"role": "user", "content": marker}],
stream=True,
extra_body={"cache": {"no-cache": True}},
)
with stream:
chunks: Final = tuple(stream)
assert {chunk.id for chunk in chunks} == {marker}
assert f"reply {marker}" in "".join(
choice.delta.content or "" for chunk in chunks for choice in chunk.choices
)
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
async def test_langfuse_otel_openai_async_sdk_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
async with _async_openai_client(candidate, key, marker) as client:
completion: Final = await client.chat.completions.create(
model=model,
messages=[{"role": "user", "content": marker}],
extra_body={"cache": {"no-cache": True}},
)
assert completion.id == marker
batches: Final[list[bytes]] = []
attributes: Final = _await_generation(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
def test_langfuse_otel_responses_api_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
with _openai_client(candidate, key, marker) as client:
completion: Final = client.responses.create(
model=model, input=marker, extra_body={"cache": {"no-cache": True}}
)
assert marker in completion.output_text, completion.output_text
batches: Final[list[bytes]] = []
attributes: Final = _await_marker_span(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
async def test_langfuse_otel_responses_api_async_streaming_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
async with _async_openai_client(candidate, key, marker) as client:
stream: Final = await client.responses.create(
model=model, input=marker, stream=True, extra_body={"cache": {"no-cache": True}}
)
events: Final = tuple([event async for event in stream])
assert events, events
batches: Final[list[bytes]] = []
attributes: Final = _await_marker_span(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
def test_langfuse_otel_anthropic_sdk_messages_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
with _anthropic_client(candidate, key, marker) as client:
message: Final = client.messages.create(
model=model, max_tokens=5, messages=[{"role": "user", "content": marker}]
)
assert marker in message.content[0].text, message.model_dump_json()
batches: Final[list[bytes]] = []
attributes: Final = _await_marker_span(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
async def test_langfuse_otel_anthropic_async_streaming_messages_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
async with _async_anthropic_client(candidate, key, marker) as client:
stream: Final = await client.messages.create(
model=model, max_tokens=5, messages=[{"role": "user", "content": marker}], stream=True
)
events: Final = tuple([event async for event in stream])
assert events, events
batches: Final[list[bytes]] = []
attributes: Final = _await_marker_span(collector, batches, marker)
assert _identity(attributes) == {"user.id": f"end-user-{marker}", "session.id": None}, attributes
def test_langfuse_otel_v2_caller_trace_user_id_wins_over_the_end_user(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_langfuse_proxy(gateway, tmp_path, collector.url, {"LITELLM_OTEL_V2": "1"}) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"metadata": {"trace_user_id": f"caller-{marker}"},
"cache": {"no-cache": True},
},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
user_spans: Final = eventually(
lambda: _trace_user_span_attributes(collector, batches, marker),
lambda spans: len(spans) >= 1,
seconds=30,
)
assert all(
_identity(attributes) == {"user.id": f"caller-{marker}", "session.id": None} for attributes in user_spans
), user_spans
def test_arize_phoenix_header_end_user_keeps_session_mapping(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_observability_proxy(
gateway,
tmp_path,
{
"PHOENIX_COLLECTOR_ENDPOINT": collector.url + "/v1/traces",
"PHOENIX_API_KEY": "phoenix-integration",
},
callbacks=("arize_phoenix",),
config_name="arize_phoenix.yaml",
) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _arize_generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) >= 1,
seconds=30,
)[0]
assert _identity(attributes) == {"user.id": owner, "session.id": f"end-user-{marker}"}, attributes
assert attributes.get("litellm.trace_id") is not None, attributes
def test_arize_phoenix_caller_session_id_stays_session(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as collector,
_observability_proxy(
gateway,
tmp_path,
{
"PHOENIX_COLLECTOR_ENDPOINT": collector.url + "/v1/traces",
"PHOENIX_API_KEY": "phoenix-integration",
},
callbacks=("arize_phoenix",),
config_name="arize_phoenix.yaml",
) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": marker}],
"user": f"end-user-{marker}",
"metadata": {"session_id": f"sess-{marker}"},
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
batches: Final[list[bytes]] = []
attributes: Final = eventually(
lambda: _arize_generation_span_attributes(collector, batches, marker),
lambda spans: len(spans) >= 1,
seconds=30,
)[0]
assert _identity(attributes) == {
"user.id": f"end-user-{marker}",
"session.id": f"end-user-{marker}",
}, attributes
assert attributes.get("litellm.trace_id") == f"sess-{marker}", attributes
def test_langfuse_otel_and_arize_phoenix_together_keep_each_mapping(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = uuid.uuid4().hex
with (
wire_server(lambda request: _upstream_reply_for(request, marker)) as provider,
wire_server(_sink) as langfuse_collector,
wire_server(_sink) as phoenix_collector,
_observability_proxy(
gateway,
tmp_path,
{
"LANGFUSE_PUBLIC_KEY": "pk-integration",
"LANGFUSE_SECRET_KEY": "sk-integration",
"LANGFUSE_HOST": langfuse_collector.url,
"PHOENIX_COLLECTOR_ENDPOINT": phoenix_collector.url + "/v1/traces",
"PHOENIX_API_KEY": "phoenix-integration",
},
callbacks=("langfuse_otel", "arize_phoenix"),
config_name="dual_callbacks.yaml",
) as candidate,
candidate.scenario() as scenario,
):
owner: Final = scenario.user()
key: Final = scenario.key(user_id=owner)
model: Final = scenario.model(model="openai/gpt-4o-mini", api_base=provider.url + "/v1")
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
headers={"x-litellm-end-user-id": f"end-user-{marker}"},
)
assert response.status_code == 200, response.text
langfuse_batches: Final[list[bytes]] = []
phoenix_batches: Final[list[bytes]] = []
langfuse_attributes: Final = _await_generation(langfuse_collector, langfuse_batches, marker)
phoenix_attributes: Final = eventually(
lambda: _arize_generation_span_attributes(phoenix_collector, phoenix_batches, marker),
lambda spans: len(spans) >= 1,
seconds=30,
)[0]
assert _identity(langfuse_attributes) == {
"user.id": f"end-user-{marker}",
"session.id": None,
}, langfuse_attributes
assert _identity(phoenix_attributes) == {
"user.id": owner,
"session.id": f"end-user-{marker}",
}, phoenix_attributes

View file

@ -1513,3 +1513,32 @@ def test_arize_mcp_emitter_is_inert_without_a_standard_logging_object():
written = {c.args[0]: c.args[1] for c in span.set_attribute.call_args_list}
assert SpanAttributes.TOOL_NAME not in written
def test_arize_session_and_user_attrs_still_emit_from_key_metadata_by_default():
from unittest.mock import MagicMock
span = MagicMock()
kwargs = {
"model": "gpt-4o",
"messages": [{"role": "user", "content": "hello"}],
"standard_logging_object": {
"call_type": "acompletion",
"model_parameters": {},
"metadata": {
"user_api_key_end_user_id": "end-1",
"user_api_key_user_id": "internal-1",
"user_api_key_team_id": "team-1",
},
"trace_id": "trace-1",
},
"optional_params": {},
"litellm_params": {"custom_llm_provider": "openai"},
}
ArizeLogger.set_arize_attributes(span, kwargs, {"id": "chatcmpl-1", "choices": [], "usage": {}})
span.set_attribute.assert_any_call(SpanAttributes.SESSION_ID, "end-1")
span.set_attribute.assert_any_call(SpanAttributes.USER_ID, "internal-1")
span.set_attribute.assert_any_call("litellm.trace_id", "trace-1")
span.set_attribute.assert_any_call("litellm.team_id", "team-1")

View file

@ -419,6 +419,31 @@ def test_langfuse_user_and_session_headers_beat_body_metadata_on_both_spans():
assert attrs["session.id"] == "from-header-s"
def test_the_proxy_end_user_fills_user_id_when_the_caller_names_no_trace_user():
logger, exporter = _logger()
root_attrs, generation_attrs = _run_named_request(
logger, exporter, {"metadata": {"user_api_key_end_user_id": "end-1"}, "proxy_server_request": {"headers": {}}}
)
assert root_attrs["user.id"] == "end-1"
def test_a_callers_trace_user_id_still_wins_over_the_proxy_end_user():
logger, exporter = _logger()
root_attrs, generation_attrs = _run_named_request(
logger,
exporter,
{
"metadata": {"trace_user_id": "caller-1", "user_api_key_end_user_id": "end-1"},
"proxy_server_request": {"headers": {}},
},
)
assert root_attrs["user.id"] == "caller-1"
def test_caller_metadata_cannot_override_the_proxy_team_identity():
logger, exporter = _logger()
response: Final = ModelResponse(choices=[Choices(message=Message(role="assistant", content="pong"))])

View file

@ -46,7 +46,11 @@ from litellm.integrations.otel.model.spans import (
root_roles,
validate_registry,
)
from litellm.integrations.otel.model.trace_controls import TraceControls, caller_trace_controls
from litellm.integrations.otel.model.trace_controls import (
TraceControls,
caller_trace_controls,
langfuse_trace_controls,
)
@pytest.fixture(autouse=True)
@ -1342,6 +1346,43 @@ def test_caller_trace_controls_carry_user_session_and_tags(request_data, expecte
assert LLMCallEvent.from_dict({"litellm_params": request_data}).trace == expected
@pytest.mark.parametrize(
("request_data", "expected"),
[
({"metadata": {"user_api_key_end_user_id": "end-1"}}, TraceControls(user_id="end-1")),
({"litellm_metadata": {"user_api_key_end_user_id": "end-2"}}, TraceControls(user_id="end-2")),
(
{"metadata": {"trace_user_id": "caller-1", "user_api_key_end_user_id": "end-1"}},
TraceControls(user_id="caller-1"),
),
(
{
"proxy_server_request": {"headers": {"langfuse_trace_user_id": "header-1"}},
"metadata": {"user_api_key_end_user_id": "end-1"},
},
TraceControls(user_id="header-1"),
),
(
{"metadata": {"user_api_key_end_user_id": "end-1", "session_id": "s-1"}},
TraceControls(user_id="end-1", session_id="s-1"),
),
({"metadata": {"user_api_key_user_id": "internal-1"}}, TraceControls()),
({}, TraceControls()),
],
ids=[
"body-end-user",
"anthropic-end-user",
"body-caller-wins",
"header-caller-wins",
"session-kept",
"internal-user-ignored",
"empty",
],
)
def test_langfuse_trace_controls_fall_back_to_the_proxy_end_user(request_data, expected):
assert langfuse_trace_controls({"litellm_params": request_data}) == expected
def test_llm_span_data_carries_the_caller_trace_controls():
controls: Final = TraceControls(name="nightly-eval", user_id="u1", session_id="s1", tags=("a", "b"))
data: Final = LLMCallSpanData.from_standard_logging_payload(_sample_payload(), trace=controls)

View file

@ -112,7 +112,7 @@ class TestLangfuseOtelIntegration:
)
mock_set_attributes.assert_called_once_with(
mock_span, mock_kwargs, mock_response, LangfuseLLMObsOTELAttributes
mock_span, mock_kwargs, mock_response, LangfuseLLMObsOTELAttributes, emit_session_and_user=False
)
mock_span.set_attribute.assert_any_call(
"langfuse.observation.type", "generation"
@ -709,7 +709,7 @@ class TestLangfuseOtelResponsesAPI:
# Verify that set_attributes was called for general attributes
mock_set_attributes.assert_called_once_with(
mock_span, kwargs, mock_response, LangfuseLLMObsOTELAttributes
mock_span, kwargs, mock_response, LangfuseLLMObsOTELAttributes, emit_session_and_user=False
)
# Verify that Langfuse-specific attributes were set
@ -1006,5 +1006,115 @@ class TestLangfuseOtelResponsesAPI:
assert output_data[0]["arguments"] == {}
class TestLangfuseOtelTraceIdentity:
def _recording_span(self):
from opentelemetry.sdk.trace import TracerProvider
return TracerProvider().get_tracer("test").start_span("generation")
def _kwargs(self, slp_metadata=None, litellm_metadata=None, model_parameters=None, slp_extra=None):
return {
"model": "gpt-4o",
"messages": [{"role": "user", "content": "hello"}],
"optional_params": {},
"litellm_params": {"metadata": litellm_metadata or {}, "custom_llm_provider": "openai"},
"standard_logging_object": {
"call_type": "acompletion",
"model_parameters": model_parameters or {},
"metadata": slp_metadata or {},
**(slp_extra or {}),
},
}
def _response_obj(self):
return {
"id": "chatcmpl-1",
"model": "gpt-4o",
"choices": [{"index": 0, "message": {"role": "assistant", "content": "hi"}, "finish_reason": "stop"}],
"usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2},
}
def _identity(self, kwargs):
span = self._recording_span()
LangfuseOtelLogger.set_langfuse_otel_attributes(span, kwargs, self._response_obj())
attributes = dict(span.attributes or {})
return {key: attributes.get(key) for key in ("user.id", "session.id")}, attributes
def test_header_end_user_beats_internal_key_owner_in_user_id(self):
identity, _ = self._identity(
self._kwargs(
slp_metadata={
"user_api_key_end_user_id": "end-1",
"user_api_key_user_id": "internal-1",
}
)
)
assert identity == {"user.id": "end-1", "session.id": None}
def test_header_end_user_lands_in_user_id_for_a_service_key(self):
identity, _ = self._identity(self._kwargs(slp_metadata={"user_api_key_end_user_id": "end-1"}))
assert identity == {"user.id": "end-1", "session.id": None}
def test_body_user_is_never_a_session(self):
identity, _ = self._identity(
self._kwargs(
slp_metadata={"user_api_key_end_user_id": "body-user"},
model_parameters={"user": "body-user"},
)
)
assert identity == {"user.id": "body-user", "session.id": None}
def test_caller_trace_user_id_wins_over_the_end_user(self):
identity, _ = self._identity(
self._kwargs(
slp_metadata={"user_api_key_end_user_id": "end-1"},
litellm_metadata={"trace_user_id": "caller-1"},
)
)
assert identity == {"user.id": "caller-1", "session.id": None}
def test_caller_trace_user_id_under_litellm_metadata_wins_over_the_end_user(self):
kwargs = self._kwargs(slp_metadata={"user_api_key_end_user_id": "end-1"})
kwargs["litellm_params"]["litellm_metadata"] = {"trace_user_id": "caller-1"}
identity, _ = self._identity(kwargs)
assert identity == {"user.id": "caller-1", "session.id": None}
def test_end_user_only_under_litellm_metadata_lands_in_user_id(self):
kwargs = self._kwargs()
kwargs["litellm_params"]["litellm_metadata"] = {"user_api_key_end_user_id": "end-1"}
identity, _ = self._identity(kwargs)
assert identity == {"user.id": "end-1", "session.id": None}
def test_caller_session_id_stays_the_session_beside_the_end_user(self):
identity, _ = self._identity(
self._kwargs(
slp_metadata={"user_api_key_end_user_id": "end-1"},
litellm_metadata={"session_id": "sess-1"},
)
)
assert identity == {"user.id": "end-1", "session.id": "sess-1"}
def test_internal_user_without_an_end_user_never_lands_in_user_id(self):
identity, _ = self._identity(self._kwargs(slp_metadata={"user_api_key_user_id": "internal-1"}))
assert identity == {"user.id": None, "session.id": None}
def test_request_context_attributes_still_emit(self):
_, attributes = self._identity(
self._kwargs(
slp_metadata={
"user_api_key_end_user_id": "end-1",
"user_api_key_team_id": "team-1",
"user_api_key_team_alias": "team-alias",
"user_api_key_alias": "key-alias",
},
slp_extra={"trace_id": "trace-1"},
)
)
assert attributes["litellm.trace_id"] == "trace-1"
assert attributes["litellm.team_id"] == "team-1"
assert attributes["litellm.team_alias"] == "team-alias"
assert attributes["litellm.key_alias"] == "key-alias"
if __name__ == "__main__":
pytest.main([__file__])