From afe33af1feacfd662685b42a8945bf38b8398d1d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=A2=A8=E8=8F=8A?= <277378677+maiphucgiang@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:33:01 +0800 Subject: [PATCH 1/2] Coalesce streamed reasoning into one client block --- converter.py | 115 +++++++++++++++++++++++++++++-- docs/advanced.md | 2 +- docs/advanced.zh-CN.md | 2 +- tests/test_realtime_streaming.py | 60 +++++++++++++--- tests/test_reasoning.py | 2 + 5 files changed, 164 insertions(+), 17 deletions(-) diff --git a/converter.py b/converter.py index 62e5936..3d5db6f 100644 --- a/converter.py +++ b/converter.py @@ -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,108 @@ 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): + """Emit one Chat reasoning delta before visible output for each stream.""" + 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 [] + 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 + 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:"): + flushed = flush() + if flushed is not None: + yield flushed + yield "" + 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) + 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 + flushed = flush() + if flushed is not None: + yield flushed + yield "" + choice = dict(choice) + clean_delta = dict(delta) + clean_delta.pop("reasoning_content", None) + choice["delta"] = clean_delta + event = dict(event) + event["choices"] = [choice] + if clean_delta or choice.get("finish_reason") is not None: + 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 +3455,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 +3877,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: diff --git a/docs/advanced.md b/docs/advanced.md index 40f3df7..228eebc 100644 --- a/docs/advanced.md +++ b/docs/advanced.md @@ -72,7 +72,7 @@ Use a source/image build and Compose configuration containing this feature; recr `stream_mode` defaults to `compatible`. Configure it through `--stream-mode compatible|realtime`, `CODEBUDDY2API_STREAM_MODE`, or the WebUI enum; explicit CLI/environment sources lock that field. Every generation request, including non-streaming, records the selected mode and freezes it with `max_collect_bytes` before routing. Hot changes affect only later requests, not in-flight failovers. Non-streaming responses remain aggregated JSON; there is no client mode override. - `compatible` preserves existing behavior: Responses streams aggregate first; Chat and Messages aggregate when tools are present and otherwise pass through upstream increments. Aggregated output is validated and replayed in fragments. Non-stream requests always use the validated aggregate path in either mode. -- `realtime` forwards reasoning, text, refusal and tool-argument increments for all three protocols. Responses assigns stable indexes when items start; Anthropic uses stable block indexes. Adapters buffer tools with missing identity until the argument phase or terminal marker, append metadata fragments without guessing from prefixes, and reject identity changes after an item starts. `max_collect_bytes` bounds retained UTF-8 output; `0` disables that limit. +- `realtime` forwards text, refusal and tool-argument increments for all three protocols, while coalescing reasoning fragments into one delta before visible output so clients keep one thinking section. Responses assigns stable indexes when items start; Anthropic uses stable block indexes. Adapters buffer tools with missing identity until the argument phase or terminal marker, append metadata fragments without guessing from prefixes, and reject identity changes after an item starts. `max_collect_bytes` bounds retained UTF-8 output; `0` disables that limit. Realtime Messages keeps one content block open at a time. The active tool remains incremental; later tool, text or thinking blocks may wait until upstream completion. Deferred event bytes share `max_collect_bytes`. A tool still awaiting its identity does not block unrelated text before its block starts. diff --git a/docs/advanced.zh-CN.md b/docs/advanced.zh-CN.md index b5aaa40..6dcf37e 100644 --- a/docs/advanced.zh-CN.md +++ b/docs/advanced.zh-CN.md @@ -72,7 +72,7 @@ Responses 投影不再修改工具定义或 schema。启用脱敏时,默认剥 `stream_mode` 默认 `compatible`,可通过 `--stream-mode compatible|realtime`、`CODEBUDDY2API_STREAM_MODE` 或 WebUI 枚举配置;显式 CLI/环境来源会锁定该项。所有生成请求(含非流式)均在选路前冻结模式及 `max_collect_bytes`,并记录所选模式。热更新只影响后续请求,不改变在途请求及换号;非流式仍返回聚合 JSON,客户端不能按请求覆盖模式。 - `compatible` 保持现有行为:Responses 流式先聚合;Chat/Messages 带工具时先聚合,无工具时沿用上游增量。聚合结果先校验,再按片段重放。两种模式下,非流式请求始终走已校验的聚合路径。 -- `realtime` 让三个协议都增量发送思考、正文、拒绝及工具参数。Responses 在输出项开始时分配稳定索引,Anthropic 使用稳定 block index。适配器在参数阶段或结束标记处确认工具身份,缺失时暂缓工具输出;元数据分片按顺序追加,不按字符串前缀猜测,输出项开始后禁止更换身份。`max_collect_bytes` 约束所保留的 UTF-8 输出,`0` 不限制。 +- `realtime` 让三个协议都增量发送正文、拒绝及工具参数,并在正文开始前合并思考增量为一个事件,避免客户端显示多个思考段。Responses 在输出项开始时分配稳定索引,Anthropic 使用稳定 block index。适配器在参数阶段或结束标记处确认工具身份,缺失时暂缓工具输出;元数据分片按顺序追加,不按字符串前缀猜测,输出项开始后禁止更换身份。`max_collect_bytes` 约束所保留的 UTF-8 输出,`0` 不限制。 实时 Messages 同时只打开一个内容块:当前工具仍增量输出,后续工具、正文或思考块可能等待上游结束,暂存事件字节计入 `max_collect_bytes`。尚未确认身份、未开始内容块的工具不会阻塞其它正文。 diff --git a/tests/test_realtime_streaming.py b/tests/test_realtime_streaming.py index ad8d4d5..4ac7ddb 100644 --- a/tests/test_realtime_streaming.py +++ b/tests/test_realtime_streaming.py @@ -644,16 +644,56 @@ async def test_realtime_http_failure_preserves_status_and_retry_after(self): self.assertIsInstance(failure, UpstreamResponseError) self.assertEqual((failure.status, failure.headers.get("Retry-After")), (429, "7")) - async def test_reasoning_delta_reaches_asgi_client_before_release(self): - before, wire = await self.drive( - "responses", "realtime", False, - first_delta={"reasoning_content": "meaningful reasoning"}) - self.assertIn(b"response.reasoning_summary_text.delta", before) - self.assertIn(b"meaningful reasoning", before) - reasoning_deltas = [event["delta"] for event in _events(wire.decode()) - if event["type"] == "response.reasoning_summary_text.delta"] - self.assertEqual("".join(reasoning_deltas), "meaningful reasoning") - self.assertIn(b"response.completed", wire) + async def test_reasoning_deltas_are_coalesced_before_visible_output(self): + rows = [ + _line({"reasoning_content": "W"}), + _line({"reasoning_content": "line"}), + _line({"reasoning_content": "478", "content": "85"}), + _line({"content": "done"}), + _line({}, "stop", {"prompt_tokens": 1, "completion_tokens": 4, "total_tokens": 5}), + "data: [DONE]\n\n", + ] + for protocol in ("chat", "responses", "messages"): + with self.subTest(protocol=protocol): + response = _FixedResponse(rows) + + @asynccontextmanager + async def backend(*args, **kwargs): + yield response + + path = {"chat": "/v1/chat/completions", "responses": "/v1/responses", + "messages": "/v1/messages"}[protocol] + sent = await self.asgi_post( + path, self.payload(protocol, False), backend, + config={"stream_mode": "realtime"}) + wire = b"".join(message.get("body", b"") for message in sent + if message["type"] == "http.response.body") + if protocol == "chat": + events = [json.loads(line[6:]) for line in wire.decode().splitlines() + if line.startswith("data: ") and line[6:] != "[DONE]"] + reasoning = [event["choices"][0]["delta"].get("reasoning_content") + for event in events + if event.get("choices") and event["choices"][0]["delta"].get("reasoning_content")] + content = [event["choices"][0]["delta"].get("content") + for event in events + if event.get("choices") and event["choices"][0]["delta"].get("content")] + self.assertEqual(reasoning, ["Wline478"]) + self.assertEqual(content, ["85", "done"]) + self.assertLess( + next(i for i, event in enumerate(events) + if event["choices"][0]["delta"].get("reasoning_content")), + next(i for i, event in enumerate(events) + if event["choices"][0]["delta"].get("content"))) + elif protocol == "responses": + events = _events(wire.decode()) + reasoning = [event["delta"] for event in events + if event["type"] == "response.reasoning_summary_text.delta"] + self.assertEqual(reasoning, ["Wline478"]) + else: + blocks = _serial_blocks(self, wire.decode()) + self.assertEqual([block["type"] for block in blocks], ["thinking", "text"]) + self.assertEqual(blocks[0]["thinking"], "Wline478") + self.assertEqual(blocks[1]["text"], "85done") async def test_hot_mode_change_affects_only_the_next_request(self): before, _ = await self.drive("responses", "realtime", True, change_to="compatible") diff --git a/tests/test_reasoning.py b/tests/test_reasoning.py index 4b9a548..abd533c 100644 --- a/tests/test_reasoning.py +++ b/tests/test_reasoning.py @@ -73,7 +73,9 @@ def test_sse_replay_reasoning_before_content(self): reasoning += delta.get("reasoning_content") or "" content += delta.get("content") or "" self.assertEqual(reasoning, "思考一思考二") + self.assertEqual(sum("reasoning_content" in line for line in lines), 1) self.assertIn("正文", content) + self.assertEqual(sum("reasoning_content" in line for line in lines), 1) # Reasoning deltas precede text deltas. first_reasoning = next(i for i, l in enumerate(lines) if "reasoning_content" in l) first_content = next(i for i, l in enumerate(lines) if '"content": "正' in l or '"content":"正' in l) From 65aa35a5bf6356f6aabdb935b0c0d2d6401204ba Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=A2=A8=E8=8F=8A?= <277378677+maiphucgiang@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:49:31 +0800 Subject: [PATCH 2/2] Harden reasoning stream coalescing --- converter.py | 52 ++++++++++++++++++++------------ tests/test_realtime_streaming.py | 1 + tests/test_reasoning.py | 30 +++++++++++++++++- 3 files changed, 62 insertions(+), 21 deletions(-) diff --git a/converter.py b/converter.py index 3d5db6f..ccff6f5 100644 --- a/converter.py +++ b/converter.py @@ -3124,11 +3124,11 @@ 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): - """Emit one Chat reasoning delta before visible output for each stream.""" + """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: @@ -3136,14 +3136,20 @@ def flush(): 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] + 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 = [] @@ -3158,10 +3164,7 @@ def flush(): yield line continue if not line.startswith("data:"): - flushed = flush() - if flushed is not None: - yield flushed - yield "" + # SSE comments and event/control lines do not end a reasoning run. yield line if embedded_boundary: yield "" @@ -3204,13 +3207,22 @@ def flush(): if flushed is not None: yield flushed yield "" - choice = dict(choice) - clean_delta = dict(delta) - clean_delta.pop("reasoning_content", None) - choice["delta"] = clean_delta + 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"] = [choice] - if clean_delta or choice.get("finish_reason") is not None: + 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 "" diff --git a/tests/test_realtime_streaming.py b/tests/test_realtime_streaming.py index 4ac7ddb..07c8771 100644 --- a/tests/test_realtime_streaming.py +++ b/tests/test_realtime_streaming.py @@ -648,6 +648,7 @@ async def test_reasoning_deltas_are_coalesced_before_visible_output(self): rows = [ _line({"reasoning_content": "W"}), _line({"reasoning_content": "line"}), + ": ping\n\n", _line({"reasoning_content": "478", "content": "85"}), _line({"content": "done"}), _line({}, "stop", {"prompt_tokens": 1, "completion_tokens": 4, "total_tokens": 5}), diff --git a/tests/test_reasoning.py b/tests/test_reasoning.py index abd533c..dc6902d 100644 --- a/tests/test_reasoning.py +++ b/tests/test_reasoning.py @@ -11,7 +11,8 @@ import httpx -from converter import _collect_stream, _merge_chat_sse_text, _chat_result_to_sse_lines +from converter import (_chat_result_to_sse_lines, _coalesce_reasoning_sse, _collect_stream, + _merge_chat_sse_text) from app.adapters.anthropic_adapter import AnthropicStreamConverter from app.adapters.responses_adapter import ResponsesStreamConverter @@ -81,6 +82,33 @@ def test_sse_replay_reasoning_before_content(self): first_content = next(i for i, l in enumerate(lines) if '"content": "正' in l or '"content":"正' in l) self.assertLess(first_reasoning, first_content) + def test_reasoning_coalescer_preserves_multiple_choices(self): + rows = [ + "data: " + json.dumps({"id": "c", "choices": [ + {"index": 0, "delta": {"reasoning_content": "思考"}, "finish_reason": None}, + {"index": 1, "delta": {"content": "旁路"}, "finish_reason": None}, + ]}), + "data: " + json.dumps({"id": "c", "choices": [ + {"index": 0, "delta": {"content": "正文"}, "finish_reason": None}, + {"index": 1, "delta": {"content": "继续"}, "finish_reason": None}, + ]}), + "data: [DONE]", + ] + + async def source(): + for row in rows: + yield row + + async def collect(): + return [line async for line in _coalesce_reasoning_sse(source())] + + events = [json.loads(line[6:]) for line in asyncio.run(collect()) + if line.startswith("data: ") and line[6:] != "[DONE]"] + self.assertTrue(events) + self.assertTrue(all(len(event["choices"]) == 2 for event in events)) + self.assertEqual(events[0]["choices"][0]["delta"]["reasoning_content"], "思考") + + class TestReplayEnvelope(unittest.TestCase): """Include complete Chat chunk envelopes in content and usage replay events."""