From 0c3f546126b98f74060d3c13b2ca66c3836f4704 Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 18:31:02 +0000 Subject: [PATCH] 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> --- tests/integration/_support/otlp_sink.py | 229 ++++ tests/integration/_support/process.py | 4 +- tests/integration/_support/upstream.py | 42 + tests/integration/contracts.json | 150 +++ .../test_otel_tenant_internal_spans.py | 1193 +++++++++++++++++ .../test_otel_tenant_internal_spans_chaos.py | 228 ++++ 6 files changed, 1844 insertions(+), 2 deletions(-) create mode 100644 tests/integration/_support/otlp_sink.py create mode 100644 tests/integration/observability/test_otel_tenant_internal_spans.py create mode 100644 tests/integration/observability/test_otel_tenant_internal_spans_chaos.py diff --git a/tests/integration/_support/otlp_sink.py b/tests/integration/_support/otlp_sink.py new file mode 100644 index 00000000000..4486c112f15 --- /dev/null +++ b/tests/integration/_support/otlp_sink.py @@ -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() diff --git a/tests/integration/_support/process.py b/tests/integration/_support/process.py index e0923c9d055..3a5e5416b76 100644 --- a/tests/integration/_support/process.py +++ b/tests/integration/_support/process.py @@ -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", ], diff --git a/tests/integration/_support/upstream.py b/tests/integration/_support/upstream.py index e9c50ea7966..2444778d42d 100644 --- a/tests/integration/_support/upstream.py +++ b/tests/integration/_support/upstream.py @@ -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), diff --git a/tests/integration/contracts.json b/tests/integration/contracts.json index cdf534bb5a1..c764aa1264b 100644 --- a/tests/integration/contracts.json +++ b/tests/integration/contracts.json @@ -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": { diff --git a/tests/integration/observability/test_otel_tenant_internal_spans.py b/tests/integration/observability/test_otel_tenant_internal_spans.py new file mode 100644 index 00000000000..30f5609ae01 --- /dev/null +++ b/tests/integration/observability/test_otel_tenant_internal_spans.py @@ -0,0 +1,1193 @@ +from __future__ import annotations + +import asyncio +import json +import os +import uuid +from collections.abc import Iterator, Mapping +from contextlib import contextmanager +from pathlib import Path +from typing import Final +from urllib.parse import urlparse + +import httpx +import pytest +import yaml +from pydantic import JsonValue + +from integration._support.client import Gateway, eventually, object_value, string_value +from integration._support.database import read_rows +from integration._support.otlp_sink import ( + configure_sink, + recorded_spans, + span_class, + spans_for_trace, +) +from integration._support.process import owned_proxy + +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") +SINK_ARIZE: Final = os.environ.get("OTEL_AUDIT_ARIZE_SINK", "http://127.0.0.1:8193") +TENANT_HOST: Final = urlparse(SINK_TENANT).netloc + +LANGFUSE_VARS: Final[dict[str, str]] = { + "langfuse_public_key": "pk-lf-audit", + "langfuse_secret_key": "sk-lf-audit", + "langfuse_host": SINK_TENANT, +} +ARIZE_VARS: Final[dict[str, str]] = {"arize_space_id": "audit-space", "arize_api_key": "audit-arize-key"} +INTERNAL_SPANS_VAR: Final = "otel_internal_spans" + + +def _nonce() -> str: + return f"otelaudit-{uuid.uuid4().hex}" + + +def _add_callback( + gateway: Gateway, + team_id: str, + callback_vars: Mapping[str, str], + *, + callback_name: str = "langfuse_otel", + callback_type: str | None = None, + key: str | None = None, +) -> httpx.Response: + body: Final[dict[str, JsonValue]] = {"callback_name": callback_name, "callback_vars": dict(callback_vars)} + if callback_type is not None: + body["callback_type"] = callback_type + return gateway.request("POST", f"/team/{team_id}/callback", body, key=key) + + +def _key_on_team(scenario: object, team_id: str, **fields: JsonValue) -> str: + return scenario.key(team_id=team_id, **fields) # type: ignore[attr-defined] # Scenario helper + + +def _audit_model(scenario: object, upstream_url: str) -> str: + return scenario.model(model="openai/audit-chat", api_base=f"{upstream_url}/v1") # type: ignore[attr-defined] + + +def _chat(gateway: Gateway, key: str, model: str, nonce: str, *, stream: bool = False) -> httpx.Response: + return gateway.request( + "POST", + "/v1/chat/completions", + {"model": model, "messages": [{"role": "user", "content": nonce}], **({"stream": True} if stream else {})}, + key=key, + ) + + +def _response_id(body: JsonValue) -> str | None: + if isinstance(body, dict): + value: Final = body.get("id") + if isinstance(value, str): + return value + return None + + +def _call_id(response: httpx.Response) -> str | None: + return response.headers.get("x-litellm-call-id") + + +def _trace_id(sink_url: str, *, call_id: str | None = None, response_id: str | None = None, seconds: float = 40) -> str: + def look() -> str | None: + _, spans = recorded_spans(sink_url) + return next( + ( + str(span["trace_id"]) + for span in spans + if (call_id is not None and (span["attributes"] or {}).get("litellm.call_id") == call_id) # type: ignore[union-attr] + or (response_id is not None and (span["attributes"] or {}).get("gen_ai.response.id") == response_id) # type: ignore[union-attr] + ), + None, + ) + + found: Final = eventually(look, lambda value: value is not None, seconds=seconds) + assert found is not None + return found + + +def _trace_spans(sink_url: str, trace_id: str, seconds: float = 30) -> tuple[dict[str, JsonValue], ...]: + def settled() -> tuple[dict[str, JsonValue], ...] | None: + _, spans = recorded_spans(sink_url) + group: Final = spans_for_trace(spans, trace_id) + classes: Final = {span_class(span) for span in group} + return group if "root" in classes and "tenant" in classes else None + + group: Final = eventually(settled, lambda value: value is not None, seconds=seconds) + assert group is not None + return group + + +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 _assert_excluded(trace_spans: tuple[dict[str, JsonValue], ...]) -> None: + counts: Final = _classes(trace_spans) + names: Final = sorted(str(span["name"]) for span in trace_spans) + assert counts["root"] == 1 and counts["internal"] == 0 and counts["tenant"] >= 1, ( + f"expected root + tenant spans only, got {counts} with {names}" + ) + + +def _assert_full(trace_spans: tuple[dict[str, JsonValue], ...]) -> None: + counts: Final = _classes(trace_spans) + names: Final = sorted(str(span["name"]) for span in trace_spans) + assert counts["root"] == 1 and counts["internal"] >= 1 and counts["tenant"] >= 1, ( + f"expected a full tree with internal spans, got {counts} with {names}" + ) + + +def _excluded_flow( + gateway: Gateway, + upstream_url: str, + send, + *, + internal_spans: str | None = "exclude", + callback_vars: Mapping[str, str] | None = None, +) -> tuple[httpx.Response, str]: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, upstream_url) + team_id: Final = scenario.team() + vars: Final = dict(LANGFUSE_VARS) + if callback_vars: + vars.update(callback_vars) + if internal_spans is not None: + vars[INTERNAL_SPANS_VAR] = internal_spans + response: Final = _add_callback(gateway, team_id, vars) + assert response.status_code == 200, f"callback setup failed: {response.status_code} {response.text}" + key: Final = _key_on_team(scenario, team_id) + nonce: Final = _nonce() + traffic: Final = send(gateway, key, model, nonce) + return traffic, nonce + + +def _assert_trace_split(traffic: httpx.Response) -> None: + assert traffic.status_code == 200, traffic.text + call_id: Final = _call_id(traffic) + response_id: Final = _response_id(traffic.json()) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + assert tenant_trace == operator_trace + + +def _assert_upstream_saw(upstream_url: str, nonce: str) -> None: + observations: Final = httpx.get(f"{upstream_url}/__observations", trust_env=False, timeout=15).json()["requests"] + matching: Final = [ + entry for entry in observations if nonce in json.dumps(entry.get("body", {})) + ] + assert len(matching) >= 1, f"upstream never saw nonce {nonce}" + + +def _openai_sync_send(gateway: Gateway, key: str, model: str, nonce: str, *, stream: bool = False) -> httpx.Response: + return _chat(gateway, key, model, nonce, stream=stream) + + +def _openai_sdk(gateway: Gateway, key: str, model: str, nonce: str) -> tuple[str | None, str | None]: + import openai + + client: Final = openai.OpenAI(base_url=f"{gateway.client.base_url}".rstrip("/"), api_key=key, timeout=30) + raw: Final = client.chat.completions.with_raw_response.create( + model=model, messages=[{"role": "user", "content": nonce}] + ) + parsed: Final = raw.parse() + return parsed.id, raw.headers.get("x-litellm-call-id") + + +def _openai_sdk_async_stream(gateway: Gateway, key: str, model: str, nonce: str) -> tuple[str | None, str | None]: + import openai + + async def run() -> tuple[str | None, str | None]: + client: Final = openai.AsyncOpenAI(base_url=f"{gateway.client.base_url}".rstrip("/"), api_key=key, timeout=30) + raw: Final = await client.chat.completions.with_raw_response.create( + model=model, messages=[{"role": "user", "content": nonce}], stream=True + ) + stream: Final = raw.parse() + last_id: str | None = None + async for chunk in stream: + if chunk.id: + last_id = chunk.id + return last_id, raw.headers.get("x-litellm-call-id") + + return asyncio.run(run()) + + +def _anthropic_sdk(gateway: Gateway, key: str, model: str, nonce: str) -> tuple[str | None, str | None]: + import anthropic + + client: Final = anthropic.Anthropic(base_url=f"{gateway.client.base_url}".rstrip("/"), api_key=key, timeout=30) + raw: Final = client.messages.with_raw_response.create( + model=model, max_tokens=16, messages=[{"role": "user", "content": nonce}] + ) + parsed: Final = raw.parse() + return parsed.id, raw.headers.get("x-litellm-call-id") + + +def _anthropic_sdk_async_stream(gateway: Gateway, key: str, model: str, nonce: str) -> tuple[str | None, str | None]: + import anthropic + + async def run() -> tuple[str | None, str | None]: + client: Final = anthropic.AsyncAnthropic(base_url=f"{gateway.client.base_url}".rstrip("/"), api_key=key, timeout=30) + last_id: str | None = None + async with client.messages.stream(model=model, max_tokens=16, messages=[{"role": "user", "content": nonce}]) as stream: + async for event in stream: + response_obj: Final = getattr(event, "message", None) + if response_obj is not None and getattr(response_obj, "id", None): + last_id = response_obj.id + return last_id, None + + return asyncio.run(run()) + + +def _responses_httpx(gateway: Gateway, key: str, model: str, nonce: str, *, stream: bool) -> httpx.Response: + return gateway.request( + "POST", + "/v1/responses", + {"model": model, "input": nonce, **({"stream": True} if stream else {})}, + key=key, + ) + + +def _responses_stream_id(response: httpx.Response) -> str | None: + for line in response.text.splitlines(): + if line.startswith("data:"): + try: + event: Final = json.loads(line[5:].strip()) + except json.JSONDecodeError: + continue + response_obj: Final = event.get("response") + if isinstance(response_obj, dict) and isinstance(response_obj.get("id"), str): + return response_obj["id"] + return None + + +def _send_and_assert_excluded( + gateway: Gateway, + upstream_url: str, + send, +) -> None: + traffic: Final[httpx.Response] + nonce: Final[str] + traffic, nonce = _excluded_flow(gateway, upstream_url, send) + _assert_trace_split(traffic) + _assert_upstream_saw(upstream_url, nonce) + + +def _audit_config(tmp_path: Path, litellm_settings: Mapping[str, JsonValue] = {}) -> Path: + 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": [TENANT_HOST], + **dict(litellm_settings), + } + config["callback_settings"] = { + "otel": {"exporter": "http/json", "endpoint": SINK_OPERATOR, "use_simple_processor": True} + } + config["general_settings"] = {**config.get("general_settings", {}), "user_api_key_cache_ttl": 2} + path: Final = tmp_path / "audit-proxy.yaml" + path.write_text(yaml.safe_dump(config)) + return path + + +@contextmanager +def _candidate( + gateway: Gateway, tmp_path: Path, env: Mapping[str, str] = {}, settings: Mapping[str, JsonValue] = {} +) -> Iterator[Gateway]: + overrides: Final = {"LITELLM_OTEL_V2": "1", **dict(env)} + with owned_proxy(gateway, tmp_path, overrides, config=_audit_config(tmp_path, settings), num_workers=2) as candidate: + yield candidate + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h1_team_exclude_chat_completions_openai_sync") +def test_team_exclude_chat_completions_openai_sync(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, f"callback setup failed: {callback.status_code} {callback.text}" + key: Final = _key_on_team(scenario, team_id) + nonce = _nonce() + response_id, call_id = _openai_sdk(gateway, key, model, nonce) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + assert tenant_trace == operator_trace + _assert_upstream_saw(gateway.upstream_url, nonce) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h2_team_exclude_chat_completions_stream_async") +def test_team_exclude_chat_completions_stream_openai_async(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + nonce: Final = _nonce() + response_id, call_id = _openai_sdk_async_stream(gateway, key, model, nonce) + assert response_id is not None, "stream produced no response id" + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + _assert_upstream_saw(gateway.upstream_url, nonce) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h3_team_exclude_messages_anthropic_sync") +def test_team_exclude_messages_anthropic_sync(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + nonce: Final = _nonce() + response_id, call_id = _anthropic_sdk(gateway, key, model, nonce) + assert call_id is not None, "no x-litellm-call-id header on /v1/messages" + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + _assert_upstream_saw(gateway.upstream_url, nonce) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h4_team_exclude_messages_stream_anthropic_async") +def test_team_exclude_messages_stream_anthropic_async(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + nonce: Final = _nonce() + response: Final = gateway.request( + "POST", + "/v1/messages", + { + "model": model, + "max_tokens": 16, + "stream": True, + "messages": [{"role": "user", "content": nonce}], + }, + key=key, + ) + assert response.status_code == 200, response.text + assert "message_stop" in response.text, response.text + call_id: Final = _call_id(response) + assert call_id is not None + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + _assert_upstream_saw(gateway.upstream_url, nonce) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h5_team_exclude_responses_openai_sync") +def test_team_exclude_responses_api(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + nonce: Final = _nonce() + response: Final = _responses_httpx(gateway, key, model, nonce, stream=False) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + response_id: Final = _response_id(response.json()) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + _assert_upstream_saw(gateway.upstream_url, nonce) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h6_team_exclude_responses_stream_httpx") +def test_team_exclude_responses_stream(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + nonce: Final = _nonce() + response: Final = _responses_httpx(gateway, key, model, nonce, stream=True) + assert response.status_code == 200, response.text + response_id: Final = _responses_stream_id(response) + call_id: Final = _call_id(response) + assert response_id is not None or call_id is not None, response.text[:400] + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + _assert_upstream_saw(gateway.upstream_url, nonce) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h7_team_include_explicit") +def test_team_include_explicit_delivers_full_trace(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "include"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + response_id: Final = _response_id(response.json()) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h8_var_absent_defaults_include") +def test_var_absent_defaults_to_include(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_full(_trace_spans(SINK_TENANT, tenant_trace)) + + +def _key_logging_entry(callback_vars: Mapping[str, str], callback_name: str = "langfuse_otel") -> list[JsonValue]: + return [{"callback_name": callback_name, "callback_vars": dict(callback_vars)}] + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h9_key_var_exclude_no_team") +def test_key_level_exclude_without_team_callbacks(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + key: Final = scenario.key( + metadata={"logging": _key_logging_entry({**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"})} + ) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h10_key_include_beats_team_exclude") +def test_key_include_wins_over_team_exclude(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = scenario.key( + team_id=team_id, + metadata={"logging": _key_logging_entry({**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "include"})}, + ) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_full(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h11_key_exclude_beats_team_include") +def test_key_exclude_wins_over_team_include(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "include"}) + assert callback.status_code == 200, callback.text + key: Final = scenario.key( + team_id=team_id, + metadata={"logging": _key_logging_entry({**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"})}, + ) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h12_two_destinations_independent_filters") +def test_langfuse_exclude_arize_include_split(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + langfuse: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert langfuse.status_code == 200, langfuse.text + arize: Final = _add_callback( + gateway, team_id, {**ARIZE_VARS, INTERNAL_SPANS_VAR: "include"}, callback_name="arize" + ) + assert arize.status_code == 200, arize.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + langfuse_trace: Final = _trace_id(SINK_TENANT, call_id=call_id) + _assert_excluded(_trace_spans(SINK_TENANT, langfuse_trace)) + arize_trace: Final = _trace_id(SINK_ARIZE, call_id=call_id) + arize_spans: Final = _trace_spans(SINK_ARIZE, arize_trace) + _assert_full(arize_spans) + assert arize_trace == langfuse_trace + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h13_llm_only_scope_with_exclude") +def test_llm_only_scope_under_exclude_keeps_model_span(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback( + gateway, + team_id, + {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude", "langfuse_span_scope": "llm_only"}, + ) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + + def settled() -> tuple[dict[str, JsonValue], ...] | None: + _, spans = recorded_spans(SINK_TENANT) + group: Final = spans_for_trace(spans, tenant_trace) + return group if group else None + + group: Final = eventually(settled, lambda value: value is not None, seconds=30) + assert group is not None + names: Final = sorted(str(span["name"]) for span in group) + assert all("gen_ai.operation.name" in (span["attributes"] or {}) for span in group), names # type: ignore[union-attr] + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h14_additive_mode_exclude") +def test_additive_mode_operator_full_tenant_excluded(gateway: Gateway, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, settings={"otel_tenant_destination_mode": "additive"}) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h15_operator_sink_untouched") +def test_operator_sink_keeps_internal_spans_under_exclude(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=_call_id(response)) + operator_spans: Final = _trace_spans(SINK_OPERATOR, operator_trace) + _assert_full(operator_spans) + internal_names: Final = sorted( + str(span["name"]) for span in operator_spans if span_class(span) == "internal" + ) + assert any("auth" in name or "redis" in name or "postgres" in name for name in internal_names), internal_names + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.h16_guardrail_span_kept_under_exclude") +def test_guardrail_span_survives_exclude(gateway: Gateway, tmp_path: Path) -> None: + guardrail_name: Final = f"audit-filter-{uuid.uuid4().hex[:8]}" + config: Final = yaml.safe_load(_audit_config(tmp_path).read_text()) + config["guardrails"] = [ + { + "guardrail_name": guardrail_name, + "litellm_params": { + "guardrail": "litellm_content_filter", + "mode": "pre_call", + "default_on": True, + "patterns": [ + { + "pattern_type": "regex", + "pattern_name": "audit_secret", + "pattern": "TOPSECRET\\d{9}", + "action": "BLOCK", + } + ], + }, + } + ] + path: Final = tmp_path / "audit-guardrail.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 = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id) + tenant_spans: Final = _trace_spans(SINK_TENANT, tenant_trace) + + def guardrail_seen() -> tuple[dict[str, JsonValue], ...] | None: + _, spans = recorded_spans(SINK_TENANT) + group: Final = spans_for_trace(spans, tenant_trace) + kept: Final = tuple( + span + for span in group + if "litellm.guardrail.name" in (span["attributes"] or {}) # type: ignore[union-attr] + ) + return kept or None + + kept: Final = eventually(guardrail_seen, lambda value: value is not None, seconds=30) + assert kept is not None, f"guardrail span missing at tenant sink: {tenant_spans}" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.d1_env_exclude_applies") +def test_env_exclude_applies_when_var_absent(gateway: Gateway, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, env={"LITELLM_OTEL_TENANT_INTERNAL_SPANS": "exclude"}) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.d2_var_include_beats_env_exclude") +def test_var_include_beats_env_exclude(gateway: Gateway, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, env={"LITELLM_OTEL_TENANT_INTERNAL_SPANS": "exclude"}) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "include"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_full(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.d3_setting_exclude_applies") +def test_litellm_setting_exclude_applies_when_var_absent(gateway: Gateway, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, settings={"otel_tenant_internal_spans": "exclude"}) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.d4_setting_include_beats_env_exclude") +def test_setting_include_beats_env_exclude(gateway: Gateway, tmp_path: Path) -> None: + with _candidate( + gateway, tmp_path, env={"LITELLM_OTEL_TENANT_INTERNAL_SPANS": "exclude"}, settings={"otel_tenant_internal_spans": "include"} + ) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_full(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.d5_env_whitespace_case_exclude") +def test_env_exclude_with_whitespace_and_case(gateway: Gateway, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, env={"LITELLM_OTEL_TENANT_INTERNAL_SPANS": " EXCLUDE "}) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.d6_env_bogus_falls_back_include") +def test_env_bogus_value_falls_back_to_include(gateway: Gateway, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, env={"LITELLM_OTEL_TENANT_INTERNAL_SPANS": "bogus"}) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_full(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.d7_empty_setting_falls_through_env") +def test_empty_setting_falls_through_to_env(gateway: Gateway, tmp_path: Path) -> None: + with _candidate( + gateway, tmp_path, env={"LITELLM_OTEL_TENANT_INTERNAL_SPANS": "exclude"}, settings={"otel_tenant_internal_spans": ""} + ) as candidate: + with candidate.scenario() as scenario: + model: Final = _audit_model(scenario, candidate.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(candidate, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(candidate, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + + +def _bad_var_rejected(response: httpx.Response) -> None: + assert response.status_code in (400, 422), f"expected rejection, got {response.status_code}: {response.text}" + assert INTERNAL_SPANS_VAR in response.text or "callback" in response.text, response.text + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s1_int_value_rejected") +def test_callback_var_int_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + response: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: 1}) # type: ignore[dict-item] # deliberately malformed input + _bad_var_rejected(response) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s2_list_value_rejected") +def test_callback_var_list_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + response: Final = _add_callback( + gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: ["exclude"]} # type: ignore[dict-item] # deliberately malformed input + ) + _bad_var_rejected(response) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s3_empty_value_rejected") +def test_callback_var_empty_string_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + response: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: ""}) + _bad_var_rejected(response) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s4_oversized_value_rejected") +def test_callback_var_oversized_string_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + response: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "x" * 5120}) + _bad_var_rejected(response) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s5_case_sensitive_rejected") +def test_callback_var_case_sensitive_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + response: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "Exclude"}) + _bad_var_rejected(response) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s6_same_value_twice_accepted") +def test_same_internal_spans_value_on_second_entry_accepted(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + first: Final = _add_callback( + gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}, callback_type="success_and_failure" + ) + assert first.status_code == 200, first.text + second: Final = _add_callback( + gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}, callback_type="success" + ) + assert second.status_code == 200, f"identical value rejected: {second.status_code} {second.text}" + listed: Final = gateway.get(f"/team/{team_id}/callback") + data: Final = object_value(listed["data"]) + assert "langfuse_otel" in data.get("success_callbacks", []), data + assert object_value(data["callback_vars"]).get(INTERNAL_SPANS_VAR) == "exclude", data + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(response)) + _assert_excluded(_trace_spans(SINK_TENANT, tenant_trace)) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s7_conflicting_value_rejected") +def test_conflicting_internal_spans_value_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + first: Final = _add_callback( + gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}, callback_type="success_and_failure" + ) + assert first.status_code == 200, first.text + second: Final = _add_callback( + gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "include"}, callback_type="success" + ) + assert second.status_code == 400, f"expected 400 conflict, got {second.status_code}: {second.text}" + assert INTERNAL_SPANS_VAR in second.text, second.text + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s8_non_otel_callback_rejected") +def test_internal_spans_on_non_otel_callback_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + response: Final = _add_callback( + gateway, + team_id, + {"langsmith_api_key": "sk-ls-audit", INTERNAL_SPANS_VAR: "exclude"}, + callback_name="langsmith", + ) + assert response.status_code == 400, f"expected 400, got {response.status_code}: {response.text}" + assert "callback" in response.text, response.text + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s9_newrelic_accepts_var") +def test_internal_spans_on_newrelic_accepted(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + response: Final = _add_callback( + gateway, + team_id, + {"newrelic_api_key": "nr-audit-key", INTERNAL_SPANS_VAR: "exclude"}, + callback_name="newrelic", + ) + assert response.status_code == 200, f"expected 200, got {response.status_code}: {response.text}" + listed: Final = gateway.get(f"/team/{team_id}/callback") + data: Final = object_value(listed["data"]) + assert object_value(data["callback_vars"]).get(INTERNAL_SPANS_VAR) == "exclude", data + assert "newrelic" in data.get("success_callbacks", []) or "newrelic" in data.get("failure_callbacks", []), data + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s10_unauthenticated_callback_post") +def test_unauthenticated_callback_post_rejected(gateway: Gateway) -> None: + with httpx.Client(base_url=str(gateway.client.base_url), timeout=15, trust_env=False) as client: + response: Final = client.post( + "/team/some-team/callback", json={"callback_name": "langfuse_otel", "callback_vars": dict(LANGFUSE_VARS)} + ) + assert response.status_code == 401, f"expected 401, got {response.status_code}: {response.text}" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s11_key_generate_bogus_rejected") +def test_key_generate_bogus_internal_spans_rejected(gateway: Gateway) -> None: + response: Final = gateway.request( + "POST", "/key/generate", {"metadata": {"logging": _key_logging_entry({**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "bogus"})}} + ) + assert response.status_code == 400, f"expected 400, got {response.status_code}: {response.text}" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s12_key_update_bogus_rejected") +def test_key_update_bogus_internal_spans_rejected(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + key: Final = scenario.key() + response: Final = gateway.request( + "POST", + "/key/update", + {"key": key, "metadata": {"logging": _key_logging_entry({**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "bogus"})}}, + ) + assert response.status_code == 400, f"expected 400, got {response.status_code}: {response.text}" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s13_team_update_unvalidated_drops_destination") +def test_team_update_bogus_internal_spans_drops_destination(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + update: Final = gateway.request( + "POST", + "/team/update", + {"team_id": team_id, "metadata": {"logging": _key_logging_entry({**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "bogus"})}}, + ) + assert update.status_code == 200, update.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + _, observed = recorded_spans(SINK_TENANT, since=0) + leaked: Final = tuple( + span for span in observed if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr] + ) + assert not leaked, f"bogus metadata.logging still reached the tenant sink: {leaked}" + + +def _sink_status_flow(gateway: Gateway, status: int) -> None: + configure_sink(SINK_TENANT, status=status) + try: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, f"caller broke on sink {status}: {response.status_code} {response.text}" + call_id: Final = _call_id(response) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + response_id: Final = _response_id(response.json()) + rows: Final = eventually( + lambda: read_rows( + 'SELECT spend FROM "LiteLLM_SpendLogs" WHERE request_id IN (%s, %s)', + (str(call_id), str(response_id)), + ), + lambda values: len(values) >= 1, + seconds=70, + ) + assert rows, "no spend row" + finally: + configure_sink(SINK_TENANT, status=200) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s14_tenant_sink_403_caller_unaffected") +def test_tenant_sink_403_does_not_break_caller(gateway: Gateway) -> None: + _sink_status_flow(gateway, 403) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s15_tenant_sink_404_caller_unaffected") +def test_tenant_sink_404_does_not_break_caller(gateway: Gateway) -> None: + _sink_status_flow(gateway, 404) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s16_upstream_500_under_exclude") +def test_upstream_500_under_exclude_keeps_internal_spans_back(gateway: Gateway) -> None: + upstream_model: Final = "audit-chat" + httpx.post( + f"{gateway.upstream_url}/__scripts/{upstream_model}", json={"statuses": [500]}, trust_env=False, timeout=15 + ).raise_for_status() + try: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code >= 500, f"expected caller 5xx, got {response.status_code}: {response.text}" + call_id: Final = _call_id(response) + + def tenant_group() -> tuple[dict[str, JsonValue], ...] | None: + _, spans = recorded_spans(SINK_TENANT) + group: Final = tuple( + span + for span in spans + if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr] + or span["trace_id"] in {s["trace_id"] for s in spans if (s["attributes"] or {}).get("litellm.call_id") == call_id} # type: ignore[union-attr] + ) + return group if group else None + + group: Final = eventually(tenant_group, lambda value: value is not None, seconds=40) + assert group is not None, "tenant sink never received the failed-request trace" + counts: Final = _classes(group) + assert counts["internal"] == 0, f"internal spans leaked on error trace: {group}" + finally: + httpx.delete(f"{gateway.upstream_url}/__scripts/{upstream_model}", trust_env=False, timeout=15) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s17_unrelated_key_unaffected") +def test_unrelated_key_unaffected_by_tenant_sink_failure(gateway: Gateway) -> None: + configure_sink(SINK_TENANT, status=403) + try: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + key: Final = scenario.key() + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, f"unrelated key broke: {response.status_code} {response.text}" + finally: + configure_sink(SINK_TENANT, status=200) + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.s18_key_health_with_exclude_team") +def test_key_health_with_excluded_team_callback(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = gateway.request("POST", "/key/health", key=key) + assert response.status_code == 200, f"/key/health failed: {response.status_code} {response.text}" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.e1_var_update_takes_effect_within_ttl") +def test_callback_var_update_include_to_exclude_takes_effect(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + first: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "include"}) + assert first.status_code == 200, first.text + key: Final = _key_on_team(scenario, team_id) + warm: Final = _chat(gateway, key, model, _nonce()) + assert warm.status_code == 200, warm.text + warm_trace: Final = _trace_id(SINK_TENANT, call_id=_call_id(warm)) + _assert_full(_trace_spans(SINK_TENANT, warm_trace)) + deleted: Final = gateway.request("DELETE", f"/team/{team_id}/callback/langfuse_otel") + assert deleted.status_code == 200, deleted.text + updated: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert updated.status_code == 200, updated.text + + issued: Final[list[str | None]] = [] + + def flipped() -> str | None: + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + issued.append(_call_id(response)) + _, spans = recorded_spans(SINK_TENANT) + group: Final = tuple( + span + for span in spans + if (span["attributes"] or {}).get("litellm.call_id") in issued # type: ignore[union-attr] + ) + if not group: + return None + newest: Final = next( + (span for span in reversed(group) if (span["attributes"] or {}).get("litellm.call_id") == issued[-1]), # type: ignore[union-attr] + group[-1], + ) + full_group: Final = spans_for_trace(spans, str(newest["trace_id"])) + if _classes(full_group)["internal"] == 0: + return str(newest["trace_id"]) + return None + + trace: Final = eventually(flipped, lambda value: value is not None, seconds=70) + assert trace is not None, "exclude never took effect within the cache TTL" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.e2_callback_delete_stops_tenant_export") +def test_callback_delete_stops_tenant_export(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, LANGFUSE_VARS) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + warm: Final = _chat(gateway, key, model, _nonce()) + assert warm.status_code == 200, warm.text + _trace_id(SINK_TENANT, call_id=_call_id(warm)) + deleted: Final = gateway.request("DELETE", f"/team/{team_id}/callback/langfuse_otel") + assert deleted.status_code == 200, deleted.text + + def drained() -> str | None: + response: Final = _chat(gateway, key, model, _nonce()) + if response.status_code != 200: + return None + call_id: Final = _call_id(response) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + return call_id if operator_trace else None + + call_id: Final = eventually(drained, lambda value: value is not None, seconds=70) + assert call_id is not None + _, spans = recorded_spans(SINK_TENANT) + leaked: Final = tuple( + span for span in spans if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr] + ) + assert not leaked, f"tenant sink still received spans after callback delete: {leaked}" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.e3_identical_requests_export_once") +def test_identical_requests_export_exactly_once(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + call_ids: Final = [] + response_ids: Final = [] + for _ in range(5): + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + call_ids.append(_call_id(response)) + response_ids.append(_response_id(response.json())) + for call_id, response_id in zip(call_ids, response_ids): + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id, response_id=response_id) + _, spans = recorded_spans(SINK_TENANT) + matching: Final = tuple( + span + for span in spans_for_trace(spans, tenant_trace) + if "gen_ai.operation.name" in (span["attributes"] or {}) # type: ignore[union-attr] + ) + assert len(matching) == 1, f"model span for {response_id} exported {len(matching)} times" + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id, response_id=response_id) + _, operator_spans = recorded_spans(SINK_OPERATOR) + operator_matching: Final = tuple( + span + for span in spans_for_trace(operator_spans, operator_trace) + if (span["attributes"] or {}).get("gen_ai.response.id") == response_id # type: ignore[union-attr] + ) + assert len(operator_matching) == 1, f"operator exported {response_id} {len(operator_matching)} times" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.e4_concurrent_requests_excluded_once") +def test_concurrent_requests_all_excluded_once(gateway: Gateway) -> None: + import threading + + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + responses: Final[list[httpx.Response]] = [] + + def hit() -> None: + responses.append(_chat(gateway, key, model, _nonce())) + + threads: Final = [threading.Thread(target=hit) for _ in range(10)] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=30) + assert len(responses) == 10 + for response in responses: + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + tenant_trace: Final = _trace_id(SINK_TENANT, call_id=call_id) + group: Final = _trace_spans(SINK_TENANT, tenant_trace) + _assert_excluded(group) + workers: Final = { + str((span["resource"] or {}).get("process.pid")) for span in group # type: ignore[union-attr] + } + assert workers, "no process attribution on tenant spans" + + +@pytest.mark.covers("other.observability.otel.tenant_internal_spans.e5_failure_only_entry_anchors_no_destination") +def test_failure_only_callback_entry_anchors_no_destination(gateway: Gateway) -> None: + with gateway.scenario() as scenario: + model: Final = _audit_model(scenario, gateway.upstream_url) + team_id: Final = scenario.team() + callback: Final = _add_callback( + gateway, + team_id, + {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}, + callback_type="failure", + ) + assert callback.status_code == 200, callback.text + key: Final = _key_on_team(scenario, team_id) + response: Final = _chat(gateway, key, model, _nonce()) + assert response.status_code == 200, response.text + call_id: Final = _call_id(response) + operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + _assert_full(_trace_spans(SINK_OPERATOR, operator_trace)) + _, spans = recorded_spans(SINK_TENANT) + leaked: Final = tuple( + span for span in spans if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr] + ) + assert not leaked, f"failure-only entry anchored a tenant destination: {leaked}" diff --git a/tests/integration/observability/test_otel_tenant_internal_spans_chaos.py b/tests/integration/observability/test_otel_tenant_internal_spans_chaos.py new file mode 100644 index 00000000000..f5131e291ef --- /dev/null +++ b/tests/integration/observability/test_otel_tenant_internal_spans_chaos.py @@ -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]