Coalesce streamed reasoning into one client block - #54
Conversation
Reviewer's GuideThe PR introduces bounded buffering around realtime Chat SSE streams so consecutive reasoning fragments are emitted as one client-visible reasoning block before incremental text, refusal, or tool output, with corresponding converter behavior, documentation updates, and cross-protocol regression tests. Sequence diagram for coalescing streamed reasoningsequenceDiagram
participant Upstream as Upstream SSE
participant Coalescer as _coalesce_reasoning_sse
participant Client as Client
Upstream->>Coalescer: reasoning_content fragment
Upstream->>Coalescer: reasoning_content fragment
Coalescer->>Coalescer: buffer and charge_text
Upstream->>Coalescer: content, refusal, tool_calls, or finish_reason
Coalescer->>Client: one reasoning_content delta
Coalescer->>Client: visible output delta
Upstream->>Coalescer: data [DONE]
Coalescer->>Client: data [DONE]
Flow diagram for bounded reasoning bufferingflowchart LR
A[Chat SSE stream] --> B[_coalesce_reasoning_sse]
B --> C{Reasoning delta?}
C -->|yes| D[Append to pending buffer]
D --> E[StreamOutputBudget.charge_text]
E --> F{Visible output or finish?}
F -->|no| D
F -->|yes| G[Flush one reasoning delta]
C -->|no| G
G --> H[Preserve incremental content refusal or tool output]
H --> I[Client SSE stream]
E --> J[max_collect_bytes bound]
File-Level Changes
Tips and commandsInteracting with Sourcery
Customizing Your ExperienceAccess your dashboard to:
Getting Help
|
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
Hey - I've found 3 issues
Prompt for AI Agents
Please address the comments from this code review:
## Individual Comments
### Comment 1
<location path="converter.py" line_range="3198" />
<code_context>
+ if template is None:
+ template = event
+ pending.append(reasoning)
+ buffer_budget.charge_text(reasoning)
+ visible = any(delta.get(key) for key in ("content", "refusal", "tool_calls", "function_call"))
+ visible = visible or bool(choice.get("finish_reason"))
</code_context>
<issue_to_address>
**issue (broader_impact):** Reasoning bytes are charged once by the upstream `ChatSSEAccumulator` budget and charged again by `_coalesce_reasoning_sse` or the adapted converter, so a valid stream can exceed `max_collect_bytes` and fail with `response_too_large` even though its retained output is within the configured limit.
**Triggers:** When realtime streaming has a nonzero `max_collect_bytes` and the stream contains reasoning or other retained text.
**Suggested fix:** Share the existing request budget with `_coalesce_reasoning_sse`, or remove the duplicate charge from one of the buffering/conversion layers.
</issue_to_address>
### Comment 2
<location path="converter.py" line_range="3140-3146" />
<code_context>
+ source = template or {}
+ payload = dict(source)
+ choices = source.get("choices") or []
+ if choices:
+ choice = dict(choices[0])
+ delta = choice.get("delta") if isinstance(choice.get("delta"), dict) else {}
+ delta = {key: value for key, value in delta.items() if key == "role"}
+ delta["reasoning_content"] = "".join(pending)
+ choice["delta"] = delta
+ choice["finish_reason"] = None
+ payload["choices"] = [choice]
+ wire = "data: " + json.dumps(payload, ensure_ascii=False)
+ template = None
</code_context>
<issue_to_address>
**issue (bug_risk):** When an upstream SSE event contains multiple choices and the first choice has reasoning, the reconstructed payload replaces the original `choices` list with `[choice]`, permanently dropping every other choice from that event.
**Triggers:** When a request or upstream provider emits more than one choice in a reasoning-bearing SSE frame.
**Suggested fix:** Preserve and transform every choice in the original `choices` list, rather than reconstructing only `choices[0]`.
```suggestion
transformed_choices = []
for choice in choices:
choice = dict(choice)
delta = choice.get("delta") if isinstance(choice.get("delta"), dict) else {}
delta = {key: value for key, value in delta.items() if key == "role"}
delta["reasoning_content"] = "".join(pending)
choice["delta"] = delta
choice["finish_reason"] = None
transformed_choices.append(choice)
payload["choices"] = transformed_choices
```
</issue_to_address>
### Comment 3
<location path="converter.py" line_range="3127" />
<code_context>
return sanitize_log_text(f"{type(error).__name__}: {str(error).strip() or 'upstream transport failed'}", 512)
+async def _coalesce_reasoning_sse(lines, *, max_bytes=0):
+ """Emit one Chat reasoning delta before visible output for each stream."""
+ pending: list[str] = []
+ template = None
</code_context>
<issue_to_address>
**nitpick:** The function docstring says it emits one Chat reasoning delta, but the function is also inserted into the Responses and Anthropic adapted streaming paths, so its documented scope does not describe its actual callers or behavior.
**Suggested fix:** Update the docstring to describe coalescing normalized Chat SSE reasoning before all downstream protocol adapters.
```suggestion
"""Coalesce normalized Chat SSE reasoning before all downstream protocol adapters."""
```
</issue_to_address>Sourcery assessment
Approval pending. 2 findings to address first.
Blocking findings: converter.py:3198, converter.py:3146
| if template is None: | ||
| template = event | ||
| pending.append(reasoning) | ||
| buffer_budget.charge_text(reasoning) |
There was a problem hiding this comment.
issue (broader_impact): Reasoning bytes are charged once by the upstream ChatSSEAccumulator budget and charged again by _coalesce_reasoning_sse or the adapted converter, so a valid stream can exceed max_collect_bytes and fail with response_too_large even though its retained output is within the configured limit.
Triggers: When realtime streaming has a nonzero max_collect_bytes and the stream contains reasoning or other retained text.
Suggested fix: Share the existing request budget with _coalesce_reasoning_sse, or remove the duplicate charge from one of the buffering/conversion layers.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: afe33af1fe
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| if not line.startswith("data:"): | ||
| flushed = flush() | ||
| if flushed is not None: | ||
| yield flushed | ||
| yield "" | ||
| yield line |
There was a problem hiding this comment.
Keep coalescing across SSE control lines
When the upstream stream contains a valid SSE comment or control line (for example a keep-alive : ping or event: line) between two reasoning deltas, this branch flushes and resets the pending buffer before forwarding that line. The following reasoning then becomes a separate client block, so Responses/Anthropic clients can observe multiple thinking sections despite the new coalescing guarantee; ignorable SSE control lines should not terminate the reasoning run.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 65aa35a5bf
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| visible = any(delta.get(key) for key in ("content", "refusal", "tool_calls", "function_call")) | ||
| visible = visible or bool(choice.get("finish_reason")) | ||
| if not visible: | ||
| continue |
There was a problem hiding this comment.
Preserve visible deltas from other choices
When an upstream chunk contains reasoning in choice 0 but visible content, refusal, or tool data in another choice, visible is computed only from choice 0 and this branch continues, dropping the entire event. The later transformed event cannot recover that other choice's delta, so multi-choice responses silently lose output; determine visibility across all choices before suppressing the frame.
Useful? React with 👍 / 👎.
Summary
max_collect_bytesSummary by Sourcery
Coalesce streamed reasoning before visible output so clients receive one continuous thinking section without losing incremental response content.
Bug Fixes:
Enhancements:
Documentation:
Tests: