fix(ccr): make StreamingCCRHandler work on OpenAI streams (#3069)
## Description `StreamingCCRHandler` (`headroom/ccr/response_handler.py`) was written against the Anthropic wire format. Constructed with `provider="openai"` it does not work: it silently drops the response, reports the wrong `finish_reason`, and emits a stream shape no OpenAI client can read. This PR fixes all three. **Reachability, stated up front:** `StreamingCCRHandler` is exported from `headroom/ccr/__init__.py` but no proxy handler instantiates it today. Every live CCR path (`handlers/openai.py:4276`, `handlers/openai.py:5936`, `handlers/anthropic.py`, `handlers/gemini.py`) calls `CCRResponseHandler.handle_response` on a non-streaming body instead. So these defects are not currently hit by proxy traffic. They bite anyone importing the public `headroom.ccr.StreamingCCRHandler` export, and they would bite the moment streaming CCR gets wired up. I would rather fix them while they are cheap than have them surface as a mysterious truncation bug later. **This PR does not fix #1026.** I found these while investigating that issue and they turned out to be unrelated to it. #1026 needs information from the reporter before anyone can say whether Headroom is even in the request path; I have asked for it there. ### The three defects **1. The whole OpenAI response was dropped.** `StreamingCCRBuffer.add_chunk` detected a tool call by scanning the accumulated bytes for the literal `"type":"tool_use"`. That is Anthropic-only. An OpenAI-compatible stream carries tool calls as a `tool_calls` array inside `choices[].delta` and never emits that marker, so `detected_ccr` could never become `True`. Independently, `process_stream` decided the stream had ended by scanning for `"stop_reason"`, another Anthropic-only field. An OpenAI stream has no such field; it terminates with the `[DONE]` sentinel. With neither marker ever matching, and nothing flushing the buffer once the source iterator ran out, the outcome was: - OpenAI stream under 10 000 bytes: **nothing at all was yielded**. The client got an empty response. - OpenAI stream over 10 000 bytes: chunks flushed in ~10 KB batches, and the final sub-threshold batch was never flushed. The response visibly stopped mid-sentence. **2. `finish_reason` was hardcoded.** `_reconstruct_openai_response` always returned `"finish_reason": "stop"`, even when it had just finished reconstructing a non-empty `tool_calls` array, where the OpenAI API requires `"tool_calls"`. A client that drives its agent loop off `finish_reason` reads `stop`, concludes the turn is over, and never executes the tool calls. The Anthropic sibling `_reconstruct_anthropic_response` does this correctly, carrying `stop_reason` through from `message_delta`. It also discarded `id`, `object`, `created`, `model`, and `usage`, returning a bare `choices` list that is not a valid `chat.completion`. **3. `_response_to_sse` emitted the wrong shape.** The OpenAI branch serialised the reconstructed **non-streaming** body into a single SSE frame. A streaming client parses `choices[].delta`; this frame has `choices[].message`. Both the text and the tool calls were invisible to it. ### Why CI did not catch it `tests/test_ccr_response_handler_extra.py` exercised `_reconstruct_openai_response` but never asserted `finish_reason`, and the one `process_stream` test that passed `provider="openai"` fed it Anthropic-shaped bytes (`"type":"tool_use"` plus `"stop_reason"`). No test had ever run a real OpenAI stream through this class. That test now uses the real OpenAI wire shape, so it actually covers the path it claims to. ## Type of Change - [x] Bug fix (non-breaking change which fixes an issue) - [ ] New feature (non-breaking change which adds functionality) - [ ] Breaking change (fix or feature that would cause existing functionality to not work as expected) - [ ] Documentation update - [ ] Refactor / internal change ## Changes Made All in `headroom/ccr/response_handler.py`: - `StreamingCCRBuffer` gained a `provider` field (defaults to `"anthropic"`, so existing construction is unchanged) and picks its tool-call marker from it: `"type":"tool_use"` for Anthropic, `"tool_calls"` for everything else. `StreamingCCRHandler.__init__` now passes its own provider down. - `process_stream` selects the end-of-stream marker by provider (`"stop_reason"` for Anthropic, `data: [DONE]` for OpenAI), and **always flushes whatever is still buffered once the source iterator is exhausted**. That second part is deliberately unconditional on the marker: upstream can truncate, a gateway can omit the sentinel, and a future stream shape may not be recognised. Buffered bytes at that point are real response data, so they get flushed rather than dropped. - Removed the dead re-iteration block that followed the detection loop. Its guard was `not detection_complete and not self.buffer.detected_ccr`, and the only `break` out of the loop above required `detected_ccr` to be `True`, so it could only ever be reached with an already-exhausted iterator. The new flush takes its place. - `_reconstruct_openai_response` derives `finish_reason`: `"tool_calls"` when the message carries tool calls, otherwise the last non-null upstream value (so a truncated turn stays reported as `"length"`), defaulting to `"stop"`. It carries `id` / `created` / `model` / `system_fingerprint` / `usage` through from the chunk envelope and stamps `"object": "chat.completion"`. It also tolerates `"delta": null` on a terminal chunk, which some OpenAI-compatible providers send instead of `{}`, in the same spirit as #2467. - New `_openai_response_to_chunks` splits a non-streaming `chat.completion` body into proper `chat.completion.chunk` frames (a role delta, a content delta, one delta per tool call, then a terminal frame carrying `finish_reason`). `_response_to_sse` uses it and then emits `[DONE]`. The Anthropic branch still delegates to `StreamingMixin._response_to_sse` and is untouched. Tests in `tests/test_ccr_response_handler_extra.py`: - Seven new tests: OpenAI CCR detection on a `tool_calls` delta (plus a non-CCR negative case), a short OpenAI stream passing through byte for byte, a stream past the 10 000-byte flush threshold keeping its tail, a stream with no `[DONE]` sentinel still flushing, `finish_reason` becoming `"tool_calls"` with the envelope preserved, the upstream `finish_reason` being kept when there are no tool calls, and `_response_to_sse` emitting parseable chunk frames. - `test_streaming_handler_falls_back_to_buffer_on_processing_error` now feeds genuine OpenAI SSE bytes instead of Anthropic ones, so it exercises the OpenAI detection path it was always meant to. - `test_response_to_sse_formats` asserts the new chunk-frame shape for OpenAI. The Anthropic half is unchanged. No behaviour change for `provider="anthropic"` beyond the end-of-iterator flush, which can only add data that was previously discarded. ## Testing - [x] Unit tests added/updated - [x] Existing tests pass - [ ] Manual testing performed - [ ] Integration tests added Each of the seven new tests was confirmed to fail against the unmodified source (`git stash` on `response_handler.py` alone, tests untouched), so they are genuine regression tests rather than assertions written to match current behaviour: ``` $ git stash push -- headroom/ccr/response_handler.py $ python -m pytest tests/test_ccr_response_handler_extra.py -q -k openai FAILED tests/test_ccr_response_handler_extra.py::test_streaming_buffer_detects_ccr_in_openai_tool_calls_delta FAILED tests/test_ccr_response_handler_extra.py::test_openai_stream_without_ccr_yields_every_chunk FAILED tests/test_ccr_response_handler_extra.py::test_openai_stream_past_flush_threshold_keeps_the_tail FAILED tests/test_ccr_response_handler_extra.py::test_openai_stream_without_done_sentinel_still_flushes FAILED tests/test_ccr_response_handler_extra.py::test_reconstruct_openai_response_marks_tool_calls_finish_reason FAILED tests/test_ccr_response_handler_extra.py::test_reconstruct_openai_response_keeps_upstream_finish_reason FAILED tests/test_ccr_response_handler_extra.py::test_response_to_sse_emits_openai_chunk_frames 7 failed, 2 passed, 13 deselected in 0.79s ``` With the fix applied, the full CCR response-handler suite passes: ``` $ python -m pytest tests/test_ccr_response_handler_extra.py tests/test_ccr_response_handler.py -q collected 57 items tests\test_ccr_response_handler_extra.py ...................... [ 38%] tests\test_ccr_response_handler.py ................................... [100%] ============================= 57 passed in 1.74s ============================== ``` Wider CCR and streaming surface: ``` $ python -m pytest tests/ -k "ccr or streaming" -q 4 failed, 696 passed, 73 skipped, 10949 deselected, 2 warnings in 175.80s (0:02:55) ``` The 4 failures are pre-existing on a clean `upstream/main` and unrelated to this change (verified by stashing both changed files and re-running exactly those four): `test_ccr_mcp_http.py::test_streamable_http_initialize_and_list_tools`, `test_cli_proxy_env.py::TestCLICompressionOnlyFlags::test_ccr_defaults_on`, and two in `test_transforms/test_smart_crusher_ccr_roundtrip.py`. Lint and types: ``` $ python -m ruff check . All checks passed! $ python -m ruff format --check . 1505 files already formatted $ python -m mypy headroom --ignore-missing-imports Found 12 errors in 3 files (checked 521 source files) ``` Zero mypy errors in `headroom/ccr/response_handler.py`. The 12 are pre-existing, in `ccr/mcp_server.py`, `memory/mcp_server.py`, and `release_version.py`, none of which this PR touches (they come from a locally installed `mcp` whose stubs differ from CI's). ## Real Behavior Proof - Environment: Windows 11, Python 3.13.11, pytest 9.1.1, ruff and mypy from the repo's pinned config, branch `fix/ccr-streaming-openai-path` off `upstream/main` at `cbb950a4`. - Exact command / steps: `python -m pytest tests/test_ccr_response_handler_extra.py tests/test_ccr_response_handler.py -q`; then `git stash push -- headroom/ccr/response_handler.py` and `python -m pytest tests/test_ccr_response_handler_extra.py -q -k openai` to confirm the new tests fail without the source fix; then `python -m pytest tests/ -k "ccr or streaming" -q`; then `python -m ruff check .`, `python -m ruff format --check .`, `python -m mypy headroom --ignore-missing-imports`. - Observed result: 57/57 pass in the CCR response-handler suites with the fix; all 7 new tests fail without it. The wider run is 696 passed with 4 failures that reproduce identically on an unmodified tree. Ruff clean, mypy clean on the changed file. In `test_openai_stream_without_ccr_yields_every_chunk` the handler now returns every input chunk byte for byte, where before it returned an empty list. - Not tested: no end-to-end run against a live OpenAI-compatible backend, because no proxy handler instantiates `StreamingCCRHandler` today, so there is no wired path to drive. Coverage is at the class level using recorded-shape SSE frames. The Anthropic path is covered only by the existing tests, which still pass unchanged. ## Runtime Rollout Safety - Rollout-managed feature(s): none. `StreamingCCRHandler` is not gated by a rollout feature and is not reachable from any proxy handler. - Minimum rollout channel: not applicable; no rollout gate is involved. - Stable/default behavior changed: no. For `provider="anthropic"` the only behavioural difference is that bytes left buffered when the source iterator ends are now flushed instead of discarded, which can only add data the client previously lost. For `provider="openai"` the class was non-functional, so there is no prior behaviour to preserve. - Kill switch / disable path: not applicable; no new configuration, env var, or feature flag is introduced. - Unsafe override required: no. - Qualification impact: none. No qualification-gated surface is touched. - Rollback path: revert this commit. It is self-contained in `headroom/ccr/response_handler.py` and `tests/test_ccr_response_handler_extra.py`, with no schema, config, or persisted-state changes. ## Review Readiness - [x] I have performed a self-review - [x] This PR is ready for human review Two judgement calls worth a reviewer's attention: 1. **Removing the dead re-iteration block** in `process_stream`. I am confident it was unreachable (the only `break` above it requires `detected_ccr`, which its own guard excludes), but it is the one deletion in this diff rather than an addition, so it is worth a second pair of eyes. 2. **The unconditional end-of-iterator flush.** I chose to flush regardless of whether an end marker matched, rather than only fixing the OpenAI marker. That makes the truncation bug unreachable even if a future provider uses a shape neither marker recognises. The cost is that a stream whose trailing bytes are genuinely not meant for the client would now be forwarded. Given the buffer only ever holds upstream response bytes, forwarding is the safer default, but flag it if you disagree. Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -571,11 +571,23 @@ class StreamingCCRBuffer:
|
||||
chunks: list[bytes] = field(default_factory=list)
|
||||
detected_ccr: bool = False
|
||||
complete_response: dict[str, Any] | None = None
|
||||
provider: str = "anthropic"
|
||||
|
||||
# Patterns to detect tool_use in stream
|
||||
# Wire markers for the start of a tool call. Anthropic streams
|
||||
# `"type":"tool_use"` content blocks; OpenAI-compatible streams carry a
|
||||
# `"tool_calls"` array inside `choices[].delta` and never emit the
|
||||
# Anthropic marker, so scanning only for the latter meant CCR was never
|
||||
# detected on an OpenAI stream.
|
||||
_tool_use_start: bytes = b'"type":"tool_use"'
|
||||
_openai_tool_use_start: bytes = b'"tool_calls"'
|
||||
_ccr_tool_pattern: bytes = f'"{CCR_TOOL_NAME}"'.encode()
|
||||
|
||||
def _tool_call_marker(self) -> bytes:
|
||||
"""The provider's on-the-wire marker for the start of a tool call."""
|
||||
if self.provider == "anthropic":
|
||||
return self._tool_use_start
|
||||
return self._openai_tool_use_start
|
||||
|
||||
def add_chunk(self, chunk: bytes) -> bool:
|
||||
"""Add a chunk and check for CCR tool calls.
|
||||
|
||||
@@ -587,7 +599,7 @@ class StreamingCCRBuffer:
|
||||
# Quick check: does accumulated content contain CCR tool?
|
||||
accumulated = b"".join(self.chunks)
|
||||
|
||||
if self._tool_use_start in accumulated and self._ccr_tool_pattern in accumulated:
|
||||
if self._tool_call_marker() in accumulated and self._ccr_tool_pattern in accumulated:
|
||||
self.detected_ccr = True
|
||||
return True
|
||||
|
||||
@@ -622,7 +634,7 @@ class StreamingCCRHandler:
|
||||
) -> None:
|
||||
self.response_handler = response_handler
|
||||
self.provider = provider
|
||||
self.buffer = StreamingCCRBuffer()
|
||||
self.buffer = StreamingCCRBuffer(provider=provider)
|
||||
|
||||
async def process_stream(
|
||||
self,
|
||||
@@ -648,8 +660,12 @@ class StreamingCCRHandler:
|
||||
Response chunks (possibly from continuation response).
|
||||
"""
|
||||
# Phase 1: Initial detection
|
||||
# Buffer chunks until we can determine if there's a CCR call
|
||||
detection_complete = False
|
||||
# Buffer chunks until we can determine if there's a CCR call.
|
||||
#
|
||||
# The end-of-stream marker is provider-specific. Anthropic signals the
|
||||
# terminal state with `stop_reason` in `message_delta`; OpenAI-compatible
|
||||
# streams have no such field and terminate with the `[DONE]` sentinel.
|
||||
end_marker = b'"stop_reason"' if self.provider == "anthropic" else b"data: [DONE]"
|
||||
|
||||
async for chunk in stream_iterator:
|
||||
self.buffer.add_chunk(chunk)
|
||||
@@ -660,9 +676,7 @@ class StreamingCCRHandler:
|
||||
accumulated = self.buffer.get_accumulated()
|
||||
|
||||
# Look for stream end markers
|
||||
if b'"stop_reason"' in accumulated:
|
||||
detection_complete = True
|
||||
|
||||
if end_marker in accumulated:
|
||||
if self.buffer.detected_ccr:
|
||||
# CCR detected - need to handle
|
||||
break
|
||||
@@ -679,13 +693,15 @@ class StreamingCCRHandler:
|
||||
yield buffered_chunk
|
||||
self.buffer.clear()
|
||||
|
||||
# Continue streaming rest of response
|
||||
if not detection_complete and not self.buffer.detected_ccr:
|
||||
async for chunk in stream_iterator:
|
||||
if self.buffer.detected_ccr:
|
||||
self.buffer.add_chunk(chunk)
|
||||
else:
|
||||
yield chunk
|
||||
# The end marker is not guaranteed to arrive: upstream can truncate, a
|
||||
# provider can omit the sentinel, or the stream can be a shape this
|
||||
# detector does not recognise. Anything still buffered once the source
|
||||
# iterator is exhausted is real response data the client has never
|
||||
# seen, so flush it instead of dropping it.
|
||||
if not self.buffer.detected_ccr and self.buffer.chunks:
|
||||
for buffered_chunk in self.buffer.chunks:
|
||||
yield buffered_chunk
|
||||
self.buffer.clear()
|
||||
|
||||
# Phase 2: Handle CCR if detected
|
||||
if self.buffer.detected_ccr:
|
||||
@@ -903,13 +919,38 @@ class StreamingCCRHandler:
|
||||
}
|
||||
|
||||
tool_calls_map: dict[int, dict[str, Any]] = {}
|
||||
finish_reason: str | None = None
|
||||
envelope: dict[str, Any] = {}
|
||||
usage: Any = None
|
||||
|
||||
for event in events:
|
||||
choices = event.get("choices", [])
|
||||
if not choices:
|
||||
# Carry the chunk envelope through. Dropping it left the
|
||||
# reconstructed body without `id`, `model`, `created` or `usage`,
|
||||
# which downstream middleware reads for routing and metering.
|
||||
for key in ("id", "created", "model", "system_fingerprint"):
|
||||
value = event.get(key)
|
||||
if value is not None:
|
||||
envelope[key] = value
|
||||
if event.get("usage") is not None:
|
||||
usage = event["usage"]
|
||||
|
||||
choices = event.get("choices")
|
||||
if not isinstance(choices, list) or not choices:
|
||||
continue
|
||||
choice = choices[0]
|
||||
if not isinstance(choice, dict):
|
||||
continue
|
||||
|
||||
delta = choices[0].get("delta", {})
|
||||
# `finish_reason` is null on every chunk but the last, so keep the
|
||||
# most recent non-null value rather than the first one seen.
|
||||
if choice.get("finish_reason") is not None:
|
||||
finish_reason = choice["finish_reason"]
|
||||
|
||||
# Some OpenAI-compatible providers send `"delta": null` on the
|
||||
# terminal chunk instead of an empty object.
|
||||
delta = choice.get("delta")
|
||||
if not isinstance(delta, dict):
|
||||
delta = {}
|
||||
|
||||
if "content" in delta and delta["content"]:
|
||||
message["content"] = (message.get("content") or "") + delta["content"]
|
||||
@@ -944,14 +985,94 @@ class StreamingCCRHandler:
|
||||
tc["function"]["arguments"] += fn["arguments"]
|
||||
|
||||
message["tool_calls"] = [tool_calls_map[i] for i in sorted(tool_calls_map.keys())]
|
||||
if not message["tool_calls"]:
|
||||
has_tool_calls = bool(message["tool_calls"])
|
||||
if not has_tool_calls:
|
||||
del message["tool_calls"]
|
||||
if not message["content"]:
|
||||
message["content"] = None
|
||||
|
||||
return {
|
||||
"choices": [{"message": message, "finish_reason": "stop"}],
|
||||
# OpenAI requires `finish_reason: "tool_calls"` whenever the message
|
||||
# carries tool calls. This was hardcoded to "stop", which tells any
|
||||
# client that drives its agent loop off `finish_reason` that the turn
|
||||
# is over, so the reconstructed tool calls were never executed.
|
||||
if has_tool_calls:
|
||||
finish_reason = "tool_calls"
|
||||
elif finish_reason is None:
|
||||
finish_reason = "stop"
|
||||
|
||||
response: dict[str, Any] = {
|
||||
"object": "chat.completion",
|
||||
**envelope,
|
||||
"choices": [{"index": 0, "message": message, "finish_reason": finish_reason}],
|
||||
}
|
||||
if usage is not None:
|
||||
response["usage"] = usage
|
||||
return response
|
||||
|
||||
def _openai_response_to_chunks(self, response: dict[str, Any]) -> list[bytes]:
|
||||
"""Split a non-streaming ``chat.completion`` body into SSE chunk frames.
|
||||
|
||||
A streaming client reads ``choices[].delta``, not ``choices[].message``.
|
||||
Serialising the reconstructed non-streaming body into a single SSE frame
|
||||
produced a stream in which both the text and the tool calls were
|
||||
invisible to the client.
|
||||
"""
|
||||
choices = response.get("choices")
|
||||
choice = choices[0] if isinstance(choices, list) and choices else {}
|
||||
if not isinstance(choice, dict):
|
||||
choice = {}
|
||||
message = choice.get("message")
|
||||
if not isinstance(message, dict):
|
||||
message = {}
|
||||
finish_reason = choice.get("finish_reason") or "stop"
|
||||
|
||||
base: dict[str, Any] = {"object": "chat.completion.chunk"}
|
||||
for key in ("id", "created", "model", "system_fingerprint"):
|
||||
if response.get(key) is not None:
|
||||
base[key] = response[key]
|
||||
|
||||
def frame(delta: dict[str, Any], reason: str | None) -> bytes:
|
||||
payload = {
|
||||
**base,
|
||||
"choices": [{"index": 0, "delta": delta, "finish_reason": reason}],
|
||||
}
|
||||
return f"data: {json.dumps(payload)}\n\n".encode()
|
||||
|
||||
frames = [frame({"role": message.get("role") or "assistant"}, None)]
|
||||
|
||||
content = message.get("content")
|
||||
if content:
|
||||
frames.append(frame({"content": content}, None))
|
||||
|
||||
tool_calls = message.get("tool_calls")
|
||||
if isinstance(tool_calls, list):
|
||||
for index, tool_call in enumerate(tool_calls):
|
||||
if not isinstance(tool_call, dict):
|
||||
continue
|
||||
function = tool_call.get("function")
|
||||
if not isinstance(function, dict):
|
||||
function = {}
|
||||
frames.append(
|
||||
frame(
|
||||
{
|
||||
"tool_calls": [
|
||||
{
|
||||
"index": index,
|
||||
"id": tool_call.get("id", ""),
|
||||
"type": tool_call.get("type", "function"),
|
||||
"function": {
|
||||
"name": function.get("name", ""),
|
||||
"arguments": function.get("arguments", ""),
|
||||
},
|
||||
}
|
||||
]
|
||||
},
|
||||
None,
|
||||
)
|
||||
)
|
||||
|
||||
frames.append(frame({}, finish_reason))
|
||||
return frames
|
||||
|
||||
async def _response_to_sse(
|
||||
self,
|
||||
@@ -968,6 +1089,7 @@ class StreamingCCRHandler:
|
||||
for chunk in StreamingMixin()._response_to_sse(response, "anthropic"):
|
||||
yield chunk
|
||||
else:
|
||||
# OpenAI SSE format
|
||||
yield f"data: {json.dumps(response)}\n\n".encode()
|
||||
# OpenAI SSE format: `chat.completion.chunk` frames, then [DONE].
|
||||
for chunk in self._openai_response_to_chunks(response):
|
||||
yield chunk
|
||||
yield b"data: [DONE]\n\n"
|
||||
|
||||
@@ -445,7 +445,15 @@ async def test_streaming_handler_falls_back_to_buffer_on_processing_error(
|
||||
lambda data: (_ for _ in ()).throw(RuntimeError("parse failed")),
|
||||
)
|
||||
|
||||
chunks = [b'{"type":"tool_use","name":"headroom_retrieve"', b',"stop_reason":"tool_use"}']
|
||||
# Real OpenAI wire shape: a `tool_calls` delta naming the CCR tool, then
|
||||
# the `[DONE]` sentinel. This test previously fed Anthropic-shaped bytes to
|
||||
# an ``openai`` handler, so it never reached the OpenAI detection path.
|
||||
chunks = [
|
||||
b'data: {"choices":[{"index":0,"delta":{"tool_calls":[{"index":0,'
|
||||
b'"id":"call_1","function":{"name":"headroom_retrieve",'
|
||||
b'"arguments":"{}"}}]}}]}\n\n',
|
||||
b"data: [DONE]\n\n",
|
||||
]
|
||||
streamed = [
|
||||
chunk
|
||||
async for chunk in handler.process_stream(_async_iter(chunks), [], None, lambda m, t: None)
|
||||
@@ -462,7 +470,14 @@ async def test_response_to_sse_formats() -> None:
|
||||
|
||||
openai = StreamingCCRHandler(CCRResponseHandler(), provider="openai")
|
||||
openai_chunks = [chunk async for chunk in openai._response_to_sse({"choices": []})]
|
||||
assert openai_chunks == [b'data: {"choices": []}\n\n', b"data: [DONE]\n\n"]
|
||||
# An empty body still produces well-formed chunk frames (role, then a
|
||||
# terminal frame carrying finish_reason) rather than a single non-streaming
|
||||
# body a streaming client cannot read.
|
||||
assert openai_chunks[-1] == b"data: [DONE]\n\n"
|
||||
frames = [json.loads(chunk.decode()[len("data: ") :]) for chunk in openai_chunks[:-1]]
|
||||
assert [frame["object"] for frame in frames] == ["chat.completion.chunk"] * 2
|
||||
assert frames[0]["choices"][0]["delta"] == {"role": "assistant"}
|
||||
assert frames[-1]["choices"][0]["finish_reason"] == "stop"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@@ -533,3 +548,229 @@ def test_reconstruct_server_tool_use_input_from_partial_json() -> None:
|
||||
assert block["type"] == "server_tool_use"
|
||||
assert block["input"] == {"query": "x"}
|
||||
assert "_partial_json" not in block
|
||||
|
||||
|
||||
def _openai_chunk(delta: dict[str, Any], finish_reason: str | None = None) -> bytes:
|
||||
"""One `chat.completion.chunk` SSE frame in the shape a real backend sends."""
|
||||
payload = {
|
||||
"id": "chatcmpl_1",
|
||||
"object": "chat.completion.chunk",
|
||||
"created": 1700000000,
|
||||
"model": "gpt-4o-mini",
|
||||
"choices": [{"index": 0, "delta": delta, "finish_reason": finish_reason}],
|
||||
}
|
||||
return f"data: {json.dumps(payload)}\n\n".encode()
|
||||
|
||||
|
||||
def test_streaming_buffer_detects_ccr_in_openai_tool_calls_delta() -> None:
|
||||
# An OpenAI-compatible stream never emits Anthropic's `"type":"tool_use"`
|
||||
# marker; its tool calls arrive as a `tool_calls` array inside
|
||||
# `choices[].delta`. Scanning only for the Anthropic marker meant CCR was
|
||||
# never detected on this provider.
|
||||
buffer = StreamingCCRBuffer(provider="openai")
|
||||
|
||||
assert buffer.add_chunk(_openai_chunk({"role": "assistant"})) is False
|
||||
detected = buffer.add_chunk(
|
||||
_openai_chunk(
|
||||
{
|
||||
"tool_calls": [
|
||||
{
|
||||
"index": 0,
|
||||
"id": "call_1",
|
||||
"function": {"name": CCR_TOOL_NAME, "arguments": ""},
|
||||
}
|
||||
]
|
||||
}
|
||||
)
|
||||
)
|
||||
|
||||
assert detected is True
|
||||
assert buffer.detected_ccr is True
|
||||
|
||||
# A non-CCR tool call on the same provider must not trip detection.
|
||||
other = StreamingCCRBuffer(provider="openai")
|
||||
assert (
|
||||
other.add_chunk(
|
||||
_openai_chunk(
|
||||
{"tool_calls": [{"index": 0, "id": "c", "function": {"name": "other_tool"}}]}
|
||||
)
|
||||
)
|
||||
is False
|
||||
)
|
||||
assert other.detected_ccr is False
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_openai_stream_without_ccr_yields_every_chunk() -> None:
|
||||
# A short OpenAI stream with no CCR call must pass through byte for byte.
|
||||
# End-of-stream was detected by scanning for Anthropic's `stop_reason`,
|
||||
# which an OpenAI stream never contains, so nothing was ever flushed and
|
||||
# the client received an empty response.
|
||||
handler = StreamingCCRHandler(CCRResponseHandler(), provider="openai")
|
||||
chunks = [
|
||||
_openai_chunk({"role": "assistant"}),
|
||||
_openai_chunk({"content": "hello "}),
|
||||
_openai_chunk({"content": "world"}),
|
||||
_openai_chunk({}, finish_reason="stop"),
|
||||
b"data: [DONE]\n\n",
|
||||
]
|
||||
|
||||
streamed = [
|
||||
chunk
|
||||
async for chunk in handler.process_stream(_async_iter(chunks), [], None, lambda m, t: None)
|
||||
]
|
||||
|
||||
assert streamed == chunks
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_openai_stream_past_flush_threshold_keeps_the_tail() -> None:
|
||||
# Past 10 000 buffered bytes the handler flushes in batches. Without an
|
||||
# end-of-stream match the final sub-threshold batch was never flushed, so
|
||||
# a long response visibly stopped mid-sentence.
|
||||
handler = StreamingCCRHandler(CCRResponseHandler(), provider="openai")
|
||||
chunks = [
|
||||
_openai_chunk({"role": "assistant"}),
|
||||
_openai_chunk({"content": "x" * 11000}),
|
||||
_openai_chunk({"content": "the tail that used to be dropped"}),
|
||||
_openai_chunk({}, finish_reason="stop"),
|
||||
b"data: [DONE]\n\n",
|
||||
]
|
||||
|
||||
streamed = [
|
||||
chunk
|
||||
async for chunk in handler.process_stream(_async_iter(chunks), [], None, lambda m, t: None)
|
||||
]
|
||||
|
||||
assert streamed == chunks
|
||||
assert b"the tail that used to be dropped" in b"".join(streamed)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_openai_stream_without_done_sentinel_still_flushes() -> None:
|
||||
# Upstream can truncate before `[DONE]`, and some gateways omit it. Bytes
|
||||
# left in the buffer when the source iterator is exhausted are real
|
||||
# response data, so they are flushed rather than discarded.
|
||||
handler = StreamingCCRHandler(CCRResponseHandler(), provider="openai")
|
||||
chunks = [_openai_chunk({"role": "assistant"}), _openai_chunk({"content": "partial answer"})]
|
||||
|
||||
streamed = [
|
||||
chunk
|
||||
async for chunk in handler.process_stream(_async_iter(chunks), [], None, lambda m, t: None)
|
||||
]
|
||||
|
||||
assert streamed == chunks
|
||||
|
||||
|
||||
def test_reconstruct_openai_response_marks_tool_calls_finish_reason() -> None:
|
||||
# OpenAI requires `finish_reason: "tool_calls"` when the message carries
|
||||
# tool calls. It was hardcoded to "stop", so a client driving its agent
|
||||
# loop off `finish_reason` ended the turn instead of running the tools.
|
||||
handler = StreamingCCRHandler(CCRResponseHandler(), provider="openai")
|
||||
|
||||
parsed = handler._reconstruct_openai_response(
|
||||
[
|
||||
{"id": "chatcmpl_1", "model": "gpt-4o-mini", "created": 1700000000},
|
||||
{
|
||||
"choices": [
|
||||
{
|
||||
"delta": {
|
||||
"tool_calls": [
|
||||
{
|
||||
"index": 0,
|
||||
"id": "call_1",
|
||||
"function": {
|
||||
"name": CCR_TOOL_NAME,
|
||||
"arguments": '{"hash":"abc"}',
|
||||
},
|
||||
}
|
||||
]
|
||||
},
|
||||
"finish_reason": None,
|
||||
}
|
||||
]
|
||||
},
|
||||
{
|
||||
"choices": [{"delta": None, "finish_reason": "tool_calls"}],
|
||||
"usage": {"prompt_tokens": 12, "completion_tokens": 3},
|
||||
},
|
||||
]
|
||||
)
|
||||
|
||||
assert parsed["choices"][0]["finish_reason"] == "tool_calls"
|
||||
# The chunk envelope is carried through so the reconstructed body is a
|
||||
# valid `chat.completion` rather than a bare `choices` list.
|
||||
assert parsed["object"] == "chat.completion"
|
||||
assert parsed["id"] == "chatcmpl_1"
|
||||
assert parsed["model"] == "gpt-4o-mini"
|
||||
assert parsed["created"] == 1700000000
|
||||
assert parsed["usage"] == {"prompt_tokens": 12, "completion_tokens": 3}
|
||||
|
||||
|
||||
def test_reconstruct_openai_response_keeps_upstream_finish_reason() -> None:
|
||||
# With no tool calls, the upstream reason is preserved instead of being
|
||||
# rewritten to "stop": a truncated turn must stay reported as truncated.
|
||||
handler = StreamingCCRHandler(CCRResponseHandler(), provider="openai")
|
||||
|
||||
parsed = handler._reconstruct_openai_response(
|
||||
[
|
||||
{"choices": [{"delta": {"content": "half an ans"}, "finish_reason": None}]},
|
||||
{"choices": [{"delta": {}, "finish_reason": "length"}]},
|
||||
]
|
||||
)
|
||||
|
||||
assert parsed["choices"][0]["finish_reason"] == "length"
|
||||
assert parsed["choices"][0]["message"]["content"] == "half an ans"
|
||||
|
||||
# And an absent reason still defaults to "stop".
|
||||
defaulted = handler._reconstruct_openai_response(
|
||||
[{"choices": [{"delta": {"content": "hi"}}]}],
|
||||
)
|
||||
assert defaulted["choices"][0]["finish_reason"] == "stop"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_response_to_sse_emits_openai_chunk_frames() -> None:
|
||||
# A streaming client reads `choices[].delta`. Re-serialising the
|
||||
# reconstructed non-streaming body (`choices[].message`) into one SSE frame
|
||||
# made both the content and the tool calls invisible to it.
|
||||
handler = StreamingCCRHandler(CCRResponseHandler(), provider="openai")
|
||||
response = {
|
||||
"id": "chatcmpl_2",
|
||||
"object": "chat.completion",
|
||||
"created": 1700000001,
|
||||
"model": "gpt-4o-mini",
|
||||
"choices": [
|
||||
{
|
||||
"index": 0,
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": "done",
|
||||
"tool_calls": [
|
||||
{
|
||||
"id": "call_9",
|
||||
"type": "function",
|
||||
"function": {"name": "do_thing", "arguments": '{"a":1}'},
|
||||
}
|
||||
],
|
||||
},
|
||||
"finish_reason": "tool_calls",
|
||||
}
|
||||
],
|
||||
}
|
||||
|
||||
chunks = [chunk async for chunk in handler._response_to_sse(response)]
|
||||
|
||||
assert chunks[-1] == b"data: [DONE]\n\n"
|
||||
frames = [json.loads(chunk.decode()[len("data: ") :]) for chunk in chunks[:-1]]
|
||||
assert all(frame["object"] == "chat.completion.chunk" for frame in frames)
|
||||
assert all("delta" in frame["choices"][0] for frame in frames)
|
||||
assert all(frame["id"] == "chatcmpl_2" for frame in frames)
|
||||
|
||||
deltas = [frame["choices"][0]["delta"] for frame in frames]
|
||||
assert deltas[0] == {"role": "assistant"}
|
||||
assert deltas[1] == {"content": "done"}
|
||||
assert deltas[2]["tool_calls"][0]["id"] == "call_9"
|
||||
assert deltas[2]["tool_calls"][0]["index"] == 0
|
||||
assert deltas[2]["tool_calls"][0]["function"] == {"name": "do_thing", "arguments": '{"a":1}'}
|
||||
assert frames[-1]["choices"][0]["finish_reason"] == "tool_calls"
|
||||
|
||||
Reference in New Issue
Block a user