This commit is contained in:
devin-ai-integration[bot] 2026-09-30 22:31:59 +00:00 • committed by GitHub
commit d2f548a521
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 4211 additions and 49 deletions

View file

@ -563,6 +563,7 @@ class OpenTelemetryV2(CustomLogger):
request_route=request_root_http_route(),
trace=call.trace,
session_id=call.session_id,
metadata_keys=tuple(self.config.baggage_metadata_keys),
)
end_time_ns: Final = to_ns(end_time)
if carrier is not None and carrier.span is not None:

View file

@ -7,7 +7,7 @@ Phoenix + any other OpenInference-aware backend simultaneously.
"""
import json
from collections.abc import Callable, Mapping, Sequence
from collections.abc import Callable, Iterator, Mapping, Sequence
from itertools import accumulate, chain, groupby
from types import MappingProxyType
from typing import Final
@ -15,10 +15,12 @@ from typing import Final
from litellm.integrations.otel.mappers.base import AttributeMap, AttrValue, SpanData
from litellm.integrations.otel.mappers.utils import (
MAX_TOOL_DEFINITION_ATTRS_PER_SPAN,
MessageToolCall,
collect,
drop_none,
drop_none_pairs,
json_if,
message_content,
message_tool_calls,
output_messages,
tool_definition_attrs,
)
@ -31,37 +33,106 @@ from litellm.integrations.otel.model.payloads import (
_INPUT_MESSAGES: Final = "llm.input_messages"
_OUTPUT_MESSAGES: Final = "llm.output_messages"
_MESSAGE_FAMILIES: Final = (_INPUT_MESSAGES, _OUTPUT_MESSAGES)
_MESSAGE_BASE: Final = -1
_ParsedMessage = tuple[object, str | None, tuple[MessageToolCall, ...]]
def _message_key_groups(attrs: Mapping[str, AttrValue]) -> Mapping[tuple[str, int], tuple[str, ...]]:
"""Per-index message keys in ``attrs`` grouped by ``(family, index)``."""
tagged: Final = sorted(
(family, int(key.split(".")[2]), key)
for key in attrs
for family in _MESSAGE_FAMILIES
if key.startswith(f"{family}.")
def _parse_message(message: object) -> _ParsedMessage:
role: Final = message.get("role") if isinstance(message, dict) else None
return role, message_content(message), message_tool_calls(message)
def _tool_call_attribute_pairs(
prefix: str, idx: int, tool_calls: tuple[MessageToolCall, ...]
) -> Iterator[tuple[str, str | None]]:
for tool_idx, tool_call in enumerate(tool_calls):
yield f"{prefix}.{idx}.message.tool_calls.{tool_idx}.tool_call.id", tool_call.id
yield f"{prefix}.{idx}.message.tool_calls.{tool_idx}.tool_call.function.name", tool_call.name
yield f"{prefix}.{idx}.message.tool_calls.{tool_idx}.tool_call.function.arguments", tool_call.arguments
def _message_attribute_pairs(
prefix: str,
messages: Sequence[_ParsedMessage],
*,
with_tool_call_attrs: bool,
) -> Iterator[tuple[str, str | None]]:
for idx, (role, content, tool_calls) in enumerate(messages):
yield f"{prefix}.{idx}.message.role", role if isinstance(role, str) else None
yield f"{prefix}.{idx}.message.content", content
if with_tool_call_attrs:
yield from _tool_call_attribute_pairs(prefix, idx, tool_calls)
def _message_value(messages: Sequence[_ParsedMessage]) -> str:
return json.dumps(
[
{
"role": role,
"content": content,
**({"tool_calls": [tool_call.to_openai_dict() for tool_call in tool_calls]} if tool_calls else {}),
}
for role, content, tool_calls in messages
]
)
def _message_key_group(key: str) -> tuple[str, int, int, str] | None:
family: Final = next(
(family for family in _MESSAGE_FAMILIES if key.startswith(f"{family}.")),
None,
)
if family is None:
return None
parts: Final = key.split(".")
message_idx: Final = int(parts[2])
tool_idx: Final = int(parts[5]) if parts[4] == "tool_calls" else _MESSAGE_BASE
return family, message_idx, tool_idx, key
def _message_key_groups(attrs: Mapping[str, AttrValue]) -> Mapping[tuple[str, int, int], tuple[str, ...]]:
"""Message and tool-call keys in ``attrs`` grouped by family, message index, and tool index."""
tagged: Final = tuple(tag for key in attrs if (tag := _message_key_group(key)) is not None)
return MappingProxyType(
{group: tuple(key for _, _, key in keys) for group, keys in groupby(tagged, key=lambda tag: tag[:2])}
{group: tuple(key for _, _, _, key in keys) for group, keys in groupby(sorted(tagged), key=lambda tag: tag[:3])}
)
def _shed_order(groups: Mapping[tuple[str, int], tuple[str, ...]]) -> tuple[tuple[str, int], ...]:
"""Message groups least valuable first: middle prompt turns, extra choices, then the opener, the newest turn
and the first choice."""
inputs: Final = sorted(idx for family, idx in groups if family == _INPUT_MESSAGES)
outputs: Final = sorted(idx for family, idx in groups if family == _OUTPUT_MESSAGES)
def _message_shed_groups(
groups: Mapping[tuple[str, int, int], tuple[str, ...]], family: str, message_idx: int
) -> Iterator[tuple[str, int, int]]:
tool_call_groups: Final = tuple(
sorted(
(group for group in groups if group[:2] == (family, message_idx) and group[2] != _MESSAGE_BASE),
key=lambda group: group[2],
reverse=True,
)
)
yield from tool_call_groups
base_group: Final = (family, message_idx, _MESSAGE_BASE)
if base_group in groups:
yield base_group
def _shed_order(groups: Mapping[tuple[str, int, int], tuple[str, ...]]) -> tuple[tuple[str, int, int], ...]:
"""Middle inputs, extra choices, pinned inputs, then the first choice, with tool calls before message keys."""
inputs: Final = sorted(frozenset(idx for family, idx, _ in groups if family == _INPUT_MESSAGES))
outputs: Final = sorted(frozenset(idx for family, idx, _ in groups if family == _OUTPUT_MESSAGES))
pinned_inputs: Final = tuple(dict.fromkeys((*inputs[:1], *inputs[-1:])))
return (
message_order: Final = (
*((_INPUT_MESSAGES, idx) for idx in inputs[1:-1]),
*((_OUTPUT_MESSAGES, idx) for idx in reversed(outputs[1:])),
*((_INPUT_MESSAGES, idx) for idx in pinned_inputs),
*((_OUTPUT_MESSAGES, idx) for idx in outputs[:1]),
)
return tuple(
chain.from_iterable(_message_shed_groups(groups, family, message_idx) for family, message_idx in message_order)
)
def fit_indexed_messages(attrs: Mapping[str, AttrValue], budget: int | None) -> Mapping[str, AttrValue]:
"""``attrs`` with whole per-index messages shed, least valuable first, until at most ``budget`` keys remain.
"""``attrs`` with indexed message attributes shed, least valuable first, until at most ``budget`` keys remain.
``None`` means the span has no attribute count limit. Every message still rides the ``input.value`` and
``output.value`` blobs, so shedding a per-index pair loses no content.
@ -85,7 +156,9 @@ class OpenInferenceMapper:
- ``llm.model_name`` / ``llm.provider`` / ``llm.invocation_parameters``
- ``llm.input_messages.{i}.message.role`` / ``...content``
- ``llm.output_messages.{i}.message.role`` / ``...content``
- ``llm.output_messages.{i}.message.tool_calls.{j}.tool_call.*``
- ``llm.token_count.prompt`` / ``...completion`` / ``...total``
- ``metadata`` — JSON object of allowlisted promoted request metadata
- ``input.value`` / ``output.value`` — JSON-serialized request / response
"""
@ -121,6 +194,7 @@ class OpenInferenceMapper:
"llm.invocation_parameters": lambda d: json_if(
collect(OpenInferenceMapper._INVOCATION_PARAMS, d.request_params)
),
"metadata": lambda d: json_if(dict(sorted(d.promoted_metadata.items()))),
}
def __init__(self, tool_attr_budget: int = MAX_TOOL_DEFINITION_ATTRS_PER_SPAN) -> None:
@ -137,27 +211,36 @@ class OpenInferenceMapper:
return {
**collect(self._LLM_CALL_ATTRS, data),
**collect(self._BLOB_ATTRS, data),
**self._messages(_INPUT_MESSAGES, "input.value", data.messages_in),
**self._messages(_OUTPUT_MESSAGES, "output.value", output_messages(data)),
**self._messages(
_INPUT_MESSAGES,
"input.value",
data.messages_in,
with_tool_call_attrs=False,
),
**self._messages(
_OUTPUT_MESSAGES,
"output.value",
output_messages(data),
with_tool_call_attrs=True,
),
**self._tools(data),
}
@staticmethod
def _messages(prefix: str, value_key: str, messages: Sequence[object]) -> AttributeMap:
def _messages(
prefix: str,
value_key: str,
messages: Sequence[object],
*,
with_tool_call_attrs: bool,
) -> AttributeMap:
"""``{prefix}.{idx}.message.*`` keys for every message + the ``value_key`` blob of all of them."""
parsed: Final = [(m.get("role") if isinstance(m, dict) else None, message_content(m)) for m in messages]
attrs: Final = drop_none(
{
key: value
for idx, (role, content) in enumerate(parsed)
for key, value in (
(f"{prefix}.{idx}.message.role", role if isinstance(role, str) else None),
(f"{prefix}.{idx}.message.content", content),
)
}
parsed: Final = tuple(_parse_message(message) for message in messages)
attrs: Final = drop_none_pairs(
_message_attribute_pairs(prefix, parsed, with_tool_call_attrs=with_tool_call_attrs)
)
if parsed:
attrs[value_key] = json.dumps([{"role": role, "content": content} for role, content in parsed])
attrs[value_key] = _message_value(parsed)
return attrs
def _tools(self, data: LLMCallSpanData) -> AttributeMap:

View file

@ -7,11 +7,28 @@ they live in one place.
import json
from collections.abc import Callable, Iterable, Mapping, Sequence
from dataclasses import dataclass
from typing import Final
from litellm.integrations.otel.mappers.base import AttributeMap, AttrValue
from litellm.integrations.otel.model.payloads import LLMCallSpanData, ToolDefinition
@dataclass(frozen=True, slots=True)
class MessageToolCall:
id: str | None
type: str
name: str | None
arguments: str | None
def to_openai_dict(self) -> dict[str, object]:
return {
"id": self.id,
"type": self.type,
"function": {"name": self.name, "arguments": self.arguments},
}
DEFAULT_SPAN_ATTRIBUTE_LIMIT: Final = 128
"""The OTel SDK's default per-span attribute count limit."""
@ -121,3 +138,36 @@ def message_content(message: object) -> str | None:
def output_messages(data: LLMCallSpanData) -> list:
"""The ``message`` payload of each response choice."""
return [c.get("message") for c in data.choices_out if isinstance(c, dict)]
def message_tool_calls(message: object) -> tuple[MessageToolCall, ...]:
if not isinstance(message, dict):
return ()
tool_calls: Final = message.get("tool_calls")
if not isinstance(tool_calls, (list, tuple)) or not tool_calls:
return ()
return tuple(tool_call for value in tool_calls if (tool_call := _message_tool_call(value)) is not None)
def _message_tool_call(value: object) -> MessageToolCall | None:
if not isinstance(value, dict):
return None
function: Final = value.get("function")
function_data: Final = function if isinstance(function, dict) else {}
raw_type: Final = value.get("type")
raw_arguments: Final = function_data.get("arguments")
arguments: Final = (
raw_arguments
if isinstance(raw_arguments, str)
else json.dumps(raw_arguments, default=str)
if raw_arguments is not None
else None
)
name: Final = function_data.get("name")
identifier: Final = value.get("id")
return MessageToolCall(
id=identifier if isinstance(identifier, str) else None,
type=raw_type if isinstance(raw_type, str) else "function",
name=name if isinstance(name, str) else None,
arguments=arguments,
)

View file

@ -18,7 +18,7 @@ from collections.abc import Callable, Mapping
from types import MappingProxyType
from typing import Final
from litellm.integrations.otel.model.metadata import REQUESTER_METADATA_PATH, RequestIdentity
from litellm.integrations.otel.model.metadata import RequestIdentity, allowlisted_metadata
from litellm.integrations.otel.model.semconv import GenAI, LiteLLM
# Attribute key -> value extractor over (identity, request_model,
@ -92,9 +92,8 @@ def promoted_metadata(metadata: Mapping[str, str], metadata_keys: tuple[str, ...
"""Allowlisted entries of a flattened metadata mapping under ``litellm.metadata.*``."""
return MappingProxyType(
{
f"{LiteLLM.METADATA_PREFIX}{meta_key.removeprefix(REQUESTER_METADATA_PATH)}": value
for meta_key in metadata_keys
if (value := metadata.get(meta_key))
f"{LiteLLM.METADATA_PREFIX}{meta_key}": value
for meta_key, value in allowlisted_metadata(metadata, metadata_keys).items()
}
)

View file

@ -53,6 +53,16 @@ REQUESTER_METADATA_KEY: Final = "requester_metadata"
REQUESTER_METADATA_PATH: Final = f"{REQUESTER_METADATA_KEY}."
def allowlisted_metadata(metadata: Mapping[str, str], metadata_keys: tuple[str, ...]) -> Mapping[str, str]:
return MappingProxyType(
{
meta_key.removeprefix(REQUESTER_METADATA_PATH): value
for meta_key in metadata_keys
if (value := metadata.get(meta_key))
}
)
@dataclass(frozen=True)
class RequestIdentity:
call_id: str | None = None

View file

@ -12,7 +12,7 @@ from urllib.parse import urlsplit
from typing_extensions import ReadOnly, TypedDict
from litellm.integrations.otel.model.metadata import RequestContext, RequestIdentity
from litellm.integrations.otel.model.metadata import RequestContext, RequestIdentity, allowlisted_metadata
from litellm.integrations.otel.model.semconv import (
GenAIOperation,
GenAIOutputType,
@ -60,6 +60,8 @@ if TYPE_CHECKING:
StandardLoggingPayload,
)
_EMPTY_METADATA: Final[Mapping[str, str]] = MappingProxyType({})
# --- typed sub-structures ---------------------------------------------------- #
@ -430,6 +432,7 @@ class LLMCallSpanData:
trace: TraceControls = field(default_factory=TraceControls)
session_id: str | None = None
embedding_output: EmbeddingOutput | None = None
promoted_metadata: Mapping[str, str] = field(default_factory=lambda: _EMPTY_METADATA)
@classmethod
def from_standard_logging_payload(
@ -440,6 +443,8 @@ class LLMCallSpanData:
request_route: str | None = None,
trace: TraceControls | None = None,
session_id: str | None = None,
*,
metadata_keys: tuple[str, ...] = (),
) -> LLMCallSpanData:
params: Final = cast(Mapping[str, object], payload.get("model_parameters") or {})
# The single parse of the request's metadata — the request-vs-provider
@ -476,6 +481,7 @@ class LLMCallSpanData:
cost=LLMCost.from_breakdown(cast("Mapping[str, object] | None", payload.get("cost_breakdown"))),
server=ServerInfo.from_api_base(context.api_base),
identity=context.identity,
promoted_metadata=allowlisted_metadata(context.identity.metadata, metadata_keys),
is_streaming=as_bool(payload.get("stream")),
tools=_extract_tools(params),
messages_in=_dicts(payload.get("messages")) if capture_content else (),

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,403 @@
from __future__ import annotations
import threading
import uuid
from concurrent.futures import ThreadPoolExecutor
from contextlib import ExitStack
from pathlib import Path
from typing import Final
from urllib.parse import urlsplit
import httpx
import psutil
import pytest
from _openinference_support import (
CHAT_TOOLS,
RESPONSES_TOOLS,
_anthropic_response,
_anthropic_stream_response,
_assert_chat_request,
_assert_messages_request,
_assert_responses_request,
_chat_caller_response,
_chat_caller_stream,
_chat_request_marker,
_chat_response,
_chat_stream_response,
_chat_tool_call,
_collect_marker_spans,
_json_object,
_llm_spans_through_markers,
_matching_marker_span,
_messages_caller_raw_stream,
_messages_caller_response,
_normalize_chat_caller_stream,
_normalize_responses_caller_body,
_normalize_responses_caller_stream,
_owned_sink_handler,
_responses_caller_response,
_responses_caller_stream,
_responses_response,
_responses_stream_response,
_rig,
_spans,
_sse_json_values,
)
from integration._support.client import Gateway, eventually
from integration._support.wire import Reply, Request, wire_server
def _call(
proxy: Gateway,
model: str,
marker: str,
*,
surface: str = "chat",
stream: bool = False,
prompt: str | None = None,
) -> httpx.Response:
match surface:
case "chat":
return proxy.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": prompt or "weather in Paris?"}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": {"trace_marker": marker},
**({"stream": True} if stream else {}),
**({"stream_options": {"include_usage": True}} if stream and surface == "chat" else {}),
"cache": {"no-cache": True},
},
)
case "responses":
return proxy.request(
"POST",
"/v1/responses",
{
"model": model,
"input": prompt or "weather in Paris?",
"tools": RESPONSES_TOOLS,
"tool_choice": {"type": "function", "name": "lookup_weather"},
"metadata": {"trace_marker": marker},
**({"stream": True} if stream else {}),
"cache": {"no-cache": True},
},
)
case "messages":
return proxy.request(
"POST",
"/v1/messages",
{
"model": model,
"max_tokens": 64,
"messages": [{"role": "user", "content": prompt or "weather in Paris?"}],
"tools": [
{
"name": "lookup_weather",
"description": "Get weather",
"input_schema": {
"type": "object",
"properties": {"city": {"type": "string"}},
},
}
],
"tool_choice": {"type": "auto"},
"metadata": {"trace_marker": marker},
**({"stream": True} if stream else {}),
"cache": {"no-cache": True},
},
)
case _:
raise AssertionError(f"Unknown endpoint: {surface}")
def _assert_response(response: httpx.Response, marker: str, surface: str, stream: bool, model: str) -> None:
assert response.status_code == 200, response.text
if stream:
observed_stream: Final = _sse_json_values(response.content)
if surface == "chat":
chat_reply: Final = _chat_stream_response(marker, (_chat_tool_call(marker),))
assert _normalize_chat_caller_stream(observed_stream) == _chat_caller_stream(chat_reply, model), (
response.text
)
return
if surface == "responses":
responses_reply: Final = _responses_stream_response(marker, (_chat_tool_call(marker),))
assert _normalize_responses_caller_stream(observed_stream) == _responses_caller_stream(
responses_reply, model
), response.text
return
if surface == "messages":
messages_reply: Final = _anthropic_stream_response(marker)
assert observed_stream == _messages_caller_raw_stream(messages_reply, model), response.text
return
raise AssertionError(f"Unknown endpoint: {surface}")
observed: Final = _json_object(response.content)
if surface == "chat":
chat_reply: Final = _chat_response(marker)
assert observed == _chat_caller_response(chat_reply, model), response.text
return
if surface == "responses":
responses_reply: Final = _responses_response(marker)
assert _normalize_responses_caller_body(observed) == _responses_caller_response(responses_reply, model), (
response.text
)
return
if surface == "messages":
messages_reply: Final = _anthropic_response(marker)
assert observed == _messages_caller_response(messages_reply, model), response.text
return
raise AssertionError(f"Unknown endpoint: {surface}")
def _call_without_worker_error(
proxy: Gateway, model: str, marker: str, *, prompt: str | None = None
) -> httpx.Response | None:
try:
return _call(proxy, model, marker, prompt=prompt)
except httpx.HTTPError:
return None
def test_arize_otel_v2_f1_sink_outage_and_recovery(gateway: Gateway, tmp_path: Path) -> None:
outage_markers: Final = tuple("f1-outage-" + uuid.uuid4().hex for _ in range(30))
recovery_markers: Final = tuple("f1-recovery-" + uuid.uuid4().hex for _ in range(30))
surfaces: Final = ("chat", "responses", "messages")
surface_streams: Final = tuple(
(surfaces[index % len(surfaces)], index % 2 == 0) for index in range(len(outage_markers))
)
calls: Final = tuple(
(marker, surface, stream) for marker, (surface, stream) in zip(outage_markers, surface_streams, strict=True)
)
recovery_calls: Final = tuple(
(marker, surface, stream) for marker, (surface, stream) in zip(recovery_markers, surface_streams, strict=True)
)
def upstream(request: Request) -> Reply:
body: Final = _json_object(request.body)
if request.target.endswith("/messages"):
messages: Final = body.get("messages")
assert isinstance(messages, list) and isinstance(messages[0], dict), body
marker: Final = messages[0].get("content")
assert isinstance(marker, str), body
_assert_messages_request(
request,
marker=marker,
prompt=marker,
stream=True if body.get("stream") is True else False,
)
return _anthropic_stream_response(marker) if body.get("stream") is True else _anthropic_response(marker)
if request.target.endswith("/responses"):
marker: Final = body.get("input")
assert isinstance(marker, str), body
_assert_responses_request(
request,
marker=marker,
input_value=marker,
stream=body.get("stream") is True,
)
return (
_responses_stream_response(marker, (_chat_tool_call(marker),))
if body.get("stream") is True
else _responses_response(marker)
)
marker: Final = _chat_request_marker(request)
_assert_chat_request(
request,
messages=[{"role": "user", "content": marker}],
stream=True if body.get("stream") is True else None,
stream_options={"include_usage": True} if body.get("stream") is True else None,
)
return (
_chat_stream_response(marker, (_chat_tool_call(marker),))
if body.get("stream") is True
else _chat_response(marker)
)
def sink(_request: Request) -> Reply:
return Reply(body=b"", content_type="application/x-protobuf")
def run_burst(
burst: tuple[tuple[str, str, bool], ...],
proxy: Gateway,
chat_model: str,
messages_model: str,
) -> None:
with ThreadPoolExecutor(max_workers=len(burst)) as executor:
futures: Final = tuple(
executor.submit(
_call,
proxy,
messages_model if surface == "messages" else chat_model,
marker,
surface=surface,
stream=stream,
prompt=marker,
)
for marker, surface, stream in burst
)
responses: Final = tuple(
(marker, surface, stream, future.result(timeout=60))
for (marker, surface, stream), future in zip(burst, futures, strict=True)
)
for marker, surface, stream, response in responses:
model_name: Final = messages_model if surface == "messages" else chat_model
_assert_response(response, marker, surface, stream, model_name)
with ExitStack() as servers:
initial_stack: Final = servers.enter_context(ExitStack())
stopped_destination: Final = initial_stack.enter_context(wire_server(_owned_sink_handler(sink)))
sink_port: Final = urlsplit(stopped_destination.url).port
assert sink_port is not None, stopped_destination.url
initial_stack.close()
with _rig(
gateway,
tmp_path,
upstream,
destination_wire=stopped_destination,
workers=1,
) as rig:
with httpx.Client(trust_env=False) as client, pytest.raises(httpx.ConnectError):
client.get(stopped_destination.url + "/health", timeout=2)
with rig.proxy.scenario() as scenario:
messages_model: Final = scenario.model(
model="anthropic/claude-opus-5-5",
api_base=rig.provider.url,
)
run_burst(calls, rig.proxy, rig.model, messages_model)
with httpx.Client(trust_env=False) as client, pytest.raises(httpx.ConnectError):
client.get(stopped_destination.url + "/health", timeout=2)
recovered_stack: Final = servers.enter_context(ExitStack())
recovered_destination: Final = recovered_stack.enter_context(
wire_server(_owned_sink_handler(sink), port=sink_port)
)
run_burst(recovery_calls, rig.proxy, rig.model, messages_model)
sentinel: Final = f"f1-sentinel-{uuid.uuid4().hex}"
sentinel_response: Final = _call(
rig.proxy,
rig.model,
sentinel,
surface="chat",
stream=False,
prompt=sentinel,
)
_assert_response(sentinel_response, sentinel, "chat", False, rig.model)
spans: Final = _llm_spans_through_markers(recovered_destination, (*recovery_markers, sentinel))
assert all(
sum(span.get("litellm.metadata.trace_marker") == marker for span in spans) == 1
for marker in recovery_markers
), spans
assert all(
sum(span.get("litellm.metadata.trace_marker") == marker for span in spans) <= 1
for marker in outage_markers
), spans
assert sum(span.get("litellm.metadata.trace_marker") == sentinel for span in spans) == 1, spans
def test_arize_otel_v2_f2_slow_sink_does_not_deadlock(gateway: Gateway, tmp_path: Path) -> None:
release: Final = threading.Event()
blocked: Final = threading.Event()
completed: Final = threading.Event()
marker: Final = "f2-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker)
def sink(request: Request) -> Reply:
if any(attributes.get("litellm.metadata.trace_marker") == marker for attributes in _spans((request,))):
blocked.set()
assert release.wait(timeout=5), "slow sink was not released"
completed.set()
return Reply(body=b"", content_type="application/x-protobuf")
with _rig(gateway, tmp_path, upstream, destination_handler=sink) as rig:
timer: Final = threading.Timer(2, release.set)
try:
response: Final = _call(rig.proxy, rig.model, marker)
_assert_response(response, marker, "chat", False, rig.model)
assert eventually(lambda: blocked.is_set(), bool, seconds=10)
timer.start()
assert eventually(lambda: completed.is_set(), bool, seconds=10)
timer.join(timeout=5)
assert not timer.is_alive(), "Slow sink timer did not finish"
finally:
release.set()
timer.cancel()
if timer.ident is not None:
timer.join(timeout=5)
_matching_marker_span(rig.destination, marker)
def test_arize_otel_v2_f3_one_proxy_worker_can_die(gateway: Gateway, tmp_path: Path) -> None:
markers: Final = tuple("f3-" + uuid.uuid4().hex for _ in range(8))
release: Final = threading.Event()
def upstream(request: Request) -> Reply:
request_marker: Final = _chat_request_marker(request)
_assert_chat_request(
request,
messages=[{"role": "user", "content": request_marker}],
)
assert release.wait(timeout=20), "F3 upstream barrier was not released"
return _chat_response(request_marker)
def sink(_request: Request) -> Reply:
return Reply(body=b"", content_type="application/x-protobuf")
with _rig(
gateway,
tmp_path,
upstream,
destination_handler=sink,
fresh_client_connections=True,
) as rig:
children: Final = psutil.Process(rig.owned.process.pid).children(recursive=True)
workers: Final = tuple(
child for child in children if child.is_running() and "resource_tracker" not in " ".join(child.cmdline())
)
assert len(workers) >= 2, tuple((worker.pid, worker.name()) for worker in workers)
with ThreadPoolExecutor(max_workers=len(markers)) as executor:
try:
futures: Final = tuple(
executor.submit(
_call_without_worker_error,
rig.proxy,
rig.model,
marker,
prompt=marker,
)
for marker in markers
)
observed: Final = eventually(
lambda: rig.provider.received.qsize(),
lambda count: count >= 2,
seconds=10,
)
assert observed >= 2, observed
workers[0].kill()
assert eventually(lambda: not workers[0].is_running(), bool, seconds=10), workers[0]
release.set()
in_flight: Final = tuple(
(marker, future.result(timeout=60)) for marker, future in zip(markers, futures, strict=True)
)
finally:
release.set()
survivor: Final = "f3-survivor-" + uuid.uuid4().hex
survivor_response: Final = _call(rig.proxy, rig.model, survivor, prompt=survivor)
_assert_response(survivor_response, survivor, "chat", False, rig.model)
served_responses: Final = tuple(
(marker, response) for marker, response in in_flight if response is not None and response.status_code == 200
)
for marker, response in served_responses:
_assert_response(response, marker, "chat", False, rig.model)
served: Final = tuple(marker for marker, _response in served_responses) + (survivor,)
collected: Final = _collect_marker_spans(rig.destination, served)
assert len(collected) == len(served), collected
exported: Final = tuple(span["litellm.metadata.trace_marker"] for span in collected)
assert len(exported) == len(served), exported
assert frozenset(exported) == frozenset(served), exported

View file

@ -0,0 +1,336 @@
from __future__ import annotations
import uuid
from collections.abc import Mapping
from pathlib import Path
from types import MappingProxyType
from typing import Final
import httpx
import pytest
from _openinference_support import (
CHAT_TOOLS,
_assert_chat_request,
_chat_caller_response,
_chat_response,
_json_messages,
_json_object,
_matching_genai_marker_span,
_matching_marker_span,
_rig,
)
from integration._support.client import Gateway
from integration._support.wire import Reply, Request
from pydantic import JsonValue
_GENAI_B3_KEYS_WITHOUT_BAGGAGE: Final = frozenset(
{
"gen_ai.input.messages",
"gen_ai.operation.name",
"gen_ai.output.messages",
"gen_ai.provider.name",
"gen_ai.request.model",
"gen_ai.response.finish_reasons",
"gen_ai.response.id",
"gen_ai.response.model",
"gen_ai.system",
"gen_ai.tool.0.description",
"gen_ai.tool.0.name",
"gen_ai.tool.0.parameters",
"gen_ai.usage.completion_tokens",
"gen_ai.usage.input_tokens",
"gen_ai.usage.output_tokens",
"gen_ai.usage.prompt_tokens",
"gen_ai.usage.total_tokens",
"litellm.api_key.hash",
"litellm.call_id",
"litellm.call_type",
"litellm.cost.discount_amount",
"litellm.cost.discount_percent",
"litellm.cost.input",
"litellm.cost.margin_fixed_amount",
"litellm.cost.margin_percent",
"litellm.cost.margin_total_amount",
"litellm.cost.original",
"litellm.cost.output",
"litellm.cost.tool_usage",
"litellm.cost.total",
"litellm.provider.model",
"litellm.request.route",
"litellm.request.tools.declared",
"llm.request.functions.0.description",
"llm.request.functions.0.name",
"llm.request.functions.0.parameters",
"server.address",
"server.port",
}
)
_GENAI_B3_KEYS_WITH_BAGGAGE: Final = _GENAI_B3_KEYS_WITHOUT_BAGGAGE | frozenset({"litellm.metadata.trace_marker"})
_LANGFUSE_B3_KEYS: Final = frozenset(
{
"gen_ai.input.messages",
"gen_ai.operation.name",
"gen_ai.output.messages",
"gen_ai.provider.name",
"gen_ai.request.model",
"gen_ai.response.finish_reasons",
"gen_ai.response.id",
"gen_ai.response.model",
"gen_ai.system",
"gen_ai.tool.0.description",
"gen_ai.tool.0.name",
"gen_ai.tool.0.parameters",
"gen_ai.usage.completion_tokens",
"gen_ai.usage.input_tokens",
"gen_ai.usage.output_tokens",
"gen_ai.usage.prompt_tokens",
"gen_ai.usage.total_tokens",
"langfuse.observation.cost_details",
"langfuse.observation.id",
"langfuse.observation.input",
"langfuse.observation.metadata.provider",
"langfuse.observation.model.name",
"langfuse.observation.output",
"langfuse.observation.type",
"langfuse.observation.usage_details",
"litellm.api_key.hash",
"litellm.call_id",
"litellm.call_type",
"litellm.cost.discount_amount",
"litellm.cost.discount_percent",
"litellm.cost.input",
"litellm.cost.margin_fixed_amount",
"litellm.cost.margin_percent",
"litellm.cost.margin_total_amount",
"litellm.cost.original",
"litellm.cost.output",
"litellm.cost.tool_usage",
"litellm.cost.total",
"litellm.provider.model",
"litellm.request.route",
"litellm.request.tools.declared",
"llm.request.functions.0.description",
"llm.request.functions.0.name",
"llm.request.functions.0.parameters",
"server.address",
"server.port",
}
)
_B3_ATTRIBUTE_KEYS: Final[Mapping[str, frozenset[str]]] = MappingProxyType(
{
"langfuse_otel": _LANGFUSE_B3_KEYS,
"langtrace": _GENAI_B3_KEYS_WITHOUT_BAGGAGE,
"signoz": _GENAI_B3_KEYS_WITH_BAGGAGE,
"newrelic": _GENAI_B3_KEYS_WITH_BAGGAGE,
"levo": _GENAI_B3_KEYS_WITH_BAGGAGE,
"agentops": _GENAI_B3_KEYS_WITH_BAGGAGE,
"otel": _GENAI_B3_KEYS_WITH_BAGGAGE,
}
)
_B3_BAGGAGE_CALLBACKS: Final = frozenset({"signoz", "newrelic", "levo", "agentops", "otel"})
_B4_ATTRIBUTE_KEYS: Final = frozenset(
{
"input.value",
"litellm.trace_id",
"llm.cost.total",
"llm.input_messages.0.message.content",
"llm.input_messages.0.message.role",
"llm.invocation_parameters",
"llm.is_streaming",
"llm.model_name",
"llm.output_messages.0.message.content",
"llm.output_messages.0.message.role",
"llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments",
"llm.output_messages.0.message.tool_calls.0.tool_call.function.name",
"llm.output_messages.0.message.tool_calls.0.tool_call.id",
"llm.provider",
"llm.request.type",
"llm.response.cost",
"llm.response.id",
"llm.response.model",
"llm.token_count.completion",
"llm.token_count.prompt",
"llm.token_count.total",
"llm.tools.0.description",
"llm.tools.0.name",
"llm.tools.0.parameters",
"metadata",
"openinference.span.kind",
"output.value",
"user.id",
}
)
def _request(proxy: Gateway, model: str, marker: str) -> httpx.Response:
return proxy.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": "weather in Paris?"}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
)
def _assert_openinference(attributes: dict[str, str], marker: str) -> None:
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.id"] == f"call_{marker}", attributes
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.name"] == "lookup_weather", (
attributes
)
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments"] == '{"city": "Paris"}'
), attributes
assert _json_object(attributes["metadata"].encode()) == {"trace_marker": marker}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
assert _json_messages(attributes["output.value"]) == [
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": f"call_{marker}",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
}
], attributes
def test_arize_otel_v2_b1_phoenix_openinference(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "b1-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker)
with _rig(
gateway,
tmp_path,
upstream,
callbacks=("arize_phoenix",),
callback_settings={"otel": {"exporter": "http/protobuf", "endpoint": "unused"}},
environment={"PHOENIX_PROJECT_NAME": "integration"},
) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker), rig.model), response.text
_assert_openinference(_matching_marker_span(rig.destination, marker), marker)
def test_arize_otel_v2_b2_weave_openinference(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "b2-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker)
with _rig(
gateway,
tmp_path,
upstream,
callbacks=("weave_otel",),
callback_settings={"otel": {"exporter": "http/protobuf", "endpoint": "unused"}},
environment={"WANDB_BASE_URL": "http://127.0.0.1"},
) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker), rig.model), response.text
_assert_openinference(_matching_marker_span(rig.destination, marker), marker)
@pytest.mark.parametrize(
"callback",
("langfuse_otel", "langtrace", "signoz", "newrelic", "levo", "agentops", "otel"),
)
def test_arize_otel_v2_b3_non_openinference_callback_family(callback: str, gateway: Gateway, tmp_path: Path) -> None:
marker: Final = f"b3-{callback}-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
body: Final = _json_object(request.body)
assert body == {
"messages": [{"role": "user", "content": "weather in Paris?"}],
"model": "gpt-4o-mini",
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"tools": CHAT_TOOLS,
}, body
return _chat_response(marker)
callback_settings: Final[dict[str, JsonValue]] = {
"otel": {"exporter": "http/protobuf", "endpoint": "unused", "mapper_names": ["genai"]}
}
with _rig(
gateway,
tmp_path,
upstream,
callbacks=(callback,),
callback_settings=callback_settings,
environment={
"LANGFUSE_HOST": "http://127.0.0.1",
**(
{
"HTTPS_PROXY": "http://127.0.0.1:0",
"https_proxy": "http://127.0.0.1:0",
"NO_PROXY": "",
"no_proxy": "",
}
if callback == "agentops"
else {}
),
},
remove_environment=(
("NEW_RELIC_LICENSE_KEY",)
if callback == "newrelic"
else ("AGENTOPS_API_KEY",)
if callback == "agentops"
else ()
),
) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker), rig.model), response.text
attributes: Final = _matching_genai_marker_span(rig.destination, marker)
assert frozenset(attributes) == _B3_ATTRIBUTE_KEYS[callback], attributes
assert "metadata" not in attributes, attributes
assert not any(".tool_calls." in key for key in attributes), attributes
baggage: Final = tuple(
sorted((key, value) for key, value in attributes.items() if key.startswith("litellm.metadata."))
)
expected_baggage: Final = (
(("litellm.metadata.trace_marker", marker),) if callback in _B3_BAGGAGE_CALLBACKS else ()
)
assert baggage == expected_baggage, attributes
def test_arize_otel_v2_b4_legacy_otel_is_unchanged(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "b4-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
body: Final = _json_object(request.body)
assert body == {
"messages": [{"role": "user", "content": "weather in Paris?"}],
"model": "gpt-4o-mini",
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"tools": CHAT_TOOLS,
}, body
return _chat_response(marker)
with _rig(
gateway,
tmp_path,
upstream,
callbacks=("arize",),
callback_settings={"otel": {"exporter": "http/protobuf", "endpoint": "unused"}},
remove_environment=("LITELLM_OTEL_V2",),
disabled_environment=("LITELLM_OTEL_V2",),
) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker), rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
assert frozenset(attributes) == _B4_ATTRIBUTE_KEYS, attributes

