fix: keep admitted heartbeats in their owning conversation
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
@@ -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:<session_id>` 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:<session_id>` 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.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user