mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-30 01:52:18 +00:00
test(otel v2): make the internal-spans audit cells own their sinks and proxy
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
0c3f546126
commit
2d1c788799
5 changed files with 539 additions and 384 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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:])}"
|
||||
|
|
|
|||
53
tests/integration/observability/conftest.py
Normal file
53
tests/integration/observability/conftest.py
Normal file
|
|
@ -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,
|
||||
}
|
||||
File diff suppressed because it is too large
Load diff
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue