fix(proxy/anthropic): stop answering a non-streaming turn with an event stream (#3142)
## Description Closes #3130. Unifies #3131 (@Joaovsales) and #3132 (@taiseii), which landed within hours of each other on the same bug. Neither is redundant — **#3131 contributed the clearest statement of the contract; #3132 contributed the reconstruction that can actually be trusted to satisfy it.** This takes both. A caller that sent `stream: false` was handed a `text/event-stream` body at HTTP 200. The reply was complete — 8756 bytes, a valid upstream `request-id` — it was simply wearing a wire format the SDK cannot parse, so the turn was lost. **On root cause.** #3130 says outright: *"I could not pin down why the upstream answered a `stream`-less request with an event stream."* I think this does. At `v0.35.0` the CCR path flips the body to `stream: false` and never touches the client's `Accept` header — I checked the tag and the count of Accept rewrites at that site is **zero**. So upstream receives a self-contradicting request: *"answer as JSON"* in the body, *"I only accept SSE"* in the headers. Both reporters (#3130, #3140) show `server: cloudflare` / `cf-ray`, and both describe it as intermittent — consistent with an edge honouring `Accept` under retry. #3102 fixed that for the CCR flip; this PR moves the rewrite to the buffered boundary **every** non-streaming request reaches, so the client's own non-streaming retry is covered too. ## Type of Change - [x] Bug fix (non-breaking change that fixes an issue) ## Changes Made **From #3131 — the contract.** `headroom/proxy/nonstream_sse_policy.py`: a pure module with a behaviour matrix and `should_recover_sse_reply` as a single predicate. The three negative arms are deliberate — a streaming caller wants SSE, a JSON content-type is already correct, a non-200 carries an upstream error the client should see verbatim. **From #3132 — the reconstruction.** `require_complete=True` demands `message_start`, a terminal `message_stop`, every opened block closed, no in-band `error` event, and no delta the reconstructor cannot replay. Anything short of that is a 502. Three things only #3132 had, each load-bearing: - **`index` is stripped from rebuilt content blocks.** The parser writes it (`streaming.py:425`) and a client persists the reconstructed turn and echoes it back — at which point Anthropic 400s with `content.0.text.index: Extra inputs are not permitted`. `_strip_streaming_only_content_fields` (`anthropic.py:185`) already documents this exact failure. That inbound stripper would mask it *while the proxy is in the path*, but the client's stored history is still polluted. - **SSE framing is normalized and `data:` no longer requires the optional space.** The old `startswith("data: ")` skipped a spec-valid stream **entirely** — zero events parsed, which is literally what the report describes (*"0 stream events received"*). - **Detection sniffs the body**, so a mislabeled or absent content-type is still caught. **Reconciled where they disagreed:** - *Headers.* #3131 hand-rolled a framing list; this uses the established `sanitize_forwarded_response_headers`. That already strips `connection`, `keep-alive` and `server` alongside the content-* family — and per the comment at `helpers.py:325`, leaving `transfer-encoding` on a rebuilt body is what produced an empty HTTP 200 in #3019. #3131's list would have left three of those on. #3132's `cf-*` filter is kept. - *Detection.* The body sniff arrives as `body_is_event_stream`, so the policy module stays pure — the sniff needs the response object and the handler owns that. - Dropped #3131's `json_reply_headers` and its test class; everything else from both PRs is retained. ## Testing - [x] Unit tests pass (`pytest`) - [x] Integration tests pass ### Test Output ```text tests/test_nonstream_sse_policy.py 18 passed (from #3131) tests/test_anthropic_buffered_sse.py 18 passed (from #3132) 36 passed Regression sweep (-k "stream or sse or ccr or anthropic or proxy or buffered or usage"): 2982 passed, 181 skipped, 0 failed in 153.41s ruff check: All checks passed ruff format --check: 527 files already formatted ``` Both contributors' suites are kept whole and both pass unmodified against the merged implementation, which is the useful signal here — they were written independently against different implementations. ## Real Behavior Proof - Environment: macOS (darwin 25.4.0), Python 3.12.13, worktree off `main`, `_core.abi3.so` copied in. - Exact command / steps: applied #3132 as the engine, layered #3131's policy module over it, rewired the decision site to the predicate, then ran both suites and a 2982-test sweep concentrated on everything touching the shared SSE parser. - Observed result: 36/36 across both contributed suites, 2982 passed / 0 failed on the sweep. The sweep matters more than usual here — `_parse_sse_to_response` is shared with the streaming path's usage accounting, and `require_complete` defaults to `False` specifically so existing callers keep the lenient reconstruction they were written against. Nothing regressed. - Not tested: no live upstream. I could not reproduce the upstream answering a `stream`-less request with SSE against real `api.anthropic.com` — that is the condition #3130 reports as intermittent and load-dependent, and the Accept explanation above remains a well-supported hypothesis rather than something I observed. The fix does not depend on it: whatever the upstream returns, a caller that did not ask for streaming is no longer handed SSE. ## Runtime Rollout Safety - Rollout-managed feature(s): None. - Minimum rollout channel: n/a - Stable/default behavior changed: Yes, deliberately, in two places. A non-streaming turn answered with SSE is now reconstructed as JSON instead of relayed; an SSE reply that cannot be faithfully reconstructed is now a 502 instead of an unparseable 200. Both are the point. `require_complete` defaults to `False`, so streaming callers of the shared parser are untouched. - Kill switch / disable path: none by design — relaying a body the client cannot parse has no legitimate mode. - Unsafe override required: No. - Qualification impact: A truncated upstream stream now surfaces as an explicit 502 rather than a short-but-successful turn. More visible failures, fewer silent ones. - Rollback path: Revert the commit. ## Review Readiness - [x] I have performed a self-review - [x] This PR is ready for human review ## Additional Notes If this lands, #3131 and #3132 should be closed as superseded — both authors are credited via `Co-authored-by:` and their tests ship intact. I would not close either before a maintainer agrees this unification is the direction, since it discards a design decision from each. **Wider context, not fixed here:** #3130 and #3140 both report against **0.35.0**, and `main` already carries a stack of fixes for this symptom class that has never shipped — #3102 (Accept), #3092, #3091, #3094, #3101, #3069, #3084, #3124, #3134. All of them are gated behind #3067 `chore: release 0.36.0`. Every closed lookalike (#3019, #3055, #3071, #3040, #2952) was fixed into that same unreleased window. Merging this PR does not help either reporter until 0.36.0 ships; **cutting that release is the higher-leverage action.** The interim workaround for anyone on 0.35.0 is `HEADROOM_NO_CCR=1` — the buffered flip is gated on `_has_headroom_retrieve_tool`, and `no_ccr` stops the tool being injected, so the flip never engages. Note `headroom wrap` has no `--no-ccr` flag in 0.35.0, so it has to be the env var. Co-authored-by: Tejas Chopra <tejas@Tejass-MacBook-Pro.local> Co-authored-by: João Souto <73318835+Joaovsales@users.noreply.github.com> Co-authored-by: taiseii <37083727+taiseii@users.noreply.github.com> Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -51,6 +51,7 @@ from headroom.proxy.image_isolation import run_image_compression_isolated
|
||||
from headroom.proxy.memory_decision import MemoryDecision
|
||||
from headroom.proxy.memory_query import MemoryQuery
|
||||
from headroom.proxy.model_router import estimate_input_tokens
|
||||
from headroom.proxy.nonstream_sse_policy import should_recover_sse_reply
|
||||
from headroom.proxy.outcome import RequestOutcome
|
||||
|
||||
logger = logging.getLogger("headroom.proxy")
|
||||
@@ -226,6 +227,66 @@ def _looks_like_sse_response(response: httpx.Response) -> bool:
|
||||
class AnthropicHandlerMixin:
|
||||
"""Mixin providing Anthropic API handler methods for HeadroomProxy."""
|
||||
|
||||
def _adapt_event_stream_to_json(
|
||||
self,
|
||||
response: httpx.Response,
|
||||
request_id: str,
|
||||
) -> httpx.Response:
|
||||
"""Rebuild an SSE reply as the JSON a non-streaming caller asked for.
|
||||
|
||||
A caller that sent ``stream: false`` cannot parse ``text/event-stream``,
|
||||
so relaying it verbatim loses a turn the upstream already charged for
|
||||
(#3130). Reconstruction is strict: a truncated stream, or one carrying
|
||||
an ``error`` event, becomes an explicit 502 rather than a successful
|
||||
HTTP 200 whose message is silently short.
|
||||
"""
|
||||
headers = {
|
||||
k: v
|
||||
for k, v in sanitize_forwarded_response_headers(
|
||||
response.headers,
|
||||
"content-type",
|
||||
).items()
|
||||
if not k.lower().startswith("cf-")
|
||||
}
|
||||
|
||||
parsed = None
|
||||
try:
|
||||
parsed = self._parse_sse_to_response(
|
||||
response.content.decode("utf-8", "replace"),
|
||||
"anthropic",
|
||||
require_complete=True,
|
||||
)
|
||||
except Exception as exc: # pragma: no cover - defensive
|
||||
logger.warning(f"[{request_id}] SSE->JSON reconstruction raised: {exc}")
|
||||
|
||||
if parsed is None:
|
||||
logger.error(
|
||||
f"[{request_id}] Upstream answered a non-streaming request with an "
|
||||
f"event stream that could not be faithfully reconstructed "
|
||||
f"(body_bytes={len(response.content)}); returning 502 rather than "
|
||||
f"a wire format the client cannot parse"
|
||||
)
|
||||
return httpx.Response(
|
||||
502,
|
||||
json={
|
||||
"type": "error",
|
||||
"error": {
|
||||
"type": "upstream_protocol_error",
|
||||
"message": (
|
||||
"Upstream answered a non-streaming request with an "
|
||||
"incomplete event stream."
|
||||
),
|
||||
},
|
||||
},
|
||||
headers=headers,
|
||||
)
|
||||
|
||||
logger.info(
|
||||
f"[{request_id}] Upstream answered a non-streaming request with an "
|
||||
f"event stream; adapted {len(response.content)} bytes of SSE to JSON"
|
||||
)
|
||||
return httpx.Response(200, json=parsed, headers=headers)
|
||||
|
||||
async def _count_tokens_offloaded(self, model, messages): # noqa: ANN001, ANN201
|
||||
from headroom.proxy.token_counting import count_tokens_offloaded
|
||||
|
||||
@@ -3509,20 +3570,9 @@ class AnthropicHandlerMixin:
|
||||
body_mutation_tracker.mark_mutated(
|
||||
"ccr_streaming_retrieve_buffered_non_stream"
|
||||
)
|
||||
# The body now asks for a non-streaming reply, so the
|
||||
# client's ``Accept: text/event-stream`` no longer describes
|
||||
# the response being requested. Forwarding it unchanged
|
||||
# sends upstream a self-contradicting request: "answer as
|
||||
# JSON" in the body, "I only accept SSE" in the headers.
|
||||
#
|
||||
# Anthropic tolerates that. Stricter Anthropic-compatible
|
||||
# gateways do not: GitHub Copilot's returns a generic
|
||||
# ``api_error``, which is why a session's first call
|
||||
# succeeded and the next one — the first to carry a
|
||||
# redeemable marker, and so the first to be buffered —
|
||||
# failed (#3078).
|
||||
_accept_key = next((k for k in headers if k.lower() == "accept"), "accept")
|
||||
headers[_accept_key] = "application/json"
|
||||
# The ``Accept`` rewrite this flip used to do lives at the
|
||||
# buffered boundary below, which every non-streaming
|
||||
# request reaches — this one and the client's own (#3130).
|
||||
logger.info(
|
||||
f"[{request_id}] CCR: stream:true request has "
|
||||
"headroom_retrieve available; using buffered stream:false "
|
||||
@@ -3681,6 +3731,29 @@ class AnthropicHandlerMixin:
|
||||
session_key=session_key,
|
||||
)
|
||||
else:
|
||||
# Whatever set it — the client's own ``stream: false`` or
|
||||
# the CCR flip above — this branch sends a non-streaming
|
||||
# request, so the client's ``Accept: text/event-stream`` no
|
||||
# longer describes what is being asked for. Forwarding it
|
||||
# unchanged puts a self-contradicting request on the wire:
|
||||
# "answer as JSON" in the body, "I only accept SSE" in the
|
||||
# headers.
|
||||
#
|
||||
# Anthropic tolerates the contradiction; stricter
|
||||
# Anthropic-compatible gateways do not — GitHub Copilot's
|
||||
# answers a generic ``api_error`` (#3078). On a
|
||||
# client-originated non-stream turn — Claude Code's retry
|
||||
# after a failed stream — an SSE answer to a JSON request
|
||||
# is the empty/malformed HTTP 200 of #3130.
|
||||
#
|
||||
# Mutated in place: ``headers`` is captured by the
|
||||
# closures defined below, and rebinding it here would
|
||||
# leave them holding the old mapping.
|
||||
if body.get("stream", False) is False:
|
||||
for _accept_key in [k for k in headers if k.lower() == "accept"]:
|
||||
headers.pop(_accept_key, None)
|
||||
headers["accept"] = "application/json"
|
||||
|
||||
# Populated once the upstream answers 200 with parseable
|
||||
# JSON, so the guard below can fall back to it (#3088).
|
||||
_salvageable_upstream: dict[str, Any] = {}
|
||||
@@ -3851,6 +3924,29 @@ class AnthropicHandlerMixin:
|
||||
f"[{request_id}] Failed to write debug dump: {dump_err}"
|
||||
)
|
||||
|
||||
# A non-streaming request answered with an event
|
||||
# stream (#3130). The turn is complete and already
|
||||
# paid for — it is just wearing the wrong wire
|
||||
# format — so adapt it to the JSON this caller asked
|
||||
# for *here*, ahead of everything that reads the
|
||||
# body: CCR retrieval, memory, turn hooks, usage and
|
||||
# cost accounting, prefix tracking, the response
|
||||
# cache, marker resolution and the security scan.
|
||||
# Adapting at the final return instead would leave
|
||||
# every one of those looking at an unparseable body.
|
||||
#
|
||||
# ``stream`` is what the *client* asked for, not what
|
||||
# went upstream: a buffered CCR turn deliberately
|
||||
# requests JSON on behalf of a streaming client and
|
||||
# re-emits SSE further down, and must keep doing so.
|
||||
if should_recover_sse_reply(
|
||||
client_requested_stream=bool(stream),
|
||||
status_code=response.status_code,
|
||||
content_type=response.headers.get("content-type"),
|
||||
body_is_event_stream=_looks_like_sse_response(response),
|
||||
):
|
||||
response = self._adapt_event_stream_to_json(response, request_id)
|
||||
|
||||
# Parse response for CCR handling
|
||||
resp_json = None
|
||||
try:
|
||||
@@ -4363,12 +4459,23 @@ class AnthropicHandlerMixin:
|
||||
)
|
||||
)
|
||||
|
||||
# Remove compression headers since httpx already decompressed the response
|
||||
response_headers = dict(response.headers)
|
||||
response_headers.pop("content-encoding", None)
|
||||
response_headers.pop(
|
||||
"content-length", None
|
||||
) # Length changed after decompression
|
||||
# Framing headers describe how the *upstream* framed
|
||||
# its body, not what this response is: httpx already
|
||||
# decompressed it, Starlette recomputes the length,
|
||||
# and uvicorn owns the connection. Replaying a stale
|
||||
# ``transfer-encoding: chunked`` over a fixed-length
|
||||
# body is what made an HTTP 200 read as empty in
|
||||
# #3019. ``cf-*`` is CDN provenance the caller has no
|
||||
# use for, and the header set clients cite as
|
||||
# evidence of an intermediary mangling a reply
|
||||
# (#3130).
|
||||
response_headers = {
|
||||
k: v
|
||||
for k, v in sanitize_forwarded_response_headers(
|
||||
response.headers
|
||||
).items()
|
||||
if not k.lower().startswith("cf-")
|
||||
}
|
||||
|
||||
# Inject Headroom compression metrics (for SaaS metering)
|
||||
response_headers["x-headroom-tokens-before"] = str(original_tokens)
|
||||
|
||||
@@ -350,7 +350,13 @@ class StreamingMixin:
|
||||
|
||||
return usage_found if usage_found else None
|
||||
|
||||
def _parse_sse_to_response(self, sse_data: str, provider: str) -> dict[str, Any] | None:
|
||||
def _parse_sse_to_response(
|
||||
self,
|
||||
sse_data: str,
|
||||
provider: str,
|
||||
*,
|
||||
require_complete: bool = False,
|
||||
) -> dict[str, Any] | None:
|
||||
"""Parse SSE data to reconstruct the API response JSON.
|
||||
|
||||
Args:
|
||||
@@ -358,6 +364,10 @@ class StreamingMixin:
|
||||
from a complete-events bytes buffer (see
|
||||
``parse_sse_events_from_byte_buffer``).
|
||||
provider: Provider type for parsing.
|
||||
require_complete: Reject anything short of a whole, replayable
|
||||
message — see the strictness note below. Off by default so
|
||||
streaming callers keep the lenient reconstruction they were
|
||||
written against.
|
||||
|
||||
Returns:
|
||||
Reconstructed response dict or None if parsing fails.
|
||||
@@ -366,11 +376,32 @@ class StreamingMixin:
|
||||
``text_delta``, ``input_json_delta``, ``thinking_delta``,
|
||||
``signature_delta``, ``citations_delta``. Also preserves
|
||||
``redacted_thinking.data`` and accumulates citations as a list.
|
||||
|
||||
Permissive mode answers with whatever blocks it managed to
|
||||
accumulate. That is right for a streaming caller salvaging a
|
||||
partial stream, and wrong for #3130, where the reconstruction is
|
||||
handed to the client *as* the turn: a truncated stream would become
|
||||
a successful — and silently short — message, and an ``error`` event
|
||||
would vanish behind an HTTP 200. ``require_complete`` demands
|
||||
``message_start``, a terminal ``message_stop``, every opened block
|
||||
closed, no ``error`` event, and no delta type this reconstructor
|
||||
cannot replay; anything else returns None so the caller can fail
|
||||
loudly instead.
|
||||
"""
|
||||
if provider != "anthropic":
|
||||
return None # Only implemented for Anthropic
|
||||
|
||||
# Event framing is CRLF in some intermediaries (and mixed after a
|
||||
# retry through one). Normalize before the line split so a
|
||||
# correctly framed stream is never read as zero events.
|
||||
sse_data = sse_data.replace("\r\n", "\n").replace("\r", "\n")
|
||||
|
||||
response: dict[str, Any] = {"content": [], "usage": {}}
|
||||
saw_message_start = False
|
||||
saw_message_stop = False
|
||||
saw_error = False
|
||||
saw_unreplayable_delta = False
|
||||
open_block_indices: set[int] = set()
|
||||
# Track blocks by their `index` field so out-of-order events
|
||||
# don't corrupt the reconstruction. The current block pointer
|
||||
# remains for backward-compat with code that walks this dict
|
||||
@@ -392,9 +423,9 @@ class StreamingMixin:
|
||||
appended_block_keys: set[int] = set()
|
||||
|
||||
for line in sse_data.split("\n"):
|
||||
if not line.startswith("data: "):
|
||||
if not line.startswith("data:"):
|
||||
continue
|
||||
data_str = line[6:].strip()
|
||||
data_str = line[5:].strip()
|
||||
if not data_str or data_str == "[DONE]":
|
||||
continue
|
||||
|
||||
@@ -406,8 +437,12 @@ class StreamingMixin:
|
||||
event_type = data.get("type", "")
|
||||
|
||||
if event_type == "message_start":
|
||||
saw_message_start = True
|
||||
msg = data.get("message", {})
|
||||
response["id"] = msg.get("id")
|
||||
response["type"] = msg.get("type", "message")
|
||||
if "stop_sequence" in msg:
|
||||
response["stop_sequence"] = msg["stop_sequence"]
|
||||
response["model"] = msg.get("model")
|
||||
response["role"] = msg.get("role", "assistant")
|
||||
response["stop_reason"] = msg.get("stop_reason")
|
||||
@@ -454,6 +489,7 @@ class StreamingMixin:
|
||||
if _k != "type":
|
||||
current_block[_k] = _v
|
||||
blocks_by_index[block_index] = current_block
|
||||
open_block_indices.add(block_index)
|
||||
|
||||
elif event_type == "content_block_delta":
|
||||
# Resolve the target block by index (preferred) or fall
|
||||
@@ -490,6 +526,10 @@ class StreamingMixin:
|
||||
citation = delta.get("citation")
|
||||
if citation is not None:
|
||||
citations.append(citation)
|
||||
else:
|
||||
# A delta this reconstructor has no rule for: the
|
||||
# accumulated block is missing whatever it carried.
|
||||
saw_unreplayable_delta = True
|
||||
|
||||
elif event_type == "content_block_stop":
|
||||
idx = data.get("index")
|
||||
@@ -523,6 +563,8 @@ class StreamingMixin:
|
||||
if block_key not in appended_block_keys:
|
||||
response["content"].append(target)
|
||||
appended_block_keys.add(block_key)
|
||||
if idx is not None:
|
||||
open_block_indices.discard(idx)
|
||||
current_block = None
|
||||
|
||||
elif event_type == "message_delta":
|
||||
@@ -531,9 +573,38 @@ class StreamingMixin:
|
||||
response["stop_reason"] = delta["stop_reason"]
|
||||
if "stop_details" in delta:
|
||||
response["stop_details"] = delta["stop_details"]
|
||||
if "stop_sequence" in delta:
|
||||
response["stop_sequence"] = delta["stop_sequence"]
|
||||
if data.get("usage"):
|
||||
response["usage"].update(data["usage"])
|
||||
|
||||
elif event_type == "message_stop":
|
||||
saw_message_stop = True
|
||||
|
||||
elif event_type == "error":
|
||||
# An in-band failure. Permissive callers keep salvaging
|
||||
# what arrived before it; a strict caller must not dress
|
||||
# the remains up as a successful turn.
|
||||
saw_error = True
|
||||
|
||||
if require_complete:
|
||||
if (
|
||||
not saw_message_start
|
||||
or not saw_message_stop
|
||||
or saw_error
|
||||
or saw_unreplayable_delta
|
||||
or open_block_indices
|
||||
):
|
||||
return None
|
||||
# ``index`` is a response-delta field. Anthropic rejects it on
|
||||
# the next request ("content.0.text.index: Extra inputs are not
|
||||
# permitted"), so it must not survive into a body the client
|
||||
# will persist and echo back.
|
||||
for block in response["content"]:
|
||||
if isinstance(block, dict):
|
||||
block.pop("index", None)
|
||||
return response
|
||||
|
||||
return response if response.get("content") else None
|
||||
|
||||
def _response_to_sse(self, response: dict[str, Any], provider: str) -> list[bytes]:
|
||||
|
||||
@@ -0,0 +1,129 @@
|
||||
"""Wire-format contract policy for the buffered (non-streaming) reply path.
|
||||
|
||||
The problem
|
||||
-----------
|
||||
|
||||
The buffered Anthropic path returns the upstream reply with its headers
|
||||
copied wholesale::
|
||||
|
||||
response_headers = dict(response.headers)
|
||||
response_headers.pop("content-encoding", None)
|
||||
response_headers.pop("content-length", None)
|
||||
...
|
||||
return Response(content=..., status_code=..., headers=response_headers)
|
||||
|
||||
``content-type`` rides along untouched. When the upstream answers a
|
||||
``stream``-less request with ``text/event-stream``, that body reaches a
|
||||
caller that asked for JSON, as a ``200`` it cannot parse. Clients report
|
||||
it as an empty or malformed response and the turn is lost — the reply is
|
||||
*present and complete*, just wearing the wrong wire format.
|
||||
|
||||
The buffered-stream (CCR) path already refuses this shape, logging the
|
||||
offending ``content-type`` and returning ``upstream_protocol_error``
|
||||
(#2952). The plain non-streaming path never got the same treatment: an
|
||||
unparseable body there was assumed to mean "no CCR handling", so it was
|
||||
logged at DEBUG and passed through.
|
||||
|
||||
The contract
|
||||
------------
|
||||
|
||||
A caller that did not set ``stream: true`` must never receive an
|
||||
event-stream body. Headroom owns both ends of that boundary, so it can
|
||||
enforce it rather than let the mismatch reach the client.
|
||||
|
||||
Behaviour matrix
|
||||
----------------
|
||||
|
||||
============================ ============== ========= ====================
|
||||
Client asked for streaming? Upstream C-T Status Result
|
||||
============================ ============== ========= ====================
|
||||
yes any any untouched
|
||||
no application/… any untouched
|
||||
no text/event-… != 200 untouched (real error)
|
||||
no text/event-… 200 recover, else refuse
|
||||
============================ ============== ========= ====================
|
||||
|
||||
Recovery reuses ``StreamingMixin._parse_sse_to_response``, the same
|
||||
reconstruction the streaming path already runs for usage accounting, so
|
||||
this adds no new parsing surface. Recovering in place — rather than
|
||||
returning early — keeps the rest of the buffered path (CCR, turn hooks,
|
||||
security scan, usage accounting) operating on a normal reply.
|
||||
|
||||
Non-200 is deliberately excluded: an error status is already actionable
|
||||
by the client, and passing it through unchanged preserves the upstream's
|
||||
own error payload.
|
||||
|
||||
Public API
|
||||
----------
|
||||
|
||||
* :func:`is_event_stream` — media-type test, parameter- and case-tolerant.
|
||||
* :func:`should_recover_sse_reply` — the gate above, as one predicate.
|
||||
|
||||
Header correction is *not* here — see the note beside the public functions.
|
||||
|
||||
Constraints (per project memory)
|
||||
--------------------------------
|
||||
|
||||
* pure: no I/O, no logging, no config — the handler owns those.
|
||||
* no regexes: media-type parsing is a single ``split``.
|
||||
* no silent fallbacks: the caller refuses loudly when recovery fails.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
SSE_MEDIA_TYPE = "text/event-stream"
|
||||
JSON_MEDIA_TYPE = "application/json"
|
||||
|
||||
|
||||
def media_type(content_type: str | None) -> str:
|
||||
"""Return the bare media type, lower-cased, with parameters dropped.
|
||||
|
||||
``"text/event-stream; charset=utf-8"`` and ``"Text/Event-Stream"`` both
|
||||
yield ``"text/event-stream"``. Returns ``""`` for a missing header.
|
||||
"""
|
||||
if not content_type:
|
||||
return ""
|
||||
return content_type.split(";", 1)[0].strip().lower()
|
||||
|
||||
|
||||
def is_event_stream(content_type: str | None) -> bool:
|
||||
"""True when ``content_type`` denotes an SSE body."""
|
||||
return media_type(content_type) == SSE_MEDIA_TYPE
|
||||
|
||||
|
||||
def should_recover_sse_reply(
|
||||
*,
|
||||
client_requested_stream: bool,
|
||||
status_code: int,
|
||||
content_type: str | None,
|
||||
body_is_event_stream: bool = False,
|
||||
) -> bool:
|
||||
"""True when a buffered reply violates the caller's non-streaming contract.
|
||||
|
||||
See the behaviour matrix in the module docstring. The three negative
|
||||
arms are all deliberate: a streaming caller *wants* SSE, a JSON
|
||||
content-type is already correct, and a non-200 carries an upstream
|
||||
error the client should see verbatim.
|
||||
|
||||
``body_is_event_stream`` covers the reply that *is* an event stream while
|
||||
saying otherwise — a mislabeled or absent ``content-type``. Trusting the
|
||||
declared type alone would let exactly the same unparseable body through,
|
||||
so the caller sniffs the payload and passes the answer in. It stays a
|
||||
parameter rather than an import because this module is pure: the sniff
|
||||
needs the response object, and the handler owns that.
|
||||
"""
|
||||
if client_requested_stream:
|
||||
return False
|
||||
if status_code != 200:
|
||||
return False
|
||||
return is_event_stream(content_type) or body_is_event_stream
|
||||
|
||||
|
||||
# Header correction deliberately lives in
|
||||
# ``helpers.sanitize_forwarded_response_headers`` rather than here. It already
|
||||
# owns ``FRAMING_RESPONSE_HEADERS`` and strips ``connection``, ``keep-alive``
|
||||
# and ``server`` alongside the content-* family — a second copy of that list
|
||||
# would drift, and the ones this module would have missed are load-bearing:
|
||||
# leaving ``transfer-encoding`` on a rebuilt body is what produced an empty
|
||||
# HTTP 200 in #3019, and ``server: cloudflare`` is one of the headers the
|
||||
# client cites as evidence of an intermediary.
|
||||
@@ -0,0 +1,318 @@
|
||||
"""A non-streaming turn must never be answered with an event stream (#3130).
|
||||
|
||||
Claude Code retries a failed streaming turn as ``stream: false``. The buffered
|
||||
Anthropic path forwarded the upstream response headers wholesale, so when the
|
||||
upstream answered that JSON request with ``content-type: text/event-stream``
|
||||
the SDK got a wire format it never asked for and lost a complete, already-paid
|
||||
turn:
|
||||
|
||||
API returned an empty or malformed response (HTTP 200) ... content-type
|
||||
event-stream, body is an event stream (the non-streaming request was
|
||||
answered with a stream), 8756 bytes
|
||||
|
||||
Two defects, fixed on both sides:
|
||||
|
||||
* the request went out contradicting itself — ``stream: false`` in the body,
|
||||
``Accept: text/event-stream`` in the headers (the narrow CCR-only rewrite
|
||||
from #3078 never covered a client-originated non-stream turn), and
|
||||
* the response was relayed verbatim instead of being adapted to the JSON the
|
||||
caller asked for.
|
||||
|
||||
Reconstruction is deliberately strict: a partial stream must fail loudly as a
|
||||
502 rather than be handed back as a successful — and silently truncated —
|
||||
message.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
|
||||
import pytest
|
||||
|
||||
fastapi = pytest.importorskip("fastapi")
|
||||
httpx = pytest.importorskip("httpx")
|
||||
|
||||
from fastapi.testclient import TestClient # noqa: E402
|
||||
|
||||
from headroom.proxy.server import ProxyConfig, create_app # noqa: E402
|
||||
|
||||
COMPLETE_SSE = (
|
||||
"event: message_start\n"
|
||||
'data: {"type":"message_start","message":{"id":"msg_1","type":"message",'
|
||||
'"role":"assistant","model":"claude-sonnet-4-6","content":[],'
|
||||
'"stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":10,'
|
||||
'"output_tokens":1,"cache_read_input_tokens":2,'
|
||||
'"cache_creation_input_tokens":3}}}\n\n'
|
||||
"event: content_block_start\n"
|
||||
'data: {"type":"content_block_start","index":0,'
|
||||
'"content_block":{"type":"text","text":""}}\n\n'
|
||||
"event: content_block_delta\n"
|
||||
'data: {"type":"content_block_delta","index":0,'
|
||||
'"delta":{"type":"text_delta","text":"hello"}}\n\n'
|
||||
"event: content_block_stop\n"
|
||||
'data: {"type":"content_block_stop","index":0}\n\n'
|
||||
"event: message_delta\n"
|
||||
'data: {"type":"message_delta","delta":{"stop_reason":"end_turn",'
|
||||
'"stop_sequence":null},"usage":{"output_tokens":5}}\n\n'
|
||||
"event: message_stop\n"
|
||||
'data: {"type":"message_stop"}\n\n'
|
||||
)
|
||||
|
||||
# Everything up to — but not including — the terminal event.
|
||||
TRUNCATED_SSE = COMPLETE_SSE.split("event: message_delta")[0]
|
||||
|
||||
ERROR_SSE = (
|
||||
COMPLETE_SSE.split("event: message_delta")[0] + "event: error\n"
|
||||
'data: {"type":"error","error":{"type":"overloaded_error",'
|
||||
'"message":"Overloaded"}}\n\n'
|
||||
)
|
||||
|
||||
JSON_REPLY = {
|
||||
"id": "msg_1",
|
||||
"type": "message",
|
||||
"role": "assistant",
|
||||
"model": "claude-sonnet-4-6",
|
||||
"content": [{"type": "text", "text": "hello"}],
|
||||
"stop_reason": "end_turn",
|
||||
"usage": {
|
||||
"input_tokens": 10,
|
||||
"output_tokens": 5,
|
||||
"cache_read_input_tokens": 0,
|
||||
"cache_creation_input_tokens": 0,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def _config() -> ProxyConfig:
|
||||
return ProxyConfig(
|
||||
optimize=False,
|
||||
cache_enabled=False,
|
||||
rate_limit_enabled=False,
|
||||
memory_enabled=False,
|
||||
ccr_inject_tool=False,
|
||||
ccr_handle_responses=False,
|
||||
ccr_context_tracking=False,
|
||||
image_optimize=False,
|
||||
)
|
||||
|
||||
|
||||
def _drive(
|
||||
*,
|
||||
upstream: httpx.Response,
|
||||
accept: str | None = "text/event-stream",
|
||||
stream: bool = False,
|
||||
) -> tuple[httpx.Response, dict[str, object]]:
|
||||
"""Run one turn against a canned upstream reply.
|
||||
|
||||
Returns the client-facing response and what went upstream.
|
||||
"""
|
||||
seen: dict[str, object] = {}
|
||||
app = create_app(_config())
|
||||
with TestClient(app) as client:
|
||||
proxy = client.app.state.proxy
|
||||
|
||||
async def _fake_retry(method, url, headers, body, stream=False, **kwargs): # noqa: ANN001
|
||||
sent = json.loads(body) if isinstance(body, (str, bytes)) else body
|
||||
seen["stream"] = sent.get("stream")
|
||||
seen["headers"] = dict(headers or {})
|
||||
return upstream
|
||||
|
||||
proxy._retry_request = _fake_retry # type: ignore[assignment]
|
||||
|
||||
headers = {"x-api-key": "test-key", "anthropic-version": "2023-06-01"}
|
||||
if accept is not None:
|
||||
headers["accept"] = accept
|
||||
resp = client.post(
|
||||
"/v1/messages",
|
||||
json={
|
||||
"model": "claude-sonnet-4-6",
|
||||
"max_tokens": 64,
|
||||
"stream": stream,
|
||||
"messages": [{"role": "user", "content": "go"}],
|
||||
},
|
||||
headers=headers,
|
||||
)
|
||||
return resp, seen
|
||||
|
||||
|
||||
def _sse_response(body: str, **extra_headers: str) -> httpx.Response:
|
||||
headers = {"content-type": "text/event-stream", **extra_headers}
|
||||
return httpx.Response(200, content=body.encode(), headers=headers)
|
||||
|
||||
|
||||
def _accepts(headers: dict) -> list[str]:
|
||||
return [v for k, v in headers.items() if k.lower() == "accept"]
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Request side: a stream:false body must not carry an SSE-only Accept
|
||||
# --------------------------------------------------------------------------- #
|
||||
def test_non_stream_turn_asks_upstream_for_json() -> None:
|
||||
_, seen = _drive(upstream=httpx.Response(200, json=JSON_REPLY))
|
||||
|
||||
assert seen["stream"] is False
|
||||
assert _accepts(seen["headers"]) == ["application/json"] # type: ignore[arg-type]
|
||||
|
||||
|
||||
def test_non_stream_turn_replaces_rather_than_appends_accept() -> None:
|
||||
_, seen = _drive(upstream=httpx.Response(200, json=JSON_REPLY), accept="TEXT/EVENT-STREAM")
|
||||
|
||||
values = _accepts(seen["headers"]) # type: ignore[arg-type]
|
||||
assert values == ["application/json"]
|
||||
assert not any("event-stream" in v.lower() for v in values)
|
||||
|
||||
|
||||
def test_non_stream_turn_without_client_accept_still_asks_for_json() -> None:
|
||||
_, seen = _drive(upstream=httpx.Response(200, json=JSON_REPLY), accept=None)
|
||||
|
||||
assert _accepts(seen["headers"]) == ["application/json"] # type: ignore[arg-type]
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Response side: SSE at 200 for a JSON request is adapted, not relayed
|
||||
# --------------------------------------------------------------------------- #
|
||||
def test_event_stream_answer_is_adapted_to_json() -> None:
|
||||
resp, _ = _drive(upstream=_sse_response(COMPLETE_SSE))
|
||||
|
||||
assert resp.status_code == 200
|
||||
assert "application/json" in resp.headers["content-type"]
|
||||
body = resp.json()
|
||||
assert body["type"] == "message"
|
||||
assert body["content"] == [{"type": "text", "text": "hello"}]
|
||||
assert body["stop_reason"] == "end_turn"
|
||||
|
||||
|
||||
def test_adapted_reply_preserves_usage_for_accounting() -> None:
|
||||
resp, _ = _drive(upstream=_sse_response(COMPLETE_SSE))
|
||||
|
||||
usage = resp.json()["usage"]
|
||||
assert usage["input_tokens"] == 10
|
||||
assert usage["output_tokens"] == 5
|
||||
assert usage["cache_read_input_tokens"] == 2
|
||||
assert usage["cache_creation_input_tokens"] == 3
|
||||
|
||||
|
||||
def test_adapted_reply_carries_no_streaming_only_index() -> None:
|
||||
"""``index`` is a response-delta field; Anthropic rejects it on replay."""
|
||||
resp, _ = _drive(upstream=_sse_response(COMPLETE_SSE))
|
||||
|
||||
assert all("index" not in block for block in resp.json()["content"])
|
||||
|
||||
|
||||
def test_adapted_reply_drops_cdn_and_framing_headers() -> None:
|
||||
resp, _ = _drive(
|
||||
upstream=_sse_response(
|
||||
COMPLETE_SSE,
|
||||
**{
|
||||
"server": "cloudflare",
|
||||
"cf-ray": "abc123",
|
||||
"cf-cache-status": "DYNAMIC",
|
||||
"request-id": "req_011CeC1JTMS8egPL3FBteQay",
|
||||
"anthropic-ratelimit-requests-remaining": "42",
|
||||
},
|
||||
)
|
||||
)
|
||||
|
||||
lowered = {k.lower() for k in resp.headers}
|
||||
assert "server" not in lowered
|
||||
assert not any(k.startswith("cf-") for k in lowered)
|
||||
# Provenance the caller legitimately needs survives.
|
||||
assert resp.headers["request-id"] == "req_011CeC1JTMS8egPL3FBteQay"
|
||||
assert resp.headers["anthropic-ratelimit-requests-remaining"] == "42"
|
||||
|
||||
|
||||
def test_truncated_event_stream_fails_loudly() -> None:
|
||||
"""A partial stream is not a successful short answer."""
|
||||
resp, _ = _drive(upstream=_sse_response(TRUNCATED_SSE))
|
||||
|
||||
assert resp.status_code == 502
|
||||
assert "application/json" in resp.headers["content-type"]
|
||||
assert resp.json()["error"]["type"] == "upstream_protocol_error"
|
||||
|
||||
|
||||
def test_error_event_is_not_reported_as_success() -> None:
|
||||
resp, _ = _drive(upstream=_sse_response(ERROR_SSE))
|
||||
|
||||
assert resp.status_code == 502
|
||||
assert resp.json()["error"]["type"] == "upstream_protocol_error"
|
||||
|
||||
|
||||
def test_plain_json_reply_is_untouched() -> None:
|
||||
resp, _ = _drive(upstream=httpx.Response(200, json=JSON_REPLY))
|
||||
|
||||
assert resp.status_code == 200
|
||||
assert "application/json" in resp.headers["content-type"]
|
||||
assert resp.json()["content"] == [{"type": "text", "text": "hello"}]
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Strict reconstruction, exercised directly
|
||||
# --------------------------------------------------------------------------- #
|
||||
@pytest.fixture()
|
||||
def proxy():
|
||||
from headroom.proxy.server import HeadroomProxy
|
||||
|
||||
return HeadroomProxy(_config())
|
||||
|
||||
|
||||
def test_strict_mode_requires_a_terminal_event(proxy) -> None:
|
||||
assert proxy._parse_sse_to_response(TRUNCATED_SSE, "anthropic", require_complete=True) is None
|
||||
|
||||
|
||||
def test_strict_mode_rejects_an_error_event(proxy) -> None:
|
||||
assert proxy._parse_sse_to_response(ERROR_SSE, "anthropic", require_complete=True) is None
|
||||
|
||||
|
||||
def test_strict_mode_rejects_an_unclosed_block(proxy) -> None:
|
||||
unclosed = COMPLETE_SSE.replace(
|
||||
'event: content_block_stop\ndata: {"type":"content_block_stop","index":0}\n\n', ""
|
||||
)
|
||||
|
||||
assert proxy._parse_sse_to_response(unclosed, "anthropic", require_complete=True) is None
|
||||
|
||||
|
||||
def test_strict_mode_rejects_an_unknown_delta_type(proxy) -> None:
|
||||
"""A future delta Headroom cannot replay must not pass as complete."""
|
||||
unknown = COMPLETE_SSE.replace('"type":"text_delta","text":"hello"', '"type":"future_delta"')
|
||||
|
||||
assert proxy._parse_sse_to_response(unknown, "anthropic", require_complete=True) is None
|
||||
|
||||
|
||||
def test_strict_mode_reads_crlf_framed_events(proxy) -> None:
|
||||
parsed = proxy._parse_sse_to_response(
|
||||
COMPLETE_SSE.replace("\n", "\r\n"), "anthropic", require_complete=True
|
||||
)
|
||||
|
||||
assert parsed is not None
|
||||
assert parsed["content"] == [{"type": "text", "text": "hello"}]
|
||||
|
||||
|
||||
def test_strict_mode_keeps_stop_sequence_and_type(proxy) -> None:
|
||||
parsed = proxy._parse_sse_to_response(COMPLETE_SSE, "anthropic", require_complete=True)
|
||||
|
||||
assert parsed is not None
|
||||
assert parsed["type"] == "message"
|
||||
assert parsed["stop_sequence"] is None
|
||||
|
||||
|
||||
def test_permissive_mode_is_unchanged_for_existing_callers(proxy) -> None:
|
||||
"""Streaming callers keep the lenient reconstruction they rely on."""
|
||||
parsed = proxy._parse_sse_to_response(TRUNCATED_SSE, "anthropic")
|
||||
|
||||
assert parsed is not None
|
||||
assert parsed["content"][0]["text"] == "hello"
|
||||
|
||||
|
||||
def test_event_stream_under_a_vague_content_type_is_still_adapted() -> None:
|
||||
"""A gateway may relay the stream without declaring it (#3130)."""
|
||||
resp, _ = _drive(
|
||||
upstream=httpx.Response(
|
||||
200,
|
||||
content=COMPLETE_SSE.encode(),
|
||||
headers={"content-type": "application/octet-stream"},
|
||||
)
|
||||
)
|
||||
|
||||
assert resp.status_code == 200
|
||||
assert "application/json" in resp.headers["content-type"]
|
||||
assert resp.json()["content"] == [{"type": "text", "text": "hello"}]
|
||||
@@ -0,0 +1,263 @@
|
||||
"""Regression tests: a non-streaming caller must never receive an SSE body.
|
||||
|
||||
The buffered Anthropic path copies the upstream response headers wholesale,
|
||||
``content-type`` included. When the upstream answers a ``stream``-less request
|
||||
with ``text/event-stream``, that body reached the caller as a ``200`` it could
|
||||
not parse — the reply was complete, just in the wrong wire format, and the turn
|
||||
was lost.
|
||||
|
||||
The buffered-stream (CCR) path already refused this shape (#2952). These tests
|
||||
pin the same protection on the plain non-streaming path, plus the recovery that
|
||||
turns a lost turn into a normal reply.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from headroom.proxy.nonstream_sse_policy import (
|
||||
is_event_stream,
|
||||
media_type,
|
||||
should_recover_sse_reply,
|
||||
)
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Fixtures
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
_SSE_REPLY = (
|
||||
"event: message_start\n"
|
||||
'data: {"type":"message_start","message":{"id":"msg_sse_recovered",'
|
||||
'"type":"message","role":"assistant","model":"claude-sonnet-4-6",'
|
||||
'"content":[],"usage":{"input_tokens":11,"output_tokens":0}}}\n'
|
||||
"\n"
|
||||
"event: content_block_start\n"
|
||||
'data: {"type":"content_block_start","index":0,'
|
||||
'"content_block":{"type":"text","text":""}}\n'
|
||||
"\n"
|
||||
"event: content_block_delta\n"
|
||||
'data: {"type":"content_block_delta","index":0,'
|
||||
'"delta":{"type":"text_delta","text":"recovered body"}}\n'
|
||||
"\n"
|
||||
"event: content_block_stop\n"
|
||||
'data: {"type":"content_block_stop","index":0}\n'
|
||||
"\n"
|
||||
"event: message_delta\n"
|
||||
'data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},'
|
||||
'"usage":{"output_tokens":4}}\n'
|
||||
"\n"
|
||||
"event: message_stop\n"
|
||||
'data: {"type":"message_stop"}\n'
|
||||
"\n"
|
||||
)
|
||||
|
||||
# Upstream headers as they actually arrive through Anthropic's edge — the
|
||||
# correlation headers here are what a client uses to report and dedup a turn,
|
||||
# so the fix must not drop them while correcting the content-type.
|
||||
_UPSTREAM_SSE_HEADERS = {
|
||||
"content-type": "text/event-stream; charset=utf-8",
|
||||
"request-id": "req_011CeC1JTMS8egPL3FBteQay",
|
||||
"anthropic-ratelimit-requests-remaining": "49",
|
||||
"cf-ray": "9a1b2c3d4e5f6789-GRU",
|
||||
"server": "cloudflare",
|
||||
}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Pure policy
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestMediaTypeParsing:
|
||||
@pytest.mark.parametrize(
|
||||
("header", "expected"),
|
||||
[
|
||||
("text/event-stream", "text/event-stream"),
|
||||
("text/event-stream; charset=utf-8", "text/event-stream"),
|
||||
("Text/Event-Stream", "text/event-stream"),
|
||||
(" text/event-stream ", "text/event-stream"),
|
||||
("application/json", "application/json"),
|
||||
(None, ""),
|
||||
("", ""),
|
||||
],
|
||||
)
|
||||
def test_parameters_and_case_are_normalized(self, header, expected) -> None:
|
||||
assert media_type(header) == expected
|
||||
|
||||
def test_is_event_stream_only_matches_sse(self) -> None:
|
||||
assert is_event_stream("text/event-stream; charset=utf-8") is True
|
||||
assert is_event_stream("application/json") is False
|
||||
assert is_event_stream(None) is False
|
||||
|
||||
|
||||
class TestShouldRecoverSseReply:
|
||||
"""The gate has three deliberate negative arms; each is a separate risk."""
|
||||
|
||||
def test_recovers_sse_200_for_a_non_streaming_caller(self) -> None:
|
||||
assert (
|
||||
should_recover_sse_reply(
|
||||
client_requested_stream=False,
|
||||
status_code=200,
|
||||
content_type="text/event-stream",
|
||||
)
|
||||
is True
|
||||
)
|
||||
|
||||
def test_streaming_caller_is_untouched(self) -> None:
|
||||
"""A streaming caller asked for SSE — rewriting it would break the turn."""
|
||||
assert (
|
||||
should_recover_sse_reply(
|
||||
client_requested_stream=True,
|
||||
status_code=200,
|
||||
content_type="text/event-stream",
|
||||
)
|
||||
is False
|
||||
)
|
||||
|
||||
def test_json_reply_is_untouched(self) -> None:
|
||||
assert (
|
||||
should_recover_sse_reply(
|
||||
client_requested_stream=False,
|
||||
status_code=200,
|
||||
content_type="application/json",
|
||||
)
|
||||
is False
|
||||
)
|
||||
|
||||
@pytest.mark.parametrize("status", [429, 500, 529])
|
||||
def test_error_status_is_passed_through(self, status) -> None:
|
||||
"""A non-200 carries an upstream error payload the client should see."""
|
||||
assert (
|
||||
should_recover_sse_reply(
|
||||
client_requested_stream=False,
|
||||
status_code=status,
|
||||
content_type="text/event-stream",
|
||||
)
|
||||
is False
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Handler end-to-end — the wiring is where the bug lived
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
pytest.importorskip("fastapi")
|
||||
|
||||
from fastapi.testclient import TestClient # noqa: E402
|
||||
|
||||
from headroom.proxy.server import ProxyConfig, create_app # noqa: E402
|
||||
|
||||
|
||||
def _make_proxy_client() -> TestClient:
|
||||
config = ProxyConfig(
|
||||
optimize=True,
|
||||
mode="token",
|
||||
cache_enabled=False,
|
||||
rate_limit_enabled=False,
|
||||
cost_tracking_enabled=False,
|
||||
log_requests=False,
|
||||
ccr_inject_tool=False,
|
||||
ccr_handle_responses=False,
|
||||
ccr_context_tracking=False,
|
||||
image_optimize=False,
|
||||
)
|
||||
return TestClient(create_app(config))
|
||||
|
||||
|
||||
def _post_non_streaming(client: TestClient):
|
||||
return client.post(
|
||||
"/v1/messages",
|
||||
headers={"x-api-key": "test-key", "anthropic-version": "2023-06-01"},
|
||||
json={
|
||||
"model": "claude-sonnet-4-6",
|
||||
"max_tokens": 64,
|
||||
"messages": [{"role": "user", "content": "hello"}],
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def _stub_upstream(proxy, response: httpx.Response) -> None:
|
||||
async def _fake_retry(method, url, headers, body, stream=False, **kwargs): # noqa: ANN001
|
||||
return response
|
||||
|
||||
proxy._retry_request = _fake_retry
|
||||
|
||||
|
||||
class TestNonStreamingCallerNeverGetsAnEventStream:
|
||||
def test_sse_reply_is_recovered_as_json(self) -> None:
|
||||
"""Before the fix this returned text/event-stream and the SDK reported
|
||||
an empty or malformed response despite a complete reply."""
|
||||
with _make_proxy_client() as client:
|
||||
_stub_upstream(
|
||||
client.app.state.proxy,
|
||||
httpx.Response(
|
||||
200,
|
||||
headers=_UPSTREAM_SSE_HEADERS,
|
||||
content=_SSE_REPLY.encode(),
|
||||
),
|
||||
)
|
||||
response = _post_non_streaming(client)
|
||||
|
||||
assert response.status_code == 200
|
||||
assert "event-stream" not in response.headers["content-type"]
|
||||
assert response.headers["content-type"].startswith("application/json")
|
||||
|
||||
payload = response.json()
|
||||
assert payload["id"] == "msg_sse_recovered"
|
||||
assert payload["content"][0]["text"] == "recovered body"
|
||||
|
||||
def test_upstream_correlation_headers_survive_recovery(self) -> None:
|
||||
with _make_proxy_client() as client:
|
||||
_stub_upstream(
|
||||
client.app.state.proxy,
|
||||
httpx.Response(
|
||||
200,
|
||||
headers=_UPSTREAM_SSE_HEADERS,
|
||||
content=_SSE_REPLY.encode(),
|
||||
),
|
||||
)
|
||||
response = _post_non_streaming(client)
|
||||
|
||||
assert response.headers["request-id"] == "req_011CeC1JTMS8egPL3FBteQay"
|
||||
|
||||
def test_unrecoverable_event_stream_is_refused_not_forwarded(self) -> None:
|
||||
"""No message_start means no message. Refuse loudly rather than hand
|
||||
the caller a 200 it cannot parse."""
|
||||
with _make_proxy_client() as client:
|
||||
_stub_upstream(
|
||||
client.app.state.proxy,
|
||||
httpx.Response(
|
||||
200,
|
||||
headers=_UPSTREAM_SSE_HEADERS,
|
||||
content=b'event: ping\ndata: {"type":"ping"}\n\n',
|
||||
),
|
||||
)
|
||||
response = _post_non_streaming(client)
|
||||
|
||||
assert response.status_code == 502
|
||||
assert "event-stream" not in response.headers["content-type"]
|
||||
assert response.json()["error"]["type"] == "upstream_protocol_error"
|
||||
|
||||
def test_ordinary_json_reply_is_unaffected(self) -> None:
|
||||
"""Control: the fix must be inert on the overwhelmingly common path."""
|
||||
with _make_proxy_client() as client:
|
||||
_stub_upstream(
|
||||
client.app.state.proxy,
|
||||
httpx.Response(
|
||||
200,
|
||||
json={
|
||||
"id": "msg_plain",
|
||||
"type": "message",
|
||||
"role": "assistant",
|
||||
"content": [{"type": "text", "text": "ok"}],
|
||||
"usage": {"input_tokens": 10, "output_tokens": 3},
|
||||
},
|
||||
),
|
||||
)
|
||||
response = _post_non_streaming(client)
|
||||
|
||||
assert response.status_code == 200
|
||||
assert json.loads(response.content)["id"] == "msg_plain"
|
||||
Reference in New Issue
Block a user