Preserve Qdrant 4xx on the Vision gateway and stop leaking upstream errors.

Map Qdrant 400-499 through as the same status with a short public detail, and keep 5xx and transport failures as sanitized 502. CLIP and other services still use the old helper default.
This commit is contained in:
2026-08-25 07:58:48 +02:00
parent bd0abab759
commit 1713f0ff79
3 changed files with 505 additions and 49 deletions
+193 -49
View File
@@ -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")