From 572446308e07f5a27e1ac01badcee99f3b5ab294 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 21:42:06 +0530 Subject: [PATCH] fix(agent): feed inline text to the live reasoning pane (#89647) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The desktop/TUI reasoning pane is driven by reasoning.delta via reasoning_callback; the scrubber-side collector only filled the final reasoning_content, which extract_reasoning already recovers from the raw content, so the pane stayed dead. - Drop the scrubber reasoning collector (_reasoning_parts, reasoning(), clear_reasoning(), \x00 sentinel, _THINK_TAG_RE) and its per-request reset hook. - StreamingThinkScrubber.feed() exposes the text it stripped from inside think blocks as last_hidden; _fire_stream_delta forwards it through _fire_reasoning_delta(inline=True) while no native reasoning delta has arrived for this model response (reset per request) — no double reasoning. CLI gating is unchanged: its reasoning_callback is None unless show_reasoning/verbose. - _finish_chat_stream fills reasoning_content from the raw content via the existing extract_reasoning when no reasoning delta arrived. - Replace the collector tests with two guards (live forwarding + native suppression; _finish_chat_stream fallback), both red on origin/main. Co-authored-by: SayHell0W0rld <852938468@qq.com> --- agent/chat_completion_helpers.py | 13 ++++--- agent/stream_delivery.py | 21 ++++++++---- agent/think_scrubber.py | 45 ++++++------------------- tests/agent/test_plugin_stream_hooks.py | 30 +++++++++++++++++ tests/agent/test_think_scrubber.py | 22 ------------ 5 files changed, 61 insertions(+), 70 deletions(-) diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index af28613307..4c1d666f0d 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -3269,13 +3269,12 @@ class _StreamingCall(StreamingWaitMonitor): args or stamping "stop".""" full_content = "".join(content_parts) or None full_reasoning = "".join(reasoning_parts) or None - if not full_reasoning: - # Providers that inline reasoning (MiniMax-M3 streams … in content) never - # send a reasoning delta — recover what the think scrubber stripped so the structured - # reasoning_content field stays populated (#89647). - think_scrubber = getattr(self.agent, "_stream_think_scrubber", None) - if think_scrubber is not None: - full_reasoning = think_scrubber.reasoning() or None + if not full_reasoning and full_content: + # Inline-reasoning providers (MiniMax-M3 streams … in content) send no + # reasoning delta; fill the structured field from the raw content (#89647). + from agent.agent_runtime_helpers import extract_reasoning + + full_reasoning = extract_reasoning(self.agent, SimpleNamespace(content=full_content)) mock_tool_calls, has_truncated_tool_args = self._assemble_tool_calls(tool_calls_acc, finish_reason) # Zero-chunk guard: nothing usable = upstream error / malformed SSE. if finish_reason is None and not content_parts and not reasoning_parts and not refusal_parts and not tool_calls_acc: diff --git a/agent/stream_delivery.py b/agent/stream_delivery.py index 9644d264e8..1027e65667 100644 --- a/agent/stream_delivery.py +++ b/agent/stream_delivery.py @@ -52,6 +52,9 @@ class StreamDeliveryMixin: """ think_scrubber = getattr(self, "_stream_think_scrubber", None) ctx_scrubber = getattr(self, "_stream_context_scrubber", None) + # Inline text is forwarded to the reasoning pane only while no native reasoning delta + # arrived for this model response (#89647). + self._native_reasoning_streamed = False # Next stream re-reads plugins.stream_reasoning_deltas (config edits land per request). self._stream_reasoning_hooks_enabled = None @@ -72,10 +75,6 @@ class StreamDeliveryMixin: if think_scrubber is not None: think_tail = think_scrubber.flush() deliver(ctx_scrubber.feed(think_tail) if think_tail and ctx_scrubber is not None else think_tail) - # Inline reasoning recovered for reasoning_content is per model response (#89647). - clear_reasoning = getattr(think_scrubber, "clear_reasoning", None) - if callable(clear_reasoning): - clear_reasoning() if ctx_scrubber is not None: deliver(ctx_scrubber.flush()) self._current_streamed_assistant_text = "" @@ -308,6 +307,11 @@ class StreamDeliveryMixin: # See #5719. scrubber = getattr(self, "_stream_context_scrubber", None) text = think_scrubber.feed(text) if think_scrubber is not None else self._strip_think_blocks(text) + # Providers that inline reasoning (MiniMax-M3 …) send no reasoning delta, so the + # live reasoning pane would stay empty; forward what the scrubber stripped instead (#89647). + hidden = think_scrubber.last_hidden if think_scrubber is not None else "" + if hidden and not getattr(self, "_native_reasoning_streamed", False): + self._fire_reasoning_delta(hidden, inline=True) text = scrubber.feed(text) if scrubber is not None else sanitize_context(text) # Only strip leading newlines on the first delta — mid-stream "\n" is legitimate markdown. # Check the parts list, not the joined property (joining per token copies the whole reply). @@ -320,8 +324,13 @@ class StreamDeliveryMixin: if delivered: self._record_streamed_assistant_text(text) - def _fire_reasoning_delta(self, text: str) -> None: - """Fire reasoning callback if registered; superseded writers are fenced like content deltas.""" + def _fire_reasoning_delta(self, text: str, *, inline: bool = False) -> None: + """Fire reasoning callback if registered; superseded writers are fenced like content deltas. + + ``inline`` marks text recovered from ```` blocks in content; any other call is a native + provider reasoning delta and stops inline forwarding for the rest of this model response.""" + if not inline: + self._native_reasoning_streamed = True if self._stream_writer_superseded(): # Single-writer guard (#65991): fence out a superseded stream's reasoning deltas the same way as # content deltas. diff --git a/agent/think_scrubber.py b/agent/think_scrubber.py index a0d9647c5c..e2dc09103f 100644 --- a/agent/think_scrubber.py +++ b/agent/think_scrubber.py @@ -23,13 +23,6 @@ THINK_TAG_NAMES: Tuple[str, ...] = ("think", "thinking", "reasoning", "thought", THINK_OPEN_TAGS: Tuple[str, ...] = tuple(f"<{name.lower()}>" for name in THINK_TAG_NAMES) THINK_CLOSE_TAGS: Tuple[str, ...] = tuple(f"" for name in THINK_TAG_NAMES) -# Tag markup stripped by reasoning(); matches every variant the scrubber -# recognizes (open, close, case-insensitive). -_THINK_TAG_RE = re.compile( - r"<\s*/?\s*(?:" + "|".join(THINK_TAG_NAMES) + r")\s*>", - re.IGNORECASE, -) - class StreamingThinkScrubber: """Stateful scrubber for streaming reasoning/thinking blocks. @@ -44,8 +37,6 @@ class StreamingThinkScrubber: _CLOSE_TAGS: Tuple[str, ...] = THINK_CLOSE_TAGS _ALL_TAGS: Tuple[str, ...] = _OPEN_TAGS + _CLOSE_TAGS _MAX_TAG_LEN: int = max(len(tag) for tag in _ALL_TAGS) - # Separates collected reasoning blocks; streamed pieces of one block rejoin verbatim. - _BLOCK_END: str = "\x00" # Orphan close tag plus trailing whitespace (matches _strip_think_blocks case 3). _ORPHAN_CLOSE_RE = re.compile( "(?:" + "|".join(re.escape(t) for t in _CLOSE_TAGS) + r")[ \t\n\r]*", re.IGNORECASE @@ -59,7 +50,8 @@ class StreamingThinkScrubber: self._in_block: bool = False self._buf: str = "" self._last_emitted_ended_newline: bool = True - self._reasoning_parts: list[str] = [] + # Reasoning text the most recent feed() stripped from inside think blocks (tags excluded). + self.last_hidden: str = "" def _emit(self, out: list[str], text: str) -> None: """Append visible prose to *out* (orphan close tags stripped) and track the newline flag.""" @@ -70,24 +62,22 @@ class StreamingThinkScrubber: def feed(self, text: str) -> str: """Feed one delta; return the scrubbed visible portion ("" when it is all reasoning or held back).""" + self.last_hidden = "" if not text: return "" buf = self._buf + text self._buf = "" out: list[str] = [] + hidden: list[str] = [] while buf: if self._in_block: close_idx, close_len = self._find_first_tag(buf, self._CLOSE_TAGS) if close_idx == -1: - # No close yet: hold back a possible partial close-tag prefix; collect - # the rest as reasoning (#89647). - discard = self._hold_partial(buf, self._CLOSE_TAGS) - if discard: - self._reasoning_parts.append(discard) + # No close yet: hold back a possible partial close-tag prefix; the rest is reasoning. + hidden.append(self._hold_partial(buf, self._CLOSE_TAGS)) break - # Found close: collect block content as reasoning (#89647). - self._reasoning_parts.append(buf[:close_idx] + self._BLOCK_END) + hidden.append(buf[:close_idx]) buf = buf[close_idx + close_len:] self._in_block = False continue @@ -99,8 +89,8 @@ class StreamingThinkScrubber: open_idx, open_len = self._find_open_at_boundary(buf, out) if pair is not None and (open_idx == -1 or pair[0] <= open_idx): self._emit(out, buf[:pair[0]]) - # Collect the stripped pair (tags removed by reasoning()) (#89647). - self._reasoning_parts.append(buf[pair[0]:pair[1]] + self._BLOCK_END) + # Pair tags are exact ````/````: inner text sits between them. + hidden.append(buf[buf.index(">", pair[0]) + 1:buf.rindex("<", pair[0], pair[1])]) buf = buf[pair[1]:] continue if open_idx != -1: @@ -114,6 +104,7 @@ class StreamingThinkScrubber: self._emit(out, self._hold_partial(buf, self._ALL_TAGS)) break + self.last_hidden = "".join(hidden) return "".join(out) def _hold_partial(self, buf: str, tags: Tuple[str, ...]) -> str: @@ -127,28 +118,12 @@ class StreamingThinkScrubber: partial reasoning is worse than a truncated answer), otherwise the tail is emitted verbatim. Always resets the boundary flag — intra-turn retries flush then stream again without ``reset()``, and a stale False flag made the new stream's opening ```` look mid-line.""" - if self._in_block: - self._reasoning_parts.append(self._buf + self._BLOCK_END) tail = "" if self._in_block else self._buf self._buf = "" self._in_block = False self._last_emitted_ended_newline = True return self._strip_orphan_close_tags(tail) if tail else "" - def reasoning(self) -> str: - """Reasoning text stripped from streamed content so far, tag markup removed ("" when none). - - Lets callers populate a structured ``reasoning_content`` field for providers (e.g. MiniMax-M3) - that inline reasoning instead of returning a reasoning delta (#89647). - """ - blocks = _THINK_TAG_RE.sub("", "".join(self._reasoning_parts)).split(self._BLOCK_END) - return "\n".join(t for t in (b.strip() for b in blocks) if t) - - def clear_reasoning(self) -> None: - """Drop collected reasoning. Called per model request so one API call's inline reasoning - does not bleed into the next call's ``reasoning_content`` within the same turn.""" - self._reasoning_parts = [] - # ── internal helpers ─────────────────────────────────────────────── @staticmethod diff --git a/tests/agent/test_plugin_stream_hooks.py b/tests/agent/test_plugin_stream_hooks.py index f44a8a9cec..414106c2b1 100644 --- a/tests/agent/test_plugin_stream_hooks.py +++ b/tests/agent/test_plugin_stream_hooks.py @@ -360,3 +360,33 @@ def test_bedrock_reasoning_delta_reaches_plugin_only_observer(monkeypatch): assert calls[0]["kind"] == "reasoning" assert calls[0]["delta"] == "bedrock reasoning" + + +def test_inline_think_reaches_reasoning_pane_unless_native_reasoning_streamed(): + """#89647: inline text stripped from content feeds reasoning_callback (the live pane), but not + once the provider streamed native reasoning for this response (no double reasoning).""" + agent = _agent() + seen = [] + agent.reasoning_callback = seen.append + agent._reset_stream_delivery_tracking() + for delta in ["", "Let me", " check config", "", "The answer is 42."]: + agent._fire_stream_delta(delta) + assert "".join(seen) == "Let me check config" + + seen.clear() + agent._reset_stream_delivery_tracking() + agent._fire_reasoning_delta("native") + agent._fire_stream_delta("dupok") + assert seen == ["native"] + + +def test_finish_chat_stream_recovers_inline_reasoning_content(): + """#89647: with no reasoning delta, reasoning_content comes from the blocks in raw content.""" + from agent import chat_completion_helpers as cch + + call = cch._StreamingCall.__new__(cch._StreamingCall) + call.agent = _agent() + deltas = ["", "Let me", " check config", "", "The answer is 42."] + resp = call._finish_chat_stream(None, "assistant", deltas, [], {}, "stop", "MiniMax-M3", None, + flush_pending=lambda: None) + assert resp.choices[0].message.reasoning_content == "Let me check config" diff --git a/tests/agent/test_think_scrubber.py b/tests/agent/test_think_scrubber.py index 38056bc6c8..8c1a2d69d9 100644 --- a/tests/agent/test_think_scrubber.py +++ b/tests/agent/test_think_scrubber.py @@ -182,25 +182,3 @@ class TestRealisticStreaming: s = StreamingThinkScrubber() deltas = ["Hello ", "world ", "how ", "are ", "you?"] assert _drive(s, deltas) == "Hello world how are you?" - - - -class TestReasoningCollection: - """The scrubber collects stripped reasoning for structured fields (#89647).""" - - def test_streamed_blocks_rejoin_verbatim_and_tags_are_removed(self) -> None: - s = StreamingThinkScrubber() - visible = _drive(s, ["", "Let me check", " their config", "", "Done. ", - "planok"]) - assert visible == "Done. ok" - assert s.reasoning() == "Let me check their config\nplan" - - def test_per_request_reset_drops_previous_response_reasoning(self) -> None: - from agent.stream_delivery import StreamDeliveryMixin - - agent = StreamDeliveryMixin.__new__(StreamDeliveryMixin) - agent._stream_think_scrubber = StreamingThinkScrubber() - agent._stream_context_scrubber = None - agent._stream_think_scrubber.feed("first callhi") - agent._reset_stream_delivery_tracking() - assert agent._stream_think_scrubber.reasoning() == ""