§
    ÷žyjÄ>  ã                  óB   — d Z ddlmZ ddlmZ dgZ G d„ d¦  «        ZdS )uí  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.

Usage::

    scrubber = StreamingThinkScrubber()
    for delta in stream:
        visible = scrubber.feed(delta)
        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.

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.
é    )Úannotations)ÚTupleÚStreamingThinkScrubberc                  óD  — e Zd ZU dZdZded<    ed„ eD ¦   «         ¦  «        Zded<    ed„ eD ¦   «         ¦  «        Zded<    e	d	„ eez   D ¦   «         ¦  «        Z
d
ed<   d"d„Zd"d„Zd#d„Zd$d„Zed%d„¦   «         Zd&d„Zd'd„Zd(d„Zed)d„¦   «         Zed#d „¦   «         Zd!S )*r   an  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.
    )ÚthinkÚthinkingÚ	reasoningÚthoughtÚREASONING_SCRATCHPADúTuple[str, ...]Ú_OPEN_TAG_NAMESc              #  ó"   K  — | ]
}d |› d�V — ŒdS )Ú<Ú>N© ©Ú.0Únames     ú:/home/ragecks/.hermes/hermes-agent/agent/think_scrubber.pyú	<genexpr>z StreamingThinkScrubber.<genexpr>Y   s*   è è € Ð'PÐ'P¸¨¨D¨¨¨Ð'PÐ'PÐ'PÐ'PÐ'PÐ'Pó    Ú
_OPEN_TAGSc              #  ó"   K  — | ]
}d |› d�V — ŒdS )ú</r   Nr   r   s     r   r   z StreamingThinkScrubber.<genexpr>Z   s*   è è € Ð(RÐ(R¸$¨¨d¨¨¨Ð(RÐ(RÐ(RÐ(RÐ(RÐ(Rr   Ú_CLOSE_TAGSc              #  ó4   K  — | ]}t          |¦  «        V — Œd S )N)Úlen)r   Útags     r   r   z StreamingThinkScrubber.<genexpr>]   s(   è è € ÐIÐI¨�C ™HœHÐIÐIÐIÐIÐIÐIr   ÚintÚ_MAX_TAG_LENÚreturnÚNonec                ó0   — d| _         d| _        d| _        d S )NFÚ T©Ú	_in_blockÚ_bufÚ_last_emitted_ended_newline©Úselfs    r   Ú__init__zStreamingThinkScrubber.__init___   s   € Ø$ˆŒØˆŒ	Ø15ˆÔ(Ð(Ð(r   c                ó0   — d| _         d| _        d| _        dS )z4Reset all state.  Call at the top of every new turn.Fr$   TNr%   r)   s    r   ÚresetzStreamingThinkScrubber.resetd   s   € àˆŒØˆŒ	Ø+/ˆÔ(Ð(Ð(r   ÚtextÚstrc                ó*  — |sdS | j         |z   }d| _         g }|�re| j        r~|                      || j        ¦  «        \  }}|dk    rD|                      || j        ¦  «        }|r|| d…         nd| _         d                     |¦  «        S |||z   d…         }d| _        �nÝ|                      |¦  «        }|                      ||¦  «        \  }}	|�u|dk    s|d         |k    rc|\  }
}|d|
…         }|rF|                      |¦  «        }|r/| 	                    |¦  «         | 
                    d¦  «        | _        ||d…         }�Œ-|dk    rh|d|…         }|rF|                      |¦  «        }|r/| 	                    |¦  «         | 
                    d¦  «        | _        d| _        |||	z   d…         }�Œ›|                      || j        ¦  «        }|                      || j        ¦  «        }t          ||¦  «        }|r|d| …         }|| d…         | _         n	|}d| _         |rF|                      |¦  «        }|r/| 	                    |¦  «         | 
                    d¦  «        | _        d                     |¦  «        S |�°ed                     |¦  «        S )zç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.
        r$   éÿÿÿÿNFr   Ú
