From 561897db59e132cd38a4bfd0134b81b9ee948326 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Mon, 7 Sep 2026 13:55:51 -0700 Subject: [PATCH] fix: keep admitted heartbeats in their owning conversation --- gateway/run_goals.py | 1 + gateway/run_heartbeat_acceptance.py | 27 +++++ gateway/run_heartbeat_restore.py | 2 +- gateway/run_turn.py | 6 ++ .../test_heartbeat_execution_ownership.py | 102 ++++++++++++++++++ website/docs/user-guide/features/heartbeat.md | 4 +- 6 files changed, 139 insertions(+), 3 deletions(-) create mode 100644 tests/gateway/test_heartbeat_execution_ownership.py diff --git a/gateway/run_goals.py b/gateway/run_goals.py index 06c01b7a09..ba7cc3eed4 100644 --- a/gateway/run_goals.py +++ b/gateway/run_goals.py @@ -153,6 +153,7 @@ class GatewayGoalsMixin: event = self._synthetic_prompt_event(source, prompt) event.metadata["gateway_session_key"] = quick_key event._heartbeat_execution_started = False + event._heartbeat_session_id = session_id # A pinned route skips topic recovery: no await between the idle # check and adapter claim. FIFO alone never wakes an idle session. try: diff --git a/gateway/run_heartbeat_acceptance.py b/gateway/run_heartbeat_acceptance.py index 7254013007..170ac42622 100644 --- a/gateway/run_heartbeat_acceptance.py +++ b/gateway/run_heartbeat_acceptance.py @@ -9,6 +9,33 @@ import logging logger = logging.getLogger("gateway.run") +async def resolve_heartbeat_owner(runner, event, entry): + """Keep normal reset/topic resolution, then admit only the original lineage.""" + expected = getattr(event, "_heartbeat_session_id", None) + if not expected: + return True + resolved = entry.session_id + if resolved != expected: + def compression_tip(): + return runner.session_store._db_for_key(entry.session_key).get_compression_tip(expected) + + tip = await runner._run_in_executor_with_context(compression_tip) + if tip != resolved: + return False + # Keep a value, not the mutable routing entry: preparation and hooks can yield + # to /new or /stop before the agent runner starts. + event._heartbeat_resolved_session_id = resolved + return heartbeat_owner_is_current(runner, event, entry.session_key) + + +def heartbeat_owner_is_current(runner, event, session_key): + expected = getattr(event, "_heartbeat_resolved_session_id", None) + if not expected: + return True + current = runner.session_store.lookup_by_session_key(session_key) + return current is not None and not current.suspended and current.session_id == expected + + def settle_heartbeat_attempt(event, manager): if not getattr(event, "_heartbeat_execution_started", False): try: diff --git a/gateway/run_heartbeat_restore.py b/gateway/run_heartbeat_restore.py index 0667b9f7d6..724572b78b 100644 --- a/gateway/run_heartbeat_restore.py +++ b/gateway/run_heartbeat_restore.py @@ -27,7 +27,7 @@ async def restore_heartbeat_watches(runner) -> None: with _profile_runtime_scope(home): entries = store.list_sessions() for entry in entries: - if entry.origin is None or not entry.session_id: + if entry.origin is None or not entry.session_id or entry.suspended: continue try: with runner._profile_scope_for_source(entry.origin): diff --git a/gateway/run_turn.py b/gateway/run_turn.py index 6bd9139c43..aba2ce11d1 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -306,6 +306,9 @@ class GatewayTurnMixin: self._cache_session_source(session_key, source) if await asyncio.to_thread(self._is_telegram_topic_lane, source): session_entry = await self._hmwa_heal_telegram_topic_binding(source, session_entry, session_key) + from gateway.run_heartbeat_acceptance import resolve_heartbeat_owner + if not await resolve_heartbeat_owner(self, event, session_entry): + return return source, session_entry, session_key async def _hmwa_heal_telegram_topic_binding(self, source, session_entry, session_key): @@ -1949,6 +1952,9 @@ class GatewayTurnMixin: # Capture the launch session id so post-run compression publication is identity-guarded # (a /new may move session_entry.session_id while the old run is still unwinding). + from gateway.run_heartbeat_acceptance import heartbeat_owner_is_current + if not heartbeat_owner_is_current(self, event, session_key): + return _run_start_session_id = session_entry.session_id _turn_started_monotonic = time.monotonic() # Admission/typing is not execution. All routing, authorization and diff --git a/tests/gateway/test_heartbeat_execution_ownership.py b/tests/gateway/test_heartbeat_execution_ownership.py new file mode 100644 index 0000000000..fb9336c413 --- /dev/null +++ b/tests/gateway/test_heartbeat_execution_ownership.py @@ -0,0 +1,102 @@ +"""Admitted heartbeat work belongs to a conversation, not a reusable route.""" +import asyncio +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest + +from evals.heartbeat_idle_wire import WireAdapter +from gateway.config import GatewayConfig, Platform, PlatformConfig +from gateway.run import GatewayRunner +from gateway.session import SessionSource, SessionStore +from hermes_cli.heartbeat import HeartbeatState, migrate_heartbeat_to_session, save_heartbeat + + +@pytest.mark.asyncio +@pytest.mark.parametrize("boundary", ["reset", "suspend", "compression", "prepare-reset", "prepare-suspend", "unchanged"]) +async def test_admitted_heartbeat_executes_only_in_own_conversation(tmp_path, boundary): + runner = GatewayRunner.__new__(GatewayRunner) + runner.config = GatewayConfig() + runner._running_agents = {} + runner._run_in_executor_with_context = asyncio.to_thread + runner.session_store = SessionStore(tmp_path / "sessions", runner.config) + adapter = WireAdapter(PlatformConfig(enabled=True, typing_indicator=False), Platform.TELEGRAM) + adapter.wire = [] + runner._adapter_for_source = lambda source: adapter + runner._is_telegram_topic_lane = lambda source: False + runner._cache_session_source = lambda *args: None + runner._clear_session_env = lambda tokens: None + source = SessionSource(platform=Platform.TELEGRAM, chat_id="boundary", user_id="owner") + entry = runner.session_store.get_or_create_session(source) + key, old = entry.session_key, entry.session_id + await runner._warm_goals_session_db("test") + save_heartbeat(old, HeartbeatState(prompt="old work", interval_seconds=60, created_at=1)) + executions = [] + + def change(): + if "reset" in boundary: + runner.session_store.reset_session(key) + elif "suspend" in boundary: + runner.session_store.suspend_session(key) + elif boundary == "compression": + child = old + "-child" + runner.session_store._db_for_key(key).publish_compression_child( + parent_session_id=old, child_session_id=child, source="telegram", + require_compression_lease=False, model="offline", model_config={}, + system_prompt="offline", messages=[{"role": "user", "content": "retained"}], + ) + assert migrate_heartbeat_to_session(old, child) + assert runner.session_store.advance_compression_session(key, old, child) + + async def prepare(*args): + if boundary.startswith("prepare-"): + change() + return runner._PreparedTurn([], "", "old work", "old work", None, None), {} + + async def model(**kwargs): + executions.append(kwargs["session_id"]) + raise asyncio.CancelledError # stop before unrelated post-turn delivery + + runner._hmwa_prepare_turn = prepare + runner._run_agent = model + runner.hooks = SimpleNamespace(emit=AsyncMock()) + adapter.set_message_handler(lambda event: runner._handle_message_with_agent(event, source, key, 1)) + try: + await runner._heartbeat_poll_once({key: (source, old)}) + assert key in adapter._session_tasks # prove admission before changing ownership + if not boundary.startswith("prepare-"): + change() + await asyncio.gather(*adapter._background_tasks, return_exceptions=True) + await asyncio.sleep(0) + expected = [old + "-child"] if boundary == "compression" else [old] if boundary == "unchanged" else [] + assert executions == expected + if boundary == "suspend": + assert runner.session_store.peek_session_id(key) != old # normal reset policy still runs + finally: + runner.session_store.close_all_db_handles() + + +@pytest.mark.asyncio +async def test_restart_does_not_restore_suspended_heartbeat(tmp_path): + from gateway.run_heartbeat_restore import restore_heartbeat_watches + + config = GatewayConfig() + store = SessionStore(tmp_path / "sessions", config) + source = SessionSource(platform=Platform.TELEGRAM, chat_id="stopped", user_id="owner") + entry = store.get_or_create_session(source) + save_heartbeat(entry.session_id, HeartbeatState(prompt="stale work", interval_seconds=60, created_at=1)) + store.suspend_session(entry.session_key) + store.close_all_db_handles() + runner = GatewayRunner.__new__(GatewayRunner) + runner.config = config + runner.session_store = SessionStore(tmp_path / "sessions", config) + runner._heartbeat_watch = {} + runner._start_heartbeat_poller = lambda: None + runner._run_in_executor_with_context = asyncio.to_thread + try: + await restore_heartbeat_watches(runner) + assert runner._heartbeat_watch == {} + restored = runner.session_store.lookup_by_session_key(entry.session_key) + assert restored.suspended and restored.session_id == entry.session_id + finally: + runner.session_store.close_all_db_handles() diff --git a/website/docs/user-guide/features/heartbeat.md b/website/docs/user-guide/features/heartbeat.md index aa1ee8ab02..f6b2043c24 100644 --- a/website/docs/user-guide/features/heartbeat.md +++ b/website/docs/user-guide/features/heartbeat.md @@ -45,8 +45,8 @@ 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 and conversation boundaries.** State lives in `SessionDB.state_meta` keyed by `heartbeat:` and follows context-compression session rotations. In the messaging gateway, leaving a conversation through reset, switch, or suspension clears its heartbeat; resuming that archived conversation does not resurrect it. Firing requires the owning process (CLI session or gateway) to be running. +- **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 and suspended conversations do not restart. No new chat message is required. +- **Persistence and conversation boundaries.** State lives in `SessionDB.state_meta` keyed by `heartbeat:` and follows context-compression session rotations. In the messaging gateway, leaving a conversation through reset, switch, or suspension clears its heartbeat; resuming that archived conversation does not resurrect it. Firing requires the owning process (CLI session or gateway) to be running. An already-admitted gateway tick is checked again after session resolution and before agent execution: it may follow a compression child, but cannot carry its old instruction into a reset or switched conversation. - **Execution accounting.** The gateway reserves a due tick at adapter admission. If that exact attempt ends before entering the agent runner (including cancellation or a routing, authorization, emergency-stop, or preparation rejection), it refunds the tick unless the schedule has since changed. Once the agent runner is entered, the fire remains counted even if execution fails or is interrupted. This count is **not** proof of a successful model response or outbound delivery; abrupt process death can prevent the refund callback. - **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.