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:
Parideboy
2026-08-18 05:55:02 +02:00
committed by GitHub
parent eeb038bc0c
commit 7ef736fb1a
2 changed files with 388 additions and 25 deletions
+145 -23
View File
@@ -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"
+243 -2
View File
@@ -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"