mirror of
https://github.com/headroomlabs-ai/headroom.git
synced 2026-08-27 14:17:10 -04:00
fix(proxy): retry passthrough on transient upstream connection close (#1513)
## Description `GET /v1/models` (and other buffered passthrough routes) returned an opaque HTTP **502** when an OpenAI-compatible upstream closed a pooled keep-alive connection mid-response, surfacing `httpx.RemoteProtocolError: peer closed connection without sending complete message body (incomplete chunked read)`. The same upstream answers a direct `curl` with 200 because curl opens a fresh connection per call, while Headroom reuses pooled keep-alive connections — so the first request issued on a stale connection fails even though the upstream is healthy. The fix makes the buffered passthrough path retry once on a fresh connection (exactly what curl does), and return a clear error only if the upstream is genuinely sending an incomplete response. Closes #1112 ## Type of Change - [x] Bug fix (non-breaking change that fixes an issue) - [ ] New feature (non-breaking change that adds functionality) - [ ] Breaking change (fix or feature that would cause existing functionality to change) - [ ] Documentation update - [ ] Performance improvement - [ ] Code refactoring (no functional changes) ## Changes Made - Add `headroom.proxy.helpers.request_with_transient_retry(client, *, request_id=None, max_retries=1, **request_kwargs)`: issues a buffered httpx request and retries on a **fresh connection** when (and only when) `httpx.RemoteProtocolError` is raised. Every other exception (`ConnectError`, timeouts, status errors) propagates immediately, so existing handling is unchanged. Documented as buffered-only (a streamed response can't be safely replayed once bytes reach the client). - Route `OpenAIHandlerMixin.handle_passthrough` through the helper, and add an `except httpx.RemoteProtocolError` arm that returns a clear `502` with error type `upstream_protocol_error` when the protocol error persists across the retry (instead of letting the raw error surface as an opaque/unhandled 502). - Add `tests/test_proxy_passthrough_transient_retry.py` (helper unit tests + handler-level tests covering the exact issue path). - Add a `CHANGELOG.md` entry under `Unreleased → Fixed`. Scope note: streaming `/v1/responses` is intentionally **out of scope** for this change — a streamed response cannot be safely retried after the first byte has been delivered to the client. The helper is written reusable so a streaming-aware follow-up can build on it. ## Testing - [x] Unit tests pass (`pytest`) - [x] Linting passes (`ruff check .`) - [x] Type checking passes (`mypy headroom`) - [x] New tests added for new functionality - [x] Manual testing performed ### Test Output ```text $ ruff check headroom/proxy/helpers.py headroom/proxy/handlers/openai.py tests/test_proxy_passthrough_transient_retry.py All checks passed! $ mypy headroom/proxy/helpers.py --ignore-missing-imports Success: no issues found in 1 source file $ pytest tests/test_proxy_passthrough_transient_retry.py -q tests/test_proxy_passthrough_transient_retry.py ....... [100%] 7 passed in 0.27s # no regressions in the surrounding passthrough/handler suites: $ pytest tests/test_proxy_passthrough_transient_retry.py tests/test_proxy_handler_helpers.py \ tests/test_proxy_byte_faithful_forwarding.py \ tests/test_proxy/test_compression_failure_action.py tests/test_proxy_copilot_auth_hooks.py -q 80 passed, 1 warning in 6.88s ``` ## Real Behavior Proof Reproduced against a **real local TCP server** (no mocks) that speaks HTTP/1.1 and, when armed, emits a chunked body then closes the socket **without** the terminating `0\r\n\r\n` — the exact condition that makes httpx raise the `incomplete chunked read` error from this issue. - Environment: macOS arm64, Python 3.12, httpx 0.28.1 (same httpx major as the report), real loopback sockets via `asyncio.start_server`. - Exact command / steps: start the local server; (1) issue a single buffered request — the pre-fix `handle_passthrough` behaviour; (2) issue the same request through `request_with_transient_retry` — the fix. Verbatim: `python repro_1112.py`. - Observed result: BEFORE the fix a single request raises `httpx.RemoteProtocolError` ("incomplete chunked read") which `handle_passthrough` surfaced as an opaque HTTP 502; AFTER the fix the same request returns **HTTP 200** (the retry opened a fresh connection, mirroring a direct `curl`). Full terminal output: ```text upstream listening on http://127.0.0.1:62374/v1/models BEFORE (single buffered request, pre-fix behaviour): raised httpx.RemoteProtocolError: peer closed connection without sending complete message body (incomplete chunked read) -> handle_passthrough surfaced this as an opaque HTTP 502 AFTER (request_with_transient_retry, the fix): HTTP 200 body={"object":"list","data":[]} -> first attempt hit the incomplete chunked read, retry on a fresh connection returned 200 (mirrors a direct curl) ``` The log line `Upstream closed connection mid-response (...incomplete chunked read); retrying on a fresh connection (attempt 1/1)` fires on the recovered request, confirming the retry path is what produced the 200. - Not tested: real third-party upstreams (LiteLLM/vLLM/etc.) — the local server reproduces the precise httpx error deterministically; the streaming `/v1/responses` path is intentionally out of scope (a streamed response cannot be safely retried after the first byte reaches the client). ## Review Readiness - [x] I have performed a self-review - [x] This PR is ready for human review ## Checklist - [x] My code follows the project's style guidelines - [x] I have performed a self-review of my code - [x] I have commented my code, particularly in hard-to-understand areas - [x] I have made corresponding changes to the documentation - [x] My changes generate no new warnings - [x] I have added tests that prove my fix is effective or that my feature works - [x] New and existing unit tests pass locally with my changes - [x] I have updated the CHANGELOG.md if applicable ## Additional Notes - No new dependencies (httpx is already a proxy dependency), so no supply-chain justification is required. - The retry is deliberately narrow: only `httpx.RemoteProtocolError` is retried, capped at one retry, so a genuinely-down upstream still fails fast via the existing `ConnectError`/timeout path. - "Documentation" checklist item refers to the `CHANGELOG.md` entry; no user-facing docs pages needed for this internal resilience fix.
This commit is contained in:
parent
32ce99e4b4
commit
5d14080c94
4 changed files with 291 additions and 2 deletions
|
|
@ -72,6 +72,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||
`min(4, cpu)` × 1). The ONNX embedder already capped its threads; this brings
|
||||
the torch path to parity
|
||||
([#198](https://github.com/headroomlabs-ai/headroom/issues/198)).
|
||||
- Buffered passthrough routes (e.g. `GET /v1/models`) no longer return an
|
||||
opaque HTTP 502 when an OpenAI-compatible upstream closes a pooled
|
||||
keep-alive connection mid-response (`httpx.RemoteProtocolError` /
|
||||
"incomplete chunked read"). Headroom now retries the request once on a
|
||||
fresh connection — mirroring a direct `curl` — and only returns a clear
|
||||
`upstream_protocol_error` 502 if the upstream is genuinely sending an
|
||||
incomplete response
|
||||
([#1112](https://github.com/chopratejas/headroom/issues/1112)).
|
||||
|
||||
### Changed
|
||||
|
||||
|
|
|
|||
|
|
@ -6746,7 +6746,11 @@ class OpenAIHandlerMixin:
|
|||
client = classify_client(headers)
|
||||
tags = extract_tags(headers)
|
||||
# PR-A5 (P5-49): strip internal x-headroom-* before forwarding upstream.
|
||||
from headroom.proxy.helpers import _strip_internal_headers, log_outbound_headers
|
||||
from headroom.proxy.helpers import (
|
||||
_strip_internal_headers,
|
||||
log_outbound_headers,
|
||||
request_with_transient_retry,
|
||||
)
|
||||
|
||||
_pre_strip_count_pt = sum(1 for k in headers if k.lower().startswith("x-headroom-"))
|
||||
headers = _strip_internal_headers(headers)
|
||||
|
|
@ -6769,7 +6773,12 @@ class OpenAIHandlerMixin:
|
|||
if _prefers_http1_passthrough(base_url):
|
||||
passthrough_client = self.http_client_h1 or self.http_client
|
||||
try:
|
||||
response = await passthrough_client.request( # type: ignore[union-attr]
|
||||
# Retry once on a transient keep-alive close (httpx
|
||||
# RemoteProtocolError / "incomplete chunked read"): the upstream
|
||||
# closed a pooled connection httpx then reused. A fresh connection
|
||||
# succeeds, mirroring what a direct curl call does. See GH #1112.
|
||||
response = await request_with_transient_retry(
|
||||
passthrough_client, # type: ignore[arg-type]
|
||||
method=request.method,
|
||||
url=url,
|
||||
headers=headers,
|
||||
|
|
@ -6795,6 +6804,37 @@ class OpenAIHandlerMixin:
|
|||
status_code=502,
|
||||
media_type="application/json",
|
||||
)
|
||||
except httpx.RemoteProtocolError as e:
|
||||
# Persisted across the retry: the upstream really is sending an
|
||||
# incomplete response. Return a clear 502 instead of letting the
|
||||
# raw protocol error surface as an opaque/unhandled 502.
|
||||
logger.warning(
|
||||
"Passthrough upstream closed connection without a complete "
|
||||
"response after retry: %s %s -> %s: %s",
|
||||
request.method,
|
||||
path,
|
||||
url,
|
||||
e,
|
||||
)
|
||||
return Response(
|
||||
content=json.dumps(
|
||||
{
|
||||
"error": {
|
||||
"type": "upstream_protocol_error",
|
||||
# Full exception detail is logged server-side above;
|
||||
# keep the client-facing message generic so upstream
|
||||
# exception/stack-trace text is not exposed (CodeQL
|
||||
# py/stack-trace-exposure).
|
||||
"message": (
|
||||
"Upstream closed the connection without sending "
|
||||
"a complete response."
|
||||
),
|
||||
}
|
||||
}
|
||||
),
|
||||
status_code=502,
|
||||
media_type="application/json",
|
||||
)
|
||||
|
||||
# Remove compression headers since httpx already decompressed the response
|
||||
response_headers = dict(response.headers)
|
||||
|
|
|
|||
|
|
@ -960,6 +960,50 @@ def retry_after_ms(response: httpx.Response, max_ms: int) -> float | None:
|
|||
RETRYABLE_OVERLOAD_STATUSES: frozenset[int] = frozenset({429, 529})
|
||||
|
||||
|
||||
async def request_with_transient_retry(
|
||||
client: httpx.AsyncClient,
|
||||
*,
|
||||
request_id: str | None = None,
|
||||
max_retries: int = 1,
|
||||
**request_kwargs: Any,
|
||||
) -> httpx.Response:
|
||||
"""Issue a buffered httpx request, retrying once on a transient close.
|
||||
|
||||
``httpx.RemoteProtocolError`` ("peer closed connection without sending
|
||||
complete message body (incomplete chunked read)") is raised when an
|
||||
upstream closes a pooled keep-alive connection that httpx then reuses for
|
||||
the next request. A direct ``curl`` never hits this because it opens a
|
||||
fresh connection per call; Headroom reuses pooled connections, so the
|
||||
first request issued on a stale connection fails even though the upstream
|
||||
is healthy (it answers a fresh connection with 200). Retrying opens a new
|
||||
connection and succeeds, mirroring curl's behaviour. See GH #1112.
|
||||
|
||||
Only ``httpx.RemoteProtocolError`` is retried — the specific stale
|
||||
keep-alive symptom; every other exception (``ConnectError``, timeouts,
|
||||
HTTP status errors) propagates immediately so existing handling is
|
||||
unchanged. Use this for buffered (non-streaming) requests only: a streamed
|
||||
response cannot be safely replayed once bytes have reached the client.
|
||||
"""
|
||||
import httpx
|
||||
|
||||
attempt = 0
|
||||
while True:
|
||||
try:
|
||||
return await client.request(**request_kwargs)
|
||||
except httpx.RemoteProtocolError as exc:
|
||||
if attempt >= max_retries:
|
||||
raise
|
||||
attempt += 1
|
||||
logger.warning(
|
||||
"Upstream closed connection mid-response (%s); retrying on a "
|
||||
"fresh connection (attempt %d/%d)%s",
|
||||
exc,
|
||||
attempt,
|
||||
max_retries,
|
||||
f" [{request_id}]" if request_id else "",
|
||||
)
|
||||
|
||||
|
||||
# Image compression availability (do not retain a global compressor instance)
|
||||
_image_compressor_available: bool | None = None
|
||||
|
||||
|
|
|
|||
197
tests/test_proxy_passthrough_transient_retry.py
Normal file
197
tests/test_proxy_passthrough_transient_retry.py
Normal file
|
|
@ -0,0 +1,197 @@
|
|||
"""Regression tests for GH #1112.
|
||||
|
||||
The Headroom proxy returned an opaque HTTP 502 when an OpenAI-compatible
|
||||
upstream closed a pooled keep-alive connection mid-response, surfacing
|
||||
``httpx.RemoteProtocolError`` ("peer closed connection without sending
|
||||
complete message body (incomplete chunked read)"). The same upstream answers
|
||||
a direct ``curl`` with 200 because curl opens a fresh connection per call;
|
||||
Headroom reuses pooled connections, so the first request on a stale connection
|
||||
fails even though the upstream is healthy.
|
||||
|
||||
The fix adds :func:`headroom.proxy.helpers.request_with_transient_retry`,
|
||||
which retries the buffered request once on a fresh connection, and wires it
|
||||
into ``OpenAIHandlerMixin.handle_passthrough`` with a clean 502 fallback when
|
||||
the protocol error persists.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import AsyncMock
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from headroom.proxy.handlers.openai import OpenAIHandlerMixin
|
||||
from headroom.proxy.helpers import request_with_transient_retry
|
||||
|
||||
_INCOMPLETE_CHUNKED = (
|
||||
"peer closed connection without sending complete message body (incomplete chunked read)"
|
||||
)
|
||||
|
||||
|
||||
def _ok_response() -> httpx.Response:
|
||||
request = httpx.Request("GET", "https://api.openai.com/v1/models")
|
||||
return httpx.Response(
|
||||
200,
|
||||
request=request,
|
||||
headers={"content-type": "application/json"},
|
||||
json={"object": "list", "data": []},
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Unit tests for the retry helper
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_helper_returns_response_without_retry_on_success() -> None:
|
||||
ok = _ok_response()
|
||||
client = SimpleNamespace(request=AsyncMock(return_value=ok))
|
||||
|
||||
result = asyncio.run(
|
||||
request_with_transient_retry(client, method="GET", url="https://up/v1/models")
|
||||
)
|
||||
|
||||
assert result is ok
|
||||
assert client.request.await_count == 1
|
||||
|
||||
|
||||
def test_helper_recovers_after_one_remote_protocol_error() -> None:
|
||||
ok = _ok_response()
|
||||
client = SimpleNamespace(
|
||||
request=AsyncMock(side_effect=[httpx.RemoteProtocolError(_INCOMPLETE_CHUNKED), ok])
|
||||
)
|
||||
|
||||
result = asyncio.run(
|
||||
request_with_transient_retry(client, method="GET", url="https://up/v1/models")
|
||||
)
|
||||
|
||||
assert result is ok
|
||||
# one initial attempt + one retry on a fresh connection
|
||||
assert client.request.await_count == 2
|
||||
|
||||
|
||||
def test_helper_reraises_persistent_remote_protocol_error() -> None:
|
||||
client = SimpleNamespace(
|
||||
request=AsyncMock(side_effect=httpx.RemoteProtocolError(_INCOMPLETE_CHUNKED))
|
||||
)
|
||||
|
||||
with pytest.raises(httpx.RemoteProtocolError):
|
||||
asyncio.run(request_with_transient_retry(client, method="GET", url="https://up/v1/models"))
|
||||
|
||||
# default max_retries=1 → exactly 2 attempts before giving up
|
||||
assert client.request.await_count == 2
|
||||
|
||||
|
||||
def test_helper_does_not_retry_other_errors() -> None:
|
||||
client = SimpleNamespace(
|
||||
request=AsyncMock(side_effect=httpx.ConnectError("connection refused"))
|
||||
)
|
||||
|
||||
with pytest.raises(httpx.ConnectError):
|
||||
asyncio.run(request_with_transient_retry(client, method="GET", url="https://up/v1/models"))
|
||||
|
||||
# ConnectError is not a transient keep-alive close; no retry
|
||||
assert client.request.await_count == 1
|
||||
|
||||
|
||||
def test_helper_respects_max_retries() -> None:
|
||||
ok = _ok_response()
|
||||
client = SimpleNamespace(
|
||||
request=AsyncMock(
|
||||
side_effect=[
|
||||
httpx.RemoteProtocolError(_INCOMPLETE_CHUNKED),
|
||||
httpx.RemoteProtocolError(_INCOMPLETE_CHUNKED),
|
||||
ok,
|
||||
]
|
||||
)
|
||||
)
|
||||
|
||||
result = asyncio.run(
|
||||
request_with_transient_retry(
|
||||
client, method="GET", url="https://up/v1/models", max_retries=2
|
||||
)
|
||||
)
|
||||
|
||||
assert result is ok
|
||||
assert client.request.await_count == 3
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Handler-level tests: the exact path from the issue traceback
|
||||
# (proxy_routes.list_models → OpenAIHandlerMixin.handle_passthrough)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class _PassthroughModelsRequest:
|
||||
method = "GET"
|
||||
headers: dict[str, str] = {}
|
||||
url = SimpleNamespace(path="/v1/models", query="")
|
||||
|
||||
async def body(self) -> bytes:
|
||||
return b""
|
||||
|
||||
|
||||
class _FlakyThenOkClient:
|
||||
"""Raises RemoteProtocolError on the first request, then succeeds —
|
||||
a stale pooled keep-alive connection followed by a fresh one."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.calls = 0
|
||||
|
||||
async def request(self, **kwargs): # noqa: ANN003, ANN201
|
||||
self.calls += 1
|
||||
if self.calls == 1:
|
||||
raise httpx.RemoteProtocolError(_INCOMPLETE_CHUNKED)
|
||||
request = httpx.Request(kwargs["method"], kwargs["url"])
|
||||
return httpx.Response(
|
||||
200,
|
||||
request=request,
|
||||
headers={"content-type": "application/json"},
|
||||
json={"object": "list", "data": []},
|
||||
)
|
||||
|
||||
|
||||
class _AlwaysProtocolErrorClient:
|
||||
def __init__(self) -> None:
|
||||
self.calls = 0
|
||||
|
||||
async def request(self, **kwargs): # noqa: ANN003, ANN201
|
||||
self.calls += 1
|
||||
raise httpx.RemoteProtocolError(_INCOMPLETE_CHUNKED)
|
||||
|
||||
|
||||
def test_passthrough_recovers_from_incomplete_chunked_read() -> None:
|
||||
handler = object.__new__(OpenAIHandlerMixin)
|
||||
client = _FlakyThenOkClient()
|
||||
handler.http_client = client
|
||||
handler.http_client_h1 = client
|
||||
|
||||
response = asyncio.run(
|
||||
handler.handle_passthrough(_PassthroughModelsRequest(), "https://api.openai.com")
|
||||
)
|
||||
|
||||
assert response.status_code == 200
|
||||
assert client.calls == 2
|
||||
assert json.loads(response.body) == {"object": "list", "data": []}
|
||||
|
||||
|
||||
def test_passthrough_returns_clean_502_on_persistent_protocol_error() -> None:
|
||||
handler = object.__new__(OpenAIHandlerMixin)
|
||||
client = _AlwaysProtocolErrorClient()
|
||||
handler.http_client = client
|
||||
handler.http_client_h1 = client
|
||||
|
||||
response = asyncio.run(
|
||||
handler.handle_passthrough(_PassthroughModelsRequest(), "https://api.openai.com")
|
||||
)
|
||||
|
||||
assert response.status_code == 502
|
||||
payload = json.loads(response.body)
|
||||
assert payload["error"]["type"] == "upstream_protocol_error"
|
||||
assert "complete response" in payload["error"]["message"]
|
||||
# initial attempt + one retry, then the clean 502
|
||||
assert client.calls == 2
|
||||
Loading…
Add table
Add a link
Reference in a new issue