fix(gateway): keep process notification routing off-loop

This commit is contained in:
dsad
2026-07-12 22:48:11 +03:00
committed by Teknium
parent 7619564fbd
commit 858a6008be
2 changed files with 42 additions and 1 deletions
+1 -1
View File
@@ -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
@@ -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."""