From be154517a87f82e50c69a8ab0198b5ef7df29f97 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Mon, 7 Sep 2026 13:02:13 -0700 Subject: [PATCH] fix: restore gateway heartbeat watches after restart Recover active watches from current persisted session origins and exact route keys, reading heartbeat state off-loop in each source's profile. Failed scans leave watches intact for the poller's next retry. Start the heartbeat poller even when startup restores no watches. Slim synthesis of #92660, #98310 and #98313, with earlier restart recovery prior art from #92594. The integration poller calls restore_heartbeat_watches on every poll, including empty registries. Co-authored-by: chelsealong Co-authored-by: liuhao1024 Co-authored-by: fangliquanflq --- gateway/run_heartbeat_restore.py | 48 ++++++++ gateway/run_startup.py | 3 + tests/gateway/test_heartbeat_watch_restore.py | 105 ++++++++++++++++++ website/docs/user-guide/features/heartbeat.md | 3 +- 4 files changed, 158 insertions(+), 1 deletion(-) create mode 100644 gateway/run_heartbeat_restore.py create mode 100644 tests/gateway/test_heartbeat_watch_restore.py diff --git a/gateway/run_heartbeat_restore.py b/gateway/run_heartbeat_restore.py new file mode 100644 index 0000000000..0667b9f7d6 --- /dev/null +++ b/gateway/run_heartbeat_restore.py @@ -0,0 +1,48 @@ +"""Recover heartbeat watches from the gateway's canonical persisted routing index.""" +from __future__ import annotations + +import logging + +logger = logging.getLogger("gateway.run") + + +async def restore_heartbeat_watches(runner) -> None: + """Retryable startup/poll scan; failed reads never prune existing watches. + + SessionStore owns one routing index across profiles. Its origin and exact key, + rather than a second heartbeat routing snapshot, also cover pre-upgrade state. + Run all storage work off-loop so a cold profile DB cannot block adapters. + """ + from gateway.run import _profile_runtime_scope + from hermes_cli.heartbeat import HeartbeatManager + from hermes_constants import get_hermes_home + + store = runner.session_store + + def scan(): + restored = [] + # The poller may have been spawned by a named profile's /heartbeat command. + # Anchor even default origins to the gateway home, not inherited context. + home = getattr(store, "_routing_home", None) or get_hermes_home() + with _profile_runtime_scope(home): + entries = store.list_sessions() + for entry in entries: + if entry.origin is None or not entry.session_id: + continue + try: + with runner._profile_scope_for_source(entry.origin): + manager = HeartbeatManager(entry.session_id) + if manager.is_active(): + restored.append((entry.session_key, entry.origin, entry.session_id)) + except Exception: + logger.debug("heartbeat restore for %s failed", entry.session_key, exc_info=True) + return restored + + try: + candidates = await runner._run_in_executor_with_context(scan) + for key, source, session_id in candidates: + # A reset/compression may have published a new owner during the executor hop. + if store.peek_session_id(key) == session_id: + runner._register_heartbeat_watch(key, source, session_id) + except Exception: + logger.debug("heartbeat restore scan failed; retrying on next poll", exc_info=True) diff --git a/gateway/run_startup.py b/gateway/run_startup.py index 818c2c3f6b..29fd2c252f 100644 --- a/gateway/run_startup.py +++ b/gateway/run_startup.py @@ -1205,6 +1205,9 @@ class GatewayStartupMixin: ) self._spawn_supervised(self._hosted_room_worker_watcher, "hosted_room_worker") self._start_loop_heartbeat_task() + from gateway.run_heartbeat_restore import restore_heartbeat_watches + self._start_heartbeat_poller() # Keep retrying even when the first scan is empty. + await restore_heartbeat_watches(self) hook_count = len(self.hooks.loaded_hooks) if hook_count: logger.info("%s hook(s) loaded", hook_count) diff --git a/tests/gateway/test_heartbeat_watch_restore.py b/tests/gateway/test_heartbeat_watch_restore.py new file mode 100644 index 0000000000..8561485b2b --- /dev/null +++ b/tests/gateway/test_heartbeat_watch_restore.py @@ -0,0 +1,105 @@ +"""Restart recovery uses current routing and never borrows another profile's heartbeat.""" +import asyncio +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest + +from gateway.config import GatewayConfig, Platform +from gateway.run import GatewayRunner, _profile_runtime_scope +from gateway.session import SessionStore, SessionSource +from hermes_cli import goals +from hermes_cli.heartbeat import HeartbeatManager +from hermes_state import SessionDB + + +@pytest.mark.asyncio +async def test_restore_retries_persisted_routes_in_their_own_profiles(tmp_path, monkeypatch): + from gateway.run_heartbeat_restore import restore_heartbeat_watches + + home = tmp_path / '.hermes' + named = home / 'profiles' / 'work' + named.mkdir(parents=True) + (named / 'config.yaml').write_text('{}') + monkeypatch.setattr(Path, 'home', lambda: tmp_path) + monkeypatch.setenv('HERMES_HOME', str(home)) + dbs = {str(p): SessionDB(db_path=p / 'state.db') for p in (home, named)} + monkeypatch.setattr(goals, '_DB_CACHE', dbs) + config = GatewayConfig(multiplex_profiles=True) + store = SessionStore(home / 'sessions', config) + entries = [] + try: + for profile, status, topic in [(None, 'active', '11'), ('work', 'active', '22'), + ('work', 'paused', '33'), (None, 'cleared', '44')]: + source = SessionSource(platform=Platform.TELEGRAM, chat_id='chat', + thread_id=topic, profile=profile, scope_id='workspace') + with _profile_runtime_scope(named if profile else home): + entry = store.get_or_create_session(source) + manager = HeartbeatManager(entry.session_id) + manager.set('check ' + topic, 60) + if status == 'paused': + manager.pause() + elif status == 'cleared': + manager.clear() + entries.append(entry) + # Reload the real persisted routing index as a fresh process would. + store.close_all_db_handles() + runner = GatewayRunner.__new__(GatewayRunner) + runner.config = config + runner.session_store = SessionStore(home / 'sessions', config) + runner._heartbeat_watch = {} + runner._start_heartbeat_poller = lambda: None + runner._profile_name_for_source = lambda source: source.profile + runner._adapter_for_source = lambda source: object() + runner._run_in_executor_with_context = asyncio.to_thread + original = runner.session_store.list_sessions + with monkeypatch.context() as patch: + patch.setattr(runner.session_store, 'list_sessions', lambda: (_ for _ in ()).throw(OSError('cold index'))) + await restore_heartbeat_watches(runner) + assert runner._heartbeat_watch == {} + assert original() + # A poller inherited from work must still restore the default profile too. + with _profile_runtime_scope(named): + await restore_heartbeat_watches(runner) + expected = {e.session_key: (e.origin, e.session_id) for e in entries[:2]} + assert runner._heartbeat_watch == expected + with monkeypatch.context() as patch: + patch.setattr('hermes_cli.heartbeat.load_heartbeat', lambda sid: None) + await restore_heartbeat_watches(runner) + assert runner._heartbeat_watch == expected + runner._heartbeat_watch.clear() + await restore_heartbeat_watches(runner) + assert runner._heartbeat_watch == expected + finally: + store.close_all_db_handles() + if 'runner' in locals(): + runner.session_store.close_all_db_handles() + for db in dbs.values(): + db.close() + + +@pytest.mark.asyncio +async def test_startup_arms_retry_poller_even_without_any_watches(monkeypatch): + runner = GatewayRunner.__new__(GatewayRunner) + runner._heartbeat_watch = {} + runner._background_tasks = set() + runner._ensure_hosted_room_worker = AsyncMock() + runner._hosted_room_worker_watcher = AsyncMock() + runner._spawn_supervised = lambda *args: None + runner._start_loop_heartbeat_task = lambda: None + runner.hooks = SimpleNamespace(loaded_hooks=[], emit=AsyncMock()) + runner.adapters = {} + runner._send_update_notification = AsyncMock(return_value=True) + runner.session_store = SimpleNamespace(list_sessions=lambda: []) + runner._run_in_executor_with_context = asyncio.to_thread + monkeypatch.setattr('gateway.channel_directory.build_channel_directory', AsyncMock(return_value={})) + await runner._start_post_connect_services(0) + task = getattr(runner, '_heartbeat_poll_task', None) + try: + assert task is not None and not task.done() + assert runner._heartbeat_watch == {} + finally: + if task: + task.cancel() + await asyncio.gather(task, return_exceptions=True) diff --git a/website/docs/user-guide/features/heartbeat.md b/website/docs/user-guide/features/heartbeat.md index 56895fb1f9..7db2ac843b 100644 --- a/website/docs/user-guide/features/heartbeat.md +++ b/website/docs/user-guide/features/heartbeat.md @@ -21,7 +21,7 @@ They look similar but serve different jobs: | | `/heartbeat` | [`hermes cron`](./cron) | |---|---|---| | Runs in | **This conversation** — full context, memory of the discussion | A fresh isolated session per tick | -| Survives process restart | State survives (SessionDB); firing resumes next time the session is driven | Yes — fully durable scheduler | +| Survives process restart | State survives (SessionDB); gateway watches resume automatically after restart | Yes — fully durable scheduler | | How many | One per session | Unlimited jobs | | Best for | "Keep an eye on X *in this thread* while we work" | Standing jobs, reports, watchdogs, deliveries | @@ -45,6 +45,7 @@ Rule of thumb: if the recurring prompt needs the conversation's context, use `/h - **Missed ticks coalesce.** If the session was busy (or the process wasn't running) through several intervals, you get **one** heartbeat turn, not a backlog. The timer re-anchors on every fire. - **User messages win.** A queued user message always takes priority; the heartbeat waits for the input queue to drain. - **Cache-safe.** The injected prompt is an ordinary user message. No system-prompt mutation, no toolset change. +- **Gateway recovery.** Startup restores active heartbeats using the current persisted conversation and thread routing, in the owning profile. Each poll retries recovery after temporary storage failures or adapter downtime; paused and cleared heartbeats do not restart. No new chat message is required. - **Persistence.** State lives in `SessionDB.state_meta` keyed by `heartbeat:` — it survives `/resume` and rides across context-compression session rotations. Firing requires the owning process (CLI session or gateway) to be running; for schedules that must survive anything, use cron. - **Don't-invent-work guard.** The injected prompt tells the agent to reply briefly and stop when nothing meaningful changed, so an idle heartbeat doesn't generate busywork.