diff --git a/docs/optimization-m12.5b-qdrant-http-semantics.md b/docs/optimization-m12.5b-qdrant-http-semantics.md new file mode 100644 index 0000000..69b5e36 --- /dev/null +++ b/docs/optimization-m12.5b-qdrant-http-semantics.md @@ -0,0 +1,39 @@ +# M12.5B — Qdrant HTTP Semantics and Safe Vision Gateway Errors + +## Old behavior + +Qdrant 4xx (including 422 for `limit > 100`) was rewritten to **Gateway 502**. Public `detail` included the internal URL and raw downstream body. + +## New behavior (Qdrant routes only) + +| Downstream | Gateway | +| --- | --- | +| 2xx | unchanged JSON | +| 400–499 | **same status**, `detail`: `Vector service rejected the request.` | +| 500–599 | **502**, `detail`: `Vector service unavailable.` | +| transport failure | **502**, sanitized | +| 2xx non-JSON | **502**, sanitized | + +CLIP / BLIP / YOLO / Maturity / Card Renderer / LLM still use the previous helper default (`preserve_client_errors=False`). + +## Scope + +Opt-in via `preserve_client_errors=True, upstream_name="qdrant"` on Qdrant `_post_json` / `_post_file` / `_get_json` / `_delete_json` calls only. + +## Deploy (do not run unless requested) + +```bash +cd /opt/docker/vision +docker compose build gateway +docker compose up -d --no-deps gateway +``` + +Do not `docker compose down`. Do not recreate qdrant-svc. + +## Rollback + +Restore previous `gateway/main.py`, then the same build + `--no-deps gateway` up. + +## Acceptance (read-only) + +`POST /vectors/search` `limit=100` → **200**. `limit=101` → **422** (not 502). Body must not contain `qdrant-svc` or Qdrant validation JSON. diff --git a/gateway/main.py b/gateway/main.py index 42e0e49..ceea5f2 100644 --- a/gateway/main.py +++ b/gateway/main.py @@ -6,7 +6,7 @@ import logging import os import time from contextlib import asynccontextmanager -from typing import Any, Dict, List, Literal, Optional +from typing import Any, Dict, List, Literal, NoReturn, Optional import httpx from fastapi import FastAPI, HTTPException, UploadFile, File, Form, Request @@ -397,54 +397,200 @@ async def _get_health(base: str) -> Dict[str, Any]: return {"status": "unreachable"} -async def _post_json(url: str, payload: Dict[str, Any]) -> Dict[str, Any]: +def _raise_for_transport_error( + exc: httpx.RequestError, + *, + url: str, + preserve_client_errors: bool, + upstream_name: str, +) -> NoReturn: + if preserve_client_errors: + logger.warning( + "downstream_transport_error upstream=%s exception_class=%s", + upstream_name, + type(exc).__name__, + ) + raise HTTPException(status_code=502, detail="Vector service unavailable.") + raise HTTPException(status_code=502, detail=f"Upstream request failed {url}: {str(exc)}") + + +def _raise_for_downstream_error( + response: httpx.Response, + *, + url: str, + preserve_client_errors: bool, + upstream_name: str, +) -> None: + status = response.status_code + if status < 400: + return + + if preserve_client_errors: + logger.warning("downstream_error upstream=%s status=%s", upstream_name, status) + if 400 <= status < 500: + raise HTTPException(status_code=status, detail="Vector service rejected the request.") + raise HTTPException(status_code=502, detail="Vector service unavailable.") + + raise HTTPException( + status_code=502, + detail=f"Upstream error {url}: {status} {response.text[:1000]}", + ) + + +def _decode_json_or_502( + response: httpx.Response, + *, + url: str, + preserve_client_errors: bool, + upstream_name: str, +) -> Dict[str, Any]: + try: + return response.json() + except Exception: + if preserve_client_errors: + logger.warning( + "downstream_non_json upstream=%s status=%s", + upstream_name, + response.status_code, + ) + raise HTTPException(status_code=502, detail="Vector service unavailable.") + raise HTTPException( + status_code=502, + detail=f"Upstream returned non-JSON at {url}: {response.status_code} {response.text[:1000]}", + ) + + +async def _post_json( + url: str, + payload: Dict[str, Any], + *, + preserve_client_errors: bool = False, + upstream_name: str = "upstream", +) -> Dict[str, Any]: t0 = time.perf_counter() try: r = await get_http_client().post(url, json=payload) except httpx.RequestError as e: - raise HTTPException(status_code=502, detail=f"Upstream request failed {url}: {str(e)}") + _raise_for_transport_error( + e, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) elapsed = (time.perf_counter() - t0) * 1000 logger.debug("POST %s status=%s elapsed_ms=%.1f", url, r.status_code, elapsed) - if r.status_code >= 400: - raise HTTPException(status_code=502, detail=f"Upstream error {url}: {r.status_code} {r.text[:1000]}") - try: - return r.json() - except Exception: - # upstream returned non-JSON (HTML error page or empty body) - raise HTTPException(status_code=502, detail=f"Upstream returned non-JSON at {url}: {r.status_code} {r.text[:1000]}") + _raise_for_downstream_error( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) + return _decode_json_or_502( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) -async def _post_file(url: str, data: bytes, fields: Dict[str, Any]) -> Dict[str, Any]: +async def _post_file( + url: str, + data: bytes, + fields: Dict[str, Any], + *, + preserve_client_errors: bool = False, + upstream_name: str = "upstream", +) -> Dict[str, Any]: files = {"file": ("image", data, "application/octet-stream")} t0 = time.perf_counter() try: r = await get_http_client().post(url, data={k: str(v) for k, v in fields.items()}, files=files) except httpx.RequestError as e: - raise HTTPException(status_code=502, detail=f"Upstream request failed {url}: {str(e)}") + _raise_for_transport_error( + e, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) elapsed = (time.perf_counter() - t0) * 1000 logger.debug("POST(file) %s status=%s elapsed_ms=%.1f", url, r.status_code, elapsed) - if r.status_code >= 400: - raise HTTPException(status_code=502, detail=f"Upstream error {url}: {r.status_code} {r.text[:1000]}") - try: - return r.json() - except Exception: - raise HTTPException(status_code=502, detail=f"Upstream returned non-JSON at {url}: {r.status_code} {r.text[:1000]}") + _raise_for_downstream_error( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) + return _decode_json_or_502( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) -async def _get_json(url: str, params: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: +async def _get_json( + url: str, + params: Optional[Dict[str, Any]] = None, + *, + preserve_client_errors: bool = False, + upstream_name: str = "upstream", +) -> Dict[str, Any]: t0 = time.perf_counter() try: r = await get_http_client().get(url, params=params) except httpx.RequestError as e: - raise HTTPException(status_code=502, detail=f"Upstream request failed {url}: {str(e)}") + _raise_for_transport_error( + e, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) elapsed = (time.perf_counter() - t0) * 1000 logger.debug("GET %s status=%s elapsed_ms=%.1f", url, r.status_code, elapsed) - if r.status_code >= 400: - raise HTTPException(status_code=502, detail=f"Upstream error {url}: {r.status_code} {r.text[:1000]}") + _raise_for_downstream_error( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) + return _decode_json_or_502( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) + + +async def _delete_json( + url: str, + *, + preserve_client_errors: bool = False, + upstream_name: str = "upstream", +) -> Dict[str, Any]: + t0 = time.perf_counter() try: - return r.json() - except Exception: - raise HTTPException(status_code=502, detail=f"Upstream returned non-JSON at {url}: {r.status_code} {r.text[:1000]}") + r = await get_http_client().delete(url) + except httpx.RequestError as e: + _raise_for_transport_error( + e, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) + elapsed = (time.perf_counter() - t0) * 1000 + logger.debug("DELETE %s status=%s elapsed_ms=%.1f", url, r.status_code, elapsed) + _raise_for_downstream_error( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) + return _decode_json_or_502( + r, + url=url, + preserve_client_errors=preserve_client_errors, + upstream_name=upstream_name, + ) @app.get("/health") @@ -618,7 +764,7 @@ async def analyze_all(payload: Dict[str, Any]): @app.post("/vectors/upsert") async def vectors_upsert(payload: Dict[str, Any]): - return await _post_json(f"{QDRANT_SVC_URL}/upsert", payload) + return await _post_json(f"{QDRANT_SVC_URL}/upsert", payload, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/upsert/file") @@ -636,17 +782,17 @@ async def vectors_upsert_file( fields["collection"] = collection if metadata_json is not None: fields["metadata_json"] = metadata_json - return await _post_file(f"{QDRANT_SVC_URL}/upsert/file", data, fields) + return await _post_file(f"{QDRANT_SVC_URL}/upsert/file", data, fields, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/upsert/vector") async def vectors_upsert_vector(payload: Dict[str, Any]): - return await _post_json(f"{QDRANT_SVC_URL}/upsert/vector", payload) + return await _post_json(f"{QDRANT_SVC_URL}/upsert/vector", payload, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/search") async def vectors_search(payload: Dict[str, Any]): - return await _post_json(f"{QDRANT_SVC_URL}/search", payload) + return await _post_json(f"{QDRANT_SVC_URL}/search", payload, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/search/file") @@ -670,32 +816,32 @@ async def vectors_search_file( fields["hnsw_ef"] = int(hnsw_ef) if filter_metadata_json is not None: fields["filter_metadata_json"] = filter_metadata_json - return await _post_file(f"{QDRANT_SVC_URL}/search/file", data, fields) + return await _post_file(f"{QDRANT_SVC_URL}/search/file", data, fields, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/search/vector") async def vectors_search_vector(payload: Dict[str, Any]): - return await _post_json(f"{QDRANT_SVC_URL}/search/vector", payload) + return await _post_json(f"{QDRANT_SVC_URL}/search/vector", payload, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/delete") async def vectors_delete(payload: Dict[str, Any]): - return await _post_json(f"{QDRANT_SVC_URL}/delete", payload) + return await _post_json(f"{QDRANT_SVC_URL}/delete", payload, preserve_client_errors=True, upstream_name="qdrant") @app.get("/vectors/collections") async def vectors_collections(): - return await _get_json(f"{QDRANT_SVC_URL}/collections") + return await _get_json(f"{QDRANT_SVC_URL}/collections", preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/collections") async def vectors_create_collection(payload: Dict[str, Any]): - return await _post_json(f"{QDRANT_SVC_URL}/collections", payload) + return await _post_json(f"{QDRANT_SVC_URL}/collections", payload, preserve_client_errors=True, upstream_name="qdrant") @app.get("/vectors/collections/{name}") async def vectors_collection_info(name: str): - return await _get_json(f"{QDRANT_SVC_URL}/collections/{name}") + return await _get_json(f"{QDRANT_SVC_URL}/collections/{name}", preserve_client_errors=True, upstream_name="qdrant") @app.get("/vectors/inspect") @@ -703,20 +849,18 @@ async def vectors_inspect(): """Full diagnostic summary for all Qdrant collections (HNSW, optimizer, payload indexes, RAM estimate).""" t0 = time.perf_counter() logger.info("vectors_inspect: start") - result = await _get_json(f"{QDRANT_SVC_URL}/inspect") + result = await _get_json(f"{QDRANT_SVC_URL}/inspect", preserve_client_errors=True, upstream_name="qdrant") logger.info("vectors_inspect: done elapsed_ms=%.1f", (time.perf_counter() - t0) * 1000) return result @app.delete("/vectors/collections/{name}") async def vectors_delete_collection(name: str): - try: - r = await get_http_client().delete(f"{QDRANT_SVC_URL}/collections/{name}") - except httpx.RequestError as exc: - raise HTTPException(status_code=502, detail=f"Upstream request failed: {exc}") - if r.status_code >= 400: - raise HTTPException(status_code=502, detail=f"Upstream error: {r.status_code}") - return r.json() + return await _delete_json( + f"{QDRANT_SVC_URL}/collections/{name}", + preserve_client_errors=True, + upstream_name="qdrant", + ) @app.get("/vectors/points/{point_id}") @@ -724,7 +868,7 @@ async def vectors_get_point(point_id: str, collection: Optional[str] = None): params = {} if collection: params["collection"] = collection - return await _get_json(f"{QDRANT_SVC_URL}/points/{point_id}", params=params) + return await _get_json(f"{QDRANT_SVC_URL}/points/{point_id}", params=params, preserve_client_errors=True, upstream_name="qdrant") @app.get("/vectors/points/by-original-id/{original_id}") @@ -732,7 +876,7 @@ async def vectors_get_point_by_original_id(original_id: str, collection: Optiona params = {} if collection: params["collection"] = collection - return await _get_json(f"{QDRANT_SVC_URL}/points/by-original-id/{original_id}", params=params) + return await _get_json(f"{QDRANT_SVC_URL}/points/by-original-id/{original_id}", params=params, preserve_client_errors=True, upstream_name="qdrant") # ---- File-based universal analyze ---- @@ -875,22 +1019,22 @@ async def cards_render_meta(payload: Dict[str, Any]): @app.get("/vectors/collections/{name}/indexes") async def vectors_collection_indexes(name: str): """List payload indexes for a collection.""" - return await _get_json(f"{QDRANT_SVC_URL}/collections/{name}/indexes") + return await _get_json(f"{QDRANT_SVC_URL}/collections/{name}/indexes", preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/collections/{name}/indexes") async def vectors_create_payload_index(name: str, payload: Dict[str, Any]): """Create a payload index on a field in a collection.""" - return await _post_json(f"{QDRANT_SVC_URL}/collections/{name}/indexes", payload) + return await _post_json(f"{QDRANT_SVC_URL}/collections/{name}/indexes", payload, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/collections/{name}/ensure-indexes") async def vectors_ensure_indexes(name: str, payload: Dict[str, Any]): """Idempotently ensure payload indexes exist for a list of fields.""" - return await _post_json(f"{QDRANT_SVC_URL}/collections/{name}/ensure-indexes", payload) + return await _post_json(f"{QDRANT_SVC_URL}/collections/{name}/ensure-indexes", payload, preserve_client_errors=True, upstream_name="qdrant") @app.post("/vectors/collections/{name}/configure") async def vectors_configure_collection(name: str, payload: Dict[str, Any]): """Update HNSW and optimizer configuration for a collection.""" - return await _post_json(f"{QDRANT_SVC_URL}/collections/{name}/configure", payload) + return await _post_json(f"{QDRANT_SVC_URL}/collections/{name}/configure", payload, preserve_client_errors=True, upstream_name="qdrant") diff --git a/tests/test_gateway_vector_errors.py b/tests/test_gateway_vector_errors.py new file mode 100644 index 0000000..007150c --- /dev/null +++ b/tests/test_gateway_vector_errors.py @@ -0,0 +1,273 @@ +from __future__ import annotations + +import unittest +from typing import Any, Dict, Optional +from unittest.mock import patch + +import httpx + +from tests.test_gateway_llm import load_gateway_module + + +QDRANT = "http://qdrant-svc:8000" +LEAK = "SECRET-QDRANT-BODY-DO-NOT-LEAK" +AUTH = {"X-API-Key": "test-key"} + + +class StubVectorClient: + def __init__( + self, + *, + post_response: Optional[httpx.Response] = None, + get_response: Optional[httpx.Response] = None, + delete_response: Optional[httpx.Response] = None, + post_exception: Optional[Exception] = None, + get_exception: Optional[Exception] = None, + delete_exception: Optional[Exception] = None, + ): + self.post_response = post_response + self.get_response = get_response + self.delete_response = delete_response + self.post_exception = post_exception + self.get_exception = get_exception + self.delete_exception = delete_exception + self.last_post: Dict[str, Any] = {} + + async def post(self, url: str, **kwargs: Any) -> httpx.Response: + self.last_post = {"url": url, **kwargs} + if self.post_exception is not None: + raise self.post_exception + if self.post_response is None: + return httpx.Response(404, json={"detail": f"No stub for POST {url}"}) + return self.post_response + + async def get(self, url: str, **kwargs: Any) -> httpx.Response: + if self.get_exception is not None: + raise self.get_exception + if self.get_response is None: + return httpx.Response(404, json={"detail": f"No stub for GET {url}"}) + return self.get_response + + async def delete(self, url: str, **kwargs: Any) -> httpx.Response: + if self.delete_exception is not None: + raise self.delete_exception + if self.delete_response is None: + return httpx.Response(404, json={"detail": f"No stub for DELETE {url}"}) + return self.delete_response + + +def _json_response(status: int, payload: Dict[str, Any], method: str = "POST", url: str = f"{QDRANT}/search") -> httpx.Response: + return httpx.Response(status, json=payload, request=httpx.Request(method, url)) + + +class GatewayVectorErrorTests(unittest.IsolatedAsyncioTestCase): + async def _request( + self, + module: Any, + method: str, + path: str, + *, + json_payload: Optional[Dict[str, Any]] = None, + files: Any = None, + data: Any = None, + ) -> httpx.Response: + transport = httpx.ASGITransport(app=module.app) + async with httpx.AsyncClient(transport=transport, base_url="http://testserver") as client: + return await client.request( + method, + path, + headers=AUTH, + json=json_payload, + files=files, + data=data, + ) + + async def _search(self, module: Any, stub: StubVectorClient, payload: Optional[Dict[str, Any]] = None) -> httpx.Response: + with patch.object(module, "get_http_client", return_value=stub): + return await self._request( + module, + "POST", + "/vectors/search", + json_payload=payload or {"url": "https://cdn.example/art.webp", "limit": 12}, + ) + + def _assert_sanitized(self, response: httpx.Response) -> None: + body = response.text + self.assertNotIn(LEAK, body) + self.assertNotIn("qdrant-svc", body) + self.assertNotIn("http://qdrant-svc:8000", body) + + async def test_search_success_unchanged(self): + module = load_gateway_module(llm_enabled=False) + payload = {"results": [{"id": 7, "score": 0.91}]} + stub = StubVectorClient(post_response=_json_response(200, payload)) + response = await self._search(module, stub) + self.assertEqual(response.status_code, 200) + self.assertEqual(response.json(), payload) + + async def test_search_preserves_client_errors(self): + module = load_gateway_module(llm_enabled=False) + cases = (400, 401, 403, 404, 408, 409, 413, 422, 429) + for status in cases: + with self.subTest(status=status): + stub = StubVectorClient( + post_response=_json_response(status, {"detail": LEAK}), + ) + response = await self._search(module, stub) + self.assertEqual(response.status_code, status) + self._assert_sanitized(response) + self.assertEqual(response.json()["detail"], "Vector service rejected the request.") + + async def test_search_maps_server_errors_to_502(self): + module = load_gateway_module(llm_enabled=False) + for status in (500, 503): + with self.subTest(status=status): + stub = StubVectorClient( + post_response=_json_response(status, {"detail": LEAK}), + ) + response = await self._search(module, stub) + self.assertEqual(response.status_code, 502) + self._assert_sanitized(response) + self.assertEqual(response.json()["detail"], "Vector service unavailable.") + + async def test_search_transport_error_is_sanitized_502(self): + module = load_gateway_module(llm_enabled=False) + stub = StubVectorClient( + post_exception=httpx.ConnectError("boom", request=httpx.Request("POST", f"{QDRANT}/search")), + ) + response = await self._search(module, stub) + self.assertEqual(response.status_code, 502) + body = response.text + self.assertNotIn("boom", body) + self.assertNotIn("qdrant-svc", body) + self.assertNotIn("http://", body) + self.assertEqual(response.json()["detail"], "Vector service unavailable.") + + async def test_search_timeout_is_sanitized_502(self): + module = load_gateway_module(llm_enabled=False) + stub = StubVectorClient( + post_exception=httpx.ReadTimeout("timed out", request=httpx.Request("POST", f"{QDRANT}/search")), + ) + response = await self._search(module, stub) + self.assertEqual(response.status_code, 502) + self.assertNotIn("timed out", response.text) + self.assertEqual(response.json()["detail"], "Vector service unavailable.") + + async def test_search_file_success_and_errors(self): + module = load_gateway_module(llm_enabled=False) + + stub_ok = StubVectorClient(post_response=_json_response(200, {"results": []}, url=f"{QDRANT}/search/file")) + with patch.object(module, "get_http_client", return_value=stub_ok): + ok = await self._request( + module, + "POST", + "/vectors/search/file", + files={"file": ("query.webp", b"image-bytes", "image/webp")}, + data={"limit": "5"}, + ) + self.assertEqual(ok.status_code, 200) + self.assertIn("file", stub_ok.last_post.get("files", {})) + self.assertEqual(stub_ok.last_post.get("data", {}).get("limit"), "5") + + stub_422 = StubVectorClient(post_response=_json_response(422, {"detail": LEAK}, url=f"{QDRANT}/search/file")) + with patch.object(module, "get_http_client", return_value=stub_422): + bad = await self._request( + module, + "POST", + "/vectors/search/file", + files={"file": ("query.webp", b"image-bytes", "image/webp")}, + data={"limit": "5"}, + ) + self.assertEqual(bad.status_code, 422) + self._assert_sanitized(bad) + + stub_500 = StubVectorClient(post_response=_json_response(500, {"detail": LEAK}, url=f"{QDRANT}/search/file")) + with patch.object(module, "get_http_client", return_value=stub_500): + server = await self._request( + module, + "POST", + "/vectors/search/file", + files={"file": ("query.webp", b"image-bytes", "image/webp")}, + data={"limit": "5"}, + ) + self.assertEqual(server.status_code, 502) + self._assert_sanitized(server) + + stub_net = StubVectorClient( + post_exception=httpx.ConnectError("boom", request=httpx.Request("POST", f"{QDRANT}/search/file")), + ) + with patch.object(module, "get_http_client", return_value=stub_net): + net = await self._request( + module, + "POST", + "/vectors/search/file", + files={"file": ("query.webp", b"image-bytes", "image/webp")}, + data={"limit": "5"}, + ) + self.assertEqual(net.status_code, 502) + self.assertNotIn("boom", net.text) + + async def test_get_collection_404_and_500(self): + module = load_gateway_module(llm_enabled=False) + stub_404 = StubVectorClient( + get_response=_json_response(404, {"detail": LEAK}, method="GET", url=f"{QDRANT}/collections/nonexistent"), + ) + with patch.object(module, "get_http_client", return_value=stub_404): + missing = await self._request(module, "GET", "/vectors/collections/nonexistent") + self.assertEqual(missing.status_code, 404) + self._assert_sanitized(missing) + + stub_500 = StubVectorClient( + get_response=_json_response(500, {"detail": LEAK}, method="GET", url=f"{QDRANT}/collections/nonexistent"), + ) + with patch.object(module, "get_http_client", return_value=stub_500): + server = await self._request(module, "GET", "/vectors/collections/nonexistent") + self.assertEqual(server.status_code, 502) + self._assert_sanitized(server) + + async def test_delete_collection_semantics(self): + module = load_gateway_module(llm_enabled=False) + stub_ok = StubVectorClient(delete_response=_json_response(200, {"ok": True}, method="DELETE", url=f"{QDRANT}/collections/test")) + with patch.object(module, "get_http_client", return_value=stub_ok): + ok = await self._request(module, "DELETE", "/vectors/collections/test") + self.assertEqual(ok.status_code, 200) + self.assertEqual(ok.json(), {"ok": True}) + + stub_404 = StubVectorClient(delete_response=_json_response(404, {"detail": LEAK}, method="DELETE", url=f"{QDRANT}/collections/test")) + with patch.object(module, "get_http_client", return_value=stub_404): + missing = await self._request(module, "DELETE", "/vectors/collections/test") + self.assertEqual(missing.status_code, 404) + self._assert_sanitized(missing) + + stub_500 = StubVectorClient(delete_response=_json_response(500, {"detail": LEAK}, method="DELETE", url=f"{QDRANT}/collections/test")) + with patch.object(module, "get_http_client", return_value=stub_500): + server = await self._request(module, "DELETE", "/vectors/collections/test") + self.assertEqual(server.status_code, 502) + self._assert_sanitized(server) + + stub_net = StubVectorClient( + delete_exception=httpx.ConnectError("boom", request=httpx.Request("DELETE", f"{QDRANT}/collections/test")), + ) + with patch.object(module, "get_http_client", return_value=stub_net): + net = await self._request(module, "DELETE", "/vectors/collections/test") + self.assertEqual(net.status_code, 502) + self.assertNotIn("boom", net.text) + + async def test_clip_422_still_mapped_to_502(self): + module = load_gateway_module(llm_enabled=False) + stub = StubVectorClient( + post_response=httpx.Response( + 422, + json={"detail": "SECRET-CLIP-BODY"}, + request=httpx.Request("POST", "http://clip:8000/analyze"), + ), + ) + with patch.object(module, "get_http_client", return_value=stub): + response = await self._request( + module, + "POST", + "/analyze/clip", + json_payload={"url": "https://cdn.example/art.webp"}, + ) + self.assertEqual(response.status_code, 502) + self.assertNotEqual(response.status_code, 422)