From 167f0ef5fa63070a4c06e374a783298ae86ef8b4 Mon Sep 17 00:00:00 2001 From: liuhao1024 Date: Tue, 18 Aug 2026 14:27:41 +0800 Subject: [PATCH] fix(streaming): stop losing the tail of commentary truncated at a tool-call boundary MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit With streaming enabled 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 interim-message path used the prefix-based _interim_content_was_streamed verdict, so the gateway's _interim_assistant_cb called on_segment_break() — finalizing the truncated bubble — and the tail ("...pick it u" vs "...pick it up") was permanently lost (#88954). That prefix semantics stays correct for the conversation-loop "previewed" marks (the streamed prefix IS on the user's screen there, and the contract test pins it per the #65919 review). The gateway decision needs the stricter test: add _interim_content_fully_streamed (exact normalized equality) and use it for the interim-message verdict. Only an exact match may skip the full-text resend; a partial prefix falls through to on_commentary() and re-delivers the complete text — a benign duplicate, never lost text. (cherry picked from commit 0b06a6660a1a4d47b4974f21ae42d7aeb1cbea15) --- agent/stream_delivery.py | 23 ++++++- tests/agent/test_run_agent_codex_responses.py | 61 +++++++++++++++++++ 2 files changed, 83 insertions(+), 1 deletion(-) diff --git a/agent/stream_delivery.py b/agent/stream_delivery.py index 1027e65667..f1e4c529a4 100644 --- a/agent/stream_delivery.py +++ b/agent/stream_delivery.py @@ -112,6 +112,27 @@ class StreamDeliveryMixin: # reverse (streamed longer) is NOT matched — it could suppress a needed resend. return bool(visible_content and streamed) and visible_content.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 + 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. @@ -202,7 +223,7 @@ class StreamDeliveryMixin: visible = "\n\n".join(undelivered_parts).strip() if commentary_parts else self._interim_assistant_visible_text(assistant_msg) if not visible or visible == "(empty)" or self._interim_text_was_delivered(visible): return - already_streamed = self._interim_content_was_streamed(visible) + already_streamed = self._interim_content_fully_streamed(visible) self._enqueue_stream_hook("on_interim_message", text=visible, already_streamed=already_streamed) self._deliver_interim(visible, already_streamed=already_streamed, record=undelivered_parts or [visible]) diff --git a/tests/agent/test_run_agent_codex_responses.py b/tests/agent/test_run_agent_codex_responses.py index d8435f044b..f0a33baa1f 100644 --- a/tests/agent/test_run_agent_codex_responses.py +++ b/tests/agent/test_run_agent_codex_responses.py @@ -1900,6 +1900,67 @@ 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.""" + 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" + 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") + ) + assert finish_reason == "tool_calls" + agent._emit_interim_assistant_message( + 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, + } + + def test_stream_delta_strips_leaked_memory_context(monkeypatch): agent = _build_agent(monkeypatch) observed = []