diff --git a/tests/integration/_support/otlp_sink.py b/tests/integration/_support/otlp_sink.py index 4486c112f15..d315d5315be 100644 --- a/tests/integration/_support/otlp_sink.py +++ b/tests/integration/_support/otlp_sink.py @@ -10,52 +10,89 @@ from __future__ import annotations import argparse import json +import os +import signal +import socket +import subprocess +import sys import threading import time -from dataclasses import dataclass +from collections.abc import Iterator, Mapping +from contextlib import ExitStack, contextmanager +from dataclasses import dataclass, field from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path from typing import Final from urllib.parse import urlparse import httpx -from pydantic import JsonValue +import psutil +from pydantic import JsonValue, TypeAdapter +from typing_extensions import ReadOnly, TypedDict 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 +class Span(TypedDict): + trace_id: ReadOnly[str] + span_id: ReadOnly[str] + parent_span_id: ReadOnly[str] + kind: ReadOnly[int] + name: ReadOnly[str] + attributes: ReadOnly[Mapping[str, JsonValue]] + resource: ReadOnly[Mapping[str, JsonValue]] - 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 + +class _SpanListing(TypedDict): + next: ReadOnly[int] + spans: ReadOnly[list[Span]] + + +_SPAN_LISTING: Final = TypeAdapter(_SpanListing) + + +def _proto_spans(body: bytes) -> list[Span]: + from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceRequest + from opentelemetry.proto.common.v1.common_pb2 import AnyValue + + def scalar(value: AnyValue) -> JsonValue: + match value.WhichOneof("value"): + case "string_value": + return value.string_value + case "bool_value": + return value.bool_value + case "int_value": + return int(value.int_value) + case "double_value": + return value.double_value + case "bytes_value": + return value.bytes_value.decode("utf-8", errors="replace") + case "array_value": + return [scalar(item) for item in value.array_value.values] + case "kvlist_value": + return {pair.key: scalar(pair.value) for pair in value.kvlist_value.values} + case _: + return None 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}, - } + Span( + 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]]: +def _json_spans(body: bytes) -> list[Span]: payload: Final = json.loads(body) def scalar(value: object) -> JsonValue: @@ -71,54 +108,51 @@ def _json_spans(body: bytes) -> list[dict[str, JsonValue]]: 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": { + Span( + trace_id=str(span.get("traceId", "")), + span_id=str(span.get("spanId", "")), + parent_span_id=str(span.get("parentSpanId", "")), + kind=int(span.get("kind", 0)), + name=str(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]]: +def decode_spans(body: bytes, content_type: str) -> list[Span]: if "protobuf" in content_type: return _proto_spans(body) return _json_spans(body) -def span_class(span: dict[str, JsonValue]) -> str: +def span_class(span: Span) -> 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 + if any(marker in span["attributes"] for marker in INTERNAL_MARKERS): return "tenant" return "internal" -def spans_for_trace(spans: tuple[dict[str, JsonValue], ...], trace_id: str) -> tuple[dict[str, JsonValue], ...]: +def spans_for_trace(spans: tuple[Span, ...], trace_id: str) -> tuple[Span, ...]: 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] + spans: list[Span] = field(default_factory=list) + requests: list[dict[str, JsonValue]] = field(default_factory=list) status: int = 200 delay_seconds: float = 0.0 - pause: threading.Event = threading.Event() + pause: threading.Event = field(default_factory=threading.Event) def __post_init__(self) -> None: - self.spans = [] - self.requests = [] self.pause.set() @@ -157,8 +191,6 @@ class _Handler(BaseHTTPRequestHandler): 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": @@ -193,11 +225,11 @@ class _Handler(BaseHTTPRequestHandler): pass -def recorded_spans(url: str, since: int = 0) -> tuple[int, tuple[dict[str, JsonValue], ...]]: +def recorded_spans(url: str, since: int = 0) -> tuple[int, tuple[Span, ...]]: 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"]) + listing: Final = _SPAN_LISTING.validate_python(response.json()) + return listing["next"], tuple(listing["spans"]) def configure_sink(url: str, **fields: JsonValue) -> None: @@ -212,6 +244,66 @@ def sink_pid(url: str) -> int: return int(httpx.get(f"{url}/__pid", trust_env=False, timeout=15).json()["pid"]) +@dataclass(frozen=True, slots=True) +class SpanSinks: + operator: str + tenant: str + arize: str + + +def _free_port() -> int: + with socket.socket() as reserve: + reserve.bind(("127.0.0.1", 0)) + return int(reserve.getsockname()[1]) + + +def _pid_reachable(url: str) -> bool: + try: + return httpx.get(f"{url}/__pid", trust_env=False, timeout=2).status_code == 200 + except httpx.TransportError: + return False + + +@contextmanager +def owned_sinks(directory: Path) -> Iterator[SpanSinks]: + from integration._support.process import group_members, signal_group, stop_root_process + + directory.mkdir(parents=True, exist_ok=True) + ports: Final = tuple(_free_port() for _ in range(3)) + root: Final = Path(__file__).resolve().parents[3] + with ExitStack() as stack: + processes: Final = tuple( + subprocess.Popen( + [sys.executable, "-m", "integration._support.otlp_sink", "--port", str(port)], + cwd=root, + stdout=stack.enter_context((directory / f"otlp-sink-{port}.log").open("w")), + stderr=subprocess.STDOUT, + start_new_session=True, + ) + for port in ports + ) + try: + urls: Final = tuple(f"http://127.0.0.1:{port}" for port in ports) + deadline: Final = time.monotonic() + 30 + while True: + alive: Final = all(process.poll() is None for process in processes) + assert alive, "OTLP sink exited before readiness" + if all(_pid_reachable(url) for url in urls): + break + assert time.monotonic() < deadline, "OTLP sink readiness deadline exceeded" + time.sleep(0.05) + yield SpanSinks(operator=urls[0], tenant=urls[1], arize=urls[2]) + finally: + for process in processes: + stopped: Final = stop_root_process(process) + residual: Final = group_members(process.pid) + if residual: + signal_group(process.pid, signal.SIGKILL) + psutil.wait_procs(residual, timeout=5) + survivors: Final = group_members(process.pid) + assert not survivors and stopped, "OTLP sink required forced cleanup" + + def main() -> None: parser: Final = argparse.ArgumentParser() parser.add_argument("--port", type=int, required=True) diff --git a/tests/integration/_support/upstream.py b/tests/integration/_support/upstream.py index 2444778d42d..5e4ff953a71 100644 --- a/tests/integration/_support/upstream.py +++ b/tests/integration/_support/upstream.py @@ -253,7 +253,6 @@ 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:])}" diff --git a/tests/integration/observability/conftest.py b/tests/integration/observability/conftest.py new file mode 100644 index 00000000000..f09a703f047 --- /dev/null +++ b/tests/integration/observability/conftest.py @@ -0,0 +1,53 @@ +from __future__ import annotations + +import uuid +from collections.abc import Callable, Iterator, Mapping +from pathlib import Path +from typing import Final +from urllib.parse import urlparse + +import pytest +import yaml +from integration._support.otlp_sink import SpanSinks, owned_sinks +from pydantic import JsonValue + +AuditConfigWriter = Callable[[Path, Mapping[str, JsonValue]], Path] + + +@pytest.fixture(scope="module") +def audit_sinks(tmp_path_factory: pytest.TempPathFactory) -> Iterator[SpanSinks]: + directory: Final = tmp_path_factory.mktemp("otel-audit-sinks") + with owned_sinks(directory) as sinks: + yield sinks + + +@pytest.fixture(scope="module") +def otel_audit_config(audit_sinks: SpanSinks) -> AuditConfigWriter: + tenant_host: Final = urlparse(audit_sinks.tenant).netloc + + def write(directory: 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": audit_sinks.operator, "use_simple_processor": True} + } + config["general_settings"] = {**config.get("general_settings", {}), "user_api_key_cache_ttl": 2} + path: Final = directory / f"otel-audit-{uuid.uuid4().hex}.yaml" + path.write_text(yaml.safe_dump(config)) + return path + + return write + + +@pytest.fixture(scope="module") +def langfuse_vars(audit_sinks: SpanSinks) -> dict[str, JsonValue]: + return { + "langfuse_public_key": "pk-lf-audit", + "langfuse_secret_key": "sk-lf-audit", + "langfuse_host": audit_sinks.tenant, + } diff --git a/tests/integration/observability/test_otel_tenant_internal_spans.py b/tests/integration/observability/test_otel_tenant_internal_spans.py index 30f5609ae01..90c0d9e11f0 100644 --- a/tests/integration/observability/test_otel_tenant_internal_spans.py +++ b/tests/integration/observability/test_otel_tenant_internal_spans.py @@ -2,43 +2,59 @@ from __future__ import annotations import asyncio import json -import os +import threading import uuid -from collections.abc import Iterator, Mapping +from collections.abc import Callable, 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.client import ( + Gateway, + Scenario, + eventually, + gateway_from_environment, + object_value, +) from integration._support.database import read_rows from integration._support.otlp_sink import ( + Span, + SpanSinks, configure_sink, recorded_spans, span_class, spans_for_trace, ) from integration._support.process import owned_proxy +from pydantic import JsonValue -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 +AuditConfigWriter = Callable[[Path, Mapping[str, JsonValue]], Path] -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" +@pytest.fixture(scope="module") +def gateway( + audit_sinks: SpanSinks, + otel_audit_config: AuditConfigWriter, + tmp_path_factory: pytest.TempPathFactory, +) -> Iterator[Gateway]: + directory: Final = tmp_path_factory.mktemp("otel-audit-proxy") + with gateway_from_environment() as base: + with owned_proxy( + base, + directory, + {"LITELLM_OTEL_V2": "1", "ARIZE_HTTP_ENDPOINT": audit_sinks.arize}, + config=otel_audit_config(directory, {}), + num_workers=2, + ) as candidate: + yield candidate + + def _nonce() -> str: return f"otelaudit-{uuid.uuid4().hex}" @@ -46,7 +62,7 @@ def _nonce() -> str: def _add_callback( gateway: Gateway, team_id: str, - callback_vars: Mapping[str, str], + callback_vars: Mapping[str, JsonValue], *, callback_name: str = "langfuse_otel", callback_type: str | None = None, @@ -58,12 +74,12 @@ def _add_callback( 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 _key_on_team(scenario: Scenario, team_id: str, **fields: JsonValue) -> str: + return scenario.key(team_id=team_id, **fields) -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 _audit_model(scenario: Scenario, upstream_url: str) -> str: + return scenario.model(model="openai/audit-chat", api_base=f"{upstream_url}/v1") def _chat(gateway: Gateway, key: str, model: str, nonce: str, *, stream: bool = False) -> httpx.Response: @@ -94,8 +110,8 @@ def _trace_id(sink_url: str, *, call_id: str | None = None, response_id: str | N ( 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] + if (call_id is not None and span["attributes"].get("litellm.call_id") == call_id) + or (response_id is not None and span["attributes"].get("gen_ai.response.id") == response_id) ), None, ) @@ -105,8 +121,8 @@ def _trace_id(sink_url: str, *, call_id: str | None = None, response_id: str | N 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: +def _trace_spans(sink_url: str, trace_id: str, seconds: float = 30) -> tuple[Span, ...]: + def settled() -> tuple[Span, ...] | None: _, spans = recorded_spans(sink_url) group: Final = spans_for_trace(spans, trace_id) classes: Final = {span_class(span) for span in group} @@ -117,11 +133,11 @@ def _trace_spans(sink_url: str, trace_id: str, seconds: float = 30) -> tuple[dic return group -def _classes(spans: tuple[dict[str, JsonValue], ...]) -> dict[str, int]: +def _classes(spans: tuple[Span, ...]) -> 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: +def _assert_excluded(trace_spans: tuple[Span, ...]) -> 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, ( @@ -129,7 +145,7 @@ def _assert_excluded(trace_spans: tuple[dict[str, JsonValue], ...]) -> None: ) -def _assert_full(trace_spans: tuple[dict[str, JsonValue], ...]) -> None: +def _assert_full(trace_spans: tuple[Span, ...]) -> 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, ( @@ -140,19 +156,20 @@ def _assert_full(trace_spans: tuple[dict[str, JsonValue], ...]) -> None: def _excluded_flow( gateway: Gateway, upstream_url: str, - send, + langfuse_vars: Mapping[str, JsonValue], + send: Callable[[Gateway, str, str, str], httpx.Response], *, internal_spans: str | None = "exclude", - callback_vars: Mapping[str, str] | None = None, + callback_vars: Mapping[str, JsonValue] | 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 + vars: Final = { + **dict(langfuse_vars), + **(dict(callback_vars) if callback_vars else {}), + **({INTERNAL_SPANS_VAR: internal_spans} if internal_spans is not None else {}), + } 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) @@ -161,14 +178,14 @@ def _excluded_flow( return traffic, nonce -def _assert_trace_split(traffic: httpx.Response) -> None: +def _assert_trace_split(audit_sinks: SpanSinks, 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)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(audit_sinks.operator, operator_trace)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) assert tenant_trace == operator_trace @@ -232,9 +249,8 @@ def _anthropic_sdk_async_stream(gateway: Gateway, key: str, model: str, nonce: s 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 + if isinstance(event, anthropic.types.RawMessageStartEvent): + last_id = event.message.id return last_id, None return asyncio.run(run()) @@ -265,101 +281,93 @@ def _responses_stream_id(response: httpx.Response) -> str | None: def _send_and_assert_excluded( gateway: Gateway, upstream_url: str, - send, + audit_sinks: SpanSinks, + langfuse_vars: Mapping[str, JsonValue], + send: Callable[[Gateway, str, str, str], httpx.Response], ) -> None: traffic: Final[httpx.Response] nonce: Final[str] - traffic, nonce = _excluded_flow(gateway, upstream_url, send) - _assert_trace_split(traffic) + traffic, nonce = _excluded_flow(gateway, upstream_url, langfuse_vars, send) + _assert_trace_split(audit_sinks, 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] = {} + gateway: Gateway, + tmp_path: Path, + audit_sinks: SpanSinks, + otel_audit_config: AuditConfigWriter, + 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: + overrides: Final = {"LITELLM_OTEL_V2": "1", "ARIZE_HTTP_ENDPOINT": audit_sinks.arize, **dict(env)} + with owned_proxy( + gateway, tmp_path, overrides, config=otel_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: +def test_team_exclude_chat_completions_openai_sync(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_team_exclude_chat_completions_stream_openai_async(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_team_exclude_messages_anthropic_sync(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_team_exclude_messages_stream_anthropic_async(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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() @@ -378,19 +386,19 @@ def test_team_exclude_messages_stream_anthropic_async(gateway: Gateway) -> None: 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_team_exclude_responses_api(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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() @@ -398,19 +406,19 @@ def test_team_exclude_responses_api(gateway: Gateway) -> None: 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_team_exclude_responses_stream(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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() @@ -419,103 +427,103 @@ def test_team_exclude_responses_stream(gateway: Gateway) -> None: 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id, response_id=response_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_team_include_explicit_delivers_full_trace(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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"}) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id, response_id=response_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_var_absent_defaults_to_include(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_full(_trace_spans(audit_sinks.tenant, tenant_trace)) -def _key_logging_entry(callback_vars: Mapping[str, str], callback_name: str = "langfuse_otel") -> list[JsonValue]: +def _key_logging_entry(callback_vars: Mapping[str, JsonValue], 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: +def test_key_level_exclude_without_team_callbacks(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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"})} + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id) + _assert_full(_trace_spans(audit_sinks.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: +def test_key_include_wins_over_team_exclude(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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"})}, + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_full(_trace_spans(audit_sinks.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: +def test_key_exclude_wins_over_team_include(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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"}) + 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"})}, + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_excluded(_trace_spans(audit_sinks.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: +def test_langfuse_exclude_arize_include_split(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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"}) + 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" @@ -525,71 +533,71 @@ def test_langfuse_exclude_arize_include_split(gateway: Gateway) -> None: 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) + langfuse_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id) + _assert_excluded(_trace_spans(audit_sinks.tenant, langfuse_trace)) + arize_trace: Final = _trace_id(audit_sinks.arize, call_id=call_id) + arize_spans: Final = _trace_spans(audit_sinks.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: +def test_llm_only_scope_under_exclude_keeps_model_span(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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"}, + {**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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) - def settled() -> tuple[dict[str, JsonValue], ...] | None: - _, spans = recorded_spans(SINK_TENANT) + def settled() -> tuple[Span, ...] | None: + _, spans = recorded_spans(audit_sinks.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] + assert all("gen_ai.operation.name" in span["attributes"] for span in group), names @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: +def test_additive_mode_operator_full_tenant_excluded(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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"}) + 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)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id) + _assert_full(_trace_spans(audit_sinks.operator, operator_trace)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id) + _assert_excluded(_trace_spans(audit_sinks.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: +def test_operator_sink_keeps_internal_spans_under_exclude(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=_call_id(response)) + operator_spans: Final = _trace_spans(audit_sinks.operator, operator_trace) _assert_full(operator_spans) internal_names: Final = sorted( str(span["name"]) for span in operator_spans if span_class(span) == "internal" @@ -598,9 +606,9 @@ def test_operator_sink_keeps_internal_spans_under_exclude(gateway: Gateway) -> N @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: +def test_guardrail_span_survives_exclude(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, 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: Final = yaml.safe_load(otel_audit_config(tmp_path, {}).read_text()) config["guardrails"] = [ { "guardrail_name": guardrail_name, @@ -627,22 +635,22 @@ def test_guardrail_span_survives_exclude(gateway: Gateway, tmp_path: Path) -> No 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"}) + 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) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id) + tenant_spans: Final = _trace_spans(audit_sinks.tenant, tenant_trace) - def guardrail_seen() -> tuple[dict[str, JsonValue], ...] | None: - _, spans = recorded_spans(SINK_TENANT) + def guardrail_seen() -> tuple[Span, ...] | None: + _, spans = recorded_spans(audit_sinks.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] + if "litellm.guardrail.name" in span["attributes"] ) return kept or None @@ -651,112 +659,110 @@ def test_guardrail_span_survives_exclude(gateway: Gateway, tmp_path: Path) -> No @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: +def test_env_exclude_applies_when_var_absent(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_excluded(_trace_spans(audit_sinks.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: +def test_var_include_beats_env_exclude(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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"}) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_full(_trace_spans(audit_sinks.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: +def test_litellm_setting_exclude_applies_when_var_absent(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_excluded(_trace_spans(audit_sinks.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"} +def test_setting_include_beats_env_exclude(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_full(_trace_spans(audit_sinks.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: +def test_env_exclude_with_whitespace_and_case(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_excluded(_trace_spans(audit_sinks.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: +def test_env_bogus_value_falls_back_to_include(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_full(_trace_spans(audit_sinks.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": ""} +def test_empty_setting_falls_through_to_env(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + with _candidate(gateway, tmp_path, audit_sinks, otel_audit_config, 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) + 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_excluded(_trace_spans(audit_sinks.tenant, tenant_trace)) def _bad_var_rejected(response: httpx.Response) -> None: @@ -765,58 +771,58 @@ def _bad_var_rejected(response: httpx.Response) -> None: @pytest.mark.covers("other.observability.otel.tenant_internal_spans.s1_int_value_rejected") -def test_callback_var_int_rejected(gateway: Gateway) -> None: +def test_callback_var_int_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> 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 + response: Final = _add_callback(gateway, team_id, {**langfuse_vars, INTERNAL_SPANS_VAR: 1}) _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: +def test_callback_var_list_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> 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 + gateway, team_id, {**langfuse_vars, INTERNAL_SPANS_VAR: ["exclude"]} ) _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: +def test_callback_var_empty_string_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> None: with gateway.scenario() as scenario: team_id: Final = scenario.team() - response: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: ""}) + 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: +def test_callback_var_oversized_string_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> 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}) + 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: +def test_callback_var_case_sensitive_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> None: with gateway.scenario() as scenario: team_id: Final = scenario.team() - response: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "Exclude"}) + 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: +def test_same_internal_spans_value_on_second_entry_accepted(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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" + 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" + 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") @@ -826,20 +832,20 @@ def test_same_internal_spans_value_on_second_entry_accepted(gateway: Gateway) -> 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)) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(response)) + _assert_excluded(_trace_spans(audit_sinks.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: +def test_conflicting_internal_spans_value_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> 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" + 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" + 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 @@ -877,72 +883,74 @@ def test_internal_spans_on_newrelic_accepted(gateway: Gateway) -> None: @pytest.mark.covers("other.observability.otel.tenant_internal_spans.s10_unauthenticated_callback_post") -def test_unauthenticated_callback_post_rejected(gateway: Gateway) -> None: +def test_unauthenticated_callback_post_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> 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)} + "/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: +def test_key_generate_bogus_internal_spans_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> None: response: Final = gateway.request( - "POST", "/key/generate", {"metadata": {"logging": _key_logging_entry({**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "bogus"})}} + "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: +def test_key_update_bogus_internal_spans_rejected(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> 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"})}}, + {"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: +def test_team_update_bogus_internal_spans_drops_destination(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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"})}}, + {"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) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id) + _assert_full(_trace_spans(audit_sinks.operator, operator_trace)) + _, observed = recorded_spans(audit_sinks.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] + span for span in observed if span["attributes"].get("litellm.call_id") == call_id ) 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) +def _sink_status_flow( + gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: Mapping[str, JsonValue], status: int +) -> None: + configure_sink(audit_sinks.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) + 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)) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id) + _assert_full(_trace_spans(audit_sinks.operator, operator_trace)) response_id: Final = _response_id(response.json()) rows: Final = eventually( lambda: read_rows( @@ -954,21 +962,25 @@ def _sink_status_flow(gateway: Gateway, status: int) -> None: ) assert rows, "no spend row" finally: - configure_sink(SINK_TENANT, status=200) + configure_sink(audit_sinks.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) +def test_tenant_sink_403_does_not_break_caller( + gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue] +) -> None: + _sink_status_flow(gateway, audit_sinks, langfuse_vars, 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) +def test_tenant_sink_404_does_not_break_caller( + gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue] +) -> None: + _sink_status_flow(gateway, audit_sinks, langfuse_vars, 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: +def test_upstream_500_under_exclude_keeps_internal_spans_back(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> None: upstream_model: Final = "audit-chat" httpx.post( f"{gateway.upstream_url}/__scripts/{upstream_model}", json={"statuses": [500]}, trust_env=False, timeout=15 @@ -977,20 +989,20 @@ def test_upstream_500_under_exclude_keeps_internal_spans_back(gateway: Gateway) 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: 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) + def tenant_group() -> tuple[Span, ...] | None: + _, spans = recorded_spans(audit_sinks.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] + if span["attributes"].get("litellm.call_id") == call_id + or span["trace_id"] in {s["trace_id"] for s in spans if s["attributes"].get("litellm.call_id") == call_id} ) return group if group else None @@ -1003,8 +1015,8 @@ def test_upstream_500_under_exclude_keeps_internal_spans_back(gateway: Gateway) @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) +def test_unrelated_key_unaffected_by_tenant_sink_failure(gateway: Gateway, audit_sinks: SpanSinks) -> None: + configure_sink(audit_sinks.tenant, status=403) try: with gateway.scenario() as scenario: model: Final = _audit_model(scenario, gateway.upstream_url) @@ -1012,14 +1024,14 @@ def test_unrelated_key_unaffected_by_tenant_sink_failure(gateway: Gateway) -> No 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) + configure_sink(audit_sinks.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: +def test_key_health_with_excluded_team_callback(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> None: with gateway.scenario() as scenario: team_id: Final = scenario.team() - callback: Final = _add_callback(gateway, team_id, {**LANGFUSE_VARS, INTERNAL_SPANS_VAR: "exclude"}) + 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) @@ -1027,20 +1039,20 @@ def test_key_health_with_excluded_team_callback(gateway: Gateway) -> None: @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: +def test_callback_var_update_include_to_exclude_takes_effect(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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"}) + 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)) + warm_trace: Final = _trace_id(audit_sinks.tenant, call_id=_call_id(warm)) + _assert_full(_trace_spans(audit_sinks.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"}) + 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]] = [] @@ -1049,16 +1061,16 @@ def test_callback_var_update_include_to_exclude_takes_effect(gateway: Gateway) - response: Final = _chat(gateway, key, model, _nonce()) assert response.status_code == 200, response.text issued.append(_call_id(response)) - _, spans = recorded_spans(SINK_TENANT) + _, spans = recorded_spans(audit_sinks.tenant) group: Final = tuple( span for span in spans - if (span["attributes"] or {}).get("litellm.call_id") in issued # type: ignore[union-attr] + if span["attributes"].get("litellm.call_id") in issued ) 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] + (span for span in reversed(group) if span["attributes"].get("litellm.call_id") == issued[-1]), group[-1], ) full_group: Final = spans_for_trace(spans, str(newest["trace_id"])) @@ -1071,16 +1083,16 @@ def test_callback_var_update_include_to_exclude_takes_effect(gateway: Gateway) - @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: +def test_callback_delete_stops_tenant_export(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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) + 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)) + _trace_id(audit_sinks.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 @@ -1089,24 +1101,24 @@ def test_callback_delete_stops_tenant_export(gateway: Gateway) -> None: if response.status_code != 200: return None call_id: Final = _call_id(response) - operator_trace: Final = _trace_id(SINK_OPERATOR, call_id=call_id) + operator_trace: Final = _trace_id(audit_sinks.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) + _, spans = recorded_spans(audit_sinks.tenant) leaked: Final = tuple( - span for span in spans if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr] + span for span in spans if span["attributes"].get("litellm.call_id") == call_id ) 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: +def test_identical_requests_export_exactly_once(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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 = [] @@ -1117,32 +1129,31 @@ def test_identical_requests_export_exactly_once(gateway: Gateway) -> None: 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) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id, response_id=response_id) + _, spans = recorded_spans(audit_sinks.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] + if "gen_ai.operation.name" in span["attributes"] ) 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_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id, response_id=response_id) + _, operator_spans = recorded_spans(audit_sinks.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] + if span["attributes"].get("gen_ai.response.id") == response_id ) 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 +def test_concurrent_requests_all_excluded_once(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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: 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]] = [] @@ -1159,24 +1170,24 @@ def test_concurrent_requests_all_excluded_once(gateway: Gateway) -> None: 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) + tenant_trace: Final = _trace_id(audit_sinks.tenant, call_id=call_id) + group: Final = _trace_spans(audit_sinks.tenant, tenant_trace) _assert_excluded(group) workers: Final = { - str((span["resource"] or {}).get("process.pid")) for span in group # type: ignore[union-attr] + str(span["resource"].get("process.pid")) for span in group } 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: +def test_failure_only_callback_entry_anchors_no_destination(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> 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_vars, INTERNAL_SPANS_VAR: "exclude"}, callback_type="failure", ) assert callback.status_code == 200, callback.text @@ -1184,10 +1195,10 @@ def test_failure_only_callback_entry_anchors_no_destination(gateway: Gateway) -> 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) + operator_trace: Final = _trace_id(audit_sinks.operator, call_id=call_id) + _assert_full(_trace_spans(audit_sinks.operator, operator_trace)) + _, spans = recorded_spans(audit_sinks.tenant) leaked: Final = tuple( - span for span in spans if (span["attributes"] or {}).get("litellm.call_id") == call_id # type: ignore[union-attr] + span for span in spans if span["attributes"].get("litellm.call_id") == call_id ) 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 index f5131e291ef..2ba017b2a8e 100644 --- a/tests/integration/observability/test_otel_tenant_internal_spans_chaos.py +++ b/tests/integration/observability/test_otel_tenant_internal_spans_chaos.py @@ -1,36 +1,48 @@ from __future__ import annotations -import json import os import signal import threading import uuid +from collections.abc import Callable, Iterator, Mapping from pathlib import Path from typing import Final from urllib.parse import urlparse import httpx import pytest +from integration._support.client import Gateway, eventually, gateway_from_environment +from integration._support.otlp_sink import Span, SpanSinks, configure_sink, recorded_spans, sink_pid, span_class +from integration._support.process import owned_proxy 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") +AuditConfigWriter = Callable[[Path, Mapping[str, JsonValue]], Path] 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, -} + + +@pytest.fixture(scope="module") +def gateway( + audit_sinks: SpanSinks, + otel_audit_config: AuditConfigWriter, + tmp_path_factory: pytest.TempPathFactory, +) -> Iterator[Gateway]: + directory: Final = tmp_path_factory.mktemp("otel-audit-chaos-proxy") + with gateway_from_environment() as base: + with owned_proxy( + base, + directory, + {"LITELLM_OTEL_V2": "1", "ARIZE_HTTP_ENDPOINT": audit_sinks.arize}, + config=otel_audit_config(directory, {}), + num_workers=2, + ) as candidate: + yield candidate def _nonce() -> str: return f"otelchaos-{uuid.uuid4().hex}" -def _classes(spans: tuple[dict[str, JsonValue], ...]) -> dict[str, int]: +def _classes(spans: tuple[Span, ...]) -> dict[str, int]: return {name: sum(1 for span in spans if span_class(span) == name) for name in ("root", "tenant", "internal")} @@ -70,13 +82,13 @@ def _send_burst(gateway: Gateway, key: str, model: str, count: int) -> list[http return responses -def _wait_trace_count(sink_url: str, call_ids: list[str | None], seconds: float) -> tuple[dict[str, JsonValue], ...]: +def _wait_trace_count(sink_url: str, call_ids: list[str | None], seconds: float) -> tuple[Span, ...]: wanted: Final = {call_id for call_id in call_ids if call_id} - def gathered() -> tuple[dict[str, JsonValue], ...] | None: + def gathered() -> tuple[Span, ...] | None: _, spans = recorded_spans(sink_url) covered: Final = { - (span["attributes"] or {}).get("litellm.call_id") for span in spans # type: ignore[union-attr] + span["attributes"].get("litellm.call_id") for span in spans } return spans if wanted <= covered else None @@ -86,13 +98,13 @@ def _wait_trace_count(sink_url: str, call_ids: list[str | None], seconds: float) @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) +def test_frozen_tenant_sink_receives_every_span_after_resume(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> None: + pid: Final = sink_pid(audit_sinks.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"}} + "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) @@ -103,66 +115,54 @@ def test_frozen_tenant_sink_receives_every_span_after_resume(gateway: Gateway) - 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) + spans: Final = _wait_trace_count(audit_sinks.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] + span for span in spans if span["attributes"].get("litellm.call_id") == call_id ) 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] + model_spans: Final = tuple(span for span in group if "gen_ai.operation.name" in span["attributes"]) 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) +def test_slow_tenant_sink_exports_each_span_once(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue]) -> None: + configure_sink(audit_sinks.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"}} + "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) + spans: Final = _wait_trace_count(audit_sinks.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] + span for span in spans if span["attributes"].get("litellm.call_id") == call_id ) 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) + configure_sink(audit_sinks.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: +def test_proxy_restart_mid_burst_keeps_serving(gateway: Gateway, audit_sinks: SpanSinks, langfuse_vars: dict[str, JsonValue], otel_audit_config: AuditConfigWriter, tmp_path: Path) -> None: + path: Final = otel_audit_config(tmp_path, {}) + overrides: Final = {"LITELLM_OTEL_V2": "1", "ARIZE_HTTP_ENDPOINT": audit_sinks.arize} + with owned_proxy(gateway, tmp_path, overrides, 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} + "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) @@ -173,19 +173,19 @@ def test_proxy_restart_mid_burst_keeps_serving(gateway: Gateway, tmp_path: Path) 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] + _, spans = recorded_spans(audit_sinks.tenant) + return any(span["attributes"].get("litellm.call_id") == first_call_id for span in spans) 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 owned_proxy(gateway, tmp_path, overrides, 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} + 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 = scenario.key(team_id=team_id) + key: Final = scenario.key(team_id=team_id) response: Final = candidate.request( "POST", "/v1/chat/completions", {"model": model, "messages": [{"role": "user", "content": _nonce()}]}, key=key ) @@ -193,14 +193,14 @@ def test_proxy_restart_mid_burst_keeps_serving(gateway: Gateway, tmp_path: Path) 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] + _, spans = recorded_spans(audit_sinks.operator) + return any(span["attributes"].get("litellm.call_id") == call_id for span in spans) 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: +def test_killing_one_worker_leaves_serving(gateway: Gateway, langfuse_vars: dict[str, JsonValue]) -> None: import psutil port: Final = int(urlparse(str(gateway.client.base_url)).port or 0) @@ -218,7 +218,7 @@ def test_killing_one_worker_leaves_serving(gateway: Gateway) -> None: 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} + "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)