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 <chelsealong@126.com> Co-authored-by: liuhao1024 <sunsky.lau@gmail.com> Co-authored-by: fangliquanflq <fangliquan@qq.com>
This commit is contained in:
48
gateway/run_heartbeat_restore.py
Normal file
48
gateway/run_heartbeat_restore.py
Normal file
@@ -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)
|
||||
@@ -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)
|
||||
|
||||
105
tests/gateway/test_heartbeat_watch_restore.py
Normal file
105
tests/gateway/test_heartbeat_watch_restore.py
Normal file
@@ -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)
|
||||
@@ -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:<session_id>` — 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.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user