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.
This commit is contained in:
Kong 2026-05-14 00:57:01 +08:00 committed by Teknium
parent d220f1fdc0
commit 808c8570a6
2 changed files with 119 additions and 2 deletions

View File

@ -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)

View File

@ -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()