T)r'   r&   Ú_find_first_tagr   Ú_max_partial_suffixÚjoinÚ_find_earliest_closed_pairÚ_find_open_at_boundaryÚ_strip_orphan_close_tagsÚappendÚendswithr(   r   Úmax)r*   r.   ÚbufÚoutÚ	close_idxÚ	close_lenÚheldÚpairÚopen_idxÚopen_lenÚ	start_idxÚend_idxÚ	precedingÚ
held_closeÚ	emit_texts                  r   ÚfeedzStreamingThinkScrubber.feedj   s  € ð ð 	Ø�2ØŒi˜$ÑˆØˆŒ	Øˆàñ Q	$ØŒ~ð P$à'+×';Ò';Ø˜Ô)ñ(ô (Ñ$�	˜9ð  ’?�?ð  ×3Ò3°C¸Ô9IÑJÔJ�DØ/3Ð ;  T E F F¤ ¸�D”IØŸ7š7 3™<œ<Ð'à˜) iÑ/Ð0Ð0Ô1�Ø!&�”‘ð ×6Ò6°sÑ;Ô;�ð &*×%@Ò%@Ø˜ñ&ô &Ñ"�˜(ð
 Ð#Ø ’N�N d¨1¤g°Ò&9Ð&9à)-Ñ&�I˜wØ # J Y J¤�IØ ð Ø$(×$AÒ$AÀ)Ñ$LÔ$L˜	Ø$ð ØŸJšJ yÑ1Ô1Ð1à )× 2Ò 2°4Ñ 8Ô 8ð !Ô<ð ˜g˜h˜hœ-�CÙà˜r’>�>ð !$ I X I¤�IØ ð Ø$(×$AÒ$AÀ)Ñ$LÔ$L˜	Ø$ð ØŸJšJ yÑ1Ô1Ð1à )× 2Ò 2°4Ñ 8Ô 8ð !Ô<ð &*�D”NØ˜h¨Ñ1Ð2Ð2Ô3�CÙð
 ×/Ò/°°T´_ÑEÔE�Ø!×5Ò5Ø˜Ô)ñô �
õ ˜4 Ñ,Ô,�Øð #Ø # F d U F¤�IØ # T E F F¤�D”I�Ià #�IØ "�D”IØð Ø $× =Ò =¸iÑ HÔ H�IØ ð ØŸ
š
 9Ñ-Ô-Ð-à%×.Ò.¨tÑ4Ô4ð Ô8ð —w’w˜s‘|”|Ð#ðc ñ Q	$ðf �wŠw�s‰|Œ|Ðr   c                óš   — | j         rd| _        d| _         d| _        dS | j        }d| _        d| _        |sdS |                      |¦  «        S )u™  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.
        r$   FT)r&   r'   r(   r8   )r*   Útails     r   ÚflushzStreamingThinkScrubber.flushÌ   sb   € ð Œ>ð 	ØˆDŒIØ"ˆDŒNà/3ˆDÔ,Ø�2ØŒyˆØˆŒ	ð ,0ˆÔ(Øð 	Ø�2Ø×,Ò,¨TÑ2Ô2Ð2r   r<   ÚtagsúTuple[int, int]c                óØ   — |                       ¦   «         }d}d}|D ]L}|                     |                      ¦   «         ¦  «        }|dk    r|dk    s||k     r|}t          |¦  «        }ŒM||fS )zfReturn (earliest_index, tag_length) over *tags*, or (-1, 0).

        Case-insensitive match.
        r1   r   )ÚlowerÚfindr   )r<   rM   Ú	buf_lowerÚbest_idxÚbest_lenr   Úidxs          r   r3   z&StreamingThinkScrubber._find_first_tagí   sx   € ð —I’I‘K”Kˆ	ØˆØˆØð 	$ð 	$ˆCØ—.’. §¢¡¤Ñ-Ô-ˆCØ�bŠyˆy˜h¨"šn˜n°°h²°Ø�Ý˜s™8œ8�øØ˜Ð!Ð!r   c                óœ  — |                      ¦   «         }d}t          | j        | j        ¦  «        D ]š\  }}|                      ¦   «         }|                      ¦   «         }|                     |¦  «        }|dk    rŒI|                     ||t          |¦  «        z   ¦  «        }	|	dk    rŒv|	t          |¦  «        z   }
