fix: publish watchdog settlement only after the abort commits

Closes the #95663 round-8 review blocker (false settlement before
commit veto): the pre-commit surface (`_surface_stall`) logged
"Force-aborting the turn and stopping lease renewal" and warned the
user "aborting it so the session can recover" BEFORE `_commit_abort`
could veto — so a turn that resumed during the warning window (or an
exceptional interrupt path that declines fail-closed) was reported as
force-aborted with lease stopped while it actually continued running.

- Split the surface: `_surface_stall` is now observational only ("no
  progress for Ns; attempting recovery"), and the definitive
  aborted/lease-stopped settlement moves to a new
  `_surface_committed_abort` that runs only after `_commit_abort`
  succeeds and the turn lease is deactivated.
- Rate-limit repeated pre-commit surfaces per observed generation: a
  turn whose aborts keep declining no longer re-logs an ERROR and
  re-warns the user every poll interval.
- Add the committed-path regression test
  (`test_watchdog_publishes_definitive_settlement_only_after_commit`)
  and extend the declined-path witness
  (`...resumes_during_warning`) to assert no committed-abort or
  definitive pre-commit claim appears when the abort is vetoed. Both
  fail on the pre-fix tree (mutation-checked).
- Document the `_interrupt_turn` lease-loss asymmetry (fires
  unconditionally, no generation claim — losing the lease means the
  process no longer owns the session).
- Trim review-round archaeology from comments/docstrings (keep the
  WHY, drop the round numbering), and drop the dead
  `cancel_event` compat note from the test fence.
- Document `agent.turn_liveness` in the configuration guide.

On top of PR #95663 by Finn763 (cherry-picked with authorship
preserved).
This commit is contained in:
kshitijk4poor
2026-09-01 02:02:13 +05:30
committed by kshitij
parent 0fe7abe37a
commit c394b005fc
4 changed files with 208 additions and 76 deletions
+72 -21
View File
@@ -34,25 +34,24 @@ startup, and a bogus value can never silently disable the watchdog
(``NaN``) or freeze the watcher thread (``Inf`` poll).
Race safety (#95663 review): the watchdog samples the activity clock and
binds the abort decision to the observed ``(generation, timestamp)``
pair. The commit callback revalidates that pair under the *same* lock
binds the abort decision to the observed ``(generation, timestamp)`` pair.
The commit callback revalidates that pair under the *same* lock
``AIAgent._touch_activity`` uses to stamp the clock, so a turn that
resumed while the stall was being logged/emitted is never hard-cancelled
— it continues and its lease keeps renewing. Round-3 carried the
revalidated generation into ``AIAgent.interrupt``
(``require_generation``). Round-4 makes that generation an actual
— it continues and its lease keeps renewing. The revalidated generation
is carried into ``AIAgent.interrupt`` (``require_generation``) as a
cancellation claim consumed at the final mutation edge: ``interrupt``
reserves the claim under the activity lock, ``_touch_activity``
invalidates the reservation the instant real progress lands, and the
claim survives every blocking boundary, including the compression
commit fence. Round-6 consumes the claim and publishes the first
observable interrupt state in ONE activity-lock critical section, so a
turn that resumes can only ever interleave before that section (the
reservation is invalidated and the abort declines) or after it (the
interrupt already committed under the lock) — never between "claim
consumed" and "state published". A turn that resumes anywhere in the
window is never hard-cancelled, and an exceptional interrupt path
declines the abort fail-closed instead of mutating interrupt state.
commit fence. Claim consumption and the first observable interrupt
state publish in ONE activity-lock critical section, so a turn that
resumes can only ever interleave before that section (the reservation
is invalidated and the abort declines) or after it (the interrupt
already committed under the lock) — never between "claim consumed" and
"state published". A turn that resumes anywhere in the window is never
hard-cancelled, and an exceptional interrupt path declines the abort
fail-closed instead of mutating interrupt state.
"""
from __future__ import annotations
@@ -217,6 +216,13 @@ class TurnLivenessWatchdog:
return
if snapshot.idle_seconds < self._timeout_s:
continue
# Pre-commit surface is OBSERVATIONAL only: it reports the
# stall and that a recovery attempt is beginning. It must not
# claim the abort or the lease withdrawal has committed — the
# next operation can still veto the outcome. The definitive
# aborted/lease-stopped settlement is published by
# _surface_committed_abort only after _commit_abort succeeds
# and the turn is deactivated (#95663 review).
self._surface_stall(snapshot)
# Commit point: bind the abort to the sampled generation/ts
# and revalidate under the lock shared with `_touch_activity`.
@@ -225,8 +231,7 @@ class TurnLivenessWatchdog:
# keeps renewing. The commit also carries the revalidated
# generation into the interrupt path, which reserves it as a
# claim, survives every blocking boundary (compression
# fence), and consumes it at the final mutation edge
# immediately before publication (#95663 round-4) — progress
# fence), and consumes it at the final mutation edge — progress
# landing anywhere in that window declines the abort.
if not self._commit_abort(snapshot, self._abort_message(snapshot)):
continue
@@ -235,6 +240,7 @@ class TurnLivenessWatchdog:
# issue's "lease keeps renewing" masking). The TTL expiry then
# lets stale-turn cleanup reclaim the row.
self._deactivate_turn()
self._surface_committed_abort(snapshot)
return
def _sample(self) -> Optional[ActivitySnapshot]:
@@ -263,18 +269,33 @@ class TurnLivenessWatchdog:
)
def _surface_stall(self, snapshot: ActivitySnapshot) -> None:
"""Log the stall loudly and emit a UI-visible warning.
"""Observationally surface the stall: log it loudly and emit a
UI-visible warning that a recovery attempt is beginning.
The lock is deliberately NOT held here: acting on the observation
happens at the commit point, after this surface window, where the
observed generation/timestamp is revalidated.
Deliberately does NOT claim the abort or the lease withdrawal has
committed: the next operation (``_commit_abort``) can still veto
the outcome when the turn resumed while this surface window was
open. The definitive aborted/lease-stopped settlement is
published by :meth:`_surface_committed_abort` only after the
abort wins and the turn is deactivated.
Rate-limited: a turn whose aborts keep declining (resumed
activity, exceptional interrupt path) must not re-log and
re-warn every poll interval — the first surface carries the
signal, repeats are suppressed until activity actually moves
again (a new generation re-arms the surface).
"""
generation = snapshot.generation
if getattr(self, "_last_surfaced_generation", None) == generation:
return
self._last_surfaced_generation = generation
session_id = getattr(self._agent, "session_id", None) or self._session_id
last_desc = getattr(self._agent, "_last_activity_desc", None)
logger.error(
"Turn liveness watchdog fired for session %s: "
"no progress for %.1fs (last activity: %r). "
"Force-aborting the turn and stopping lease renewal (#95548).",
"Attempting recovery: force-interrupting the turn and "
"stopping lease renewal if it cannot resume (#95548).",
session_id,
snapshot.idle_seconds,
last_desc,
@@ -286,7 +307,37 @@ class TurnLivenessWatchdog:
emit_warning(
"⚠️ This turn stopped making progress "
f"({int(snapshot.idle_seconds)}s without activity); "
"aborting it so the session can recover."
"attempting recovery so the session can continue."
)
except Exception:
logger.debug("Failed to emit turn liveness warning", exc_info=True)
def _surface_committed_abort(self, snapshot: ActivitySnapshot) -> None:
"""Publish the definitive settlement AFTER the abort has authority.
Runs only once ``_commit_abort`` succeeded (the interrupt was
published) and the turn lease was deactivated: the turn IS
force-aborted and lease renewal IS stopped, so stating that is
now true. Separated from the pre-commit surface so a declined
abort never reports a committed outcome (#95663 review).
"""
session_id = getattr(self._agent, "session_id", None) or self._session_id
logger.error(
"Turn liveness watchdog aborted turn for session %s: "
"no progress for %.1fs; turn interrupted and lease renewal "
"stopped (#95548).",
session_id,
snapshot.idle_seconds,
)
emit_warning = getattr(self._agent, "_emit_warning", None)
if not callable(emit_warning):
return
try:
emit_warning(
"⚠️ Turn aborted by the liveness watchdog "
f"({int(snapshot.idle_seconds)}s without activity); "
"lease renewal stopped so the session can be reclaimed. "
"You can retry your message."
)
except Exception:
logger.debug("Failed to emit committed-abort warning", exc_info=True)
+52 -51
View File
@@ -3422,17 +3422,16 @@ class AIAgent:
atomic signal even while ordinary interrupts are masked.
tool_reason: Trusted fixed category safe to expose in tool output.
Arbitrary diagnostic or caller text belongs in message.
require_generation: Optional activity-generation claim (#95663
round-3/round-4/round-6 review). When set, the
interrupt is published only if the turn's activity
generation still equals this value at the final
mutation edge. The claim is RESERVED under the
activity lock — ``_touch_activity`` invalidates the
reservation the instant real progress lands —
survives every blocking boundary in between
(including the compression commit fence), and is
CONSUMED in ONE lock critical section together with
the first observable publication
require_generation: Optional activity-generation claim (#95663).
When set, the interrupt is published only if the
turn's activity generation still equals this value
at the final mutation edge. The claim is RESERVED
under the activity lock — ``_touch_activity``
invalidates the reservation the instant real
progress lands — survives every blocking boundary in
between (including the compression commit fence),
and is CONSUMED in ONE lock critical section
together with the first observable publication
(``_interrupt_requested`` / ``_interrupt_message`` /
``_tool_interrupt_reason`` and the hard-cancel
event). If the turn resumed in the window, the call
@@ -3455,14 +3454,13 @@ class AIAgent:
running_agent.interrupt(new_message.text)
"""
if require_generation is not None:
# Round-4 (#95663): RESERVE the abort's generation claim under
# the SAME lock `_touch_activity` stamps the clock with. Real
# progress invalidates the reservation the instant it lands,
# and the claim is CONSUMED at the final mutation edge — after
# every blocking boundary — in ONE critical section with the
# first observable publication (#95663 round-6). A resumed turn
# therefore abandons the abort instead of being hard-cancelled
# by a stale proof.
# RESERVE the abort's generation claim under the SAME lock
# `_touch_activity` stamps the clock with. Real progress
# invalidates the reservation the instant it lands, and the
# claim is CONSUMED at the final mutation edge — after every
# blocking boundary — in ONE critical section with the first
# observable publication. A resumed turn therefore abandons
# the abort instead of being hard-cancelled by a stale proof.
with self._liveness_activity_lock():
if (
getattr(self, "_turn_liveness_activity_generation", 0)
@@ -3478,7 +3476,7 @@ class AIAgent:
# in-flight commit to finish. This call deliberately publishes
# NOTHING observable: the generation claim and the first
# interrupt state (including the hard-stop event) commit
# together at the final mutation edge below (#95663 round-6).
# together at the final mutation edge below.
fence = vars(self).get("_active_compression_commit_fence")
cancel_before_commit = getattr(
type(fence), "cancel_before_commit", None
@@ -3507,30 +3505,25 @@ class AIAgent:
_hard_event.set()
def _consume_claim_and_publish_first_state() -> bool:
# Final mutation edge (#95663 round-6): when a generation claim
# is in play, claim consumption and the FIRST observable
# interrupt publication are ONE activity-lock critical section.
# Round-4 consumed the claim under the lock, released it, and
# only then assigned
# `_interrupt_requested` / `_interrupt_message` /
# `_tool_interrupt_reason` — a turn that resumed in that
# consume→publication window (real progress, generation G+1)
# was still hard-cancelled by an already-consumed claim. Now
# the abort's commit point and its first publication share the
# lock `_touch_activity` stamps the clock with, so the
# generation winner is total: either the claim survives and
# the interrupt state commits under the lock BEFORE any later
# activity stamp, or the stamp landed first and the abort
# declines without publishing anything.
# Final mutation edge: when a generation claim is in play,
# claim consumption and the FIRST observable interrupt
# publication are ONE activity-lock critical section — the
# same lock `_touch_activity` stamps the clock with. The
# generation winner is therefore total: either the claim
# survives and the interrupt state commits under the lock
# BEFORE any later activity stamp, or the stamp landed first
# and the abort declines without publishing anything. (A
# consume-then-release-then-publish split would let a turn
# that resumed in the consume→publication window be
# hard-cancelled by an already-consumed claim.)
if require_generation is None:
# No claim to race the activity clock against: publish
# exactly as round-4 did, WITHOUT touching the liveness
# lock. ``AIAgent`` stand-ins used by unrelated suites
# (e.g. the start-order gate `_Stub`) do not carry the
# liveness seam, and an unconditional
# WITHOUT touching the liveness lock. ``AIAgent``
# stand-ins used by unrelated suites (e.g. the
# start-order gate `_Stub`) do not carry the liveness
# seam, and an unconditional
# ``_liveness_activity_lock()`` acquisition here
# regressed them (AttributeError caught by round-6
# exact-head CI).
# regresses them with AttributeError.
_publish_interrupt_state()
return True
with self._liveness_activity_lock():
@@ -3556,9 +3549,9 @@ class AIAgent:
if _redirect_lock is not None:
with _redirect_lock:
# The (potentially blocking) compression fence runs BEFORE
# the atomic claim/publication edge (#95663 round-6); the
# redirect lock is still held across the fence, exactly as
# before, so /stop cannot race with an accepted correction.
# the atomic claim/publication edge; the redirect lock is
# still held across the fence, exactly as before, so /stop
# cannot race with an accepted correction.
if hard_cancel:
_admit_hard_cancel()
if not _consume_claim_and_publish_first_state():
@@ -4337,7 +4330,9 @@ class AIAgent:
)
# Lazy per-instance lock (inline so bare doubles like
# types.SimpleNamespace fixtures keep working — see
# types.SimpleNamespace fixtures keep working — they bind
# _touch_activity without the class, so they cannot call
# self._liveness_activity_lock(); see
# tests/run_agent/test_session_activity_persist.py).
_clock_lock = getattr(self, "_turn_liveness_activity_lock", None)
if _clock_lock is None:
@@ -4350,12 +4345,11 @@ class AIAgent:
self._last_activity_ts = time.time()
self._last_activity_desc = bound_activity_description(desc)
self._last_activity_provenance = normalize_activity_provenance(provenance)
# Round-4 (#95663): real progress invalidates any reserved
# abort claim. A watchdog interrupt that is still in flight
# (e.g. parked inside the compression commit fence) must
# abandon itself at the final mutation edge instead of
# publishing against a generation the turn has already left
# behind.
# Real progress invalidates any reserved abort claim. A watchdog
# interrupt that is still in flight (e.g. parked inside the
# compression commit fence) must abandon itself at the final
# mutation edge instead of publishing against a generation the
# turn has already left behind.
self._turn_liveness_abort_claim = None
if os.environ.get("HERMES_KANBAN_TASK"):
try:
@@ -9229,6 +9223,13 @@ class AIAgent:
)
def _interrupt_turn(message: str) -> None:
# Lease-loss interrupts fire UNCONDITIONALLY (no
# require_generation claim): losing the durable lease
# means this process no longer owns the session, so
# the turn must stop regardless of activity-clock
# progress. The generation-claim machinery is the
# liveness watchdog's only — its stalls can be
# spuriously stale, a lost lease cannot.
nonlocal durable_turn_lease_interrupt_message
with durable_turn_lease_activity_lock:
if (
+71 -4
View File
@@ -71,7 +71,7 @@ class _DB:
class _BlockingCommitFence:
"""Controllable compression commit fence for the round-4 witness.
"""Controllable compression commit fence for the stall-abort witness.
``cancel_before_commit`` parks the interrupt thread AFTER it has passed
the internal ``require_generation`` comparison, simulating the unbounded
@@ -86,9 +86,9 @@ class _BlockingCommitFence:
self.calls = 0
def cancel_before_commit(self, cancel_event=None):
# The round-3 code passes the hard-stop Event; the round-4 code
# passes nothing (publication moved to the final claim edge).
# Accept both so the same witness is valid across the transition.
# `cancel_event` is accepted (and ignored) to mirror the production
# fence signature; publication happens at the final claim edge, so
# nothing observable is set here.
self.calls += 1
self.entered.set()
assert self.release.wait(10.0), "fence was never released"
@@ -415,11 +415,78 @@ def test_watchdog_declines_abort_when_activity_resumes_during_warning(
and "stalled-session" in record.getMessage()
for record in caplog.records
)
# …but the surface was OBSERVATIONAL only: the declined abort must
# never publish a committed-abort or lease-stop settlement
# (#95663 review — false settlement before commit veto). On the
# pre-fix tree the pre-commit surface itself logged the definitive
# "Force-aborting … stopping lease renewal" outcome — this assertion
# is what made that witness red.
assert not any(
"watchdog aborted turn" in record.getMessage()
for record in caplog.records
), "declined abort published a committed-abort settlement"
assert not any(
"Force-aborting" in record.getMessage()
for record in caplog.records
), "pre-commit surface published the definitive abort outcome"
# …and the lease kept renewing through the resumed turn.
assert len(db.refresh_times) >= 1
assert db.events[-1][0] == "release"
def test_watchdog_publishes_definitive_settlement_only_after_commit(
watchdog_config, monkeypatch, caplog
):
"""#95663 review: the definitive aborted/lease-stopped settlement is
published only AFTER the abort has authority (commit succeeded and
the turn lease was deactivated). The committed path must show the
settlement; the pre-commit surface must not claim it."""
db = _DB()
agent = _agent_with_db(db)
warnings = []
def stalled_loop(_agent, _message, _system, history, *_args, **_kwargs):
while not _agent._interrupt_requested:
time.sleep(0.005)
return {
"final_response": "aborted",
"messages": history,
"api_calls": 0,
"completed": False,
"interrupted": True,
}
agent._emit_warning = lambda msg: warnings.append(msg)
with caplog.at_level(logging.ERROR, logger="agent.turn_liveness"):
result = _run_turn(agent, stalled_loop, monkeypatch)
assert result["interrupted"] is True
# Pre-commit surface: observational, recovery-attempt language.
assert any(
"Turn liveness watchdog fired" in record.getMessage()
and "Attempting recovery" in record.getMessage()
for record in caplog.records
)
# Definitive settlement: only present because the abort committed.
assert any(
"watchdog aborted turn" in record.getMessage()
and "lease renewal stopped" in record.getMessage()
for record in caplog.records
), "committed abort did not publish the definitive settlement"
# User-visible warnings follow the same split: first observational,
# then (and only then) the committed outcome.
assert any("attempting recovery" in w for w in warnings)
assert any("Turn aborted by the liveness watchdog" in w for w in warnings)
# Ordering: the committed-abort warning came after the recovery one.
recovery_idx = next(i for i, w in enumerate(warnings) if "attempting recovery" in w)
aborted_idx = next(
i for i, w in enumerate(warnings) if "Turn aborted by the liveness watchdog" in w
)
assert aborted_idx > recovery_idx
assert db.events[-1][0] == "release"
def test_watchdog_declines_abort_when_activity_resumes_after_revalidation(
watchdog_config, monkeypatch, caplog
):
+13
View File
@@ -1796,6 +1796,19 @@ agent:
The same gate also enables **result-reference stubbing**: when a re-issued identical tool call returns a byte-identical fresh result, the duplicate payload enters context as a short reference stub pointing at the earlier result (tool name, `tool_call_id`, an args summary, and — if the first result was persisted to disk — its spillover path) instead of repeating the full output. The tool still executes every time, so polling semantics are preserved: a changed result always flows through whole. Results under 512 characters, error results, and multimodal results are never stubbed, and pollers *are* stubbed (an unchanged poll is exactly the case where the duplicate payload carries no information).
### Turn liveness watchdog
`agent.turn_liveness` bounds how long a conversation turn may make **no observable progress** before Hermes force-recovers it. The watchdog keys off the activity clock (the same signal that stamps API waits, stream tokens, and tool heartbeats — lease renewal never counts), so a turn that silently wedges mid-flight (observed as issue #95548: no tool execution, no API call, no error, but the session stays "busy" indefinitely) is surfaced loudly, interrupted so it unwinds as a retriable interrupted turn, and — when the interrupt cannot unwind the wedge — its durable turn lease stops renewing so stale-turn cleanup can reclaim the session instead of it hanging until the process is killed.
```yaml
agent:
turn_liveness:
timeout_s: 600.0 # idle bound; <= 0 disables the watchdog
poll_s: 15.0 # sampling interval (seconds)
```
Legitimately slow work is not penalized: streaming responses, tool heartbeats (every 30s while a tool runs), and approval waits all keep touching the clock, so only a turn making *zero* progress for the full bound fires the watchdog. Invalid values (a typo, `NaN`, `Inf`, non-positive `poll_s`) log a warning and fall back to the defaults — they never crash startup or silently disable the watchdog. A fired abort reports the stall as it begins recovery, and publishes the definitive aborted/lease-stopped outcome only once the interrupt has actually committed.
## TTS Configuration
```yaml