refactor(agent/prompt): remove dead code, unify duplicated helpers, compact docstrings across prompt/skill/redaction modules
Dead (zero refs): coding_system_blocks, get_friendly_tool_labels, get_scan_ordered_skills_dirs, _project_quarantine_cache_clear, clear_stable_prefixes, _redact_http_request_target_query_params, _has_http_method_substring, PromptCachePlan.marker_count, display _diff_* colour thunks (-> _diff_ansi), pass-through RedactingFormatter.__init__. Unified: _slugify -> slugify_skill_name; reload diff -> diff_command_snapshots; _is_summary_item -> is_compaction_summary_message alias; sanitizer walkers -> _sanitize_messages/_sanitize_structure; assignment redaction passes -> _redact_assignments/_should_redact_assignment; quiet-mode tool lines -> _CUTE_LINES table.
This commit is contained in:
+92
-226
@@ -1,29 +1,10 @@
|
||||
"""Stateful scrubber for reasoning/thinking blocks in streamed assistant text.
|
||||
|
||||
``run_agent._strip_think_blocks`` is regex-based and correct for a complete
|
||||
string, but when it runs *per-delta* in ``_fire_stream_delta`` it destroys
|
||||
the state that downstream consumers (CLI ``_stream_delta``, gateway
|
||||
``GatewayStreamConsumer._filter_and_accumulate``) rely on.
|
||||
|
||||
Concretely, when MiniMax-M2.7 streams
|
||||
|
||||
delta1 = "<think>"
|
||||
delta2 = "Let me check their config"
|
||||
delta3 = "</think>"
|
||||
|
||||
the per-delta regex erases delta1 entirely (case 2: unterminated-open at
|
||||
boundary matches ``^<think>...``), so the downstream state machine never
|
||||
sees the open tag, treats delta2 as regular content, and leaks reasoning
|
||||
to the user. Consumers that don't run their own state machine (ACP,
|
||||
api_server, TTS) never had any defence at all — they just emitted
|
||||
whatever survived the upstream regex.
|
||||
|
||||
This module centralises the tag-suppression state machine at the
|
||||
upstream layer so every stream_delta_callback sees text that has
|
||||
already had reasoning blocks removed. Partial tags at delta
|
||||
boundaries are held back until the next delta resolves them, and
|
||||
end-of-stream flushing surfaces any held-back prose that turned out
|
||||
not to be a real tag.
|
||||
The regex ``run_agent._strip_think_blocks`` is correct for a complete string but,
|
||||
run per-delta, erases an opening ``<think>`` that arrives alone in one delta, so
|
||||
downstream state machines never see the open tag and leak reasoning. This class
|
||||
centralises tag suppression upstream: partial tags at delta boundaries are held
|
||||
back until resolved, and ``flush()`` releases held-back prose that was not a tag.
|
||||
|
||||
Usage::
|
||||
|
||||
@@ -33,25 +14,15 @@ Usage::
|
||||
if visible:
|
||||
emit(visible)
|
||||
tail = scrubber.flush() # at end of stream
|
||||
if tail:
|
||||
emit(tail)
|
||||
|
||||
The scrubber is re-entrant per agent instance. Call ``reset()`` at
|
||||
the top of each new turn so a hung block from an interrupted prior
|
||||
stream cannot taint the next turn's output.
|
||||
Call ``reset()`` at the top of each turn so an interrupted block cannot taint
|
||||
the next turn. Tags handled (case-insensitive): ``<think>``, ``<thinking>``,
|
||||
``<reasoning>``, ``<thought>``, ``<REASONING_SCRATCHPAD>``.
|
||||
|
||||
Tag variants handled (case-insensitive):
|
||||
``<think>``, ``<thinking>``, ``<reasoning>``, ``<thought>``,
|
||||
``<REASONING_SCRATCHPAD>``.
|
||||
|
||||
Block-boundary rule for opens: an opening tag is only treated as a
|
||||
reasoning-block opener when it appears at the start of the stream,
|
||||
after a newline (optionally followed by whitespace), or when only
|
||||
whitespace has been emitted on the current line. This prevents prose
|
||||
that *mentions* the tag name (e.g. ``"use <think> tags here"``) from
|
||||
being incorrectly suppressed. Closed pairs (``<think>X</think>``) are
|
||||
always suppressed regardless of boundary; a closed pair is an
|
||||
intentional, bounded construct.
|
||||
Boundary rule: an opening tag only starts a block at a block boundary (stream
|
||||
start, after a newline, or with only whitespace emitted on the current line), so
|
||||
prose that *mentions* ``<think>`` is not suppressed. Closed pairs
|
||||
(``<think>X</think>``) are always suppressed — a closed pair is intentional.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -64,16 +35,10 @@ __all__ = ["StreamingThinkScrubber"]
|
||||
class StreamingThinkScrubber:
|
||||
"""Stateful scrubber for streaming reasoning/thinking blocks.
|
||||
|
||||
State machine:
|
||||
- ``_in_block``: True while inside an opened block, waiting for
|
||||
a close tag. All text inside is discarded.
|
||||
- ``_buf``: held-back partial-tag tail. Emitted / discarded on
|
||||
the next ``feed()`` call or by ``flush()``.
|
||||
- ``_last_emitted_ended_newline``: True iff the most recent
|
||||
emission to the consumer ended with ``\\n``, or nothing has
|
||||
been emitted yet (start-of-stream counts as a boundary). Used
|
||||
to decide whether an open tag at buffer position 0 is at a
|
||||
block boundary.
|
||||
State: ``_in_block`` (inside an open block; text discarded), ``_buf``
|
||||
(held-back partial-tag tail), ``_last_emitted_ended_newline`` (True iff the
|
||||
last emission ended with ``\\n`` or nothing has been emitted yet — decides
|
||||
whether an open tag at buffer position 0 sits at a block boundary).
|
||||
"""
|
||||
|
||||
_OPEN_TAG_NAMES: Tuple[str, ...] = (
|
||||
@@ -84,31 +49,33 @@ class StreamingThinkScrubber:
|
||||
"REASONING_SCRATCHPAD",
|
||||
)
|
||||
|
||||
# Materialise literal tag strings so the hot path does string
|
||||
# operations, not regex compilation per feed().
|
||||
# Literal tag strings so the hot path does string ops, not regex per feed().
|
||||
_OPEN_TAGS: Tuple[str, ...] = tuple(f"<{name}>" for name in _OPEN_TAG_NAMES)
|
||||
_CLOSE_TAGS: Tuple[str, ...] = tuple(f"</{name}>" for name in _OPEN_TAG_NAMES)
|
||||
|
||||
# Pre-compute the longest tag (for partial-tag hold-back bound).
|
||||
_MAX_TAG_LEN: int = max(len(tag) for tag in _OPEN_TAGS + _CLOSE_TAGS)
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.reset()
|
||||
|
||||
def reset(self) -> None:
|
||||
"""Reset all state. Call at the top of every new turn."""
|
||||
self._in_block: bool = False
|
||||
self._buf: str = ""
|
||||
self._last_emitted_ended_newline: bool = True
|
||||
|
||||
def reset(self) -> None:
|
||||
"""Reset all state. Call at the top of every new turn."""
|
||||
self._in_block = False
|
||||
self._buf = ""
|
||||
self._last_emitted_ended_newline = True
|
||||
def _emit(self, out: list[str], text: str) -> None:
|
||||
"""Append visible prose to *out* (orphan close tags stripped) and track the newline flag."""
|
||||
if text:
|
||||
text = self._strip_orphan_close_tags(text)
|
||||
if text:
|
||||
out.append(text)
|
||||
self._last_emitted_ended_newline = text.endswith("\n")
|
||||
|
||||
def feed(self, text: str) -> str:
|
||||
"""Feed one delta; return the scrubbed visible portion.
|
||||
|
||||
May return an empty string when the entire delta is reasoning
|
||||
content or is being held back pending resolution of a partial
|
||||
tag at the boundary.
|
||||
Returns "" when the whole delta is reasoning content or is held back
|
||||
pending resolution of a partial tag at the boundary.
|
||||
"""
|
||||
if not text:
|
||||
return ""
|
||||
@@ -118,130 +85,67 @@ class StreamingThinkScrubber:
|
||||
|
||||
while buf:
|
||||
if self._in_block:
|
||||
# Hunt for the earliest close tag.
|
||||
close_idx, close_len = self._find_first_tag(
|
||||
buf, self._CLOSE_TAGS,
|
||||
)
|
||||
close_idx, close_len = self._find_first_tag(buf, self._CLOSE_TAGS)
|
||||
if close_idx == -1:
|
||||
# No close yet — hold back a potential partial
|
||||
# close-tag prefix; discard everything else.
|
||||
# No close yet: hold back a possible partial close-tag prefix, drop the rest.
|
||||
held = self._max_partial_suffix(buf, self._CLOSE_TAGS)
|
||||
self._buf = buf[-held:] if held else ""
|
||||
return "".join(out)
|
||||
# Found close: discard block content + tag, continue.
|
||||
buf = buf[close_idx + close_len:]
|
||||
self._in_block = False
|
||||
continue
|
||||
|
||||
# Priority 1: closed <tag>X</tag> pair anywhere (no boundary gating —
|
||||
# even inline pairs are almost certainly leaked reasoning).
|
||||
# Priority 2: unterminated open tag at a block boundary (gated so
|
||||
# prose that mentions '<think>' isn't over-stripped). Earliest wins.
|
||||
pair = self._find_earliest_closed_pair(buf)
|
||||
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]])
|
||||
buf = buf[pair[1]:]
|
||||
continue
|
||||
if open_idx != -1:
|
||||
self._emit(out, buf[:open_idx])
|
||||
self._in_block = True
|
||||
buf = buf[open_idx + open_len:]
|
||||
continue
|
||||
|
||||
# No resolvable tag: hold back any partial-tag prefix at the tail
|
||||
# so a tag split across deltas isn't missed, then emit the rest.
|
||||
held = max(
|
||||
self._max_partial_suffix(buf, self._OPEN_TAGS),
|
||||
self._max_partial_suffix(buf, self._CLOSE_TAGS),
|
||||
)
|
||||
if held:
|
||||
self._emit(out, buf[:-held])
|
||||
self._buf = buf[-held:]
|
||||
else:
|
||||
# Priority 1 — closed <tag>X</tag> pair anywhere in
|
||||
# buf. Closed pairs are always an intentional,
|
||||
# bounded construct (even mid-line prose containing
|
||||
# an open/close pair is almost certainly a model
|
||||
# leaking reasoning inline), so no boundary gating.
|
||||
pair = self._find_earliest_closed_pair(buf)
|
||||
# Priority 2 — unterminated open tag at a block
|
||||
# boundary. Boundary-gated so prose that mentions
|
||||
# '<think>' isn't over-stripped.
|
||||
open_idx, open_len = self._find_open_at_boundary(
|
||||
buf, out,
|
||||
)
|
||||
|
||||
# Pick whichever match comes earliest in the buffer.
|
||||
if pair is not None and (
|
||||
open_idx == -1 or pair[0] <= open_idx
|
||||
):
|
||||
start_idx, end_idx = pair
|
||||
preceding = buf[:start_idx]
|
||||
if preceding:
|
||||
preceding = self._strip_orphan_close_tags(preceding)
|
||||
if preceding:
|
||||
out.append(preceding)
|
||||
self._last_emitted_ended_newline = (
|
||||
preceding.endswith("\n")
|
||||
)
|
||||
buf = buf[end_idx:]
|
||||
continue
|
||||
|
||||
if open_idx != -1:
|
||||
# Unterminated open at boundary — emit preceding,
|
||||
# enter block, continue loop with remainder.
|
||||
preceding = buf[:open_idx]
|
||||
if preceding:
|
||||
preceding = self._strip_orphan_close_tags(preceding)
|
||||
if preceding:
|
||||
out.append(preceding)
|
||||
self._last_emitted_ended_newline = (
|
||||
preceding.endswith("\n")
|
||||
)
|
||||
self._in_block = True
|
||||
buf = buf[open_idx + open_len:]
|
||||
continue
|
||||
|
||||
# No resolvable tag structure in buf. Hold back any
|
||||
# partial-tag prefix at the tail so a split tag
|
||||
# across deltas isn't missed, then emit the rest.
|
||||
held = self._max_partial_suffix(buf, self._OPEN_TAGS)
|
||||
held_close = self._max_partial_suffix(
|
||||
buf, self._CLOSE_TAGS,
|
||||
)
|
||||
held = max(held, held_close)
|
||||
if held:
|
||||
emit_text = buf[:-held]
|
||||
self._buf = buf[-held:]
|
||||
else:
|
||||
emit_text = buf
|
||||
self._buf = ""
|
||||
if emit_text:
|
||||
emit_text = self._strip_orphan_close_tags(emit_text)
|
||||
if emit_text:
|
||||
out.append(emit_text)
|
||||
self._last_emitted_ended_newline = (
|
||||
emit_text.endswith("\n")
|
||||
)
|
||||
return "".join(out)
|
||||
self._emit(out, buf)
|
||||
return "".join(out)
|
||||
|
||||
return "".join(out)
|
||||
|
||||
def flush(self) -> str:
|
||||
"""End-of-stream flush.
|
||||
|
||||
If still inside an unterminated block, held-back content is
|
||||
discarded — leaking partial reasoning is worse than a
|
||||
truncated answer. Otherwise the held-back partial-tag tail is
|
||||
emitted verbatim (it turned out not to be a real tag prefix).
|
||||
|
||||
Always treats the next ``feed()`` as a fresh stream boundary.
|
||||
Intra-turn retries (thinking-only prefill, empty-response
|
||||
retry) flush then stream again without calling ``reset()``;
|
||||
leaving ``_last_emitted_ended_newline`` False made a new
|
||||
stream's opening ``<think>`` look mid-line and leak into the
|
||||
visible reply.
|
||||
Inside an unterminated block the held-back content is discarded (leaking
|
||||
partial reasoning is worse than a truncated answer); otherwise the
|
||||
held-back 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._buf = ""
|
||||
self._in_block = False
|
||||
# Next feed() is a new stream — start-of-stream is a boundary.
|
||||
self._last_emitted_ended_newline = True
|
||||
return ""
|
||||
tail = self._buf
|
||||
tail = "" if self._in_block else self._buf
|
||||
self._buf = ""
|
||||
# Same for the non-block path: do NOT derive the boundary flag
|
||||
# from the flushed tail (e.g. a held-back '<'). End-of-stream
|
||||
# means the next feed() starts a new model response.
|
||||
self._in_block = False
|
||||
self._last_emitted_ended_newline = True
|
||||
if not tail:
|
||||
return ""
|
||||
return self._strip_orphan_close_tags(tail)
|
||||
return self._strip_orphan_close_tags(tail) if tail else ""
|
||||
|
||||
# ── internal helpers ───────────────────────────────────────────────
|
||||
|
||||
@staticmethod
|
||||
def _find_first_tag(
|
||||
buf: str, tags: Tuple[str, ...],
|
||||
) -> Tuple[int, int]:
|
||||
"""Return (earliest_index, tag_length) over *tags*, or (-1, 0).
|
||||
|
||||
Case-insensitive match.
|
||||
"""
|
||||
def _find_first_tag(buf: str, tags: Tuple[str, ...]) -> Tuple[int, int]:
|
||||
"""Return (earliest_index, tag_length) over *tags* (case-insensitive), or (-1, 0)."""
|
||||
buf_lower = buf.lower()
|
||||
best_idx = -1
|
||||
best_len = 0
|
||||
@@ -253,14 +157,10 @@ class StreamingThinkScrubber:
|
||||
return best_idx, best_len
|
||||
|
||||
def _find_earliest_closed_pair(self, buf: str):
|
||||
"""Return (start_idx, end_idx) of the earliest closed pair, else None.
|
||||
"""Return (start_idx, end_idx) of the earliest ``<tag>...</tag>`` pair, else None.
|
||||
|
||||
A closed pair is ``<tag>...</tag>`` of any variant. Matches are
|
||||
case-insensitive and non-greedy (the closest close tag after
|
||||
an open tag wins), matching the regex ``<tag>.*?</tag>``
|
||||
semantics of ``_strip_think_blocks`` case 1. When two tag
|
||||
variants could both match, the one whose open tag appears
|
||||
earlier wins.
|
||||
Case-insensitive and non-greedy (closest close after the open wins),
|
||||
matching ``_strip_think_blocks`` case 1; the earliest open tag wins.
|
||||
"""
|
||||
buf_lower = buf.lower()
|
||||
best: "tuple[int, int] | None" = None
|
||||
@@ -270,23 +170,15 @@ class StreamingThinkScrubber:
|
||||
open_idx = buf_lower.find(open_lower)
|
||||
if open_idx == -1:
|
||||
continue
|
||||
close_idx = buf_lower.find(
|
||||
close_lower, open_idx + len(open_lower),
|
||||
)
|
||||
close_idx = buf_lower.find(close_lower, open_idx + len(open_lower))
|
||||
if close_idx == -1:
|
||||
continue
|
||||
end_idx = close_idx + len(close_lower)
|
||||
if best is None or open_idx < best[0]:
|
||||
best = (open_idx, end_idx)
|
||||
best = (open_idx, close_idx + len(close_lower))
|
||||
return best
|
||||
|
||||
def _find_open_at_boundary(
|
||||
self, buf: str, already_emitted: list[str],
|
||||
) -> Tuple[int, int]:
|
||||
"""Return the earliest block-boundary open-tag (idx, len).
|
||||
|
||||
Returns (-1, 0) if no boundary-legal opener is present.
|
||||
"""
|
||||
def _find_open_at_boundary(self, buf: str, already_emitted: list[str]) -> Tuple[int, int]:
|
||||
"""Return the earliest block-boundary open-tag (idx, len), or (-1, 0)."""
|
||||
buf_lower = buf.lower()
|
||||
best_idx = -1
|
||||
best_len = 0
|
||||
@@ -305,50 +197,30 @@ class StreamingThinkScrubber:
|
||||
search_start = idx + 1
|
||||
return best_idx, best_len
|
||||
|
||||
def _is_block_boundary(
|
||||
self, buf: str, idx: int, already_emitted: list[str],
|
||||
) -> bool:
|
||||
def _is_block_boundary(self, buf: str, idx: int, already_emitted: list[str]) -> bool:
|
||||
"""True iff position *idx* in *buf* is a block boundary.
|
||||
|
||||
A block boundary is:
|
||||
- buf position 0 AND the most recent emission ended with
|
||||
a newline (or nothing has been emitted yet)
|
||||
- any position whose preceding text on the current line
|
||||
(since the last newline in buf) is whitespace-only, AND
|
||||
if there is no newline in the preceding buf portion, the
|
||||
most recent prior emission ended with a newline
|
||||
Boundary = position 0 with the prior emission ending in a newline (or
|
||||
nothing emitted yet), or any position whose preceding text on the current
|
||||
line is whitespace-only (when no newline precedes it in *buf*, the prior
|
||||
emission must also have ended with a newline).
|
||||
"""
|
||||
prior_newline = (
|
||||
already_emitted[-1].endswith("\n") if already_emitted else self._last_emitted_ended_newline
|
||||
)
|
||||
if idx == 0:
|
||||
# Check whether the last already-emitted chunk in THIS
|
||||
# feed() call ended with a newline, otherwise fall back
|
||||
# to the cross-feed flag.
|
||||
if already_emitted:
|
||||
return already_emitted[-1].endswith("\n")
|
||||
return self._last_emitted_ended_newline
|
||||
return prior_newline
|
||||
preceding = buf[:idx]
|
||||
last_nl = preceding.rfind("\n")
|
||||
if last_nl == -1:
|
||||
# No newline in buf before the tag — boundary only if the
|
||||
# prior emission ended with a newline AND everything since
|
||||
# is whitespace.
|
||||
if already_emitted:
|
||||
prior_newline = already_emitted[-1].endswith("\n")
|
||||
else:
|
||||
prior_newline = self._last_emitted_ended_newline
|
||||
return prior_newline and preceding.strip() == ""
|
||||
# Newline present — text between it and the tag must be
|
||||
# whitespace-only.
|
||||
return preceding[last_nl + 1:].strip() == ""
|
||||
|
||||
@classmethod
|
||||
def _max_partial_suffix(
|
||||
cls, buf: str, tags: Tuple[str, ...],
|
||||
) -> int:
|
||||
"""Return the longest buf-suffix that is a prefix of any tag.
|
||||
def _max_partial_suffix(cls, buf: str, tags: Tuple[str, ...]) -> int:
|
||||
"""Longest buf-suffix that is a strict prefix of any tag (case-insensitive).
|
||||
|
||||
Only prefixes strictly shorter than the tag itself count
|
||||
(full-length suffixes are the tag and are handled as matches,
|
||||
not held-back partials). Case-insensitive.
|
||||
Full-length matches are real tags handled elsewhere, not held-back partials.
|
||||
"""
|
||||
if not buf:
|
||||
return 0
|
||||
@@ -364,12 +236,7 @@ class StreamingThinkScrubber:
|
||||
|
||||
@classmethod
|
||||
def _strip_orphan_close_tags(cls, text: str) -> str:
|
||||
"""Remove any close tags from *text* (orphan-close handling).
|
||||
|
||||
An orphan close tag has no matching open in the current
|
||||
scrubber state; it's always noise, stripped with any trailing
|
||||
whitespace so the surrounding prose flows naturally.
|
||||
"""
|
||||
"""Remove close tags with no matching open (always noise) plus trailing whitespace."""
|
||||
if "</" not in text:
|
||||
return text
|
||||
text_lower = text.lower()
|
||||
@@ -382,8 +249,7 @@ class StreamingThinkScrubber:
|
||||
tag_lower = tag.lower()
|
||||
tag_len = len(tag_lower)
|
||||
if text_lower[i:i + tag_len] == tag_lower:
|
||||
# Skip the tag and any trailing whitespace,
|
||||
# matching _strip_think_blocks case 3.
|
||||
# Skip the tag and trailing whitespace (matches _strip_think_blocks case 3).
|
||||
j = i + tag_len
|
||||
while j < len(text) and text[j] in " \t\n\r":
|
||||
j += 1
|
||||
|
||||
Reference in New Issue
Block a user