Skip to content
Closed
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
64 changes: 58 additions & 6 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -12653,7 +12653,7 @@ async def _deliver_media_from_response(
# extract_local_files scanned text that still contained MEDIA: tags,
# producing false-positive bare-path matches with the MEDIA: prefix
# glued on. This matches the chain order in gateway/platforms/base.py.
_, cleaned = adapter.extract_images(cleaned)
remote_images, cleaned = adapter.extract_images(cleaned)
local_files, _ = adapter.extract_local_files(cleaned)
local_files = BasePlatformAdapter.filter_local_delivery_paths(local_files)

Expand Down Expand Up @@ -12685,12 +12685,13 @@ async def _deliver_media_from_response(
else:
non_image_local.append(file_path)

if image_paths:
image_batch = list(remote_images or [])
image_batch.extend((f"file://{_quote(p)}", "") for p in image_paths)
if image_batch:
try:
images = [(f"file://{_quote(p)}", "") for p in image_paths]
await adapter.send_multiple_images(
chat_id=event.source.chat_id,
images=images,
images=image_batch,
metadata=_thread_meta,
)
except Exception as e:
Expand Down Expand Up @@ -12742,6 +12743,50 @@ async def _deliver_media_from_response(
logger.warning("Post-stream media extraction failed: %s", e)


async def _send_queued_first_response(
self,
response: str,
event: MessageEvent,
adapter,
metadata: Optional[Dict[str, Any]] = None,
) -> None:
"""Deliver a queued-turn first response with normal media handling.

The queued follow-up path cannot hand the response back to
``BasePlatformAdapter._process_message_background()``, so sending the raw
text through ``adapter.send()`` would expose MEDIA directives and skip
attachments. Mirror the normal response pipeline: send only display text,
then deliver extracted MEDIA/local files separately.
"""
if not response:
return

text_content = response
try:
Comment thread
Qwinty marked this conversation as resolved.
from gateway.platforms.base import _strip_media_directives

_, text_content = adapter.extract_media(text_content)
_, text_content = adapter.extract_images(text_content)
_, text_content = adapter.extract_local_files(text_content)
text_content = _strip_media_directives(text_content).strip()
except Exception as exc:
logger.warning(
"[%s] Queued follow-up response media cleanup failed: %s",
getattr(adapter, "name", "adapter"),
exc,
)
text_content = response

if text_content:
await adapter.send(
event.source.chat_id,
text_content,
metadata=metadata,
)

await self._deliver_media_from_response(response, event, adapter)



async def _run_background_task(
self,
Expand Down Expand Up @@ -19057,9 +19102,16 @@ def _stream_confirmed_final_delivery(
"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,
response_event = MessageEvent(
text="",
message_type=MessageType.TEXT,
source=source,
message_id=event_message_id,
)
await self._send_queued_first_response(
first_response,
response_event,
adapter,
metadata=_status_thread_metadata,
)
except Exception as e:
Expand Down
115 changes: 115 additions & 0 deletions tests/gateway/test_queued_followup_media_delivery.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
from types import SimpleNamespace
from unittest.mock import AsyncMock

import pytest

from gateway.config import Platform
from gateway.platforms.base import BasePlatformAdapter, MessageEvent
from gateway.run import GatewayRunner
from gateway.session import SessionSource


def _runner() -> GatewayRunner:
runner = object.__new__(GatewayRunner)
runner._reply_anchor_for_event = lambda event: event.message_id
runner._thread_metadata_for_source = lambda source, _reply=None: (
{"thread_id": source.thread_id} if source.thread_id else None
)
return runner


def _event() -> MessageEvent:
return MessageEvent(
text="next",
source=SessionSource(
platform=Platform.TELEGRAM,
chat_id="273403055",
chat_type="dm",
user_id="273403055",
thread_id="374476",
),
message_id="25741",
)


def _adapter():
return SimpleNamespace(
name="Telegram",
extract_media=BasePlatformAdapter.extract_media,
extract_images=lambda text: ([], text),
extract_local_files=BasePlatformAdapter.extract_local_files,
send=AsyncMock(),
send_voice=AsyncMock(),
send_video=AsyncMock(),
send_document=AsyncMock(),
send_multiple_images=AsyncMock(),
)


@pytest.mark.asyncio
async def test_queued_first_response_strips_media_and_sends_document(tmp_path):
html = tmp_path / "pivin-scenario.html"
html.write_text("<!doctype html><title>Pivin</title>", encoding="utf-8")

adapter = _adapter()
response = f"Done.\n\nMEDIA:{html}"

await _runner()._send_queued_first_response(
response,
_event(),
adapter,
metadata={"thread_id": "374476"},
)

adapter.send.assert_awaited_once_with(
"273403055",
"Done.",
metadata={"thread_id": "374476"},
)
adapter.send_document.assert_awaited_once()
assert adapter.send_document.await_args.kwargs["chat_id"] == "273403055"
assert adapter.send_document.await_args.kwargs["file_path"] == str(html)
assert adapter.send_document.await_args.kwargs["metadata"] == {"thread_id": "374476"}


@pytest.mark.asyncio
async def test_queued_first_response_all_media_sends_no_empty_text(tmp_path):
html = tmp_path / "handoff.html"
html.write_text("<!doctype html><title>Only media</title>", encoding="utf-8")

adapter = _adapter()

await _runner()._send_queued_first_response(
f"MEDIA:{html}",
_event(),
adapter,
metadata={"thread_id": "374476"},
)

adapter.send.assert_not_awaited()
adapter.send_document.assert_awaited_once()
assert adapter.send_document.await_args.kwargs["file_path"] == str(html)


@pytest.mark.asyncio
async def test_queued_first_response_preserves_remote_image_delivery():
adapter = _adapter()
adapter.extract_images = BasePlatformAdapter.extract_images

await _runner()._send_queued_first_response(
"Chart ready.\n\n![usage chart](https://example.com/usage.png)",
_event(),
adapter,
metadata={"thread_id": "374476"},
)

adapter.send.assert_awaited_once_with(
"273403055",
"Chart ready.",
metadata={"thread_id": "374476"},
)
adapter.send_multiple_images.assert_awaited_once_with(
chat_id="273403055",
images=[("https://example.com/usage.png", "usage chart")],
metadata={"thread_id": "374476"},
)
Loading