Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 35 additions & 20 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This changes transformed finals from edits to fresh sends, but the unchanged tests/gateway/test_run_progress_topics.py:1043-1071 still requires the appended text in adapter.edits. Update that integration test to assert the new fresh-send contract, including the failed-delivery already_sent outcome.

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
Expand Down
47 changes: 45 additions & 2 deletions gateway/stream_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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(
Expand Down
29 changes: 29 additions & 0 deletions plugins/platforms/telegram/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
75 changes: 63 additions & 12 deletions tests/gateway/test_run_progress_topics.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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,
Expand All @@ -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
)


Expand Down
85 changes: 85 additions & 0 deletions tests/gateway/test_stream_consumer_fresh_final.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
Loading