From e52acf76a1527ca8b94cc202a947e042bb777800 Mon Sep 17 00:00:00 2001 From: HenryG <314296805+HenryGHJ@users.noreply.github.com> Date: Sat, 8 Aug 2026 01:42:59 +0800 Subject: [PATCH] fix(gateway): keep media history reads off event loop --- gateway/platforms/base.py | 38 +++++-- .../gateway/test_73771_media_resend_dedup.py | 100 ++++++++++++++++++ 2 files changed, 128 insertions(+), 10 deletions(-) diff --git a/gateway/platforms/base.py b/gateway/platforms/base.py index 39de048616502..61376f4de9b10 100644 --- a/gateway/platforms/base.py +++ b/gateway/platforms/base.py @@ -45,6 +45,10 @@ _AUDIO_EXTS = frozenset(_AUDIO_MIME_TYPES) _TELEGRAM_AUDIO_ATTACHMENT_EXTS = frozenset({'.mp3', '.m4a'}) _TELEGRAM_VOICE_EXTS = frozenset({'.ogg', '.opus'}) _POST_DELIVERY_CALLBACK_TIMEOUT_SECONDS = 30.0 +# Delivery-time history is best-effort dedup metadata, not canonical state. +# Keep this comfortably below the Discord heartbeat watchdog window and fail +# open rather than withholding a legitimate attachment. +_HISTORY_MEDIA_LOOKUP_TIMEOUT_SECONDS = 5.0 def _platform_name(platform) -> str: @@ -5892,16 +5896,6 @@ class BasePlatformAdapter(ABC): media_files, response = self.extract_media(response) media_files = self.filter_media_delivery_paths(media_files) - # Do NOT deduplicate MEDIA tags against prior turns here. - # The auto-append path in GatewayRunner._run_agent_inner already - # deduplicates auto-appended tags via _collect_auto_append_media_tags - # with history_media_paths, so this filter would only catch explicit - # MEDIA tags the model deliberately included in its response — which - # must be preserved (user asked to resend an image, the model echoed - # a path intentionally, etc.). Bare-file-path dedup still applies - # to local_files below via the same _history_media_paths set. - _history_media_paths = self._history_media_paths_for_session(session_key) - # Extract image URLs and send them as native platform attachments images, text_content = self.extract_images(response) # Strip any remaining internal directives from message body (fixes #1561). @@ -5920,6 +5914,30 @@ class BasePlatformAdapter(ABC): # instead of becoming native uploads. local_files, text_content = self.extract_local_files(text_content) local_files = self.filter_local_delivery_paths(local_files) + # Do NOT load the full SQLite transcript for ordinary text or + # explicit MEDIA tags. History is needed only for bare local + # paths auto-detected above. Run that synchronous DB/decode + # work off the platform event loop so a slow state.db read + # cannot block Discord heartbeats and trigger the liveness + # watchdog. On lookup failure the helper returns None and we + # fail open by delivering the candidate file. + _history_media_paths = None + if local_files: + try: + _history_media_paths = await asyncio.wait_for( + asyncio.to_thread( + self._history_media_paths_for_session, + session_key, + ), + timeout=_HISTORY_MEDIA_LOOKUP_TIMEOUT_SECONDS, + ) + except asyncio.TimeoutError: + logger.warning( + "[%s] Timed out loading media-delivery history for %s; " + "delivering bare local file path(s) without history dedup", + self.name, + session_key, + ) if _history_media_paths: _suppressed = [p for p in local_files if p in _history_media_paths] if _suppressed: diff --git a/tests/gateway/test_73771_media_resend_dedup.py b/tests/gateway/test_73771_media_resend_dedup.py index 84e8385adf781..dbed4ed554a81 100644 --- a/tests/gateway/test_73771_media_resend_dedup.py +++ b/tests/gateway/test_73771_media_resend_dedup.py @@ -23,6 +23,7 @@ sibling): import asyncio import logging +import time from types import SimpleNamespace from unittest.mock import AsyncMock @@ -218,6 +219,105 @@ async def test_bare_local_path_history_dedup_survives_and_logs(tmp_path, monkeyp ) +@pytest.mark.asyncio +async def test_plain_text_response_does_not_load_transcript(): + """Ordinary text delivery must not touch SQLite-backed history at all.""" + adapter = _DummyAdapter() + adapter._keep_typing = _hold_typing + + class _ExplodingStore: + calls = 0 + + def peek_session_id(self, _session_key): + self.calls += 1 + raise AssertionError("plain text delivery loaded session history") + + store = _ExplodingStore() + adapter.set_session_store(store) + + async def handler(_event): + return "Plain response with no local attachment path." + + adapter.set_message_handler(handler) + event = _make_event() + await adapter._process_message_background(event, build_session_key(event.source)) + + assert store.calls == 0 + assert any("Plain response" in item["content"] for item in adapter.sent) + + +@pytest.mark.asyncio +async def test_bare_path_history_lookup_does_not_block_event_loop(tmp_path, monkeypatch): + """A slow transcript read must run outside the platform event loop.""" + pdf = _allowed_file(tmp_path, monkeypatch, "slow-history.pdf") + monkeypatch.setattr("gateway.platforms.base.LOCAL_DELIVERY_SAFE_ROOTS", (pdf.parent,), raising=False) + adapter = _DummyAdapter() + adapter._keep_typing = _hold_typing + monkeypatch.setattr( + type(adapter), "filter_local_delivery_paths", staticmethod(lambda paths: list(paths)) + ) + + class _SlowStore: + def peek_session_id(self, _session_key): + return "sess-slow" + + def load_transcript(self, _session_id): + time.sleep(0.15) + return [{"role": "user", "content": "current"}] + + adapter.set_session_store(_SlowStore()) + + async def handler(_event): + return f"Generated file: {pdf}" + + adapter.set_message_handler(handler) + event = _make_event() + delivery = asyncio.create_task( + adapter._process_message_background(event, build_session_key(event.source)) + ) + await asyncio.sleep(0.02) + assert not delivery.done() + # If load_transcript ran on the event-loop thread, this 20 ms sleep would + # not resume until after the 150 ms blocking read completed. + assert not adapter.documents + await delivery + assert adapter.documents == [str(pdf)] + + +@pytest.mark.asyncio +async def test_bare_path_history_lookup_timeout_fails_open(tmp_path, monkeypatch): + """A wedged transcript read must not hold response delivery indefinitely.""" + pdf = _allowed_file(tmp_path, monkeypatch, "timeout-history.pdf") + monkeypatch.setattr("gateway.platforms.base.LOCAL_DELIVERY_SAFE_ROOTS", (pdf.parent,), raising=False) + monkeypatch.setattr("gateway.platforms.base._HISTORY_MEDIA_LOOKUP_TIMEOUT_SECONDS", 0.02) + adapter = _DummyAdapter() + adapter._keep_typing = _hold_typing + monkeypatch.setattr( + type(adapter), "filter_local_delivery_paths", staticmethod(lambda paths: list(paths)) + ) + + class _WedgedStore: + def peek_session_id(self, _session_key): + return "sess-wedged" + + def load_transcript(self, _session_id): + time.sleep(0.2) + return [] + + adapter.set_session_store(_WedgedStore()) + + async def handler(_event): + return f"Generated file: {pdf}" + + adapter.set_message_handler(handler) + event = _make_event() + started = time.monotonic() + await adapter._process_message_background(event, build_session_key(event.source)) + + assert time.monotonic() - started < 0.15 + assert adapter.documents == [str(pdf)] + + # --------------------------------------------------------------------------- # Streaming sibling (run.py _deliver_media_from_response) # ---------------------------------------------------------------------------