diff --git a/evals/heartbeat_idle_wire.py b/evals/heartbeat_idle_wire.py index 9603ae5285..b83404321d 100644 --- a/evals/heartbeat_idle_wire.py +++ b/evals/heartbeat_idle_wire.py @@ -60,6 +60,7 @@ async def main(base_poller): release = asyncio.Event() async def handler(event): + event._heartbeat_execution_started = True # fake agent execution boundary received.append(event.text) await release.wait() return "wire reply" @@ -108,6 +109,16 @@ async def main(base_poller): assert key in adapter._active_sessions and recovery == [] await drain() print(json.dumps({"phase": "pinned route", "recovery_calls": len(recovery)})) + # Admission can be cancelled before the fake agent boundary is reached. + clock.now += 60 + before = heartbeat.HeartbeatManager("wire-session").state.to_json() + turns_before = len(received) + await runner._heartbeat_poll_once(watch) + await adapter.cancel_session_processing(key) + await drain() + assert len(received) == turns_before + assert heartbeat.HeartbeatManager("wire-session").state.to_json() == before + print(json.dumps({"phase": "cancelled admission", "claim_refunded": True})) runner.config = GatewayConfig() runner.session_store = SessionStore( sessions_dir=Path(os.environ["HERMES_HOME"]) / "sessions", config=runner.config) @@ -122,10 +133,13 @@ async def main(base_poller): "recovery_calls": len(recovery)})) # Route mismatch is a real adapter rejection, not a fabricated return value. adapter._session_key_profile = lambda source: "different-profile" + route_sid = resolved[1].session_id + heartbeat.HeartbeatManager(route_sid).set("check status", 60) clock.now += 60 - before = heartbeat.HeartbeatManager("wire-session").state.to_json() + before = heartbeat.HeartbeatManager(route_sid).state.to_json() await runner._heartbeat_poll_once(watch) - assert heartbeat.HeartbeatManager("wire-session").state.to_json() == before + assert watch[key][1] == route_sid + assert heartbeat.HeartbeatManager(route_sid).state.to_json() == before assert snapshot()["queue_depth"] == 0 print(json.dumps({"phase": "rejected route", "claim_refunded": True})) await adapter.disconnect() diff --git a/gateway/run_goals.py b/gateway/run_goals.py index 2d54e93606..06c01b7a09 100644 --- a/gateway/run_goals.py +++ b/gateway/run_goals.py @@ -127,9 +127,11 @@ class GatewayGoalsMixin: store = getattr(self, "session_store", None) if store is not None: current = store.peek_session_id(quick_key) - if current: - session_id = current - watch[quick_key] = (source, session_id) + if not current: + watch.pop(quick_key, None) + return + session_id = current + watch[quick_key] = (source, session_id) adapter = self._adapter_for_source(source) if adapter is None or not adapter._message_handler: return @@ -150,12 +152,17 @@ class GatewayGoalsMixin: return event = self._synthetic_prompt_event(source, prompt) event.metadata["gateway_session_key"] = quick_key + event._heartbeat_execution_started = False # A pinned route skips topic recovery: no await between the idle # check and adapter claim. FIFO alone never wakes an idle session. try: await adapter.handle_message(event) finally: - if quick_key not in adapter._active_sessions: + task = getattr(adapter, "_session_tasks", {}).get(quick_key) + if task is not None: + from gateway.run_heartbeat_acceptance import settle_heartbeat_attempt + task.add_done_callback(lambda done: settle_heartbeat_attempt(event, mgr)) + elif quick_key not in adapter._active_sessions: mgr.abandon_fire() def _start_heartbeat_poller(self) -> None: diff --git a/gateway/run_heartbeat_acceptance.py b/gateway/run_heartbeat_acceptance.py new file mode 100644 index 0000000000..7254013007 --- /dev/null +++ b/gateway/run_heartbeat_acceptance.py @@ -0,0 +1,17 @@ +"""Accounting at completion of an exact heartbeat admission attempt. + +Execution means entering the gateway's agent runner after turn preparation, +not reserving an adapter slot, successful model completion, or outbound delivery. +The done callback retains the watch's profile ContextVars and manager claim. +""" +import logging + +logger = logging.getLogger("gateway.run") + + +def settle_heartbeat_attempt(event, manager): + if not getattr(event, "_heartbeat_execution_started", False): + try: + manager.abandon_fire() + except Exception: + logger.warning("Failed to refund unexecuted heartbeat", exc_info=True) diff --git a/gateway/run_turn.py b/gateway/run_turn.py index b42ba50eb5..6bd9139c43 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -1951,6 +1951,9 @@ class GatewayTurnMixin: # (a /new may move session_entry.session_id while the old run is still unwinding). _run_start_session_id = session_entry.session_id _turn_started_monotonic = time.monotonic() + # Admission/typing is not execution. All routing, authorization and + # turn preparation gates have passed when the agent runner is entered. + event._heartbeat_execution_started = True agent_result = await self._run_agent( message=message_text, context_prompt=prepared.context_prompt, history=history, source=source, session_id=_run_start_session_id, session_key=session_key, diff --git a/tests/gateway/test_heartbeat_acceptance.py b/tests/gateway/test_heartbeat_acceptance.py new file mode 100644 index 0000000000..7471dfb192 --- /dev/null +++ b/tests/gateway/test_heartbeat_acceptance.py @@ -0,0 +1,88 @@ +"""An adapter slot is not proof a heartbeat reached execution.""" +import asyncio +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest + +from gateway.config import Platform, PlatformConfig +from gateway.run import GatewayRunner +from gateway.session import SessionSource, build_session_key +from hermes_cli.heartbeat import HeartbeatManager, HeartbeatState, save_heartbeat +from evals.heartbeat_idle_wire import WireAdapter + + +@pytest.mark.asyncio +async def test_cancelled_admission_is_refunded_but_started_execution_is_not(): + runner = object.__new__(GatewayRunner) + runner._running_agents = {} + runner._run_in_executor_with_context = asyncio.to_thread + adapter = WireAdapter(PlatformConfig(enabled=True, typing_indicator=False), Platform.TELEGRAM) + adapter.wire = [] + runner._adapter_for_source = lambda source: adapter + source = SessionSource(platform=Platform.TELEGRAM, chat_id='42', user_id='42') + key = build_session_key(source) + watch = {key: (source, 'cancel-test')} + started = asyncio.Event() + + async def model(**kwargs): + started.set() + await asyncio.Event().wait() + + runner._run_agent = model + runner.hooks = SimpleNamespace(emit=AsyncMock()) + runner._hmwa_resolve_session = AsyncMock(return_value=(source, SimpleNamespace(session_id="cancel-test"), key)) + runner._hmwa_prepare_turn = AsyncMock(return_value=( + runner._PreparedTurn([], "", "check", "check", None, None), {})) + runner._clear_session_env = lambda tokens: None + + async def handler(event): + return await runner._handle_message_with_agent(event, source, key, 1) + + adapter.set_message_handler(handler) + await runner._warm_goals_session_db('test') + for cancel_before_start in (True, False): + save_heartbeat('cancel-test', HeartbeatState(prompt='check', interval_seconds=60, created_at=1)) + await runner._heartbeat_poll_once(watch) + task = adapter._session_tasks[key] + if not cancel_before_start: + await asyncio.wait_for(started.wait(), 2) + await adapter.cancel_session_processing(key) + await asyncio.gather(task, return_exceptions=True) + await asyncio.sleep(0) # task callbacks settle the exact attempt + assert HeartbeatManager('cancel-test').state.fire_count == int(not cancel_before_start) + + +@pytest.mark.asyncio +async def test_runner_rejection_does_not_consume_tick_or_overwrite_replacement(): + runner = object.__new__(GatewayRunner) + runner._running_agents = {} + runner._run_in_executor_with_context = asyncio.to_thread + adapter = WireAdapter(PlatformConfig(enabled=True, typing_indicator=False), Platform.TELEGRAM) + adapter.wire = [] + runner._adapter_for_source = lambda source: adapter + source = SessionSource(platform=Platform.TELEGRAM, chat_id='42', user_id='42') + key = build_session_key(source) + watch = {key: (source, 'rejected-test')} + await runner._warm_goals_session_db('test') + + async def denied(event): + return 'not authorized' + + adapter.set_message_handler(denied) + save_heartbeat('rejected-test', HeartbeatState(prompt='check', interval_seconds=60, created_at=1)) + await runner._heartbeat_poll_once(watch) + await asyncio.gather(*adapter._background_tasks) + await asyncio.sleep(0) + assert HeartbeatManager('rejected-test').state.fire_count == 0 + + async def replaced(event): + HeartbeatManager('rejected-test').set('replacement', 120) + return None + + adapter.set_message_handler(replaced) + await runner._heartbeat_poll_once(watch) + await asyncio.gather(*adapter._background_tasks) + await asyncio.sleep(0) + state = HeartbeatManager('rejected-test').state + assert state.prompt == 'replacement' and state.fire_count == 0 diff --git a/tests/gateway/test_heartbeat_poller.py b/tests/gateway/test_heartbeat_poller.py index 07d21af506..70b2914297 100644 --- a/tests/gateway/test_heartbeat_poller.py +++ b/tests/gateway/test_heartbeat_poller.py @@ -57,6 +57,7 @@ async def test_idle_wake_coalesces_intervals_while_adapter_owns_turn(poller): received = [] async def handler(event): + event._heartbeat_execution_started = True # fake agent execution boundary received.append(event) started.set() await release.wait() @@ -96,6 +97,7 @@ async def test_unavailable_or_busy_session_leaves_persisted_tick_due(poller): await runner._heartbeat_poll_once(watch) async def handler(event): + event._heartbeat_execution_started = True # fake agent execution boundary return None adapter.set_message_handler(handler) diff --git a/tests/gateway/test_heartbeat_watch_lifecycle.py b/tests/gateway/test_heartbeat_watch_lifecycle.py index f09bf34830..1d8ecd30cf 100644 --- a/tests/gateway/test_heartbeat_watch_lifecycle.py +++ b/tests/gateway/test_heartbeat_watch_lifecycle.py @@ -74,3 +74,6 @@ async def test_watch_tracks_rotated_route_owner_without_reviving_parent(): assert len(events) == 1 assert watch['route'][1] == 'child' assert HeartbeatManager('parent').state is None + runner.session_store.peek_session_id = lambda key: None + await runner._heartbeat_poll_once(watch) + assert not watch # a departed route never falls back to its stale heartbeat diff --git a/website/docs/user-guide/features/heartbeat.md b/website/docs/user-guide/features/heartbeat.md index 7db2ac843b..aa1ee8ab02 100644 --- a/website/docs/user-guide/features/heartbeat.md +++ b/website/docs/user-guide/features/heartbeat.md @@ -46,7 +46,8 @@ Rule of thumb: if the recurring prompt needs the conversation's context, use `/h - **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. +- **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. +- **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. ## Example