From 25568999ee3ed75df5194375561f9da8a2ade134 Mon Sep 17 00:00:00 2001 From: yucheng Date: Sat, 26 Sep 2026 02:38:04 +0000 Subject: [PATCH] test(integration): keep upstream support unchanged Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- tests/integration/_support/upstream.py | 49 +++---------------- .../test_otel_excluded_services.py | 19 +++---- 2 files changed, 16 insertions(+), 52 deletions(-) diff --git a/tests/integration/_support/upstream.py b/tests/integration/_support/upstream.py index 27a18915ff3..e645126032b 100644 --- a/tests/integration/_support/upstream.py +++ b/tests/integration/_support/upstream.py @@ -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"]), diff --git a/tests/integration/observability/test_otel_excluded_services.py b/tests/integration/observability/test_otel_excluded_services.py index ee9455fcac8..8beaecf95f1 100644 --- a/tests/integration/observability/test_otel_excluded_services.py +++ b/tests/integration/observability/test_otel_excluded_services.py @@ -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(