Repository navigation
Coalesce streamed reasoning into one client block #54
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -3100,8 +3100,8 @@ def _line(delta: dict, fr=None) -> str: | |
| return "data: " + json.dumps(payload, ensure_ascii=False) | ||
|
|
||
| lines = [_line({"role": "assistant", "content": ""})] | ||
| for i in range(0, len(reasoning), 48): | ||
| lines.append(_line({"reasoning_content": reasoning[i:i + 48]})) | ||
| if reasoning: | ||
| lines.append(_line({"reasoning_content": reasoning})) | ||
| for i in range(0, len(content), 48): | ||
| lines.append(_line({"content": content[i:i + 48]})) | ||
| refusal = m.get("refusal") or "" | ||
|
|
@@ -3123,6 +3123,120 @@ def _line(delta: dict, fr=None) -> str: | |
| def _network_error_text(error: Exception) -> str: | ||
| 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): | ||
| """Coalesce normalized Chat SSE reasoning before downstream protocol adapters.""" | ||
| pending: list[str] = [] | ||
| template = None | ||
| buffer_budget = StreamOutputBudget(max_bytes) | ||
|
|
||
| def flush(): | ||
| nonlocal template, pending | ||
| if not pending: | ||
| return None | ||
| source = template or {} | ||
| payload = dict(source) | ||
| choices = source.get("choices") or [] | ||
| transformed = [] | ||
| for index, raw_choice in enumerate(choices): | ||
| if not isinstance(raw_choice, dict): | ||
| transformed.append(raw_choice) | ||
| continue | ||
| choice = dict(raw_choice) | ||
| if index == 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 | ||
| transformed.append(choice) | ||
| payload["choices"] = transformed | ||
| wire = "data: " + json.dumps(payload, ensure_ascii=False) | ||
| template = None | ||
| pending = [] | ||
| return wire | ||
|
|
||
| async for raw_line in lines: | ||
| embedded_boundary = isinstance(raw_line, str) and raw_line.endswith(("\n", "\r")) | ||
| line = raw_line.rstrip("\r\n") if isinstance(raw_line, str) else raw_line | ||
| if not line or not line.strip(): | ||
| if pending: | ||
| continue | ||
| yield line | ||
| continue | ||
| if not line.startswith("data:"): | ||
| # SSE comments and event/control lines do not end a reasoning run. | ||
| yield line | ||
| if embedded_boundary: | ||
| yield "" | ||
| continue | ||
| data = line[5:].strip() | ||
| if data == "[DONE]": | ||
| flushed = flush() | ||
| if flushed is not None: | ||
| yield flushed | ||
| yield "" | ||
| yield line | ||
| if embedded_boundary: | ||
| yield "" | ||
| continue | ||
| try: | ||
| event = json.loads(data) | ||
| except (TypeError, ValueError): | ||
| flushed = flush() | ||
| if flushed is not None: | ||
| yield flushed | ||
| yield "" | ||
| yield line | ||
| if embedded_boundary: | ||
| yield "" | ||
| continue | ||
| choices = event.get("choices") if isinstance(event, dict) else None | ||
| choice = choices[0] if isinstance(choices, list) and choices else None | ||
| delta = choice.get("delta") if isinstance(choice, dict) else None | ||
| reasoning = delta.get("reasoning_content") if isinstance(delta, dict) else None | ||
| if isinstance(reasoning, str) and reasoning: | ||
| if template is None: | ||
| template = event | ||
| pending.append(reasoning) | ||
| buffer_budget.charge_text(reasoning) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. issue (broader_impact): Reasoning bytes are charged once by the upstream Triggers: When realtime streaming has a nonzero Suggested fix: Share the existing request budget with |
||
| 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 | ||
|
Comment on lines
+3202
to
+3205
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
When an upstream chunk contains reasoning in choice 0 but visible content, refusal, or tool data in another choice, Useful? React with 👍 / 👎. |
||
| flushed = flush() | ||
| if flushed is not None: | ||
| yield flushed | ||
| yield "" | ||
| transformed_choices = [] | ||
| for index, raw_choice in enumerate(choices): | ||
| if not isinstance(raw_choice, dict): | ||
| transformed_choices.append(raw_choice) | ||
| continue | ||
| clean_choice = dict(raw_choice) | ||
| if index == 0: | ||
| clean_delta = dict(delta) | ||
| clean_delta.pop("reasoning_content", None) | ||
| clean_choice["delta"] = clean_delta | ||
| transformed_choices.append(clean_choice) | ||
| event = dict(event) | ||
| event["choices"] = transformed_choices | ||
| if any(isinstance(item, dict) and | ||
| (item.get("delta") or item.get("finish_reason") is not None) | ||
| for item in transformed_choices): | ||
| yield "data: " + json.dumps(event, ensure_ascii=False) | ||
| if embedded_boundary: | ||
| yield "" | ||
| continue | ||
| flushed = flush() | ||
| if flushed is not None: | ||
| yield flushed | ||
| yield "" | ||
| yield line | ||
| if embedded_boundary: | ||
| yield "" | ||
|
|
||
|
|
||
|
|
||
| def _public_sse_line(line, model_name): | ||
| if CONFIG.get("control_store") is not None and line.startswith("data:"): | ||
| try: | ||
|
|
@@ -3353,7 +3467,9 @@ async def _stream_upstream(url: str, headers: dict, body: dict, | |
| model_name: str = "?", t0: float = 0.0, rid: str = "", cred=None): | ||
| policy = body.pop(_REQUEST_POLICY_KEY, None) or _snapshot_stream_policy("chat", body) | ||
| sent = False | ||
| upstream = _chat_sse_lines(url, headers, body, model_name, t0, rid, cred, policy=policy) | ||
| upstream = _coalesce_reasoning_sse(_chat_sse_lines( | ||
| url, headers, body, model_name, t0, rid, cred, policy=policy), | ||
| max_bytes=policy.max_collect_bytes) | ||
| try: | ||
| try: | ||
| async for line in upstream: | ||
|
|
@@ -3773,9 +3889,10 @@ async def _stream_adapted(url, headers, body, model_name, t0, rid, cred=None, *, | |
| parallel_tool_calls=body.get("parallel_tool_calls", True), | ||
| tool_registry=tool_registry)) | ||
| sent = False | ||
| upstream = _chat_sse_lines( | ||
| upstream = _coalesce_reasoning_sse(_chat_sse_lines( | ||
| url, headers, body, model_name, t0, rid, cred, | ||
| policy=policy, tracker=tracker, state=state) | ||
| policy=policy, tracker=tracker, state=state), | ||
| max_bytes=policy.max_collect_bytes) | ||
| try: | ||
| try: | ||
| async for line in upstream: | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
When the upstream stream contains a valid SSE comment or control line (for example a keep-alive
: pingorevent: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 👍 / 👎.