diff --git a/gateway/run.py b/gateway/run.py index 487846ddc2..464637e246 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -23638,7 +23638,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew is not a transactional boundary: a process crash after adapter acceptance can still cause durable at-least-once replay. """ - source = self._build_process_event_source(evt) + source = await asyncio.to_thread(self._build_process_event_source, evt) if not source: # API-server-originated sessions bind a RAW session key (the # X-Hermes-Session-Id value — see _bind_api_server_session), not a diff --git a/tests/gateway/test_background_process_notifications.py b/tests/gateway/test_background_process_notifications.py index 3dfff7a206..76941bb710 100644 --- a/tests/gateway/test_background_process_notifications.py +++ b/tests/gateway/test_background_process_notifications.py @@ -9,6 +9,7 @@ Contributed by @PeterFile (PR #593), reimplemented on current main. import asyncio import queue +import threading from types import SimpleNamespace from unittest.mock import AsyncMock @@ -292,6 +293,46 @@ async def test_inject_watch_notification_carries_message_id_reply_anchor(monkeyp assert synth_event.source.thread_id == "24296" +@pytest.mark.asyncio +async def test_inject_watch_notification_loads_session_store_off_loop(monkeypatch, tmp_path): + from gateway.session import SessionSource + + runner = _build_runner(monkeypatch, tmp_path, "all") + adapter = runner.adapters[Platform.TELEGRAM] + runner.session_store._entries["agent:main:telegram:dm:123:24296"] = SimpleNamespace( + origin=SessionSource( + platform=Platform.TELEGRAM, + chat_id="123", + chat_type="dm", + thread_id="24296", + user_id="1", + user_name="Fabio", + ) + ) + loop_thread = threading.get_ident() + load_threads = [] + real_ensure_loaded = runner.session_store._ensure_loaded + + def spy_ensure_loaded(): + load_threads.append(threading.get_ident()) + return real_ensure_loaded() + + monkeypatch.setattr(runner.session_store, "_ensure_loaded", spy_ensure_loaded) + + await runner._inject_watch_notification( + "[SYSTEM: Background process matched]", + { + "session_id": "proc_watch", + "session_key": "agent:main:telegram:dm:123:24296", + "message_id": "777", + }, + ) + + adapter.handle_message.assert_awaited_once() + assert load_threads + assert all(thread_id != loop_thread for thread_id in load_threads) + + @pytest.mark.asyncio async def test_inject_watch_notification_ignores_foreground_event_source(monkeypatch, tmp_path): """Negative test: watch notification must NOT route to the foreground thread."""