|�||d         k     r||
f}Œ›|S )a·  Return (start_idx, end_idx) of the earliest closed 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.
        Nr1   r   )rP   Úzipr   r   rQ   r   )r*   r<   rR   ÚbestÚopen_tagÚ	close_tagÚ
open_lowerÚclose_lowerrB   r>   rE   s              r   r6   z1StreamingThinkScrubber._find_earliest_closed_pairÿ   sÚ   € ð —I’I‘K”Kˆ	Ø)-ˆÝ#& t¤¸Ô8HÑ#IÔ#Ið 	+ð 	+ÑˆH�iØ!ŸšÑ)Ô)ˆJØ#Ÿ/š/Ñ+Ô+ˆKØ —~’~ jÑ1Ô1ˆHØ˜2Š~ˆ~ØØ!ŸšØ˜X­¨J©¬Ñ7ñô ˆIð ˜BŠˆØØ¥# kÑ"2Ô"2Ñ2ˆGØˆ|˜x¨$¨q¬'Ò1Ð1Ø  'Ð*�øØˆr   Úalready_emittedú	list[str]c                ó,  — |                      ¦   «         }d}d}| j        D ]q}|                      ¦   «         }d}	 |                     ||¦  «        }	|	dk    rn;|                      ||	|¦  «        r|dk    s|	|k     r|	}t	          |¦  «        }n|	dz   }ŒXŒr||fS )z�Return the earliest block-boundary open-tag (idx, len).

        Returns (-1, 0) if no boundary-legal opener is present.
        r1   r   Té   )rP   r   rQ   Ú_is_block_boundaryr   )
r*   r<   r]   rR   rS   rT   r   Ú	tag_lowerÚsearch_startrU   s
             r   r7   z-StreamingThinkScrubber._find_open_at_boundary  s¼   € ð —I’I‘K”Kˆ	ØˆØˆØ”?ð 	'ð 	'ˆCØŸ	š	™œˆIØˆLð	'Ø—n’n Y°Ñ=Ô=�Ø˜"’9�9ØØ×*Ò*¨3°°_ÑEÔEð Ø 2’~�~¨¨xª¨Ø#&˜Ý#& s¡8¤8˜ØØ" Q™w�ð	'øð ˜Ð!Ð!r   rU   Úboolc                ód  — |dk    r$|r|d                               d¦  «        S | j        S |d|…         }|                     d¦  «        }|dk    r?|r|d                               d¦  «        }n| j        }|o|                     ¦   «         dk    S ||dz   d…                              ¦   «         dk    S )aÞ  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
        r   r1   r2   Nr$   r`   )r:   r(   ÚrfindÚstrip)r*   r<   rU   r]   rF   Úlast_nlÚprior_newlines          r   ra   z)StreamingThinkScrubber._is_block_boundary4  sÍ   € ð �!Š8ˆ8ð ð :Ø& rÔ*×3Ò3°DÑ9Ô9Ð9ØÔ3Ð3Ø˜˜˜”Iˆ	Ø—/’/ $Ñ'Ô'ˆØ�bŠ=ˆ=ð ð AØ /°Ô 3× <Ò <¸TÑ BÔ B��à $Ô @�Ø Ð< Y§_¢_Ñ%6Ô%6¸"Ò%<Ð<ð ˜ 1™˜˜Ô&×,Ò,Ñ.Ô.°"Ò4Ð4r   c                óL  — |sdS |                      ¦   «         }t          t          |¦  «        | j        dz
  ¦  «        }t	          |dd¦  «        D ]T}|| d…         }|D ]D}|                      ¦   «         }t          |¦  «        |k    r|                     |¦  «        r|c c S ŒEŒUdS )zÿReturn the longest buf-suffix that is a prefix of any tag.

        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.
        r   r`   r1   N)rP   Úminr   r    ÚrangeÚ
