refactor(agent): pin session activity heartbeat cadence + harden best-effort write
Heartbeat write discipline for the durable SessionDB activity projection: - Pin the cadence in a named constant (SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS = 60s, contract >= 30s, deliberately config-independent so no compression.*/agent.* setting can turn the heartbeat into a high-frequency writer on the contended SessionDB write path). - The write already rides the standard _execute_write patience path via SessionDB.touch_session_activity — verified, now documented in the docstring. - Best-effort hardening: a failed heartbeat write never raises into the agent loop; the bare 'pass' becomes an explicit debug log with traceback. - Tests: direct proof that a heartbeat DB failure doesn't propagate, the cadence constant is pinned >= 30s, and the rate limiter keys off the shared constant (boundary tested on both sides of the window).
This commit is contained in:
@@ -18,6 +18,16 @@ from typing import Any, Mapping, Optional
|
||||
|
||||
ACTIVITY_DESCRIPTION_MAX = 120
|
||||
|
||||
# Durable SessionDB activity heartbeat cadence (seconds between writes per
|
||||
# session). Contract: MUST stay >= 30s — the SessionDB write path is
|
||||
# contended (deadline/patience retry, compression-lock patience), and the
|
||||
# heartbeat is an observation-only projection that never justifies extra
|
||||
# write pressure. This cadence is deliberately a code constant, independent
|
||||
# of any compression.* or agent.* config, so no configuration can turn the
|
||||
# heartbeat into a high-frequency writer. Matches the kanban auto-heartbeat
|
||||
# cadence. force_persist (terminal stamps) is the only bypass.
|
||||
SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS = 60.0
|
||||
|
||||
|
||||
class ActivityProvenance(str, Enum):
|
||||
"""Where a durable/in-memory activity stamp came from."""
|
||||
|
||||
+18
-7
@@ -3680,8 +3680,11 @@ class AIAgent:
|
||||
def _persist_session_activity_if_due(self) -> None:
|
||||
"""Best-effort durable activity heartbeat for SessionDB consumers.
|
||||
|
||||
Rate-limited to one write per 60s per agent (same cadence as the
|
||||
kanban auto-heartbeat). Fail-open: never raises into the agent loop.
|
||||
Cadence is pinned by SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS
|
||||
(>=30s per session, config-independent — see agent/session_activity.py).
|
||||
The write rides the standard SessionDB ``_execute_write`` patience
|
||||
path via ``touch_session_activity``. Fail-open: a failed heartbeat
|
||||
write must NEVER raise into the agent loop (swallow + debug-log).
|
||||
"""
|
||||
session_id = getattr(self, "session_id", None)
|
||||
session_db = getattr(self, "_session_db", None)
|
||||
@@ -3690,14 +3693,17 @@ class AIAgent:
|
||||
touch = getattr(session_db, "touch_session_activity", None)
|
||||
if not callable(touch):
|
||||
return
|
||||
from agent.session_activity import (
|
||||
SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS,
|
||||
normalize_activity_provenance,
|
||||
)
|
||||
|
||||
now_mono = time.monotonic()
|
||||
last_mono = getattr(self, "_session_activity_last_persist_mono", 0.0)
|
||||
if (now_mono - last_mono) < 60.0:
|
||||
if (now_mono - last_mono) < SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS:
|
||||
return
|
||||
self._session_activity_last_persist_mono = now_mono
|
||||
try:
|
||||
from agent.session_activity import normalize_activity_provenance
|
||||
|
||||
touch(
|
||||
session_id,
|
||||
getattr(self, "_last_activity_ts", None),
|
||||
@@ -3707,8 +3713,13 @@ class AIAgent:
|
||||
),
|
||||
)
|
||||
except Exception:
|
||||
# Never let durable heartbeat I/O break the agent loop.
|
||||
pass
|
||||
# Never let durable heartbeat I/O break the agent loop. The
|
||||
# heartbeat is an observation-only projection; the next due
|
||||
# window retries naturally.
|
||||
logger.debug(
|
||||
"session activity heartbeat write failed (ignored)",
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
def _reset_activity_labels_after_turn(self) -> None:
|
||||
"""Drop mid-turn activity labels once the turn is no longer running.
|
||||
|
||||
@@ -116,6 +116,63 @@ def test_touch_activity_persist_errors_are_swallowed(monkeypatch):
|
||||
assert agent._last_activity_desc == "tool completed: terminal (1.0s)"
|
||||
|
||||
|
||||
def test_heartbeat_write_failure_never_propagates_direct(monkeypatch):
|
||||
"""_persist_session_activity_if_due itself must swallow DB failures.
|
||||
|
||||
The heartbeat is best-effort by contract: a SessionDB write failure
|
||||
(locked db, disk error, closed connection) must never raise into the
|
||||
agent loop — it debug-logs and retries naturally on the next due window.
|
||||
"""
|
||||
agent = _agent_with_db()
|
||||
agent._session_db.touch_session_activity.side_effect = OSError("disk gone")
|
||||
monkeypatch.setattr(run_agent.time, "monotonic", lambda: 5000.0)
|
||||
agent._last_activity_ts = 1.0
|
||||
agent._last_activity_desc = "x"
|
||||
|
||||
# Must not raise, despite the DB write blowing up.
|
||||
agent._persist_session_activity_if_due()
|
||||
assert agent._session_db.touch_session_activity.called
|
||||
|
||||
|
||||
def test_heartbeat_cadence_constant_pinned():
|
||||
"""Heartbeat cadence is a config-independent constant and >= 30s.
|
||||
|
||||
The SessionDB write path is contended; the heartbeat must stay
|
||||
low-frequency regardless of compression/agent config.
|
||||
"""
|
||||
from agent.session_activity import (
|
||||
SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS,
|
||||
)
|
||||
|
||||
assert SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS >= 30.0
|
||||
|
||||
|
||||
def test_heartbeat_respects_cadence_constant(monkeypatch):
|
||||
"""The rate limiter must key off the shared cadence constant."""
|
||||
from agent.session_activity import (
|
||||
SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS as INTERVAL,
|
||||
)
|
||||
|
||||
agent = _agent_with_db()
|
||||
mono = {"t": 1000.0}
|
||||
monkeypatch.setattr(run_agent.time, "time", lambda: 1_700_000_000.0)
|
||||
monkeypatch.setattr(run_agent.time, "monotonic", lambda: mono["t"])
|
||||
monkeypatch.delenv("HERMES_KANBAN_TASK", raising=False)
|
||||
|
||||
agent._touch_activity("first")
|
||||
assert agent._session_db.touch_session_activity.call_count == 1
|
||||
|
||||
# Just inside the window: no write.
|
||||
mono["t"] = 1000.0 + INTERVAL - 0.5
|
||||
agent._touch_activity("inside window")
|
||||
assert agent._session_db.touch_session_activity.call_count == 1
|
||||
|
||||
# Just past the window: write.
|
||||
mono["t"] = 1000.0 + INTERVAL + 0.5
|
||||
agent._touch_activity("past window")
|
||||
assert agent._session_db.touch_session_activity.call_count == 2
|
||||
|
||||
|
||||
def test_get_activity_summary_exposes_shared_activity_contract(monkeypatch):
|
||||
agent = _agent_with_db()
|
||||
monkeypatch.setattr(run_agent.time, "time", lambda: 1_700_000_010.0)
|
||||
|
||||
Reference in New Issue
Block a user