test(otel): audit Arize OTel v2 OpenInference spans across endpoints, modes and chaos

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-29 23:41:17 +00:00
parent 6fd0390bef
commit 4b0d28f704
6 changed files with 3671 additions and 265 deletions

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,377 @@
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,
_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:
markers: Final = tuple("f1-" + uuid.uuid4().hex for _ in range(30))
surfaces: Final = ("chat", "responses", "messages")
calls: Final = tuple(
(marker, surfaces[index % len(surfaces)], index % 2 == 0) for index, marker in enumerate(markers)
)
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")
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,
environment={"OTEL_BSP_SCHEDULE_DELAY": "20000"},
destination_wire=stopped_destination,
) 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,
)
with ThreadPoolExecutor(max_workers=len(calls)) as executor:
futures: Final = tuple(
executor.submit(
_call,
rig.proxy,
messages_model if surface == "messages" else rig.model,
marker,
surface=surface,
stream=stream,
prompt=marker,
)
for marker, surface, stream in calls
)
responses: Final = tuple(
(marker, surface, stream, future.result(timeout=60))
for (marker, surface, stream), future in zip(calls, futures, strict=True)
)
for marker, surface, stream, response in responses:
model: Final = messages_model if surface == "messages" else rig.model
_assert_response(response, marker, surface, stream, 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)
)
spans: Final = _collect_marker_spans(recovered_destination, markers, timeout_seconds=90)
assert len(spans) == len(markers), spans
recorded: Final = tuple(span["litellm.metadata.trace_marker"] for span in spans)
assert len(recorded) == len(markers), recorded
assert frozenset(recorded) == frozenset(markers), recorded
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)
requests: Final = rig.destination.drain()
spans: Final = tuple(
attributes
for attributes in _spans(requests)
if attributes.get("openinference.span.kind") == "LLM"
and attributes.get("litellm.metadata.trace_marker") == marker
)
assert len(spans) == 1, spans
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,438 @@
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,
_spans,
)
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)
disabled_spans: Final = tuple(
attributes
for attributes in _spans(rig.destination.drain())
if attributes.get("openinference.span.kind") == "LLM"
)
assert disabled_spans == (), disabled_spans
@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