91 lines
3.3 KiB
Python
91 lines
3.3 KiB
Python
"""Safe helpers for provider text streams.
|
|
|
|
Some OpenAI-compatible providers place reasoning in the normal ``content``
|
|
stream as ``<think>...</think>`` instead of using LangChain's dedicated
|
|
reasoning fields. Those fragments must be separated from report/file content
|
|
while callers can still surface them in an explicitly marked live UI trace.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
class InlineThinkTagFilter:
|
|
"""Separate inline ``<think>`` blocks while preserving token streaming.
|
|
|
|
Tags can be split across arbitrary provider chunks, so regex-per-chunk is
|
|
not safe: it can leak a partial tag or the model's reasoning into the
|
|
visible draft. This small state machine holds only the suffix that might
|
|
still form a tag and emits all confirmed answer text immediately.
|
|
"""
|
|
|
|
_OPEN_TAG = "<think>"
|
|
_CLOSE_TAG = "</think>"
|
|
|
|
def __init__(self) -> None:
|
|
self._inside_think = False
|
|
self._pending = ""
|
|
|
|
@staticmethod
|
|
def _possible_tag_prefix_size(content: str, tag: str) -> int:
|
|
"""Return the longest suffix which may become the next tag."""
|
|
lowered = content.lower()
|
|
for size in range(min(len(content), len(tag) - 1), 0, -1):
|
|
if lowered[-size:] == tag[:size]:
|
|
return size
|
|
return 0
|
|
|
|
def push_parts(self, chunk: str) -> tuple[str, str]:
|
|
"""Return ``(visible_text, thinking_text)`` for one provider chunk."""
|
|
if not chunk:
|
|
return "", ""
|
|
|
|
remaining = self._pending + chunk
|
|
self._pending = ""
|
|
visible: list[str] = []
|
|
thinking: list[str] = []
|
|
|
|
while remaining:
|
|
tag = self._CLOSE_TAG if self._inside_think else self._OPEN_TAG
|
|
marker_at = remaining.lower().find(tag)
|
|
if marker_at >= 0:
|
|
(thinking if self._inside_think else visible).append(remaining[:marker_at])
|
|
remaining = remaining[marker_at + len(tag) :]
|
|
self._inside_think = not self._inside_think
|
|
continue
|
|
|
|
suffix_size = self._possible_tag_prefix_size(remaining, tag)
|
|
confirmed = remaining[:-suffix_size] if suffix_size else remaining
|
|
(thinking if self._inside_think else visible).append(confirmed)
|
|
self._pending = remaining[-suffix_size:] if suffix_size else ""
|
|
break
|
|
|
|
return "".join(visible), "".join(thinking)
|
|
|
|
def push(self, chunk: str) -> str:
|
|
"""Compatibility helper returning only confirmed visible text."""
|
|
return self.push_parts(chunk)[0]
|
|
|
|
def finish_parts(self) -> tuple[str, str]:
|
|
"""Flush the remaining visible or thinking text at end of stream."""
|
|
if self._inside_think:
|
|
thinking = self._pending
|
|
self._pending = ""
|
|
self._inside_think = False
|
|
return "", thinking
|
|
tail = self._pending
|
|
self._pending = ""
|
|
return tail, ""
|
|
|
|
def finish(self) -> str:
|
|
"""Compatibility helper returning only the visible stream tail."""
|
|
return self.finish_parts()[0]
|
|
|
|
|
|
def strip_inline_think_tags(content: str) -> str:
|
|
"""Return complete text with inline model-thinking blocks removed."""
|
|
text_filter = InlineThinkTagFilter()
|
|
return text_filter.push(content) + text_filter.finish()
|
|
|
|
|
|
__all__ = ["InlineThinkTagFilter", "strip_inline_think_tags"]
|