refactor(streaming): share interim visible/streamed normalization between prefix and exact verdicts

Both interim predicates now compare one (visible, streamed) pair from
_interim_visible_and_streamed so their normalization cannot drift; the two
#88954 codex tests collapse into one parametrized test.
This commit is contained in:
kshitijk4poor
2026-09-24 19:22:41 +05:30
committed by kshitij
parent 1f61f7087b
commit c43eb76d86
2 changed files with 28 additions and 64 deletions

View File

@@ -103,35 +103,23 @@ class StreamDeliveryMixin:
def _normalize_interim_visible_text(text: str) -> str:
return re.sub(r"\s+", " ", text).strip() if isinstance(text, str) else ""
def _interim_visible_and_streamed(self, content: str) -> tuple[str, str]:
"""(visible content, streamed text), both think-stripped and whitespace-normalized."""
normalize = self._normalize_interim_visible_text
streamed = getattr(self, "_current_streamed_assistant_text", "") or ""
return normalize(self._strip_think_blocks(content or "")), normalize(self._strip_think_blocks(streamed))
def _interim_content_was_streamed(self, content: str) -> bool:
visible_content = self._normalize_interim_visible_text(self._strip_think_blocks(content or ""))
streamed = self._normalize_interim_visible_text(
self._strip_think_blocks(getattr(self, "_current_streamed_assistant_text", "") or "")
)
# Prefix match, not equality: the final may be streamed text plus a trailing delta. The
# reverse (streamed longer) is NOT matched — it could suppress a needed resend.
return bool(visible_content and streamed) and visible_content.startswith(streamed)
visible, streamed = self._interim_visible_and_streamed(content)
return bool(visible and streamed) and visible.startswith(streamed)
def _interim_content_fully_streamed(self, content: str) -> bool:
"""Exact-equality variant for the gateway interim-message decision.
``_interim_content_was_streamed`` keeps its prefix semantics for the conversation-loop
"previewed" marks (the streamed prefix IS on the user's screen there). The gateway interim
path is different: a True verdict makes ``_interim_assistant_cb`` call
``on_segment_break()``, which finalizes the streaming bubble as-is — anything after the
streamed prefix is never delivered. On Telegram the stream is routinely truncated at the
text→tool_calls boundary, so the prefix match marked a truncated bubble complete and the
tail was lost (#88954). Only an exact match may skip the full-text resend; a partial prefix
falls through to ``on_commentary`` and re-delivers the complete text (benign duplicate,
never lost text).
"""
visible_content = self._normalize_interim_visible_text(self._strip_think_blocks(content or ""))
if not visible_content:
return False
streamed = self._normalize_interim_visible_text(
self._strip_think_blocks(getattr(self, "_current_streamed_assistant_text", "") or "")
)
return bool(streamed) and streamed == visible_content
"""Exact-match variant for the gateway interim path: a True verdict finalizes the bubble
as-is, so a truncated prefix must fall through to a full-text resend (#88954)."""
visible, streamed = self._interim_visible_and_streamed(content)
return bool(visible) and visible == streamed
def _extract_codex_interim_visible_parts(self, assistant_msg: Dict[str, Any]) -> List[str]:
"""Visible Codex commentary (``phase=commentary`` items), one string per message item.

View File

@@ -1900,23 +1900,27 @@ def test_interim_content_was_streamed_matches_prefix_not_exact(monkeypatch):
assert agent._interim_content_was_streamed("hello") is False
def test_interim_commentary_partial_stream_receives_full_text(monkeypatch):
"""A stream truncated at the text→tool_calls boundary must not be marked
already_streamed (#88954).
On Telegram the streamed commentary is routinely cut mid-text when the
tool_calls finish arrives; the partial stream IS a prefix of the full
commentary. The old prefix-based verdict set already_streamed=True, the
gateway called on_segment_break() finalizing the truncated bubble, and
the tail ("...pick it u" vs "...pick it up") was permanently lost. Only
an exact match may skip the full-text resend."""
@pytest.mark.parametrize(
("streamed", "expected_already_streamed"),
[
# Truncated at the text→tool_calls boundary (#88954): prefix only → full-text resend.
("checking the queue to pick it u", False),
# Fully streamed → gateway settles the bubble without a duplicate resend.
("checking the queue to pick it up", True),
],
)
def test_interim_commentary_already_streamed_requires_exact_match(
monkeypatch, streamed, expected_already_streamed
):
"""Only an exact stream match may mark commentary already_streamed; a prefix-only match used
to finalize the truncated bubble and permanently lose the tail (#88954)."""
agent = _build_agent(monkeypatch)
observed = {}
agent.interim_assistant_callback = lambda text, *, already_streamed=False: observed.update(
{"text": text, "already_streamed": already_streamed}
)
agent._current_streamed_assistant_text = "checking the queue to pick it u"
agent._current_streamed_assistant_text = streamed
from agent.codex_responses_adapter import _normalize_codex_response
normalized, finish_reason = _normalize_codex_response(
@@ -1927,37 +1931,9 @@ def test_interim_commentary_partial_stream_receives_full_text(monkeypatch):
agent._build_assistant_message(normalized, finish_reason)
)
# Prefix-only match: the gateway must receive the FULL text with
# already_streamed=False so on_commentary() re-delivers it.
assert observed == {
"text": "checking the queue to pick it up",
"already_streamed": False,
}
def test_interim_commentary_exact_stream_still_marks_streamed(monkeypatch):
"""The exact-equality fast path is preserved: a fully streamed commentary
keeps already_streamed=True so the gateway settles the bubble without a
duplicate resend."""
agent = _build_agent(monkeypatch)
observed = {}
agent.interim_assistant_callback = lambda text, *, already_streamed=False: observed.update(
{"text": text, "already_streamed": already_streamed}
)
agent._current_streamed_assistant_text = "checking the queue to pick it up"
from agent.codex_responses_adapter import _normalize_codex_response
normalized, finish_reason = _normalize_codex_response(
_codex_commentary_final_tool_response("checking the queue to pick it up")
)
agent._emit_interim_assistant_message(
agent._build_assistant_message(normalized, finish_reason)
)
assert observed == {
"text": "checking the queue to pick it up",
"already_streamed": True,
"already_streamed": expected_already_streamed,
}