diff --git a/gateway/run.py b/gateway/run.py index 37dfd9d92cde5..6dd9e88fbe09e 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -23663,27 +23663,42 @@ def _stream_confirmed_final_delivery( ) response["already_sent"] = True elif not _is_empty_sentinel and _transformed and _sc is not None: - # Plugin hooks transformed the response after streaming — edit the - # existing streamed message instead of sending a duplicate. - _sc_msg_id = _sc.message_id - if _sc_msg_id: - try: - await _sc.adapter.edit_message( - chat_id=source.chat_id, - message_id=_sc_msg_id, - content=response["final_response"], - finalize=True, - ) - response["already_sent"] = True - logger.info( - "Edited streamed message %s for session %s to include plugin-transformed content.", - _sc_msg_id, session_key or "?", - ) - except Exception as _edit_err: - logger.warning( - "Failed to edit streamed message for session %s: %s", - session_key or "?", _edit_err, + # The plugin-transformed final no longer matches the streamed + # preview. Commit it through the consumer's fresh-final path so + # topic metadata is preserved, stale previews are removed only + # after confirmed delivery, and a failed SendResult leaves the + # caller's normal final-send fallback enabled. + try: + _deliver_transformed = getattr( + _sc, "deliver_transformed_final", None + ) + if callable(_deliver_transformed): + _delivery_result = _deliver_transformed( + response["final_response"] ) + if inspect.isawaitable(_delivery_result): + _delivery_result = await _delivery_result + _transformed_delivered = bool(_delivery_result) + else: + _transformed_delivered = False + except Exception as _send_err: + _transformed_delivered = False + logger.warning( + "Failed to send transformed final for session %s: %s", + session_key or "?", _send_err, + ) + if _transformed_delivered: + response["already_sent"] = True + logger.info( + "Delivered transformed final for session %s through fresh send.", + session_key or "?", + ) + else: + logger.warning( + "Transformed final delivery was not confirmed for session %s; " + "leaving normal final-send fallback enabled.", + session_key or "?", + ) # Schedule deletion of tracked temporary progress bubbles after the # final response lands. Failed runs skip this so bubbles remain as diff --git a/gateway/stream_consumer.py b/gateway/stream_consumer.py index 7d3e31cf15355..17d631c462c9e 100644 --- a/gateway/stream_consumer.py +++ b/gateway/stream_consumer.py @@ -247,6 +247,7 @@ def __init__( self._last_edit_overflowed = False self._fallback_final_send = False self._fallback_prefix = "" + self._fallback_resend_full = False # True when fallback is sending only the missing tail after a partial # Telegram overflow delivery. In that case the already-visible prefix # is intentional content, not a stale preview to delete. @@ -481,6 +482,7 @@ def _reset_segment_state(self, *, preserve_no_edit: bool = False) -> None: self._last_sent_text = "" self._fallback_final_send = False self._fallback_prefix = "" + self._fallback_resend_full = False self._fallback_preserve_partial_messages = False self._segment_preview_message_ids = set() # #29346: a tool/segment boundary means what we delivered was an interim @@ -1173,6 +1175,8 @@ def _visible_prefix(self) -> str: def _continuation_text(self, final_text: str) -> str: """Return only the part of final_text the user has not already seen.""" + if self._fallback_resend_full: + return final_text prefix = self._fallback_prefix or self._visible_prefix() if prefix and final_text.startswith(prefix): return final_text[len(prefix):].lstrip() @@ -1371,8 +1375,18 @@ async def _send_fallback_final(self, text: str) -> None: _len_fn = self.adapter.message_len_fn_for_chat(self.chat_id) except Exception as e: logger.debug("per-chat limit resolution failed: %s", e) - safe_limit = max(500, raw_limit - 100) + safe_limit = ( + max(1, raw_limit - 100) + if self._fallback_resend_full + else max(500, raw_limit - 100) + ) chunks = self._split_text_chunks(continuation, safe_limit, len_fn=_len_fn) + if self._fallback_resend_full and len(chunks) > 1: + total = len(chunks) + chunks = [ + f"{chunk} ({index}/{total})" + for index, chunk in enumerate(chunks, 1) + ] stale_message_id = self._message_id # partial message to clean up last_message_id: Optional[str] = None @@ -1458,6 +1472,7 @@ async def _send_fallback_final(self, text: str) -> None: self._final_content_delivered = True self._last_sent_text = chunks[-1] self._fallback_prefix = "" + self._fallback_resend_full = False self._fallback_preserve_partial_messages = False async def _send_empty_fallback_final(self, final_text: str) -> str: @@ -1898,8 +1913,29 @@ async def _try_fresh_final(self, text: str, *, is_turn_final: bool = True) -> bo self._last_sent_text = text if is_turn_final: self._final_response_sent = True + self._final_content_delivered = True return True + async def deliver_transformed_final(self, text: str) -> bool: + """Commit a plugin-transformed final reply as numbered fresh chunks. + + A transformed response no longer matches the streamed preview. Force the + consumer-owned fallback path so every adapter ``send`` stays below the + platform limit, preserves routing metadata, and is confirmed separately. + The stale preview is removed only after every numbered chunk succeeds. + False leaves the gateway's normal full-response fallback enabled. + """ + text = self._clean_for_display(text) + if not text.strip(): + return False + self._fallback_resend_full = True + self._fallback_prefix = "" + self._fallback_preserve_partial_messages = False + self._final_response_sent = False + self._final_content_delivered = False + await self._send_fallback_final(text) + return bool(self._final_response_sent and self._final_content_delivered) + async def _suppress_silence_marker(self) -> None: """Retract any streamed preview when the final reply is a silence marker. @@ -2156,7 +2192,14 @@ async def _send_or_edit( or self._message_id ) delivered_prefix = raw_response.get("delivered_prefix") - if isinstance(delivered_prefix, str) and delivered_prefix: + if raw_response.get("resend_full_final"): + # The visible preview could not be finalized and + # therefore is not a reliable numbered chunk 1. + # Replace it with the complete final response. + self._fallback_prefix = "" + self._fallback_resend_full = True + self._fallback_preserve_partial_messages = False + elif isinstance(delivered_prefix, str) and delivered_prefix: self._last_sent_text = delivered_prefix self._fallback_prefix = delivered_prefix self._fallback_preserve_partial_messages = text.startswith( diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index f9f9e21d0b61f..fbac1921b75f1 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -4956,6 +4956,35 @@ async def _edit_overflow_split( # First chunk identical to current text — fall through to # send continuations. pass + elif getattr(e, "retry_after", None) is not None or "retry after" in err_str: + # The saturated preview is visible, but the rejected finalize + # edit means it has no trustworthy (1/N) marker. Report a + # partial overflow that requires a complete fresh resend rather + # than treating the preview as a valid delivered prefix. + logger.warning( + "[%s] Overflow split: first-chunk edit rate-limited; " + "falling back to full numbered resend: %s", + self.name, e, + ) + delivered_prefix = re.sub(r" \(\d+/\d+\)$", "", first_chunk) + retry_after = getattr(e, "retry_after", None) + return SendResult( + success=False, + message_id=message_id, + error="overflow_first_chunk_edit_rate_limited", + retryable=True, + retry_after=float(retry_after) if retry_after is not None else None, + raw_response={ + "partial_overflow": True, + "delivered_chunks": 1, + "total_chunks": len(chunks), + "last_message_id": message_id, + "delivered_prefix": delivered_prefix, + "resend_full_final": True, + "continuation_message_ids": (), + }, + continuation_message_ids=(), + ) else: logger.error( "[%s] Overflow split: first-chunk edit failed: %s", diff --git a/tests/gateway/test_run_progress_topics.py b/tests/gateway/test_run_progress_topics.py index 822cc0fb904d0..1fc14ff936360 100644 --- a/tests/gateway/test_run_progress_topics.py +++ b/tests/gateway/test_run_progress_topics.py @@ -168,6 +168,21 @@ async def edit_message(self, chat_id, message_id, content) -> SendResult: return await super().edit_message(chat_id, message_id, content) +class FailingTransformedFinalAdapter(MetadataEditProgressCaptureAdapter): + """Fail only the plugin-transformed fresh final delivery attempt.""" + + async def send(self, chat_id, content, reply_to=None, metadata=None) -> SendResult: + if "[plugin appended this]" not in content: + return await super().send(chat_id, content, reply_to, metadata) + self.sent.append({ + "chat_id": chat_id, + "content": content, + "reply_to": reply_to, + "metadata": metadata, + }) + return SendResult(success=False, error="simulated transformed delivery failure") + + class NonEditingProgressCaptureAdapter(ProgressCaptureAdapter): SUPPORTS_MESSAGE_EDITING = False @@ -1225,12 +1240,10 @@ def run_conversation(self, message, conversation_history=None, task_id=None): @pytest.mark.asyncio -async def test_transformed_response_edits_streamed_message_in_place(monkeypatch, tmp_path): - """When a transform_llm_output hook modifies the response after streaming, - the gateway must edit the existing streamed message in place with the full - transformed content (so plugins like content filters / appenders reach the - user) and still mark already_sent=True (no duplicate send). - """ +async def test_transformed_response_uses_confirmed_fresh_delivery( + monkeypatch, tmp_path +): + """A plugin-transformed final is freshly sent with topic metadata.""" adapter, result = await _run_with_agent( monkeypatch, tmp_path, @@ -1247,13 +1260,51 @@ async def test_transformed_response_edits_streamed_message_in_place(monkeypatch, adapter_cls=MetadataEditProgressCaptureAdapter, ) - # Final delivery happened (no duplicate send fallback). assert result.get("already_sent") is True - # The transformed final text reached the user — appended portion is present - # in an edit_message call (not just in the streamed sends). - edited_texts = [e["content"] for e in adapter.edits] - assert any("[plugin appended this]" in text for text in edited_texts), ( - f"expected transformed text in adapter.edits, got: {edited_texts!r}" + transformed_sends = [ + call for call in adapter.sent + if "[plugin appended this]" in call["content"] + ] + assert len(transformed_sends) == 1 + assert transformed_sends[0]["metadata"]["thread_id"] == "$thread" + assert not any( + "[plugin appended this]" in edit["content"] + for edit in adapter.edits + ) + + +@pytest.mark.asyncio +async def test_failed_transformed_fresh_delivery_leaves_fallback_enabled( + monkeypatch, tmp_path +): + """A failed transformed send must not suppress normal final delivery.""" + adapter, result = await _run_with_agent( + monkeypatch, + tmp_path, + TransformedStreamAgent, + session_id="sess-transformed-stream-failure", + config_data={ + "display": {"tool_progress": "off", "interim_assistant_messages": False}, + "streaming": {"enabled": True, "edit_interval": 0.01, "buffer_threshold": 1}, + }, + platform=Platform.MATRIX, + chat_id="!room:matrix.example.org", + chat_type="group", + thread_id="$thread", + adapter_cls=FailingTransformedFinalAdapter, + ) + + assert result.get("already_sent") is not True + assert result["final_response"].endswith("[plugin appended this]") + transformed_attempts = [ + call for call in adapter.sent + if "[plugin appended this]" in call["content"] + ] + assert len(transformed_attempts) == 1 + assert transformed_attempts[0]["metadata"]["thread_id"] == "$thread" + assert not any( + "[plugin appended this]" in edit["content"] + for edit in adapter.edits ) diff --git a/tests/gateway/test_stream_consumer_fresh_final.py b/tests/gateway/test_stream_consumer_fresh_final.py index f8270cfd86dca..d7c47318fbb7e 100644 --- a/tests/gateway/test_stream_consumer_fresh_final.py +++ b/tests/gateway/test_stream_consumer_fresh_final.py @@ -174,6 +174,91 @@ async def test_no_edit_sentinel_is_not_affected(self): assert consumer._should_send_fresh_final() is False +class TestTransformedFinalDelivery: + """Plugin-transformed finals replace previews through the safe fresh path.""" + + @pytest.mark.asyncio + async def test_success_preserves_topic_metadata_and_deletes_preview(self): + adapter = _make_adapter() + adapter.send.return_value = SimpleNamespace( + success=True, message_id="transformed-final", + ) + consumer = GatewayStreamConsumer( + adapter=adapter, + chat_id="chat", + metadata={"thread_id": "77", "reply_to_message_id": "41"}, + ) + consumer._message_id = "preview" + consumer._track_preview_id("preview") + + delivered = await consumer.deliver_transformed_final("complete transformed reply") + + assert delivered is True + adapter.send.assert_awaited_once_with( + chat_id="chat", + content="complete transformed reply", + metadata={ + "thread_id": "77", + "reply_to_message_id": "41", + "notify": True, + }, + ) + adapter.delete_message.assert_awaited_once_with("chat", "preview") + assert consumer.final_response_sent is True + assert consumer.final_content_delivered is True + + @pytest.mark.asyncio + async def test_failure_keeps_preview_and_does_not_confirm_delivery(self): + adapter = _make_adapter() + adapter.send.return_value = SimpleNamespace(success=False, error="flood_control") + consumer = GatewayStreamConsumer( + adapter=adapter, + chat_id="chat", + metadata={"thread_id": "77"}, + ) + consumer._message_id = "preview" + consumer._track_preview_id("preview") + + delivered = await consumer.deliver_transformed_final("complete transformed reply") + + assert delivered is False + adapter.delete_message.assert_not_awaited() + assert consumer.final_response_sent is False + assert consumer.final_content_delivered is False + + @pytest.mark.asyncio + async def test_later_chunk_failure_never_confirms_or_deletes_preview(self): + adapter = _make_adapter() + adapter.send.side_effect = [ + SimpleNamespace(success=True, message_id="chunk-1"), + SimpleNamespace(success=False, error="continuation failed"), + ] + consumer = GatewayStreamConsumer( + adapter=adapter, + chat_id="chat", + metadata={"thread_id": "77", "reply_to_message_id": "41"}, + ) + consumer._message_id = "preview" + consumer._track_preview_id("preview") + final_text = "word " * 1800 + + delivered = await consumer.deliver_transformed_final(final_text) + + assert delivered is False + assert adapter.send.await_count == 2 + first = adapter.send.await_args_list[0].kwargs + second = adapter.send.await_args_list[1].kwargs + assert first["content"].endswith(" (1/3)") + assert second["content"].endswith(" (2/3)") + assert len(first["content"]) <= adapter.MAX_MESSAGE_LENGTH + assert len(second["content"]) <= adapter.MAX_MESSAGE_LENGTH + assert first["metadata"]["thread_id"] == "77" + assert second["metadata"]["thread_id"] == "77" + adapter.delete_message.assert_not_awaited() + assert consumer.final_response_sent is False + assert consumer.final_content_delivered is False + + class TestSegmentBreakDoesNotMarkFinalSent: """Regression for #29346 — silent response loss after tool calls. diff --git a/tests/gateway/test_telegram_overflow_partial.py b/tests/gateway/test_telegram_overflow_partial.py index 663d1c83af04c..a4936c3a71eaa 100644 --- a/tests/gateway/test_telegram_overflow_partial.py +++ b/tests/gateway/test_telegram_overflow_partial.py @@ -95,6 +95,36 @@ async def test_edit_overflow_split_reports_partial_failure_when_continuation_fai assert result.continuation_message_ids == () +@pytest.mark.asyncio +async def test_edit_overflow_split_first_edit_rate_limit_requires_full_resend(telegram_adapter): + """A rate-limited finalize edit cannot establish a numbered chunk 1.""" + class RetryAfterLike(RuntimeError): + retry_after = 167 + + content = "word " * 120 + telegram_adapter._bot.edit_message_text = AsyncMock( + side_effect=RetryAfterLike("Flood control exceeded. Retry in 167 seconds") + ) + telegram_adapter._bot.send_message = AsyncMock() + + result = await telegram_adapter._edit_overflow_split( + "12345", "201", content, finalize=True, metadata={"thread_id": "77"} + ) + + assert result.success is False + assert result.retryable is True + assert result.error == "overflow_first_chunk_edit_rate_limited" + assert result.message_id == "201" + assert result.raw_response["partial_overflow"] is True + assert result.raw_response["delivered_chunks"] == 1 + assert result.raw_response["total_chunks"] > 1 + assert result.raw_response["last_message_id"] == "201" + assert result.raw_response["delivered_prefix"] + assert result.raw_response["resend_full_final"] is True + assert result.continuation_message_ids == () + telegram_adapter._bot.send_message.assert_not_awaited() + + @pytest.mark.asyncio async def test_stream_consumer_fallback_sends_tail_after_partial_overflow(): """A partial overflow edit enters fallback instead of marking final delivered.""" @@ -138,3 +168,57 @@ async def test_stream_consumer_fallback_sends_tail_after_partial_overflow(): adapter.delete_message.assert_not_awaited() assert consumer.final_response_sent is True assert consumer.final_content_delivered is True + + +@pytest.mark.asyncio +async def test_first_chunk_rate_limit_resends_full_numbered_final_and_deletes_preview(): + """Recovery replaces an unnumbered preview with a complete numbered reply.""" + adapter = MagicMock() + adapter.MAX_MESSAGE_LENGTH = 160 + adapter.edit_message = AsyncMock( + return_value=SendResult( + success=False, + message_id="preview-1", + error="overflow_first_chunk_edit_rate_limited", + retryable=True, + raw_response={ + "partial_overflow": True, + "delivered_chunks": 1, + "total_chunks": 4, + "last_message_id": "preview-1", + "delivered_prefix": "word " * 25, + "resend_full_final": True, + }, + ) + ) + next_message_id = 0 + + async def send_numbered_chunk(**_kwargs): + nonlocal next_message_id + next_message_id += 1 + return SendResult(success=True, message_id=f"final-{next_message_id}") + + adapter.send = AsyncMock(side_effect=send_numbered_chunk) + adapter.delete_message = AsyncMock(return_value=True) + + consumer = GatewayStreamConsumer(adapter, "chat-1", metadata={"thread_id": "77"}) + consumer._message_id = "preview-1" + consumer._last_sent_text = "word " * 25 + final_text = "word " * 180 + + ok = await consumer._send_or_edit(final_text, finalize=True) + assert ok is False + + await consumer._send_fallback_final(final_text) + + delivered = [call.kwargs["content"] for call in adapter.send.await_args_list] + assert len(delivered) > 1 + assert all( + chunk.endswith(f" ({index}/{len(delivered)})") + for index, chunk in enumerate(delivered, 1) + ) + assert all(len(chunk) <= adapter.MAX_MESSAGE_LENGTH for chunk in delivered) + assert "".join(chunk.rsplit(" (", 1)[0] for chunk in delivered) == final_text + adapter.delete_message.assert_awaited_once_with("chat-1", "preview-1") + assert consumer.final_response_sent is True + assert consumer.final_content_delivered is True