fix: refund heartbeat admissions that never enter agent execution
Bind settlement to the exact adapter task and event. Rejected or cancelled preparation refunds the existing claim; cancellation after the agent runner starts remains counted. Keep profile-scoped callback context and the manager replacement guard. A fire count is not outbound delivery proof. Drop departed routes rather than executing their stale schedule. Slim accounting-invariant salvage of #93174; preserve current direct adapter dispatch instead of reviving its FIFO/inflight implementation. Prior art #92858. Co-authored-by: Finn763 <165816600+Finn763@users.noreply.github.com> Co-authored-by: fangliquanflq <fangliquan@qq.com>
This commit is contained in:
@@ -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()
|
||||
|
||||
+11
-4
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:<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.
|
||||
- **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.
|
||||
- **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
|
||||
|
||||
Reference in New Issue
Block a user