From 808c8570a6dcef4f56857a6b29865affb5a8e784 Mon Sep 17 00:00:00 2001 From: Kong Date: Thu, 14 May 2026 00:57:01 +0800 Subject: [PATCH] fix(gateway): preserve queued follow-up media delivery Ensure queued follow-up resends keep MEDIA-backed attachments by replaying the first response through the gateway's text-plus-media delivery flow instead of a plain adapter text send. --- gateway/run.py | 44 +++++++++++++- tests/gateway/test_tts_media_routing.py | 77 +++++++++++++++++++++++++ 2 files changed, 119 insertions(+), 2 deletions(-) diff --git a/gateway/run.py b/gateway/run.py index 2c6de596f896f..c1b3803393ace 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -3041,6 +3041,22 @@ def _is_control_interrupt_message(message: Optional[str]) -> bool: return normalized in _CONTROL_INTERRUPT_MESSAGES +def _strip_response_attachments_for_direct_send(response: str, adapter) -> str: + """Return the visible text portion of a response before direct send(). + + Queued follow-up resends only replay attachments we can deliver natively in + this path: ``MEDIA:`` tags, bare local files, and internal directives. Keep + ordinary image URLs in the visible text until the queued path grows native + image-URL delivery too. + """ + _, cleaned = adapter.extract_media(response) + cleaned = cleaned.replace("[[audio_as_voice]]", "").strip() + cleaned = cleaned.replace("[[as_document]]", "").strip() + cleaned = re.sub(r"MEDIA:\s*\S+", "", cleaned).strip() + _, cleaned = adapter.extract_local_files(cleaned) + return cleaned.strip() + + def _skill_slug_from_frontmatter(skill_md: Path) -> tuple[str | None, str | None]: """Derive the /command slug and declared frontmatter name from a SKILL.md. @@ -19912,7 +19928,29 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew except Exception as e: logger.warning("Post-stream media extraction failed: %s", e) + async def _deliver_queued_first_response( + self, + response: str, + source: SessionSource, + adapter, + metadata: Optional[Dict[str, Any]] = None, + event_message_id: Optional[str] = None, + ) -> None: + """Deliver a queued response using the normal text+attachment split.""" + text_content = _strip_response_attachments_for_direct_send(response, adapter) + if text_content: + await adapter.send( + source.chat_id, + text_content, + metadata=metadata, + ) + synthetic_event = MessageEvent( + text="", + source=source, + message_id=event_message_id, + ) + await self._deliver_media_from_response(response, synthetic_event, adapter) async def _run_background_task( self, @@ -26297,10 +26335,12 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew "Queued follow-up for session %s: final stream delivery not confirmed; sending first response before continuing.", session_key or "?", ) - await adapter.send( - source.chat_id, + await self._deliver_queued_first_response( first_response, + source=source, + adapter=adapter, metadata=_status_thread_metadata, + event_message_id=event_message_id, ) except Exception as e: logger.warning("Failed to send first response before queued message: %s", e) diff --git a/tests/gateway/test_tts_media_routing.py b/tests/gateway/test_tts_media_routing.py index 50bc6368f16bd..f17d49870c013 100644 --- a/tests/gateway/test_tts_media_routing.py +++ b/tests/gateway/test_tts_media_routing.py @@ -196,3 +196,80 @@ class _DiscordMediaFailureAdapter(BasePlatformAdapter): async def get_chat_info(self, chat_id): return {"id": chat_id, "type": "dm"} + + +@pytest.mark.asyncio +async def test_queued_followup_delivery_strips_media_tag_from_text_and_sends_image(): + event = _event(thread_id="topic-1") + runner = object.__new__(GatewayRunner) + runner._thread_metadata_for_source = lambda source, anchor=None: {"thread_id": "topic-1"} + runner._reply_anchor_for_event = lambda event: event.message_id + + adapter = SimpleNamespace( + name="test", + extract_media=BasePlatformAdapter.extract_media, + extract_images=BasePlatformAdapter.extract_images, + extract_local_files=BasePlatformAdapter.extract_local_files, + send=AsyncMock(return_value=SendResult(success=True, message_id="text")), + send_multiple_images=AsyncMock(return_value=None), + send_voice=AsyncMock(return_value=SendResult(success=True, message_id="voice")), + send_document=AsyncMock(return_value=SendResult(success=True, message_id="doc")), + send_video=AsyncMock(return_value=SendResult(success=True, message_id="video")), + ) + + await GatewayRunner._deliver_queued_first_response( + runner, + "Quote here\nMEDIA:/tmp/pricelist.png", + source=event.source, + adapter=adapter, + metadata={"thread_id": "topic-1"}, + event_message_id=event.message_id, + ) + + adapter.send.assert_awaited_once_with( + "chat-1", + "Quote here", + metadata={"thread_id": "topic-1"}, + ) + adapter.send_multiple_images.assert_awaited_once_with( + chat_id="chat-1", + images=[("file:///tmp/pricelist.png", "")], + metadata={"thread_id": "topic-1"}, + ) + + +@pytest.mark.asyncio +async def test_queued_followup_delivery_keeps_remote_image_url_in_text(): + event = _event(thread_id="topic-1") + runner = object.__new__(GatewayRunner) + runner._thread_metadata_for_source = lambda source, anchor=None: {"thread_id": "topic-1"} + runner._reply_anchor_for_event = lambda event: event.message_id + + adapter = SimpleNamespace( + name="test", + extract_media=BasePlatformAdapter.extract_media, + extract_images=BasePlatformAdapter.extract_images, + extract_local_files=BasePlatformAdapter.extract_local_files, + send=AsyncMock(return_value=SendResult(success=True, message_id="text")), + send_multiple_images=AsyncMock(return_value=None), + send_voice=AsyncMock(return_value=SendResult(success=True, message_id="voice")), + send_document=AsyncMock(return_value=SendResult(success=True, message_id="doc")), + send_video=AsyncMock(return_value=SendResult(success=True, message_id="video")), + ) + + response = "See this mockup\nhttps://example.com/mockup.png" + await GatewayRunner._deliver_queued_first_response( + runner, + response, + source=event.source, + adapter=adapter, + metadata={"thread_id": "topic-1"}, + event_message_id=event.message_id, + ) + + adapter.send.assert_awaited_once_with( + "chat-1", + response, + metadata={"thread_id": "topic-1"}, + ) + adapter.send_multiple_images.assert_not_awaited()