test(integration): keep upstream support unchanged

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-26 02:38:04 +00:00
parent 661b1dbcf4
commit 25568999ee
2 changed files with 16 additions and 52 deletions

View file

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

View file

@ -4,6 +4,7 @@ import time
import uuid
from collections.abc import Callable, Iterator, Mapping
from pathlib import Path
from types import MappingProxyType
from typing import Final
import httpx
@ -11,7 +12,6 @@ import pytest
import yaml
from integration._support.client import (
Gateway,
Scenario,
eventually,
gateway_from_environment,
)
@ -39,8 +39,8 @@ def _config_with(
directory: Path,
otel_audit_config: AuditConfigWriter,
*,
otel: Mapping[str, JsonValue] = {},
extra: Callable[[dict], None] | None = None,
otel: Mapping[str, JsonValue] = MappingProxyType({}),
extra: Callable[[dict[str, JsonValue]], None] | None = None,
) -> Path:
config: Final = yaml.safe_load(otel_audit_config(directory, {}).read_text())
config["callback_settings"]["otel"].update(dict(otel))
@ -128,12 +128,7 @@ def _await_db_span(sink_url: str, trace_id: str | None, needle: str, seconds: fl
def _db_systems(spans: tuple[Span, ...]) -> set[str]:
return {
str(span["attributes"][key])
for span in spans
for key in DB_SYSTEM_KEYS
if key in span["attributes"]
}
return {str(span["attributes"][key]) for span in spans for key in DB_SYSTEM_KEYS if key in span["attributes"]}
def _assert_core_spans_present(spans: tuple[Span, ...]) -> None:
@ -202,7 +197,7 @@ def test_env_excluded_services_drops_only_redis(
gateway, tmp_path, {"LITELLM_OTEL_V2": "1", "LITELLM_OTEL_EXCLUDED_SERVICES": "redis"}, config=config, workers=2
) as candidate:
start, _ = recorded_spans(audit_sinks.tenant)
traffic: Final = _drive(candidate, langfuse_vars)
_drive(candidate, langfuse_vars)
_await_db_span(audit_sinks.tenant, None, "postgresql", seconds=60, since=start)
_, tenant_spans = recorded_spans(audit_sinks.tenant, start)
systems: Final = _db_systems(tenant_spans)
@ -229,7 +224,9 @@ def test_config_excluded_services_wins_over_env(
_, all_tenant = recorded_spans(audit_sinks.tenant, start)
systems: Final = _db_systems(tenant_spans)
assert "redis" in systems, f"redis spans missing at tenant: {systems}"
assert "postgresql" not in _db_systems(all_tenant), f"postgresql spans reached tenant: {_db_systems(all_tenant)}"
assert "postgresql" not in _db_systems(all_tenant), (
f"postgresql spans reached tenant: {_db_systems(all_tenant)}"
)
def test_bogus_excluded_service_fails_proxy_start(