test(otel v2): integration audit cells for per-destination otel_internal_spans

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-23 18:31:02 +00:00
parent 3a835bbf3c
commit 0c3f546126
6 changed files with 1844 additions and 2 deletions

View file

@ -0,0 +1,229 @@
"""OTLP/HTTP trace sink: records exported spans and exposes them over a control API.
Accepts ``application/x-protobuf`` ``ExportTraceServiceRequest`` bodies and OTLP
``http/json`` bodies on any path. Tests read spans through ``recorded_spans`` and
steer the sink through ``configure``; the process can also be frozen with
``SIGSTOP``/``SIGCONT`` after reading its pid from ``/__pid``.
"""
from __future__ import annotations
import argparse
import json
import threading
import time
from dataclasses import dataclass
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from typing import Final
from urllib.parse import urlparse
import httpx
from pydantic import JsonValue
INTERNAL_MARKERS: Final = ("gen_ai.operation.name", "mcp.method.name", "litellm.guardrail_name")
def _proto_spans(body: bytes) -> list[dict[str, JsonValue]]:
from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceRequest
def scalar(value: object) -> JsonValue:
which: Final = value.WhichOneof("value") # type: ignore[attr-defined] # protobuf AnyValue
if which is None:
return None
raw: Final = getattr(value, which)
if which == "array_value":
return [scalar(item) for item in raw.values]
if which == "kvlist_value":
return {pair.key: scalar(pair.value) for pair in raw.values}
return raw
request: Final = ExportTraceServiceRequest()
request.ParseFromString(body)
return [
{
"trace_id": span.trace_id.hex(),
"span_id": span.span_id.hex(),
"parent_span_id": span.parent_span_id.hex(),
"kind": span.kind,
"name": span.name,
"attributes": {attribute.key: scalar(attribute.value) for attribute in span.attributes},
"resource": {attribute.key: scalar(attribute.value) for attribute in resource.resource.attributes},
}
for resource in request.resource_spans
for scope in resource.scope_spans
for span in scope.spans
]
def _json_spans(body: bytes) -> list[dict[str, JsonValue]]:
payload: Final = json.loads(body)
def scalar(value: object) -> JsonValue:
if not isinstance(value, dict):
return value if isinstance(value, (str, int, float, bool)) or value is None else str(value)
for key in ("stringValue", "intValue", "doubleValue", "boolValue", "bytesValue"):
if key in value:
return value[key]
if "arrayValue" in value:
return [scalar(item) for item in value["arrayValue"].get("values", [])]
if "kvlistValue" in value:
return {pair["key"]: scalar(pair["value"]) for pair in value["kvlistValue"].get("values", [])}
return None
return [
{
"trace_id": span.get("traceId", ""),
"span_id": span.get("spanId", ""),
"parent_span_id": span.get("parentSpanId", ""),
"kind": span.get("kind", 0),
"name": span.get("name", ""),
"attributes": {attribute["key"]: scalar(attribute.get("value")) for attribute in span.get("attributes", [])},
"resource": {
attribute["key"]: scalar(attribute.get("value"))
for attribute in resource.get("resource", {}).get("attributes", [])
},
}
for resource in payload.get("resourceSpans", [])
for scope in resource.get("scopeSpans", [])
for span in scope.get("spans", [])
]
def decode_spans(body: bytes, content_type: str) -> list[dict[str, JsonValue]]:
if "protobuf" in content_type:
return _proto_spans(body)
return _json_spans(body)
def span_class(span: dict[str, JsonValue]) -> str:
if span["kind"] == 2:
return "root"
attributes: Final = span.get("attributes") or {}
if any(marker in attributes for marker in INTERNAL_MARKERS): # type: ignore[operator] # attributes is a dict
return "tenant"
return "internal"
def spans_for_trace(spans: tuple[dict[str, JsonValue], ...], trace_id: str) -> tuple[dict[str, JsonValue], ...]:
return tuple(span for span in spans if span["trace_id"] == trace_id)
@dataclass(slots=True)
class _State:
spans: list[dict[str, JsonValue]] = None # type: ignore[assignment] # initialized in __post_init__
requests: list[dict[str, JsonValue]] = None # type: ignore[assignment]
status: int = 200
delay_seconds: float = 0.0
pause: threading.Event = threading.Event()
def __post_init__(self) -> None:
self.spans = []
self.requests = []
self.pause.set()
class _Handler(BaseHTTPRequestHandler):
state: _State
protocol_version = "HTTP/1.1"
def _read_body(self) -> bytes:
return self.rfile.read(int(self.headers.get("content-length", "0")))
def _send_json(self, payload: object, status: int = 200) -> None:
body: Final = json.dumps(payload).encode()
self.send_response(status)
self.send_header("content-type", "application/json")
self.send_header("content-length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def _record(self) -> None:
body: Final = self._read_body()
self.state.pause.wait(timeout=120)
if self.state.delay_seconds > 0:
time.sleep(self.state.delay_seconds)
recorded: Final = decode_spans(body, self.headers.get("content-type", ""))
self.state.spans.extend(recorded)
self.state.requests.append({"path": self.path, "count": len(recorded)})
self._send_json({"recorded": len(recorded)}, status=self.state.status)
do_POST = _record
do_PUT = _record
def do_GET(self) -> None:
parsed: Final = urlparse(self.path)
if parsed.path == "/__spans":
since: Final = int(dict(part.split("=", 1) for part in parsed.query.split("&") if part).get("since", "0"))
self._send_json({"next": len(self.state.spans), "spans": self.state.spans[since:]})
return
if parsed.path == "/__pid":
import os
self._send_json({"pid": os.getpid()})
return
if parsed.path == "/__requests":
self._send_json({"requests": self.state.requests})
return
self._send_json({"error": "unknown"}, status=404)
def do_DELETE(self) -> None:
if urlparse(self.path).path == "/__spans":
self.state.spans.clear()
self.state.requests.clear()
self._send_json({"cleared": True})
return
self._send_json({"error": "unknown"}, status=404)
def do_PATCH(self) -> None:
if urlparse(self.path).path != "/__control":
self._send_json({"error": "unknown"}, status=404)
return
fields: Final = json.loads(self._read_body() or b"{}")
if "status" in fields:
self.state.status = int(fields["status"])
if "delay_seconds" in fields:
self.state.delay_seconds = float(fields["delay_seconds"])
if fields.get("paused") is True:
self.state.pause.clear()
if fields.get("paused") is False:
self.state.pause.set()
self._send_json({"status": self.state.status, "delay_seconds": self.state.delay_seconds})
def log_message(self, format: str, *args: object) -> None:
pass
def recorded_spans(url: str, since: int = 0) -> tuple[int, tuple[dict[str, JsonValue], ...]]:
response: Final = httpx.get(f"{url}/__spans", params={"since": since}, trust_env=False, timeout=15)
response.raise_for_status()
payload: Final = response.json()
return int(payload["next"]), tuple(payload["spans"])
def configure_sink(url: str, **fields: JsonValue) -> None:
httpx.request("PATCH", f"{url}/__control", json=dict(fields), trust_env=False, timeout=15).raise_for_status()
def reset_sink(url: str) -> None:
httpx.delete(f"{url}/__spans", trust_env=False, timeout=15).raise_for_status()
def sink_pid(url: str) -> int:
return int(httpx.get(f"{url}/__pid", trust_env=False, timeout=15).json()["pid"])
def main() -> None:
parser: Final = argparse.ArgumentParser()
parser.add_argument("--port", type=int, required=True)
arguments: Final = parser.parse_args()
class BoundHandler(_Handler):
state = _State()
server: Final = ThreadingHTTPServer(("127.0.0.1", arguments.port), BoundHandler)
server.daemon_threads = True
server.serve_forever()
if __name__ == "__main__":
main()

View file

@ -46,7 +46,7 @@ def stop_root_process(process: subprocess.Popen[bytes]) -> bool:
@contextmanager
def owned_proxy(gateway: Gateway, directory: Path, overrides: Mapping[str, str], *, config: Path | None = None, remove_environment: tuple[str, ...] = ()) -> Iterator[Gateway]:
def owned_proxy(gateway: Gateway, directory: Path, overrides: Mapping[str, str], *, config: Path | None = None, num_workers: int = 1, remove_environment: tuple[str, ...] = ()) -> Iterator[Gateway]:
with socket.socket() as reserve:
reserve.bind(("127.0.0.1", 0))
port: Final = reserve.getsockname()[1]
@ -73,7 +73,7 @@ def owned_proxy(gateway: Gateway, directory: Path, overrides: Mapping[str, str],
"--port",
str(port),
"--num_workers",
"1",
str(num_workers),
"--use_prisma_db_push",
"--enforce_prisma_migration_check",
],

View file

@ -168,6 +168,46 @@ class Provider:
self.scripts[name] = deque(int(str(value)) for value in statuses)
return JSONResponse({"configured": len(statuses)})
async def responses(self, request: Request) -> Response:
body: Final = JSON_OBJECT.validate_json(await request.body())
self.observations.put(Observation(request.url.path, request.headers.get("authorization", ""), body))
response_id: Final = f"resp-{uuid.uuid4().hex[:24]}"
model: Final = body.get("model") if isinstance(body.get("model"), str) else "audit-chat"
def envelope(with_usage: bool) -> dict[str, JsonValue]:
return {
"id": response_id,
"object": "response",
"created_at": 1,
"status": "completed",
"model": model,
"output": [
{
"type": "message",
"id": f"msg-{response_id}",
"status": "completed",
"role": "assistant",
"content": [
{"type": "output_text", "text": "integration upstream response", "annotations": []}
],
}
],
"usage": (
{"input_tokens": 20, "output_tokens": 20, "total_tokens": 40} if with_usage else None
),
}
if body.get("stream") is True:
events: Final = (
{"type": "response.created", "response": envelope(False)},
{"type": "response.output_text.delta", "delta": "integration upstream "},
{"type": "response.output_text.delta", "delta": "response"},
{"type": "response.completed", "response": envelope(True)},
)
stream_body: Final = "".join(f"event: {event['type']}\ndata: {json.dumps(event)}\n\n" for event in events)
return Response(content=stream_body.encode(), media_type="text/event-stream")
return JSONResponse(envelope(True))
async def observed(self, _request: Request) -> Response:
values: Final = tuple(self.observations.get() for _ in range(self.observations.qsize()))
return JSONResponse(
@ -213,6 +253,7 @@ class Provider:
)
response: Final = self.scenario_store.get(scenario_id)
if response is None:
print(f"upstream scripted 404: {request.method} {request.url.path}", flush=True)
return JSONResponse({"error": "Unknown scenario"}, status_code=404)
if isinstance(response, RoutedResponse):
route_key: Final = f"{request.method} /{'/'.join(segments[1:])}"
@ -338,6 +379,7 @@ class Provider:
Route("/v1/completions", completions, methods=["POST"]),
Route("/v1/embeddings", embeddings, methods=["POST"]),
Route("/v1/moderations", moderations, methods=["POST"]),
Route("/v1/responses", self.responses, methods=["POST"]),
Route("/{path:path}", self.scripted, methods=["POST"]),
Route("/{path:path}", self.scripted, methods=["GET"]),
WebSocketRoute("/v1/realtime", self.realtime),

View file

@ -1767,6 +1767,156 @@
],
"tests/integration/mcp/test_mcp_lifecycle.py::test_same_url_server_grants_scope_discovery_and_direct_or_virtual_execution[bearer]": [
"other.mcp.permissions.same_url_servers_enforce_discovery_and_execution"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_exclude_chat_completions_openai_sync": [
"other.observability.otel.tenant_internal_spans.h1_team_exclude_chat_completions_openai_sync"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_exclude_chat_completions_stream_openai_async": [
"other.observability.otel.tenant_internal_spans.h2_team_exclude_chat_completions_stream_async"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_exclude_messages_anthropic_sync": [
"other.observability.otel.tenant_internal_spans.h3_team_exclude_messages_anthropic_sync"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_exclude_messages_stream_anthropic_async": [
"other.observability.otel.tenant_internal_spans.h4_team_exclude_messages_stream_anthropic_async"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_exclude_responses_api": [
"other.observability.otel.tenant_internal_spans.h5_team_exclude_responses_openai_sync"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_exclude_responses_stream": [
"other.observability.otel.tenant_internal_spans.h6_team_exclude_responses_stream_httpx"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_include_explicit_delivers_full_trace": [
"other.observability.otel.tenant_internal_spans.h7_team_include_explicit"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_var_absent_defaults_to_include": [
"other.observability.otel.tenant_internal_spans.h8_var_absent_defaults_include"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_key_level_exclude_without_team_callbacks": [
"other.observability.otel.tenant_internal_spans.h9_key_var_exclude_no_team"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_key_include_wins_over_team_exclude": [
"other.observability.otel.tenant_internal_spans.h10_key_include_beats_team_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_key_exclude_wins_over_team_include": [
"other.observability.otel.tenant_internal_spans.h11_key_exclude_beats_team_include"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_langfuse_exclude_arize_include_split": [
"other.observability.otel.tenant_internal_spans.h12_two_destinations_independent_filters"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_llm_only_scope_under_exclude_keeps_model_span": [
"other.observability.otel.tenant_internal_spans.h13_llm_only_scope_with_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_additive_mode_operator_full_tenant_excluded": [
"other.observability.otel.tenant_internal_spans.h14_additive_mode_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_operator_sink_keeps_internal_spans_under_exclude": [
"other.observability.otel.tenant_internal_spans.h15_operator_sink_untouched"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_guardrail_span_survives_exclude": [
"other.observability.otel.tenant_internal_spans.h16_guardrail_span_kept_under_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_env_exclude_applies_when_var_absent": [
"other.observability.otel.tenant_internal_spans.d1_env_exclude_applies"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_var_include_beats_env_exclude": [
"other.observability.otel.tenant_internal_spans.d2_var_include_beats_env_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_litellm_setting_exclude_applies_when_var_absent": [
"other.observability.otel.tenant_internal_spans.d3_setting_exclude_applies"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_setting_include_beats_env_exclude": [
"other.observability.otel.tenant_internal_spans.d4_setting_include_beats_env_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_env_exclude_with_whitespace_and_case": [
"other.observability.otel.tenant_internal_spans.d5_env_whitespace_case_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_env_bogus_value_falls_back_to_include": [
"other.observability.otel.tenant_internal_spans.d6_env_bogus_falls_back_include"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_empty_setting_falls_through_to_env": [
"other.observability.otel.tenant_internal_spans.d7_empty_setting_falls_through_env"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_callback_var_int_rejected": [
"other.observability.otel.tenant_internal_spans.s1_int_value_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_callback_var_list_rejected": [
"other.observability.otel.tenant_internal_spans.s2_list_value_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_callback_var_empty_string_rejected": [
"other.observability.otel.tenant_internal_spans.s3_empty_value_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_callback_var_oversized_string_rejected": [
"other.observability.otel.tenant_internal_spans.s4_oversized_value_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_callback_var_case_sensitive_rejected": [
"other.observability.otel.tenant_internal_spans.s5_case_sensitive_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_same_internal_spans_value_on_second_entry_accepted": [
"other.observability.otel.tenant_internal_spans.s6_same_value_twice_accepted"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_conflicting_internal_spans_value_rejected": [
"other.observability.otel.tenant_internal_spans.s7_conflicting_value_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_internal_spans_on_non_otel_callback_rejected": [
"other.observability.otel.tenant_internal_spans.s8_non_otel_callback_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_internal_spans_on_newrelic_accepted": [
"other.observability.otel.tenant_internal_spans.s9_newrelic_accepts_var"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_unauthenticated_callback_post_rejected": [
"other.observability.otel.tenant_internal_spans.s10_unauthenticated_callback_post"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_key_generate_bogus_internal_spans_rejected": [
"other.observability.otel.tenant_internal_spans.s11_key_generate_bogus_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_key_update_bogus_internal_spans_rejected": [
"other.observability.otel.tenant_internal_spans.s12_key_update_bogus_rejected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_team_update_bogus_internal_spans_drops_destination": [
"other.observability.otel.tenant_internal_spans.s13_team_update_unvalidated_drops_destination"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_tenant_sink_403_does_not_break_caller": [
"other.observability.otel.tenant_internal_spans.s14_tenant_sink_403_caller_unaffected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_tenant_sink_404_does_not_break_caller": [
"other.observability.otel.tenant_internal_spans.s15_tenant_sink_404_caller_unaffected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_upstream_500_under_exclude_keeps_internal_spans_back": [
"other.observability.otel.tenant_internal_spans.s16_upstream_500_under_exclude"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_unrelated_key_unaffected_by_tenant_sink_failure": [
"other.observability.otel.tenant_internal_spans.s17_unrelated_key_unaffected"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_key_health_with_excluded_team_callback": [
"other.observability.otel.tenant_internal_spans.s18_key_health_with_exclude_team"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_callback_var_update_include_to_exclude_takes_effect": [
"other.observability.otel.tenant_internal_spans.e1_var_update_takes_effect_within_ttl"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_callback_delete_stops_tenant_export": [
"other.observability.otel.tenant_internal_spans.e2_callback_delete_stops_tenant_export"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_identical_requests_export_exactly_once": [
"other.observability.otel.tenant_internal_spans.e3_identical_requests_export_once"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_concurrent_requests_all_excluded_once": [
"other.observability.otel.tenant_internal_spans.e4_concurrent_requests_excluded_once"
],
"tests/integration/observability/test_otel_tenant_internal_spans.py::test_failure_only_callback_entry_anchors_no_destination": [
"other.observability.otel.tenant_internal_spans.e5_failure_only_entry_anchors_no_destination"
],
"tests/integration/observability/test_otel_tenant_internal_spans_chaos.py::test_frozen_tenant_sink_receives_every_span_after_resume": [
"other.observability.otel.tenant_internal_spans.c1_frozen_sink_delivers_after_resume"
],
"tests/integration/observability/test_otel_tenant_internal_spans_chaos.py::test_slow_tenant_sink_exports_each_span_once": [
"other.observability.otel.tenant_internal_spans.c3_slow_sink_no_duplicates"
],
"tests/integration/observability/test_otel_tenant_internal_spans_chaos.py::test_proxy_restart_mid_burst_keeps_serving": [
"other.observability.otel.tenant_internal_spans.c4_proxy_restart_keeps_serving"
],
"tests/integration/observability/test_otel_tenant_internal_spans_chaos.py::test_killing_one_worker_leaves_serving": [
"other.observability.otel.tenant_internal_spans.c5_worker_kill_survivor_serves"
]
},
"browser": {

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,228 @@
from __future__ import annotations
import json
import os
import signal
import threading
import uuid
from pathlib import Path
from typing import Final
from urllib.parse import urlparse
import httpx
import pytest
from pydantic import JsonValue
from integration._support.client import Gateway, eventually
from integration._support.otlp_sink import configure_sink, recorded_spans, sink_pid, span_class, spans_for_trace
SINK_OPERATOR: Final = os.environ.get("OTEL_AUDIT_OPERATOR_SINK", "http://127.0.0.1:8191")
SINK_TENANT: Final = os.environ.get("OTEL_AUDIT_TENANT_SINK", "http://127.0.0.1:8192")
INTERNAL_SPANS_VAR: Final = "otel_internal_spans"
LANGFUSE_VARS: Final[dict[str, str]] = {
"langfuse_public_key": "pk-lf-audit",
"langfuse_secret_key": "sk-lf-audit",
"langfuse_host": SINK_TENANT,
}
def _nonce() -> str:
return f"otelchaos-{uuid.uuid4().hex}"
def _classes(spans: tuple[dict[str, JsonValue], ...]) -> dict[str, int]:
return {name: sum(1 for span in spans if span_class(span) == name) for name in ("root", "tenant", "internal")}
def _send(gateway: Gateway, key: str, model: str, nonce: str, index: int) -> httpx.Response:
if index % 3 == 0:
return gateway.request("POST", "/v1/chat/completions", {"model": model, "messages": [{"role": "user", "content": nonce}]}, key=key)
if index % 3 == 1:
return gateway.request(
"POST", "/v1/messages", {"model": model, "max_tokens": 16, "messages": [{"role": "user", "content": nonce}]}, key=key
)
return gateway.request("POST", "/v1/responses", {"model": model, "input": nonce}, key=key)
def _send_burst(gateway: Gateway, key: str, model: str, count: int) -> list[httpx.Response]:
responses: Final[list[httpx.Response]] = []
lock: Final = threading.Lock()
def hit(index: int) -> None:
nonce: Final = _nonce()
if index % 2 == 0:
response: Final = _send(gateway, key, model, nonce, index)
else:
response = gateway.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": nonce}], "stream": True},
key=key,
)
with lock:
responses.append(response)
threads: Final = [threading.Thread(target=hit, args=(i,)) for i in range(count)]
for thread in threads:
thread.start()
for thread in threads:
thread.join(timeout=90)
return responses
def _wait_trace_count(sink_url: str, call_ids: list[str | None], seconds: float) -> tuple[dict[str, JsonValue], ...]:
wanted: Final = {call_id for call_id in call_ids if call_id}
def gathered() -> tuple[dict[str, JsonValue], ...] | None:
_, spans = recorded_spans(sink_url)
covered: Final = {
(span["attributes"] or {}).get("litellm.call_id") for span in spans # type: ignore[union-attr]
}
return spans if wanted <= covered else None
result: Final = eventually(gathered, lambda value: value is not None, seconds=seconds)
assert result is not None
return result
@pytest.mark.covers("other.observability.otel.tenant_internal_spans.c1_frozen_sink_delivers_after_resume")
def test_frozen_tenant_sink_receives_every_span_after_resume(gateway: Gateway) -> None:
pid: Final = sink_pid(SINK_TENANT)
with gateway.scenario() as scenario:
model: Final = scenario.model(model="openai/audit-chat", api_base=f"{gateway.upstream_url}/v1")
team_id: Final = scenario.team()
callback: Final = gateway.request(
"POST", f"/team/{team_id}/callback", {"callback_name": "langfuse_otel", "callback_vars": {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}}
)
assert callback.status_code == 200, callback.text
key: Final = scenario.key(team_id=team_id)
os.kill(pid, signal.SIGSTOP)
try:
responses: Final = _send_burst(gateway, key, model, 30)
finally:
os.kill(pid, signal.SIGCONT)
assert all(response.status_code == 200 for response in responses), [r.status_code for r in responses]
call_ids: Final = [response.headers.get("x-litellm-call-id") for response in responses]
spans: Final = _wait_trace_count(SINK_TENANT, call_ids, seconds=120)
for call_id in call_ids:
group: Final = tuple(
span for span in spans if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr]
)
assert group, f"call {call_id} never reached the tenant sink"
model_spans: Final = tuple(span for span in group if "gen_ai.operation.name" in (span["attributes"] or {})) # type: ignore[union-attr]
assert len(model_spans) == 1, f"call {call_id} exported {len(model_spans)} times"
assert _classes(group)["internal"] == 0, f"internal spans leaked for {call_id}"
@pytest.mark.covers("other.observability.otel.tenant_internal_spans.c3_slow_sink_no_duplicates")
def test_slow_tenant_sink_exports_each_span_once(gateway: Gateway) -> None:
configure_sink(SINK_TENANT, delay_seconds=2.0)
try:
with gateway.scenario() as scenario:
model: Final = scenario.model(model="openai/audit-chat", api_base=f"{gateway.upstream_url}/v1")
team_id: Final = scenario.team()
callback: Final = gateway.request(
"POST", f"/team/{team_id}/callback", {"callback_name": "langfuse_otel", "callback_vars": {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}}
)
assert callback.status_code == 200, callback.text
key: Final = scenario.key(team_id=team_id)
responses: Final = _send_burst(gateway, key, model, 20)
assert all(response.status_code == 200 for response in responses), [r.status_code for r in responses]
call_ids: Final = [response.headers.get("x-litellm-call-id") for response in responses]
spans: Final = _wait_trace_count(SINK_TENANT, call_ids, seconds=120)
for call_id in call_ids:
group: Final = tuple(
span for span in spans if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr]
)
assert group, f"call {call_id} never reached the slow sink"
span_ids: Final = [span["span_id"] for span in group]
assert len(span_ids) == len(set(span_ids)), f"duplicate spans for {call_id}"
finally:
configure_sink(SINK_TENANT, delay_seconds=0.0)
@pytest.mark.covers("other.observability.otel.tenant_internal_spans.c4_proxy_restart_keeps_serving")
def test_proxy_restart_mid_burst_keeps_serving(gateway: Gateway, tmp_path: Path) -> None:
import yaml
from integration._support.otlp_sink import spans_for_trace as _trace
from integration._support.process import owned_proxy
config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())
config["litellm_settings"] = {
**config.get("litellm_settings", {}),
"callbacks": ["otel"],
"provider_url_destination_allowed_hosts": [urlparse(SINK_TENANT).netloc],
}
config["callback_settings"] = {"otel": {"exporter": "http/json", "endpoint": SINK_OPERATOR, "use_simple_processor": True}}
path: Final = tmp_path / "audit-restart.yaml"
path.write_text(yaml.safe_dump(config))
with owned_proxy(gateway, tmp_path, {"LITELLM_OTEL_V2": "1"}, config=path, num_workers=2) as candidate:
with candidate.scenario() as scenario:
model: Final = scenario.model(model="openai/audit-chat", api_base=f"{candidate.upstream_url}/v1")
team_id: Final = scenario.team()
callback: Final = candidate.request(
"POST", f"/team/{team_id}/callback", {"callback_name": "langfuse_otel", "callback_vars": LANGFUSE_VARS}
)
assert callback.status_code == 200, callback.text
key: Final = scenario.key(team_id=team_id)
first: Final = candidate.request(
"POST", "/v1/chat/completions", {"model": model, "messages": [{"role": "user", "content": _nonce()}]}, key=key
)
assert first.status_code == 200, first.text
first_call_id: Final = first.headers.get("x-litellm-call-id")
def first_landed() -> bool:
_, spans = recorded_spans(SINK_TENANT)
return any((span["attributes"] or {}).get("litellm.call_id") == first_call_id for span in spans) # type: ignore[union-attr]
assert eventually(first_landed, bool, seconds=40), "pre-restart trace never reached the tenant sink"
with owned_proxy(gateway, tmp_path, {"LITELLM_OTEL_V2": "1"}, config=path, num_workers=2) as candidate:
with candidate.scenario() as scenario:
model = scenario.model(model="openai/audit-chat", api_base=f"{candidate.upstream_url}/v1")
team_id = scenario.team()
callback = candidate.request(
"POST", f"/team/{team_id}/callback", {"callback_name": "langfuse_otel", "callback_vars": LANGFUSE_VARS}
)
assert callback.status_code == 200, callback.text
key = scenario.key(team_id=team_id)
response: Final = candidate.request(
"POST", "/v1/chat/completions", {"model": model, "messages": [{"role": "user", "content": _nonce()}]}, key=key
)
assert response.status_code == 200, f"post-restart request failed: {response.status_code} {response.text}"
call_id: Final = response.headers.get("x-litellm-call-id")
def landed() -> bool:
_, spans = recorded_spans(SINK_OPERATOR)
return any((span["attributes"] or {}).get("litellm.call_id") == call_id for span in spans) # type: ignore[union-attr]
assert eventually(landed, bool, seconds=40), "post-restart trace never reached the operator sink"
@pytest.mark.covers("other.observability.otel.tenant_internal_spans.c5_worker_kill_survivor_serves")
def test_killing_one_worker_leaves_serving(gateway: Gateway) -> None:
import psutil
port: Final = int(urlparse(str(gateway.client.base_url)).port or 0)
listeners: Final = {
connection.laddr.port: connection.pid
for connection in psutil.net_connections(kind="tcp")
if connection.status == "LISTEN" and connection.pid
}
owner: Final = listeners.get(port)
assert owner is not None, f"no process listening on {port}"
process: Final = psutil.Process(owner)
children: Final = process.children(recursive=True)
assert children, "leg proxy has no worker children to kill"
with gateway.scenario() as scenario:
model: Final = scenario.model(model="openai/audit-chat", api_base=f"{gateway.upstream_url}/v1")
team_id: Final = scenario.team()
callback: Final = gateway.request(
"POST", f"/team/{team_id}/callback", {"callback_name": "langfuse_otel", "callback_vars": LANGFUSE_VARS}
)
assert callback.status_code == 200, callback.text
key: Final = scenario.key(team_id=team_id)
victim: Final = children[-1]
victim.terminate()
responses: Final = _send_burst(gateway, key, model, 10)
assert all(response.status_code == 200 for response in responses), [r.status_code for r in responses]