startswith)	Úclsr<   rM   rR   Ú	max_checkÚiÚsuffixr   rb   s	            r   r4   z*StreamingThinkScrubber._max_partial_suffixW  s¿   € ð ð 	Ø�1Ø—I’I‘K”Kˆ	Ý�˜I™œ¨Ô(8¸1Ñ(<Ñ=Ô=ˆ	Ý�y ! RÑ(Ô(ð 	ð 	ˆAØ ˜r˜s˜s”^ˆFØð ð �ØŸIšI™KœK�	Ý�y‘>”> AÒ%Ð%¨)×*>Ò*>¸vÑ*FÔ*FÐ%Ø�H�H�H�H�Høðð ˆqr   c                ó.  — d|vr|S |                      ¦   «         }g }d}|t          |¦  «        k     rÐd}|||dz   …         dk    rˆ| j        D ]€}|                      ¦   «         }t          |¦  «        }||||z   …         |k    rJ||z   }	|	t          |¦  «        k     r,||	         dv r"|	dz  }	|	t          |¦  «        k     r
||	         dv °"|	}d} nŒ�|s |                     ||         ¦  «         |dz  }|t          |¦  «        k     °Ðd                     |¦  «        S )	a  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.
        r   r   Fé   z 	
r`   Tr$   )rP   r   r   r9   r5   )
rn   r.   Ú
text_lowerr=   rp   Úmatchedr   rb   Útag_lenÚjs
             r   r8   z/StreamingThinkScrubber._strip_orphan_close_tagsm  sH  € ð �tÐÐØˆKØ—Z’Z‘\”\ˆ
ØˆØˆØ•#�d‘)”)ŠmˆmØˆGØ˜!˜A ™E˜'Ô" dÒ*Ð*Øœ?ð ð �CØ #§	¢	¡¤�IÝ! )™nœn�GØ! ! A¨¡K -Ô0°IÒ=Ð=ð  ™K˜Ø¥# d¡)¤)šm˜m°°Q´¸9Ð0DÐ0DØ ™F˜Að  ¥# d¡)¤)šm˜m°°Q´¸9Ð0DÐ0Dà˜Ø"&˜Ø˜ð >ð ð Ø—
’
˜4 œ7Ñ#Ô#Ð#Ø�Q‘�ð# •#�d‘)”)Šmˆmð$ �wŠw�s‰|Œ|Ðr   N)r!   r"   )r.   r/   r!   r/   )r!   r/   )r<   r/   rM   r   r!   rN   )r<   r/   )r<   r/   r]   r^   r!   rN   )r<   r/   rU   r   r]   r^   r!   rd   )r<   r/   rM   r   r!   r   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   Ú__annotations__Útupler   r   r;   r    r+   r-   rI   rL   Ústaticmethodr3   r6   r7   ra   Úclassmethodr4   r8   r   r   r   r   r   @   sµ  € € € € € € ðð ð(€Oð ð ð ñ ð #( %Ð'PÐ'PÀÐ'PÑ'PÔ'PÑ"PÔ"P€JÐPÐPÐPÑPØ#( 5Ð(RÐ(RÀ/Ð(RÑ(RÔ(RÑ#RÔ#R€KÐRÐRÐRÑRð ˜ÐIÐI°
¸[Ñ0HÐIÑIÔIÑIÔI€LÐIÐIÐIÑIð6ð 6ð 6ð 6ð
0ð 0ð 0ð 0ð`ð `ð `ð `ðD3ð 3ð 3ð 3ðB ð"ð "ð "ñ „\ð"ð"ð ð ð ð8"ð "ð "ð "ð2!5ð !5ð !5ð !5ðF ðð ð ñ „[ðð* ðð ð ñ „[ðð ð r   N)r{   Ú
__future__r   Útypingr   Ú__all__r   r   r   r   ú<module>rƒ      sz   ðð6ð 6ðp #Ð "Ð "Ð "Ð "Ð "à Ð Ð Ð Ð Ð à#Ð
$€ðLð Lð Lð Lð Lñ Lô Lð Lð Lð Lr   