fix(agent): feed inline <think> text to the live reasoning pane (#89647)
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>
This commit is contained in:
@@ -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 <think>…</think> 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 <think>…</think> 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:
|
||||
|
||||
@@ -52,6 +52,9 @@ class StreamDeliveryMixin:
|
||||
"""
|
||||
think_scrubber = getattr(self, "_stream_think_scrubber", None)
|
||||
ctx_scrubber = getattr(self, "_stream_context_scrubber", None)
|
||||
# Inline <think> 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 <think>…</think>) 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 ``<think>`` 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.
|
||||
|
||||
@@ -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"</{name.lower()}>" 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 ``<name>``/``</name>``: 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 ``<think>`` 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
|
||||
|
||||
@@ -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 <think> 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 ["<think>", "Let me", " check config", "</think>", "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("<think>dup</think>ok")
|
||||
assert seen == ["native"]
|
||||
|
||||
|
||||
def test_finish_chat_stream_recovers_inline_reasoning_content():
|
||||
"""#89647: with no reasoning delta, reasoning_content comes from the <think> blocks in raw content."""
|
||||
from agent import chat_completion_helpers as cch
|
||||
|
||||
call = cch._StreamingCall.__new__(cch._StreamingCall)
|
||||
call.agent = _agent()
|
||||
deltas = ["<think>", "Let me", " check config", "</think>", "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"
|
||||
|
||||
@@ -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, ["<think>", "Let me check", " their config", "</think>", "Done. ",
|
||||
"<thinking>plan</thinking>ok"])
|
||||
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("<think>first call</think>hi")
|
||||
agent._reset_stream_delivery_tracking()
|
||||
assert agent._stream_think_scrubber.reasoning() == ""
|
||||
|
||||
Reference in New Issue
Block a user