View file

@ -0,0 +1,431 @@
from __future__ import annotations
import uuid
from collections.abc import Callable
from pathlib import Path
from typing import Final
import httpx
import pytest
from _openinference_support import (
CHAT_TOOLS,
_assert_chat_request,
_chat_caller_response,
_chat_plain_response,
_chat_request_marker,
_chat_response,
_json_messages,
_json_object,
_matching_marker_span,
_matching_output_value_span,
_matching_span,
_response_tool_calls,
_rig,
)
from integration._support.client import Gateway
from integration._support.wire import Reply, Request
_DEFAULT_METADATA: Final = {
"requester_ip_address": "127.0.0.1",
"user_api_key_user_id": "default_user_id",
}
_DEFAULT_METADATA_BAGGAGE: Final = frozenset({("litellm.metadata.user_api_key_user_id", "default_user_id")})
def _request(
proxy: Gateway,
model: str,
marker: str,
*,
prompt: str = "weather in Paris?",
headers: dict[str, str] | None = None,
key: str | None = None,
) -> httpx.Response:
return proxy.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": prompt}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
key=key,
headers=headers,
)
def _assert_success_body(response: httpx.Response, marker: str, model: str) -> None:
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker), model), response.text
def _upstream(marker: str) -> Callable[[Request], Reply]:
def reply(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker)
return reply
def _assert_output_tool_call(attributes: dict[str, str], marker: str) -> None:
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.id"] == f"call_{marker}", attributes
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.name"] == "lookup_weather", (
attributes
)
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments"] == '{"city": "Paris"}'
), attributes
assert _json_messages(attributes["output.value"]) == [
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": f"call_{marker}",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
}
], attributes
def _assert_default_allowlist_attributes(attributes: dict[str, str], marker: str) -> None:
observed_baggage: Final = frozenset(
(key, value) for key, value in attributes.items() if key.startswith("litellm.metadata.")
)
assert observed_baggage == _DEFAULT_METADATA_BAGGAGE, attributes
metadata: Final = _json_object(attributes["metadata"].encode())
assert metadata == _DEFAULT_METADATA, attributes
assert "trace_marker" not in metadata, attributes
assert "litellm.metadata.trace_marker" not in attributes, attributes
_assert_output_tool_call(attributes, marker)
def test_arize_otel_v2_c1_absent_allowlist(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "c1-" + uuid.uuid4().hex
with _rig(
gateway,
tmp_path,
_upstream(marker),
remove_environment=("LITELLM_OTEL_BAGGAGE_METADATA_KEYS",),
disabled_environment=("LITELLM_OTEL_BAGGAGE_METADATA_KEYS",),
) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
_assert_success_body(response, marker, rig.model)
attributes: Final = _matching_span(rig.destination, marker)
_assert_default_allowlist_attributes(attributes, marker)
def test_arize_otel_v2_c2_empty_allowlist(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "c2-" + uuid.uuid4().hex
with _rig(gateway, tmp_path, _upstream(marker), environment={"LITELLM_OTEL_BAGGAGE_METADATA_KEYS": ""}) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
_assert_success_body(response, marker, rig.model)
attributes: Final = _matching_span(rig.destination, marker)
assert "metadata" not in attributes, attributes
assert not any(key.startswith("litellm.metadata.") for key in attributes), attributes
_assert_output_tool_call(attributes, marker)
def test_arize_otel_v2_c3_absent_allowlisted_key(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "c3-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker)
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = rig.proxy.request(
"POST",
"/v1/chat/completions",
{
"model": rig.model,
"messages": [{"role": "user", "content": "weather in Paris?"}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": {"other": "value"},
"cache": {"no-cache": True},
},
)
_assert_success_body(response, marker, rig.model)
attributes: Final = _matching_span(rig.destination, marker)
assert "metadata" not in attributes, attributes
assert "litellm.metadata.trace_marker" not in attributes, attributes
_assert_output_tool_call(attributes, marker)
def test_arize_otel_v2_c4_promotes_marker_and_alias(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "c4-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker)
with (
_rig(
gateway,
tmp_path,
upstream,
environment={"LITELLM_OTEL_BAGGAGE_METADATA_KEYS": "requester_metadata.trace_marker,user_api_key_alias"},
) as rig,
rig.proxy.scenario() as scenario,
):
key: Final = scenario.key(key_alias="alias-c4")
response: Final = rig.proxy.request(
"POST",
"/v1/chat/completions",
{
"model": rig.model,
"messages": [{"role": "user", "content": "weather in Paris?"}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
key=key,
)
_assert_success_body(response, marker, rig.model)
attributes: Final = _matching_marker_span(rig.destination, marker)
assert _json_object(attributes["metadata"].encode()) == {
"trace_marker": marker,
"user_api_key_alias": "alias-c4",
}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
assert attributes["litellm.metadata.user_api_key_alias"] == "alias-c4", attributes
def test_arize_otel_v2_c5_yaml_allowlist_does_not_reach_preset_so_default_applies(
gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = "c5-" + uuid.uuid4().hex
with _rig(
gateway,
tmp_path,
_upstream(marker),
callback_settings={"otel": {"baggage_metadata_keys": ["requester_metadata.trace_marker"]}},
remove_environment=("LITELLM_OTEL_BAGGAGE_METADATA_KEYS",),
disabled_environment=("LITELLM_OTEL_BAGGAGE_METADATA_KEYS",),
) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
_assert_success_body(response, marker, rig.model)
attributes: Final = _matching_marker_span(rig.destination, marker)
_assert_default_allowlist_attributes(attributes, marker)
def test_arize_otel_v2_c6_content_capture_disabled(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "c6-" + uuid.uuid4().hex
with _rig(
gateway,
tmp_path,
_upstream(marker),
environment={
"LITELLM_OTEL_BAGGAGE_METADATA_KEYS": "requester_metadata.trace_marker",
},
remove_environment=("OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT",),
disabled_environment=("OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT",),
) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
_assert_success_body(response, marker, rig.model)
attributes: Final = _matching_marker_span(rig.destination, marker)
assert _json_object(attributes["metadata"].encode()) == {"trace_marker": marker}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
assert "gen_ai.input.messages" not in attributes, attributes
assert "gen_ai.output.messages" not in attributes, attributes
assert not any(
key.startswith("llm.input_messages.") or key.startswith("llm.output_messages.") for key in attributes
), attributes
assert "input.value" not in attributes, attributes
assert "output.value" not in attributes, attributes
assert not any(".tool_calls." in key for key in attributes), attributes
def test_arize_otel_v2_c7_key_and_team_logging_callbacks(gateway: Gateway, tmp_path: Path) -> None:
key_marker: Final = "c7-key-" + uuid.uuid4().hex
team_marker: Final = "c7-team-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
marker: Final = _chat_request_marker(request)
assert marker in (key_marker, team_marker), request
_assert_chat_request(request, messages=[{"role": "user", "content": marker}])
return _chat_response(marker)
logging_metadata: Final = {"logging": [{"callback_name": "arize", "callback_type": "success"}]}
with _rig(gateway, tmp_path, upstream) as rig, rig.proxy.scenario() as scenario:
key: Final = scenario.key(key_alias="key-c7", metadata=logging_metadata)
team: Final = scenario.team(metadata=logging_metadata)
team_key: Final = scenario.key(team_id=team)
key_response: Final = _request(rig.proxy, rig.model, key_marker, prompt=key_marker, key=key)
_assert_success_body(key_response, key_marker, rig.model)
key_attributes: Final = _matching_marker_span(rig.destination, key_marker)
team_response: Final = _request(rig.proxy, rig.model, team_marker, prompt=team_marker, key=team_key)
_assert_success_body(team_response, team_marker, rig.model)
team_attributes: Final = _matching_marker_span(rig.destination, team_marker)
assert _json_object(key_attributes["metadata"].encode()) == {"trace_marker": key_marker}, key_attributes
assert _json_object(team_attributes["metadata"].encode()) == {"trace_marker": team_marker}, team_attributes
assert key_attributes["litellm.metadata.trace_marker"] == key_marker, key_attributes
assert team_attributes["litellm.metadata.trace_marker"] == team_marker, team_attributes
_assert_output_tool_call(key_attributes, key_marker)
_assert_output_tool_call(team_attributes, team_marker)
def test_arize_otel_v2_c8_request_callback_disable(gateway: Gateway, tmp_path: Path) -> None:
control_marker: Final = "c8-control-" + uuid.uuid4().hex
disabled_marker: Final = "c8-disabled-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
marker: Final = _chat_request_marker(request)
assert marker in (control_marker, disabled_marker), marker
return _chat_response(marker)
with _rig(
gateway,
tmp_path,
upstream,
litellm_settings={"allow_dynamic_callback_disabling": True},
) as rig:
control_response: Final = _request(rig.proxy, rig.model, control_marker, prompt=control_marker)
_assert_success_body(control_response, control_marker, rig.model)
control_attributes: Final = _matching_marker_span(rig.destination, control_marker)
_assert_output_tool_call(control_attributes, control_marker)
disabled_response: Final = _request(
rig.proxy,
rig.model,
disabled_marker,
prompt=disabled_marker,
headers={"x-litellm-disable-callbacks": "arize"},
)
_assert_success_body(disabled_response, disabled_marker, rig.model)
@pytest.mark.parametrize("failure_status", (401, 500))
def test_arize_otel_v2_c9_upstream_failures_are_recorded(failure_status: int, gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "c9-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
assert request.method == "POST", request.method
assert request.body, f"{request.method} {request.target}"
body: Final = _json_object(request.body)
assert body == {
"messages": [{"role": "user", "content": "weather in Paris?"}],
"model": "gpt-4o-mini",
}, body
return Reply(
status=failure_status,
body=b'{"error":{"message":"upstream failure"}}',
content_type="application/json",
)
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = rig.proxy.request(
"POST",
"/v1/chat/completions",
{
"model": rig.model,
"messages": [{"role": "user", "content": "weather in Paris?"}],
"metadata": {"trace_marker": marker, "failure_status": failure_status},
"cache": {"no-cache": True},
},
)
assert response.status_code == failure_status, response.text
error_body: Final = _json_object(response.content)
error_name: Final = {401: "AuthenticationError", 500: "InternalServerError"}[failure_status]
error_type: Final = {401: "authentication_error", 500: "internal_server_error"}[failure_status]
provider_message: Final = f"litellm.{error_name}: {error_name}: OpenAIException - upstream failure"
caller_message: Final = (
f"{provider_message}\n\nLiteLLM: model group '{rig.model}' failed with the error above. "
"No fallback was attempted."
)
assert error_body == {
"error": {
"message": caller_message,
"type": error_type,
"param": None,
"code": str(failure_status),
}
}, error_body
attributes: Final = _matching_marker_span(rig.destination, marker)
metadata: Final = attributes.get("metadata")
assert metadata is not None, attributes
assert _json_object(metadata.encode()) == {"trace_marker": marker}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
assert attributes["error.message"] == provider_message, attributes
assert attributes["error.type"] == error_name, attributes
assert not any(".tool_calls." in key for key in attributes), attributes
def test_arize_otel_v2_c10_attribute_limit_keeps_tool_prefix(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "c10-" + uuid.uuid4().hex
calls: Final = _response_tool_calls(marker, ("Paris", "Berlin", "Rome", "Tokyo", "Oslo", "Lima", "Accra", "Delhi"))
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker, calls)
with _rig(gateway, tmp_path, upstream, environment={"OTEL_SPAN_ATTRIBUTE_COUNT_LIMIT": "57"}) as rig:
response: Final = _request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker, calls), rig.model), (
response.text
)
attributes: Final = _matching_output_value_span(rig.destination, marker)
indexes: Final = tuple(
int(key.split(".tool_calls.")[1].split(".")[0])
for key in attributes
if ".tool_calls." in key and key.endswith(".tool_call.id")
)
assert indexes == (0,), attributes
assert attributes["llm.output_messages.0.message.role"] == "assistant", attributes
fields: Final = ("id", "function.name", "function.arguments")
expected_tool_call_keys: Final = frozenset().union(
*(
frozenset(f"llm.output_messages.0.message.tool_calls.{index}.tool_call.{field}" for field in fields)
for index in indexes
)
)
observed_tool_call_keys: Final = frozenset(key for key in attributes if ".tool_calls." in key)
assert observed_tool_call_keys == expected_tool_call_keys, attributes
assert _json_messages(attributes["output.value"])[0]["tool_calls"] == calls, attributes
assert _json_object(attributes["metadata"].encode()) == {"trace_marker": marker}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
history: Final = [{"role": "user", "content": f"history-{index}"} for index in range(40)]
history_marker: Final = marker + "-history"
def history_upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=history, include_tools=False)
return _chat_plain_response(history_marker, "history retained")
with _rig(gateway, tmp_path, history_upstream) as history_rig:
history_response: Final = history_rig.proxy.request(
"POST",
"/v1/chat/completions",
{
"model": history_rig.model,
"messages": history,
"metadata": {"trace_marker": history_marker},
"cache": {"no-cache": True},
},
)
assert history_response.status_code == 200, history_response.text
assert _json_object(history_response.content) == _chat_caller_response(
_chat_plain_response(history_marker, "history retained"), history_rig.model
), history_response.text
history_attributes: Final = _matching_marker_span(history_rig.destination, history_marker)
assert _json_object(history_attributes["metadata"].encode()) == {"trace_marker": history_marker}, (
history_attributes
)
assert history_attributes["litellm.metadata.trace_marker"] == history_marker, history_attributes
assert tuple(history_attributes[f"llm.input_messages.{index}.message.role"] for index in range(40)) == tuple(
str(message["role"]) for message in history
)
assert tuple(history_attributes[f"llm.input_messages.{index}.message.content"] for index in range(40)) == tuple(
str(message["content"]) for message in history
)
assert _json_messages(history_attributes["output.value"]) == [
{"role": "assistant", "content": "history retained"}
], history_attributes

View file

@ -0,0 +1,512 @@
from __future__ import annotations
import json
import uuid
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
from typing import Final
import httpx
import pytest
from _openinference_support import (
CHAT_TOOLS,
_assert_chat_request,
_chat_caller_response,
_chat_output_value,
_chat_request_marker,
_collect_marker_spans,
_json_messages,
_json_object,
_json_object_value,
_matching_marker_span,
_rig,
_span_attributes,
_spans,
)
from integration._support.client import Gateway
from integration._support.wire import Reply, Request
from pydantic import JsonValue
def _call(
proxy: Gateway,
model: str,
marker: str,
metadata: JsonValue | None = None,
key: str | None = None,
*,
prompt: str = "weather in Paris?",
) -> httpx.Response:
return proxy.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": prompt}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": metadata if metadata is not None else {"trace_marker": marker},
"cache": {"no-cache": True},
},
key=key,
)
def _success(marker: str) -> Reply:
return Reply(
body=json.dumps(
{
"id": marker,
"object": "chat.completion",
"created": 1,
"model": "gpt-4o-mini",
"choices": [
{
"index": 0,
"finish_reason": "tool_calls",
"message": {
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": "call_" + marker,
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
},
}
],
"usage": {"prompt_tokens": 10, "completion_tokens": 2, "total_tokens": 12},
}
).encode()
)
def _chat_message(body: bytes) -> dict[str, JsonValue]:
response: Final = _json_object(body)
choices: Final = response["choices"]
assert isinstance(choices, list) and len(choices) == 1 and isinstance(choices[0], dict), response
message: Final = choices[0]["message"]
assert isinstance(message, dict), response
return message
@pytest.mark.parametrize(
("value", "expected_trace"),
(
(7, "7"),
(["one", 2], None),
("", None),
("x" * 5000, "x" * 5000),
({"enabled": True}, None),
),
)
def test_arize_otel_v2_d1_metadata_value_shapes(
value: JsonValue, expected_trace: str | None, gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = "d1-" + uuid.uuid4().hex
metadata: Final = {"trace_marker": value}
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _success(marker)
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = _call(rig.proxy, rig.model, marker, metadata)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_success(marker), rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
if expected_trace is None:
assert "metadata" not in attributes, attributes
assert "litellm.metadata.trace_marker" not in attributes, attributes
else:
assert json.loads(attributes["metadata"]) == {"trace_marker": expected_trace}, attributes
assert attributes["litellm.metadata.trace_marker"] == expected_trace, attributes
def test_arize_otel_v2_d2_duplicate_json_metadata_keys(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "d2-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _success(marker)
with _rig(gateway, tmp_path, upstream) as rig:
request_body: Final = (
'{"model":"'
+ rig.model
+ '","messages":'
+ json.dumps([{"role": "user", "content": "weather in Paris?"}])
+ ',"tools":'
+ json.dumps(CHAT_TOOLS)
+ ","
+ '"tool_choice":{"type":"function","function":{"name":"lookup_weather"}},'
+ '"metadata":{"trace_marker":"'
+ marker
+ '","trace_marker":"'
+ marker
+ '"},"cache":{"no-cache":true}}'
)
response: Final = rig.proxy.client.post(
"/v1/chat/completions",
content=request_body,
headers={
"authorization": f"Bearer {rig.proxy.key}",
"content-type": "application/json",
},
)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_success(marker), rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
assert _json_object(attributes["metadata"].encode()) == {"trace_marker": marker}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
def test_arize_otel_v2_d4_unauthenticated_request_has_no_span(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "d4-" + uuid.uuid4().hex
with _rig(gateway, tmp_path, lambda _request: _success(marker)) as rig:
with httpx.Client(base_url=str(rig.proxy.client.base_url), trust_env=False) as client:
response: Final = client.post(
"/v1/chat/completions",
json={"model": rig.model, "messages": [{"role": "user", "content": marker}]},
)
assert response.status_code == 401, response.text
body: Final = _json_object(response.content)
assert body == {
"error": {
"message": "Authentication Error, No api key passed in.",
"type": "auth_error",
"param": "None",
"code": "401",
}
}, body
assert rig.provider.received.qsize() == 0
spans: Final = tuple(
attributes
for attributes in _spans(rig.destination.drain())
if attributes.get("openinference.span.kind") == "LLM"
)
assert spans == (), "unauthenticated request exported an LLM span"
def test_arize_otel_v2_d5_unknown_model_leaves_proxy_ready(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "d5-" + uuid.uuid4().hex
unknown_model: Final = "unknown-model-" + uuid.uuid4().hex
with _rig(gateway, tmp_path, lambda _request: _success(marker)) as rig:
response: Final = rig.proxy.request(
"POST",
"/v1/chat/completions",
{"model": unknown_model, "messages": [{"role": "user", "content": marker}]},
)
assert response.status_code == 400, response.text
body: Final = _json_object(response.content)
error_message: Final = (
f"/chat/completions: Invalid model name passed in model={unknown_model}. "
"Call `/v1/models` to view available models for your key."
)
assert body == {
"error": {
"message": error_message,
"type": "invalid_request_error",
"param": None,
"code": "400",
"provider_specific_fields": {"error": error_message},
}
}, body
assert rig.provider.received.qsize() == 0
readiness: Final = rig.proxy.client.get("/health/readiness")
assert readiness.status_code == 200, readiness.text
assert _json_object(readiness.content) == {"status": "healthy", "db": "connected"}, readiness.text
spans: Final = tuple(
attributes
for attributes in _spans(rig.destination.drain())
if attributes.get("openinference.span.kind") == "LLM"
)
assert spans == (), "unknown model exported an LLM span"
def test_arize_otel_v2_d6_sink_rejections_do_not_change_caller_response(gateway: Gateway, tmp_path: Path) -> None:
markers: Final = tuple(f"d6-{status}-" + uuid.uuid4().hex for status in (403, 404))
unrelated_marker: Final = "d6-unrelated-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
response_marker: Final = _chat_request_marker(request)
_assert_chat_request(request, messages=[{"role": "user", "content": response_marker}])
return _success(response_marker)
def sink(request: Request) -> Reply:
exported_markers: Final = tuple(
attributes.get("litellm.metadata.trace_marker")
for attributes in _span_attributes(request)
if "litellm.metadata.trace_marker" in attributes
)
if markers[0] in exported_markers:
return Reply(status=403, body=b"rejected")
if markers[1] in exported_markers:
return Reply(status=404, body=b"rejected")
return Reply(body=b"", content_type="application/x-protobuf")
with (
_rig(
gateway,
tmp_path,
upstream,
destination_handler=sink,
) as rig,
rig.proxy.scenario() as scenario,
):
unrelated_key: Final = scenario.key(key_alias="unrelated-d6")
first: Final = _call(rig.proxy, rig.model, markers[0], prompt=markers[0])
assert first.status_code == 200, first.text
assert _json_object(first.content) == _chat_caller_response(_success(markers[0]), rig.model), first.text
assert _matching_marker_span(rig.destination, markers[0])["litellm.metadata.trace_marker"] == markers[0]
second: Final = _call(rig.proxy, rig.model, markers[1], prompt=markers[1])
assert second.status_code == 200, second.text
assert _json_object(second.content) == _chat_caller_response(_success(markers[1]), rig.model), second.text
assert _matching_marker_span(rig.destination, markers[1])["litellm.metadata.trace_marker"] == markers[1]
unrelated: Final = _call(
rig.proxy,
rig.model,
unrelated_marker,
key=unrelated_key,
prompt=unrelated_marker,
)
assert unrelated.status_code == 200, unrelated.text
assert _json_object(unrelated.content) == _chat_caller_response(_success(unrelated_marker), rig.model), (
unrelated.text
)
assert (
_matching_marker_span(rig.destination, unrelated_marker)["litellm.metadata.trace_marker"]
== unrelated_marker
)
def test_arize_otel_v2_d7_missing_space_id_is_stable(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "d7-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _success(marker)
with _rig(
gateway,
tmp_path,
upstream,
remove_environment=("ARIZE_SPACE_ID",),
disabled_environment=("ARIZE_SPACE_ID",),
) as rig:
response: Final = _call(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_success(marker), rig.model), response.text
assert _matching_marker_span(rig.destination, marker)["litellm.metadata.trace_marker"] == marker
def test_arize_otel_v2_e1_uncached_request_exports_one_llm_span(gateway: Gateway, tmp_path: Path) -> None:
markers: Final = tuple("e1-" + uuid.uuid4().hex for _ in range(3))
def upstream(request: Request) -> Reply:
marker: Final = _chat_request_marker(request)
assert marker in markers, request
_assert_chat_request(request, messages=[{"role": "user", "content": marker}])
return _success(marker)
with _rig(gateway, tmp_path, upstream) as rig:
responses: Final = tuple(_call(rig.proxy, rig.model, marker, prompt=marker) for marker in markers)
assert all(response.status_code == 200 for response in responses), tuple(
response.text for response in responses
)
assert tuple(_json_object(response.content) for response in responses) == tuple(
_chat_caller_response(_success(marker), rig.model) for marker in markers
), responses
spans: Final = _collect_marker_spans(rig.destination, markers)
assert len(spans) == len(markers), spans
spans_by_id: Final = {span["gen_ai.response.id"]: span for span in spans}
assert frozenset(spans_by_id) == frozenset(markers), spans
assert all(
spans_by_id[marker]["gen_ai.response.id"] in response.text
for marker, response in zip(markers, responses, strict=True)
), spans
def test_arize_otel_v2_e2_concurrent_unique_markers(gateway: Gateway, tmp_path: Path) -> None:
markers: Final = tuple("e2-" + uuid.uuid4().hex for _ in range(20))
def upstream(request: Request) -> Reply:
marker: Final = _chat_request_marker(request)
assert marker.startswith("e2-"), request
_assert_chat_request(request, messages=[{"role": "user", "content": marker}])
return _success(marker)
with _rig(gateway, tmp_path, upstream) as rig:
with ThreadPoolExecutor(max_workers=20) as executor:
responses: Final = tuple(
executor.map(
lambda marker: _call(rig.proxy, rig.model, marker, prompt=marker),
markers,
)
)
assert all(response.status_code == 200 for response in responses), tuple(
response.text for response in responses
)
assert tuple(_json_object(response.content) for response in responses) == tuple(
_chat_caller_response(_success(marker), rig.model) for marker in markers
), responses
spans: Final = _collect_marker_spans(rig.destination, markers)
assert len(spans) == len(markers), spans
assert tuple(
_json_object(next(span for span in spans if marker in span.values())["metadata"].encode())
for marker in markers
) == tuple({"trace_marker": marker} for marker in markers), spans
assert (
tuple(
next(span for span in spans if marker in span.values())["litellm.metadata.trace_marker"]
for marker in markers
)
== markers
), spans
@pytest.mark.parametrize(
"shape",
("object-arguments", "missing-name", "non-dict-call", "null-tool-calls", "integer-id"),
)
def test_arize_otel_v2_d3_malformed_tool_calls_are_normalized(shape: str, gateway: Gateway, tmp_path: Path) -> None:
marker: Final = f"d3-{shape}-" + uuid.uuid4().hex
def malformed_reply() -> Reply:
call: Final = {
"id": 17 if shape == "integer-id" else f"call_{marker}",
"type": "function",
"function": {
**({} if shape == "missing-name" else {"name": "lookup_weather"}),
"arguments": {"city": "Paris"} if shape == "object-arguments" else '{"city": "Paris"}',
},
}
tool_calls: Final[JsonValue] = (
None if shape == "null-tool-calls" else ["not-a-call"] if shape == "non-dict-call" else [call]
)
message: Final = {
"role": "assistant",
"content": None,
"tool_calls": tool_calls,
}
return Reply(
body=json.dumps(
{
"id": marker,
"object": "chat.completion",
"created": 1,
"model": "gpt-4o-mini",
"choices": [{"index": 0, "finish_reason": "tool_calls", "message": message}],
"usage": {"prompt_tokens": 10, "completion_tokens": 2, "total_tokens": 12},
}
).encode()
)
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return malformed_reply()
with _rig(gateway, tmp_path, upstream) as rig:
proxy_log: Final = rig.owned.log
response: Final = _call(rig.proxy, rig.model, marker)
if shape == "non-dict-call":
assert response.status_code == 400, response.text
error_body: Final = _json_object(response.content)
assert frozenset(error_body) == frozenset({"error"}), error_body
error: Final = _json_object_value(error_body["error"])
assert frozenset(error) == frozenset({"type", "code", "param", "message"}), error
assert error["type"] == "invalid_request_error", error
assert error["code"] == "400", error
assert error["param"] is None, error
message: Final = error["message"]
assert isinstance(message, str), error
assert "AttributeError: 'str' object has no attribute 'get'" in message, error
spans: Final = tuple(_spans(rig.destination.drain()))
assert all(not any(".tool_calls." in key for key in attributes) for attributes in spans), spans
readiness: Final = rig.proxy.client.get("/health/readiness")
assert readiness.status_code == 200, readiness.text
assert _json_object(readiness.content) == {"status": "healthy", "db": "connected"}, readiness.text
else:
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(malformed_reply(), rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
indexed_prefix: Final = "llm.output_messages.0.message.tool_calls.0.tool_call."
expected_fields: Final = (
("id", "function.name", "function.arguments"),
("id", "function.arguments"),
(),
("function.name", "function.arguments"),
)[("object-arguments", "missing-name", "null-tool-calls", "integer-id").index(shape)]
expected_keys: Final = frozenset(indexed_prefix + field for field in expected_fields)
observed_keys: Final = frozenset(key for key in attributes if ".tool_calls." in key)
assert observed_keys == expected_keys, attributes
if shape == "object-arguments":
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.id"] == f"call_{marker}", (
attributes
)
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.name"] == "lookup_weather"
), attributes
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments"]
== '{"city": "Paris"}'
), attributes
elif shape == "missing-name":
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.id"] == f"call_{marker}", (
attributes
)
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments"]
== '{"city": "Paris"}'
), attributes
elif shape == "integer-id":
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.name"] == "lookup_weather"
), attributes
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments"]
== '{"city": "Paris"}'
), attributes
assert _json_messages(attributes["output.value"]) == _json_messages(
_chat_output_value(malformed_reply())
), attributes
assert "Exception while exporting Span batch" not in proxy_log.read_text(), proxy_log.read_text()
@pytest.mark.parametrize("shape", ("empty", "null", "missing"))
def test_arize_otel_v2_e3_empty_or_missing_tool_calls_never_indexed(
shape: str, gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = f"e3-{shape}-" + uuid.uuid4().hex
message: Final = (
{"role": "assistant", "content": None, "tool_calls": []}
if shape == "empty"
else {"role": "assistant", "content": None, "tool_calls": None}
if shape == "null"
else {"role": "assistant", "content": None}
)
expected_response: Final[dict[str, JsonValue]] = {
"id": marker,
"object": "chat.completion",
"created": 1,
"model": "gpt-4o-mini",
"choices": [{"index": 0, "finish_reason": "stop", "message": message}],
"usage": {"prompt_tokens": 10, "completion_tokens": 2, "total_tokens": 12},
}
expected_reply: Final = Reply(body=json.dumps(expected_response).encode())
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return expected_reply
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = _call(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
response_body: Final = _json_object(response.content)
assert response_body == _chat_caller_response(expected_reply, rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
assert not any(".tool_calls." in key for key in attributes), attributes
assert "tool_calls" not in _json_messages(attributes["output.value"])[0], attributes

View file

@ -0,0 +1,926 @@
from __future__ import annotations
import asyncio
import json
import uuid
from itertools import chain
from pathlib import Path
from typing import Final
import anthropic
import httpx
import openai
import pytest
from _openinference_support import (
CHAT_TOOLS,
RESPONSES_TOOLS,
Rig,
_anthropic_stream_response,
_assert_chat_request,
_assert_messages_request,
_assert_responses_request,
_assert_tool_span,
_chat_cache_hit_caller_stream,
_chat_caller_response,
_chat_caller_stream,
_chat_plain_response,
_chat_request_marker,
_chat_response,
_chat_stream_response,
_chat_tool_call,
_json_messages,
_json_object,
_llm_spans_through_markers,
_matching_marker_span,
_messages_caller_response,
_messages_caller_stream,
_messages_caller_stream_response,
_normalize_chat_caller_stream,
_normalize_responses_caller_body,
_normalize_responses_caller_stream,
_response_tool_calls,
_responses_caller_response,
_responses_caller_stream,
_responses_response,
_responses_stream_response,
_rig,
)
from integration._support.client import Gateway, object_value, string_value
from integration._support.wire import Reply, Request
from pydantic import JsonValue
def _chat_upstream(request: Request) -> None:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
def _responses_upstream(request: Request, marker: str, *, stream: bool = False) -> None:
_assert_responses_request(request, marker=marker, stream=stream)
def _messages_upstream(request: Request, marker: str, *, stream: bool = False) -> None:
_assert_messages_request(request, marker=marker, stream=stream)
def _messages_response(identity: str) -> Reply:
return Reply(
body=json.dumps(
{
"id": identity,
"type": "message",
"role": "assistant",
"model": "claude-opus-5-5",
"content": [
{
"type": "tool_use",
"id": "call_" + identity,
"name": "lookup_weather",
"input": {"city": "Paris"},
}
],
"stop_reason": "tool_use",
"stop_sequence": None,
"usage": {"input_tokens": 11, "output_tokens": 4},
}
).encode()
)
def _assert_tool_span_for_marker(attributes: dict[str, str], marker: str, *, content: str | None = None) -> None:
_assert_tool_span(
attributes,
marker=marker,
output=[
{
"role": "assistant",
"content": content,
"tool_calls": [
{
"id": "call_" + marker,
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
}
],
calls=[("call_" + marker, "lookup_weather", {"city": "Paris"})],
metadata={"trace_marker": marker},
baggage={"trace_marker": marker},
)
def _chat_request(
proxy: Gateway, model: str, marker: str, *, stream: bool = False, no_cache: bool = True
) -> httpx.Response:
return proxy.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": "weather in Paris?"}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": {"trace_marker": marker},
**({"stream": True} if stream else {}),
**({"cache": {"no-cache": True}} if no_cache else {}),
},
)
def _responses_request(proxy: Gateway, model: str, marker: str) -> httpx.Response:
return proxy.request(
"POST",
"/v1/responses",
{
"model": model,
"input": "weather in Paris?",
"tools": RESPONSES_TOOLS,
"tool_choice": {"type": "function", "name": "lookup_weather"},
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
)
def _messages_request(proxy: Gateway, model: str, marker: str) -> httpx.Response:
return proxy.request(
"POST",
"/v1/messages",
{
"model": model,
"max_tokens": 64,
"messages": [{"role": "user", "content": "weather in Paris?"}],
"tools": [
{
"name": "lookup_weather",
"description": "Get weather",
"input_schema": {
"type": "object",
"properties": {"city": {"type": "string"}},
},
}
],
"tool_choice": {"type": "auto"},
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
)
def _openai_client(proxy: Gateway) -> openai.OpenAI:
return openai.OpenAI(base_url=str(proxy.client.base_url) + "/v1", api_key=proxy.key, max_retries=0)
def _async_openai_client(proxy: Gateway) -> openai.AsyncOpenAI:
return openai.AsyncOpenAI(base_url=str(proxy.client.base_url) + "/v1", api_key=proxy.key, max_retries=0)
def test_arize_otel_v2_a1_chat_sync_sdk(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a1-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_chat_upstream(request)
body: Final = _json_object(request.body)
assert body == {
"messages": [{"role": "user", "content": "weather in Paris?"}],
"model": "gpt-4o-mini",
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"tools": CHAT_TOOLS,
}, body
return _chat_response(marker)
with _rig(gateway, tmp_path, upstream) as rig:
client: Final = _openai_client(rig.proxy)
response: Final = client.chat.completions.create(
model=rig.model,
messages=[{"role": "user", "content": "weather in Paris?"}],
tools=CHAT_TOOLS,
tool_choice={"type": "function", "function": {"name": "lookup_weather"}},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": True}},
)
assert response.id == marker, response
assert response.model_dump(mode="json", exclude_unset=True) == _chat_caller_response(
_chat_response(marker), rig.model
), response
calls: Final = response.choices[0].message.tool_calls
assert calls is not None and len(calls) == 1, response
assert (calls[0].id, calls[0].function.name, calls[0].function.arguments) == (
f"call_{marker}",
"lookup_weather",
'{"city": "Paris"}',
), response
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
def test_arize_otel_v2_a2_chat_async_sdk(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a2-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_chat_upstream(request)
return _chat_response(marker)
async def call() -> None:
with _rig(gateway, tmp_path, upstream) as rig:
client: Final = _async_openai_client(rig.proxy)
response: Final = await client.chat.completions.create(
model=rig.model,
messages=[{"role": "user", "content": "weather in Paris?"}],
tools=CHAT_TOOLS,
tool_choice={"type": "function", "function": {"name": "lookup_weather"}},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": True}},
)
assert response.id == marker, response
assert response.model_dump(mode="json", exclude_unset=True) == _chat_caller_response(
_chat_response(marker), rig.model
), response
calls: Final = response.choices[0].message.tool_calls
assert calls is not None and len(calls) == 1, response
assert (calls[0].id, calls[0].function.name, calls[0].function.arguments) == (
f"call_{marker}",
"lookup_weather",
'{"city": "Paris"}',
), response
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
asyncio.run(call())
def test_arize_otel_v2_a5_responses_sync_sdk(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a5-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_responses_upstream(request, marker)
return _responses_response(marker)
with _rig(gateway, tmp_path, upstream) as rig:
client: Final = _openai_client(rig.proxy)
response: Final = client.responses.create(
model=rig.model,
input="weather in Paris?",
tools=RESPONSES_TOOLS,
tool_choice={"type": "function", "name": "lookup_weather"},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": True}},
)
assert response.id.startswith("resp_"), response
assert _normalize_responses_caller_body(response.model_dump(mode="json", exclude_unset=True)) == (
_responses_caller_response(_responses_response(marker), rig.model)
), response
assert (response.status, response.model) == ("completed", rig.model), response
assert (
response.output[0].type,
response.output[0].call_id,
response.output[0].name,
response.output[0].arguments,
) == ("function_call", f"call_{marker}", "lookup_weather", '{"city": "Paris"}'), response
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
def test_arize_otel_v2_a6_responses_async_streaming_sdk(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a6-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_responses_upstream(request, marker, stream=True)
return _responses_stream_response(marker, (_chat_tool_call(marker),))
async def call() -> None:
with _rig(gateway, tmp_path, upstream) as rig:
client: Final = _async_openai_client(rig.proxy)
stream: Final = await client.responses.create(
model=rig.model,
input="weather in Paris?",
tools=RESPONSES_TOOLS,
tool_choice={"type": "function", "name": "lookup_weather"},
stream=True,
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": True}},
)
events: Final = tuple([event async for event in stream])
assert events[-1].type == "response.completed", events
assert events[-1].response.id.startswith("resp_"), events[-1]
assert _normalize_responses_caller_stream(
tuple(event.model_dump(mode="json", exclude_unset=True) for event in events)
) == _responses_caller_stream(_responses_stream_response(marker, (_chat_tool_call(marker),)), rig.model), (
events
)
call: Final = events[-1].response.output[0]
assert (call.type, call.call_id, call.name, call.arguments) == (
"function_call",
f"call_{marker}",
"lookup_weather",
'{"city": "Paris"}',
), events[-1]
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
asyncio.run(call())
def test_arize_otel_v2_a7_messages_sync_sdk(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a7-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_messages_upstream(request, marker)
return _messages_response(marker)
with _rig(gateway, tmp_path, upstream, model_name="anthropic/claude-opus-5-5", api_base_suffix="") as rig:
client: Final = anthropic.Anthropic(
base_url=str(rig.proxy.client.base_url), api_key=rig.proxy.key, max_retries=0
)
response: Final = client.messages.create(
model=rig.model,
max_tokens=64,
messages=[{"role": "user", "content": "weather in Paris?"}],
tools=[
{
"name": "lookup_weather",
"description": "Get weather",
"input_schema": {"type": "object", "properties": {"city": {"type": "string"}}},
}
],
tool_choice={"type": "auto"},
metadata={"trace_marker": marker},
)
assert response.id == marker, response
assert response.model_dump(mode="json", exclude_unset=True) == _messages_caller_response(
_messages_response(marker), rig.model
), response
call: Final = response.content[0]
assert (call.type, call.id, call.name, call.input) == (
"tool_use",
f"call_{marker}",
"lookup_weather",
{"city": "Paris"},
), response
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
def test_arize_otel_v2_a3_chat_streaming(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a3-" + uuid.uuid4().hex
stream_reply: Final = _chat_stream_response(marker, (_chat_tool_call(marker),), include_usage=False)
def upstream(request: Request) -> Reply:
_assert_chat_request(
request,
messages=[{"role": "user", "content": "weather in Paris?"}],
stream=True,
stream_options={"include_usage": False},
)
return stream_reply
with _rig(gateway, tmp_path, upstream, general_settings={"always_include_stream_usage": False}) as rig:
client: Final = _openai_client(rig.proxy)
stream: Final = client.chat.completions.create(
model=rig.model,
messages=[{"role": "user", "content": "weather in Paris?"}],
tools=CHAT_TOOLS,
tool_choice={"type": "function", "function": {"name": "lookup_weather"}},
stream=True,
stream_options={"include_usage": False},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": True}},
)
chunks: Final = tuple(stream)
assert chunks[0].id == marker and chunks[-1].id == marker, chunks
assert _normalize_chat_caller_stream(
tuple(chunk.model_dump(mode="json", exclude_unset=True) for chunk in chunks)
) == _chat_caller_stream(stream_reply, rig.model), chunks
assert chunks[-1].choices[0].finish_reason == "tool_calls", chunks
tool_call_deltas: Final = tuple(
chain.from_iterable(chunk.choices[0].delta.tool_calls or () for chunk in chunks)
)
assert len(tool_call_deltas) == 2, chunks
assert (
tool_call_deltas[0].id,
tool_call_deltas[0].function.name,
tool_call_deltas[1].function.arguments,
) == (f"call_{marker}", "lookup_weather", '{"city": "Paris"}'), chunks
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
def test_arize_otel_v2_a4_chat_async_streaming(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a4-" + uuid.uuid4().hex
stream_reply: Final = _chat_stream_response(marker, (_chat_tool_call(marker),), include_usage=False)
def upstream(request: Request) -> Reply:
_assert_chat_request(
request,
messages=[{"role": "user", "content": "weather in Paris?"}],
stream=True,
stream_options={"include_usage": False},
)
return stream_reply
async def call() -> None:
with _rig(gateway, tmp_path, upstream, general_settings={"always_include_stream_usage": False}) as rig:
client: Final = _async_openai_client(rig.proxy)
stream: Final = await client.chat.completions.create(
model=rig.model,
messages=[{"role": "user", "content": "weather in Paris?"}],
tools=CHAT_TOOLS,
tool_choice={"type": "function", "function": {"name": "lookup_weather"}},
stream=True,
stream_options={"include_usage": False},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": True}},
)
chunks: Final = tuple([chunk async for chunk in stream])
assert chunks[0].id == marker and chunks[-1].id == marker, chunks
assert _normalize_chat_caller_stream(
tuple(chunk.model_dump(mode="json", exclude_unset=True) for chunk in chunks)
) == _chat_caller_stream(stream_reply, rig.model), chunks
assert chunks[-1].choices[0].finish_reason == "tool_calls", chunks
tool_call_deltas: Final = tuple(
chain.from_iterable(chunk.choices[0].delta.tool_calls or () for chunk in chunks)
)
assert len(tool_call_deltas) == 2, chunks
assert (
tool_call_deltas[0].id,
tool_call_deltas[0].function.name,
tool_call_deltas[1].function.arguments,
) == (f"call_{marker}", "lookup_weather", '{"city": "Paris"}'), chunks
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
asyncio.run(call())
def test_arize_otel_v2_a8_messages_streaming(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a8-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_messages_upstream(request, marker, stream=True)
return _anthropic_stream_response(marker)
async def call() -> None:
with _rig(gateway, tmp_path, upstream, model_name="anthropic/claude-opus-5-5", api_base_suffix="") as rig:
client: Final = anthropic.AsyncAnthropic(
base_url=str(rig.proxy.client.base_url), api_key=rig.proxy.key, max_retries=0
)
async with client.messages.stream(
model=rig.model,
max_tokens=64,
messages=[{"role": "user", "content": "weather in Paris?"}],
tools=[
{
"name": "lookup_weather",
"description": "Get weather",
"input_schema": {"type": "object", "properties": {"city": {"type": "string"}}},
}
],
tool_choice={"type": "auto"},
metadata={"trace_marker": marker},
) as stream:
events: Final = tuple([event async for event in stream])
response: Final = await stream.get_final_message()
assert events[-1].type == "message_stop", events
assert response.id == marker, response
assert response.model_dump(mode="json", exclude_unset=True) == _messages_caller_stream_response(
_messages_response(marker), rig.model
), response
assert tuple(event.model_dump(mode="json", exclude_unset=True) for event in events) == (
_messages_caller_stream(
_anthropic_stream_response(marker),
rig.model,
final_message=_messages_caller_stream_response(_messages_response(marker), rig.model),
)
), events
call: Final = response.content[0]
assert (call.type, call.id, call.name, call.input) == (
"tool_use",
f"call_{marker}",
"lookup_weather",
{"city": "Paris"},
), response
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker, content="")
asyncio.run(call())
def test_arize_otel_v2_llm_span_carries_openinference_tool_calls_and_metadata(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a9-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_chat_upstream(request)
return _chat_response(marker)
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = _chat_request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker), rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
_assert_tool_span_for_marker(attributes, marker)
def test_arize_otel_v2_responses_span_carries_openinference_tool_calls_and_metadata(
gateway: Gateway, tmp_path: Path
) -> None:
marker: Final = "a10-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_responses_upstream(request, marker)
return _responses_response(marker)
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = _responses_request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
response_body: Final = _json_object(response.content)
assert _normalize_responses_caller_body(response_body) == _responses_caller_response(
_responses_response(marker), rig.model
), response.text
_assert_tool_span_for_marker(_matching_marker_span(rig.destination, marker), marker)
def test_arize_otel_v2_a11_parallel_output_tool_calls(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a11-" + uuid.uuid4().hex
calls: Final = _response_tool_calls(marker, ("Paris", "Berlin"))
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=[{"role": "user", "content": "weather in Paris?"}])
return _chat_response(marker, calls)
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = _chat_request(rig.proxy, rig.model, marker)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker, calls), rig.model), (
response.text
)
attributes: Final = _matching_marker_span(rig.destination, marker)
expected_calls: Final = (
("call_" + marker + "-Paris", "lookup_weather", {"city": "Paris"}),
("call_" + marker + "-Berlin", "lookup_weather", {"city": "Berlin"}),
)
_assert_tool_span(
attributes,
marker=marker,
output=[{"role": "assistant", "content": None, "tool_calls": calls}],
calls=expected_calls,
metadata={"trace_marker": marker},
baggage={"trace_marker": marker},
)
def test_arize_otel_v2_a12_plain_text_has_metadata_without_tool_calls(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a12-" + uuid.uuid4().hex
def upstream(request: Request) -> Reply:
_assert_chat_request(
request,
messages=[{"role": "user", "content": "weather in Paris?"}],
include_tools=False,
)
return _chat_plain_response(marker, "The weather is clear")
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = rig.proxy.request(
"POST",
"/v1/chat/completions",
{
"model": rig.model,
"messages": [{"role": "user", "content": "weather in Paris?"}],
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(
_chat_plain_response(marker, "The weather is clear"), rig.model
), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
assert _json_object(attributes["metadata"].encode()) == {"trace_marker": marker}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
assert not any(".tool_calls." in key for key in attributes), attributes
assert _json_messages(attributes["output.value"]) == [
{"role": "assistant", "content": "The weather is clear"}
], attributes
def test_arize_otel_v2_a13_multiturn_input_tool_calls(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a13-" + uuid.uuid4().hex
messages: Final = [
{"role": "user", "content": "weather in Paris?"},
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": "call-prior",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
},
{"role": "tool", "tool_call_id": "call-prior", "content": '{"temperature": 20}'},
]
upstream_messages: Final = [
{key: value for key, value in message.items() if value is not None} for message in messages
]
def upstream(request: Request) -> Reply:
_assert_chat_request(request, messages=upstream_messages)
return _chat_response(marker)
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = rig.proxy.request(
"POST",
"/v1/chat/completions",
{
"model": rig.model,
"messages": messages,
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(_chat_response(marker), rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
_assert_tool_span_for_marker(attributes, marker)
assert not any(key.startswith("llm.input_messages.") and ".tool_calls." in key for key in attributes), (
attributes
)
assert tuple(attributes[f"llm.input_messages.{index}.message.role"] for index in range(3)) == (
"user",
"assistant",
"tool",
), attributes
assert tuple(attributes[f"llm.input_messages.{index}.message.content"] for index in (0, 2)) == (
"weather in Paris?",
'{"temperature": 20}',
), attributes
expected_input_value: Final = [
messages[0],
messages[1],
{"role": "tool", "content": '{"temperature": 20}'},
]
assert _json_messages(attributes["input.value"]) == expected_input_value, attributes
def test_arize_otel_v2_a14_two_choices_each_with_tool_calls(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "a14-" + uuid.uuid4().hex
first_call: Final = _response_tool_calls(marker, ("Paris",))[0]
second_call: Final = _response_tool_calls(marker, ("Berlin",))[0]
def choice(index: int, call: dict[str, JsonValue]) -> dict[str, JsonValue]:
return {
"index": index,
"finish_reason": "tool_calls",
"message": {"role": "assistant", "content": None, "tool_calls": [call]},
}
expected_response: Final = Reply(
body=json.dumps(
{
"id": marker,
"object": "chat.completion",
"created": 1,
"model": "gpt-4o-mini",
"choices": [choice(0, first_call), choice(1, second_call)],
"usage": {"prompt_tokens": 11, "completion_tokens": 4, "total_tokens": 15},
}
).encode()
)
def upstream(request: Request) -> Reply:
_assert_chat_request(
request,
messages=[{"role": "user", "content": "weather in Paris and Berlin?"}],
n=2,
)
return expected_response
with _rig(gateway, tmp_path, upstream) as rig:
response: Final = rig.proxy.request(
"POST",
"/v1/chat/completions",
{
"model": rig.model,
"messages": [{"role": "user", "content": "weather in Paris and Berlin?"}],
"tools": CHAT_TOOLS,
"tool_choice": {"type": "function", "function": {"name": "lookup_weather"}},
"n": 2,
"metadata": {"trace_marker": marker},
"cache": {"no-cache": True},
},
)
assert response.status_code == 200, response.text
assert _json_object(response.content) == _chat_caller_response(expected_response, rig.model), response.text
attributes: Final = _matching_marker_span(rig.destination, marker)
assert _json_messages(attributes["output.value"]) == [
{"role": "assistant", "content": None, "tool_calls": [first_call]},
{"role": "assistant", "content": None, "tool_calls": [second_call]},
], attributes
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.id"] == str(first_call["id"]), (
attributes
)
assert attributes["llm.output_messages.1.message.tool_calls.0.tool_call.id"] == str(second_call["id"]), (
attributes
)
assert attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.name"] == "lookup_weather", (
attributes
)
assert (
attributes["llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments"] == '{"city": "Paris"}'
), attributes
assert attributes["llm.output_messages.1.message.tool_calls.0.tool_call.function.name"] == "lookup_weather", (
attributes
)
assert (
attributes["llm.output_messages.1.message.tool_calls.0.tool_call.function.arguments"]
== '{"city": "Berlin"}'
), attributes
assert _json_object(attributes["metadata"].encode()) == {"trace_marker": marker}, attributes
assert attributes["litellm.metadata.trace_marker"] == marker, attributes
def _cache_call(
rig: Rig, surface: str, marker: str, *, cache_hit: bool = False
) -> tuple[tuple[str, str, str], str, httpx.Headers]:
match surface:
case "chat":
client: Final = _openai_client(rig.proxy)
raw: Final = client.chat.completions.with_raw_response.create(
model=rig.model,
messages=[{"role": "user", "content": marker}],
tools=CHAT_TOOLS,
tool_choice={"type": "function", "function": {"name": "lookup_weather"}},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": False}},
)
response: Final = raw.parse()
assert response.model_dump(mode="json", exclude_unset=True) == _chat_caller_response(
_chat_response(marker), rig.model
), response
assert response.choices[0].message.tool_calls is not None, response
call: Final = response.choices[0].message.tool_calls[0]
return (call.id, call.function.name, call.function.arguments), response.id, raw.headers
case "chat-stream":
client: Final = _openai_client(rig.proxy)
with client.chat.completions.with_streaming_response.create(
model=rig.model,
messages=[{"role": "user", "content": marker}],
tools=CHAT_TOOLS,
tool_choice={"type": "function", "function": {"name": "lookup_weather"}},
stream=True,
stream_options={"include_usage": False},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": False}},
) as raw:
chunks: Final = tuple(raw.parse())
expected_chunks: Final = (
_chat_cache_hit_caller_stream(marker, rig.model, (_chat_tool_call(marker),))
if cache_hit
else _chat_caller_stream(
_chat_stream_response(marker, (_chat_tool_call(marker),), include_usage=False), rig.model
)
)
assert (
_normalize_chat_caller_stream(
tuple(chunk.model_dump(mode="json", exclude_unset=True) for chunk in chunks)
)
== expected_chunks
), chunks
calls: Final = tuple(chain.from_iterable(chunk.choices[0].delta.tool_calls or () for chunk in chunks))
if cache_hit:
assert len(calls) == 1, chunks
call: Final = calls[0]
assert (
call.id is not None and call.function.name is not None and call.function.arguments is not None
), chunks
return (
(call.id, call.function.name, call.function.arguments),
chunks[0].id,
raw.headers,
)
assert len(calls) == 2, chunks
call: Final = calls[0]
arguments_call: Final = calls[1]
assert (
call.id is not None
and call.function.name is not None
and arguments_call.function.arguments is not None
), chunks
return (
(call.id, call.function.name, arguments_call.function.arguments),
chunks[0].id,
raw.headers,
)
case "responses":
client: Final = _openai_client(rig.proxy)
raw: Final = client.responses.with_raw_response.create(
model=rig.model,
input=marker,
tools=RESPONSES_TOOLS,
tool_choice={"type": "function", "name": "lookup_weather"},
extra_body={"metadata": {"trace_marker": marker}, "cache": {"no-cache": False}},
)
response: Final = raw.parse()
assert _normalize_responses_caller_body(response.model_dump(mode="json", exclude_unset=True)) == (
_responses_caller_response(_responses_response(marker), rig.model)
), response
call: Final = response.output[0]
assert call.type == "function_call", response
return (call.call_id, call.name, call.arguments), response.id, raw.headers
case "messages":
client: Final = anthropic.Anthropic(
base_url=str(rig.proxy.client.base_url), api_key=rig.proxy.key, max_retries=0
)
raw: Final = client.messages.with_raw_response.create(
model=rig.model,
max_tokens=64,
messages=[{"role": "user", "content": marker}],
tools=[
{
"name": "lookup_weather",
"description": "Get weather",
"input_schema": {"type": "object", "properties": {"city": {"type": "string"}}},
}
],
tool_choice={"type": "auto"},
extra_body={"cache": {"no-cache": False}, "metadata": {"trace_marker": marker}},
)
response: Final = raw.parse()
assert response.model_dump(mode="json", exclude_unset=True) == _messages_caller_response(
_messages_response(marker), rig.model
), response
call: Final = response.content[0]
assert call.type == "tool_use", response
return (call.id, call.name, json.dumps(call.input)), response.id, raw.headers
case _:
raise AssertionError(f"Unknown cache surface: {surface}")
@pytest.mark.parametrize("surface", ("chat", "chat-stream", "responses", "messages"))
def test_arize_otel_v2_a_cache(surface: str, gateway: Gateway, tmp_path: Path) -> None:
marker: Final = f"a-cache-{surface}-" + uuid.uuid4().hex
sentinel: Final = f"{marker}-sentinel"
def upstream(request: Request) -> Reply:
body: Final = _json_object(request.body)
request_marker: Final = (
_chat_request_marker(request)
if surface in ("chat", "chat-stream", "messages")
else string_value(object_value(body["metadata"])["trace_marker"])
)
assert request_marker in (marker, sentinel), request
if surface in ("chat", "chat-stream"):
_assert_chat_request(
request,
messages=[{"role": "user", "content": request_marker}],
stream=True if surface == "chat-stream" else None,
stream_options={"include_usage": False} if surface == "chat-stream" else None,
)
return (
_chat_stream_response(request_marker, (_chat_tool_call(request_marker),), include_usage=False)
if surface == "chat-stream"
else _chat_response(request_marker)
)
if surface == "responses":
assert body.get("metadata") == {"trace_marker": request_marker}, body
_assert_responses_request(request, marker=request_marker, input_value=request_marker)
return _responses_response(request_marker)
_assert_messages_request(request, marker=request_marker, prompt=request_marker)
return _messages_response(request_marker)
with _rig(
gateway,
tmp_path,
upstream,
general_settings={"always_include_stream_usage": False} if surface == "chat-stream" else None,
model_name="anthropic/claude-opus-5-5" if surface == "messages" else "gpt-4o-mini",
api_base_suffix="" if surface == "messages" else "/v1",
workers=1,
) as rig:
expected_arguments: Final = '{"city": "Paris"}'
first: Final = _cache_call(rig, surface, marker)
assert first[0] == (f"call_{marker}", "lookup_weather", expected_arguments), first
assert first[1].startswith("resp_") if surface == "responses" else first[1] == marker, first
assert not first[2].get("x-litellm-cache-key"), first[2]
first_span: Final = _matching_marker_span(rig.destination, marker)
_assert_tool_span_for_marker(first_span, marker)
second: Final = _cache_call(rig, surface, marker, cache_hit=True)
assert second[0] == first[0], second
assert second[1].startswith("resp_") if surface == "responses" else second[1] == first[1], second
forwarded: Final = tuple(
request for request in rig.provider.drain() if request.method == "POST" and marker.encode() in request.body
)
assert len(forwarded) == 1, forwarded
if surface == "messages":
assert not second[2].get("x-litellm-cache-key"), second[2]
else:
assert second[2].get("x-litellm-cache-key"), second[2]
forwarded_body: Final = _json_object(forwarded[0].body)
if surface in ("chat", "chat-stream"):
assert "metadata" not in forwarded_body, forwarded[0]
elif surface == "responses":
assert forwarded_body["metadata"] == {"trace_marker": marker}, forwarded[0]
else:
assert forwarded_body["metadata"] == {}, forwarded[0]
sentinel_response: Final = _cache_call(rig, surface, sentinel)
assert not sentinel_response[2].get("x-litellm-cache-key"), sentinel_response[2]
sentinel_forwarded: Final = tuple(
request
for request in rig.provider.drain()
if request.method == "POST" and sentinel.encode() in request.body
)
assert len(sentinel_forwarded) == 1, sentinel_forwarded
spans: Final = _llm_spans_through_markers(rig.destination, (sentinel,))
sentinel_spans: Final = tuple(
attributes for attributes in spans if attributes.get("litellm.metadata.trace_marker") == sentinel
)
assert len(sentinel_spans) == 1, spans
_assert_tool_span_for_marker(sentinel_spans[0], sentinel)
assert not any(attributes.get("litellm.metadata.trace_marker") == marker for attributes in spans), spans

View file

@ -7,22 +7,22 @@ import pytest
pytest.importorskip("opentelemetry")
from litellm.integrations.otel import ( # noqa: E402
GenAI,
HTTP,
GenAI,
LiteLLM,
OpenTelemetryV2Config,
promoted_baggage,
)
from litellm.integrations.otel.plumbing import context as ctx_mod # noqa: E402
from litellm.integrations.otel.plumbing import providers # noqa: E402
from litellm.integrations.otel.emitter import SpanEmitter # noqa: E402
from litellm.integrations.otel.model.baggage import BAGGAGE_PROMOTED_KEYS # noqa: E402
from litellm.integrations.otel.model.payloads import ( # noqa: E402
GuardrailSpanData,
LLMCallSpanData,
ServiceSpanData,
)
from litellm.integrations.otel.model.baggage import BAGGAGE_PROMOTED_KEYS # noqa: E402
from litellm.integrations.otel.model.spans import SpanRole # noqa: E402
from litellm.integrations.otel.plumbing import context as ctx_mod # noqa: E402
from litellm.integrations.otel.plumbing import providers # noqa: E402
def _payload():
@ -63,9 +63,7 @@ def test_identity_promoted_onto_every_span():
root = engine.start_span(SpanRole.PROXY_REQUEST, "POST /chat/completions", ctx)
root_ctx = ctx_mod.context_from_span(root, ctx)
engine.emit(SpanRole.LLM_CALL, data, parent_context=root_ctx)
engine.emit(
SpanRole.GUARDRAIL, GuardrailSpanData("presidio", status="success"), root_ctx
)
engine.emit(SpanRole.GUARDRAIL, GuardrailSpanData("presidio", status="success"), root_ctx)
engine.emit(SpanRole.SERVICE, ServiceSpanData("redis", call_type="set"), root_ctx)
root.end()
@ -161,9 +159,7 @@ def test_allowlisted_metadata_subkey_promoted_blob_excluded():
engine.emit(SpanRole.SERVICE, ServiceSpanData("redis", call_type="set"), ctx)
(span,) = exporter.get_finished_spans()
# allowlisted metadata sub-key is promoted
assert (
span.attributes.get(f"{LiteLLM.METADATA_PREFIX}user_api_key_org_id") == "org1"
)
assert span.attributes.get(f"{LiteLLM.METADATA_PREFIX}user_api_key_org_id") == "org1"
# non-allowlisted metadata is NOT promoted (no full-blob dumping)
assert all("private_note" not in k for k in span.attributes)
@ -208,6 +204,17 @@ def test_nested_metadata_key_promoted_under_caller_path():
assert not any(k.startswith(f"{LiteLLM.METADATA_PREFIX}requester_metadata") for k in span.attributes)
def test_llm_call_promoted_metadata_strips_requester_prefix_and_uses_allowlist():
payload = _payload()
payload["metadata"]["requester_metadata"] = {"trace_id": "trace-123"}
data = LLMCallSpanData.from_standard_logging_payload(
payload,
metadata_keys=("requester_metadata.trace_id", "user_api_key_org_id", "missing"),
)
assert data.promoted_metadata == {"trace_id": "trace-123", "user_api_key_org_id": "org1"}
assert LLMCallSpanData.from_standard_logging_payload(payload).promoted_metadata == {}
def test_http_attributes_never_promoted():
"""Even if http.* is present in baggage, the processor must not stamp it on
child spans (it belongs on the SERVER span only)."""
@ -228,9 +235,7 @@ def test_http_attributes_never_promoted():
def test_arbitrary_upstream_baggage_not_promoted():
engine, exporter = _engine_and_exporter()
ctx = ctx_mod.set_request_baggage(
{LiteLLM.TEAM_ID: "t1", "some.upstream.key": "leak"}
)
ctx = ctx_mod.set_request_baggage({LiteLLM.TEAM_ID: "t1", "some.upstream.key": "leak"})
engine.emit(SpanRole.SERVICE, ServiceSpanData("redis", call_type="set"), ctx)
(span,) = exporter.get_finished_spans()
assert span.attributes.get(LiteLLM.TEAM_ID) == "t1"

View file

@ -7,6 +7,7 @@ backends, so one trace lights up every configured destination.
import json
from collections.abc import Mapping
from itertools import chain
from typing import Final
import pytest
@ -19,6 +20,7 @@ from litellm.integrations.otel.mappers import (
WeaveMapper,
resolve_mappers,
)
from litellm.integrations.otel.mappers.openinference import fit_indexed_messages
from litellm.integrations.otel.model.payloads import (
EmbeddingOutput,
LLMCallSpanData,
@ -29,6 +31,7 @@ from litellm.integrations.otel.model.payloads import (
ToolDefinition,
)
from litellm.integrations.otel.model.trace_controls import TraceControls
from tests.unit.integrations.otel.test_otel_v2_sources_of_truth import _responses_payload
def _llm_call(**overrides):
@ -118,6 +121,232 @@ def test_openinference_multimodal_content_text_only():
assert attrs["llm.input_messages.0.message.content"] == "hi there"
def test_openinference_output_tool_calls_preserve_calls_in_attributes_and_value():
tool_calls: Final = [
{
"id": "call_paris",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
},
{
"id": "call_search",
"function": {"name": "search", "arguments": {"q": 1}},
"index": 0,
},
"ignored",
]
data: Final = _llm_call(
choices_out=(
{
"finish_reason": "tool_calls",
"message": {"role": "assistant", "content": None, "tool_calls": tool_calls},
},
)
)
attrs: Final = OpenInferenceMapper().map(data)
assert {key: value for key, value in attrs.items() if ".tool_calls." in key} == {
"llm.output_messages.0.message.tool_calls.0.tool_call.id": "call_paris",
"llm.output_messages.0.message.tool_calls.0.tool_call.function.name": "lookup_weather",
"llm.output_messages.0.message.tool_calls.0.tool_call.function.arguments": '{"city": "Paris"}',
"llm.output_messages.0.message.tool_calls.1.tool_call.id": "call_search",
"llm.output_messages.0.message.tool_calls.1.tool_call.function.name": "search",
"llm.output_messages.0.message.tool_calls.1.tool_call.function.arguments": '{"q": 1}',
}
assert json.loads(attrs["output.value"]) == [
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": "call_paris",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
},
{
"id": "call_search",
"type": "function",
"function": {"name": "search", "arguments": '{"q": 1}'},
},
],
}
]
def test_openinference_responses_tool_calls_are_emitted_as_output_attributes():
data: Final = LLMCallSpanData.from_standard_logging_payload(
_responses_payload(
[
{
"type": "function_call",
"call_id": "call_resp",
"name": "lookup_weather",
"arguments": '{"city": "Paris"}',
}
]
),
capture_content=True,
)
attrs: Final = OpenInferenceMapper().map(data)
tool_call: Final = "llm.output_messages.0.message.tool_calls.0.tool_call."
assert {key: value for key, value in attrs.items() if ".tool_calls." in key} == {
tool_call + "id": "call_resp",
tool_call + "function.name": "lookup_weather",
tool_call + "function.arguments": '{"city": "Paris"}',
}
assert json.loads(attrs["output.value"]) == [
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": "call_resp",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
}
]
def test_openinference_input_tool_calls_stay_in_value_only():
data: Final = _llm_call(
messages_in=(
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": "call_weather",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
},
)
)
attrs: Final = OpenInferenceMapper().map(data)
assert all(".tool_calls." not in key for key in attrs)
assert json.loads(attrs["input.value"]) == [
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": "call_weather",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
}
]
def test_openinference_output_tool_calls_do_not_shed_input_roles_under_budget():
message_groups: Final = tuple(
(
{"role": "user", "content": f"Question {index}"},
{
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": f"call_{index}",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
},
{"role": "tool", "tool_call_id": f"call_{index}", "content": f"Result {index}"},
)
for index in range(13)
)
messages_in: Final = tuple(chain.from_iterable(message_groups)) + ({"role": "user", "content": "Final request"},)
tools: Final = tuple(
ToolDefinition(name=name, description="Tool", parameters_json='{"type":"object"}')
for name in ("lookup_weather", "search", "get_location", "convert_units")
)
data: Final = _llm_call(
messages_in=messages_in,
tools=tools,
choices_out=(
{
"finish_reason": "tool_calls",
"message": {
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": "call_output",
"type": "function",
"function": {"name": "lookup_weather", "arguments": '{"city": "Paris"}'},
}
],
},
},
),
)
attrs: Final = fit_indexed_messages(OpenInferenceMapper().map(data), 128)
tool_call: Final = "llm.output_messages.0.message.tool_calls.0.tool_call."
assert {
f"llm.input_messages.{index}.message.role": attrs.get(f"llm.input_messages.{index}.message.role")
for index in range(40)
} == {f"llm.input_messages.{index}.message.role": message["role"] for index, message in enumerate(messages_in)}
assert {key: value for key, value in attrs.items() if ".tool_calls." in key} == {
tool_call + "id": "call_output",
tool_call + "function.name": "lookup_weather",
tool_call + "function.arguments": '{"city": "Paris"}',
}
def test_openinference_budget_sheds_trailing_output_tool_calls_before_the_message():
tool_calls: Final = tuple(
{
"id": f"call_{index}",
"type": "function",
"function": {
"name": "lookup_weather",
"arguments": f'{{"city": "C{index}"}}',
},
}
for index in range(60)
)
data: Final = _llm_call(
choices_out=(
{
"finish_reason": "tool_calls",
"message": {"role": "assistant", "content": None, "tool_calls": tool_calls},
},
)
)
mapped: Final = OpenInferenceMapper().map(data)
attrs: Final = fit_indexed_messages(mapped, len(mapped) - 34)
retained_tool_call_keys: Final = tuple(key for key in attrs if ".tool_calls." in key)
retained_tool_call_indices: Final = frozenset(int(key.split(".")[5]) for key in retained_tool_call_keys)
assert retained_tool_call_indices == frozenset(range(50))
assert len(retained_tool_call_keys) == 150
assert attrs["llm.output_messages.0.message.role"] == "assistant"
assert not any(key.startswith("llm.input_messages.0.") for key in attrs)
assert len(attrs) == len(mapped) - 34
assert json.loads(attrs["output.value"]) == [{"role": "assistant", "content": None, "tool_calls": list(tool_calls)}]
def test_openinference_plain_output_messages_keep_the_existing_value_shape():
attrs: Final = OpenInferenceMapper().map(_llm_call())
assert all(".tool_calls." not in key for key in attrs)
assert json.loads(attrs["output.value"]) == [{"role": "assistant", "content": "Sunny."}]
def test_openinference_metadata_contains_only_promoted_metadata():
attrs: Final = OpenInferenceMapper().map(
_llm_call(promoted_metadata={"trace_marker": "m", "user_api_key_alias": "k"})
)
assert json.loads(attrs["metadata"]) == {"trace_marker": "m", "user_api_key_alias": "k"}
assert "metadata" not in OpenInferenceMapper().map(_llm_call())
# --------------------------------------------------------------------------- #
# Langfuse
# --------------------------------------------------------------------------- #