diff --git a/CHANGELOG.md b/CHANGELOG.md index d44291c26..5c8d74125 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## Unreleased ### Fixed +- Concurrent large requests no longer 502 on a transient HTTP/2 stream reset. + A single upstream `StreamReset` poisons the shared h2 connection and raises + `RemoteProtocolError` / `LocalProtocolError` on every in-flight request; those + transport errors weren't in the proxy's retry paths, so they collapsed + straight to a 502 with no reconnect. The Anthropic non-streaming and streaming + retry paths now treat any `httpx.TransportError` (including h2 protocol + errors) as retryable before the first client byte, so the bad connection + is dropped and the request re-sent on a fresh one + ([#1639](https://github.com/headroomlabs-ai/headroom/issues/1639)). - `headroom wrap claude` no longer leaves a dead `ANTHROPIC_BASE_URL` in a project's `.claude/settings.local.json` after an unclean exit (`SIGKILL`, OOM, reboot, or terminal/tmux close via `SIGHUP`, which was not caught). diff --git a/headroom/proxy/handlers/streaming.py b/headroom/proxy/handlers/streaming.py index 3723e6352..7581b420b 100644 --- a/headroom/proxy/handlers/streaming.py +++ b/headroom/proxy/handlers/streaming.py @@ -1078,7 +1078,12 @@ class StreamingMixin: await asyncio.sleep(delay_with_jitter / 1000) continue break - except (httpx.ConnectError, httpx.ConnectTimeout, httpx.PoolTimeout) as e: + # Retry any transport-level failure while opening the upstream + # stream — including HTTP/2 protocol errors (Local/RemoteProtocol + # `StreamReset`) from a poisoned shared h2 connection. This runs + # before any body byte is forwarded to the client, so re-sending + # on a fresh connection is safe and avoids a 502. (#1639) + except httpx.TransportError as e: last_connect_error = e if attempt >= retry_attempts - 1: raise @@ -1097,7 +1102,10 @@ class StreamingMixin: if upstream_response is None: raise last_connect_error or RuntimeError("upstream connection did not start") - except (httpx.ConnectError, httpx.ConnectTimeout, httpx.PoolTimeout) as e: + # Retries exhausted (or a transport failure escaped the loop): emit a + # clean SSE error instead of letting an h2 StreamReset bubble up as a + # 502. Covers ConnectError/timeouts and Local/RemoteProtocolError. (#1639) + except httpx.TransportError as e: error_msg = str(e) or repr(e) logger.error(f"[{request_id}] Connection error to upstream API: {error_msg}") diff --git a/headroom/proxy/server.py b/headroom/proxy/server.py index 28c320896..062c43fe5 100644 --- a/headroom/proxy/server.py +++ b/headroom/proxy/server.py @@ -1875,7 +1875,12 @@ class HeadroomProxy( return response - except (httpx.ConnectError, httpx.TimeoutException, httpx.HTTPStatusError) as e: + # httpx.TransportError covers ConnectError, the timeout family, and — + # crucially — the protocol errors (Local/RemoteProtocolError, e.g. an + # HTTP/2 `StreamReset`) that a poisoned shared h2 connection raises on + # every in-flight request. Retrying drops the bad connection and + # re-sends on a fresh one instead of collapsing to a 502. (#1639) + except (httpx.TransportError, httpx.HTTPStatusError) as e: last_error = e if not self.config.retry_enabled or attempt >= self.config.retry_max_attempts - 1: diff --git a/tests/test_h2_stream_reset_retry.py b/tests/test_h2_stream_reset_retry.py new file mode 100644 index 000000000..f1774999c --- /dev/null +++ b/tests/test_h2_stream_reset_retry.py @@ -0,0 +1,144 @@ +"""HTTP/2 stream-reset resilience (issue #1639). + +Under concurrent load a single upstream HTTP/2 stream reset poisons the shared +h2 connection and surfaces as `RemoteProtocolError` / `LocalProtocolError` on +every in-flight request. Those are transport errors, so the proxy must retry +them (dropping the bad connection and re-sending on a fresh one) instead of +collapsing to a 502. These tests drive the real `_retry_request` and +`_stream_response` paths. +""" + +from __future__ import annotations + +from unittest.mock import AsyncMock, MagicMock + +import httpx +import pytest + +from headroom.proxy.server import HeadroomProxy + + +def _mock_proxy(): + proxy = object.__new__(HeadroomProxy) + proxy.http_client = MagicMock(spec=httpx.AsyncClient) + proxy._config = MagicMock() + proxy._config.memory_enabled = False + proxy._config.ccr_inject_tool = False + proxy._config.retry_enabled = True + proxy._config.retry_max_attempts = 2 + proxy._config.retry_base_delay_ms = 0 + proxy._config.retry_max_delay_ms = 0 + proxy.config = proxy._config + proxy.memory_handler = None + proxy._parse_sse_usage_from_buffer = MagicMock(return_value=None) + proxy._finalize_stream_response = AsyncMock(return_value=None) + return proxy + + +def _good_stream_response(chunks): + resp = AsyncMock() + resp.headers = httpx.Headers({"content-type": "text/event-stream"}) + resp.status_code = 200 + + async def aiter_bytes(): + for chunk in chunks: + yield chunk + + resp.aiter_bytes = aiter_bytes + resp.aclose = AsyncMock() + return resp + + +async def _run_stream(proxy, session_key="k"): + return await proxy._stream_response( + url="https://api.anthropic.com/v1/messages", + headers={"x-api-key": "sk-test"}, + body={ + "model": "claude-sonnet-4-20250514", + "max_tokens": 100, + "stream": True, + "messages": [{"role": "user", "content": "hi"}], + }, + provider="anthropic", + model="claude-sonnet-4-20250514", + request_id="test-1639", + original_tokens=10, + optimized_tokens=10, + tokens_saved=0, + transforms_applied=[], + tags={}, + optimization_latency=0.0, + session_key=session_key, + ) + + +@pytest.mark.asyncio +async def test_retry_request_retries_remote_protocol_error(): + proxy = _mock_proxy() + good = MagicMock() + good.status_code = 200 + good.request = MagicMock() + proxy.http_client.post = AsyncMock( + side_effect=[httpx.RemoteProtocolError(""), good] + ) + + result = await proxy._retry_request( + "POST", + "https://api.anthropic.com/v1/messages", + {"x-api-key": "sk-test"}, + {"model": "claude-sonnet-4-20250514", "messages": []}, + ) + + assert result is good + assert proxy.http_client.post.await_count == 2 + + +@pytest.mark.asyncio +async def test_retry_request_reraises_after_exhaustion(): + proxy = _mock_proxy() + proxy.http_client.post = AsyncMock(side_effect=httpx.RemoteProtocolError("reset")) + + with pytest.raises(httpx.RemoteProtocolError): + await proxy._retry_request( + "POST", + "https://api.anthropic.com/v1/messages", + {"x-api-key": "sk-test"}, + {"model": "claude-sonnet-4-20250514", "messages": []}, + ) + assert proxy.http_client.post.await_count == 2 + + +@pytest.mark.asyncio +async def test_stream_retries_h2_stream_reset_then_succeeds(): + proxy = _mock_proxy() + good = _good_stream_response( + [ + b'event: message_start\ndata: {"type":"message_start"}\n\n', + b'event: message_stop\ndata: {"type":"message_stop"}\n\n', + ] + ) + proxy.http_client.build_request = MagicMock(return_value=MagicMock()) + proxy.http_client.send = AsyncMock( + side_effect=[httpx.RemoteProtocolError(""), good] + ) + + result = await _run_stream(proxy) + body = b"".join([chunk async for chunk in result.body_iterator]) + + assert proxy.http_client.send.await_count == 2 + assert b"message_start" in body + assert b"connection_error" not in body + + +@pytest.mark.asyncio +async def test_stream_reset_exhaustion_yields_sse_error_not_crash(): + proxy = _mock_proxy() + proxy.http_client.build_request = MagicMock(return_value=MagicMock()) + proxy.http_client.send = AsyncMock(side_effect=httpx.RemoteProtocolError("reset")) + + result = await _run_stream(proxy) + body = b"".join([chunk async for chunk in result.body_iterator]) + + assert proxy.http_client.send.await_count == 2 + assert b"event: error" in body + assert b"connection_error" in body