Skip to content
Merged
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
2 changes: 1 addition & 1 deletion gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -14600,7 +14600,7 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g
if _media_adapter:
await self._deliver_media_from_response(
response, event, _media_adapter,
history_media_paths=_history_media_paths,
history_media_paths=_collect_history_media_paths(history),
)
# Streaming already delivered the body text, but the footer was
# intentionally held back (see the `not already_sent` gate above).
Expand Down
48 changes: 48 additions & 0 deletions tests/gateway/test_42039_duplicate_user_message.py
Original file line number Diff line number Diff line change
Expand Up @@ -212,6 +212,54 @@ async def test_not_new_messages_skip_db_when_agent_has_session_db(
)


# ── Post-stream MEDIA delivery keeps prior-turn deduplication ──────────


@pytest.mark.asyncio
async def test_streamed_response_receives_prior_turn_media_paths(
monkeypatch, tmp_path
):
"""An ordinary streamed reply completes the post-stream delivery branch.

The history-derived dedup set is part of that branch's contract, rather
than an optional best-effort hint: passing an undefined local crashes the
entire reply, while passing an empty set reintroduces duplicate MEDIA
attachments on later streamed responses.
"""
runner = _bootstrap(monkeypatch, tmp_path)
prior_path = "/tmp/already-delivered.png"
runner.session_store.load_transcript.return_value = [
{"role": "assistant", "content": f"MEDIA:{prior_path}"},
]
runner.adapters = {Platform.TELEGRAM: MagicMock()}
runner._deliver_media_from_response = AsyncMock()
runner._run_agent = AsyncMock(
return_value={
"final_response": "the streamed reply completed normally",
"messages": [
{"role": "assistant", "content": f"MEDIA:{prior_path}"},
{"role": "user", "content": "what is my status?"},
{"role": "assistant", "content": "the streamed reply completed normally"},
],
"tools": [],
"history_offset": 1,
"last_prompt_tokens": 0,
"already_sent": True,
"failed": False,
}
)

response = await runner._handle_message_with_agent(
_event(), _source(), "agent:main:telegram:group:-1001:12345", 1
)

assert response is None
runner._deliver_media_from_response.assert_awaited_once()
assert runner._deliver_media_from_response.await_args.kwargs[
"history_media_paths"
] == {prior_path}


# ── Test 4: normal path (new_messages found) uses skip_db=True ────────


Expand Down
Loading