diff --git a/agent/turn_liveness.py b/agent/turn_liveness.py index ce2070aecb..939fc9b917 100644 --- a/agent/turn_liveness.py +++ b/agent/turn_liveness.py @@ -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) diff --git a/run_agent.py b/run_agent.py index d8501e0ed2..36c90e2f78 100644 --- a/run_agent.py +++ b/run_agent.py @@ -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 ( diff --git a/tests/run_agent/test_turn_liveness_watchdog.py b/tests/run_agent/test_turn_liveness_watchdog.py index d189640176..cfba7b7076 100644 --- a/tests/run_agent/test_turn_liveness_watchdog.py +++ b/tests/run_agent/test_turn_liveness_watchdog.py @@ -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 ): diff --git a/website/docs/user-guide/configuration.md b/website/docs/user-guide/configuration.md index 3989fc20f9..7e087add37 100644 --- a/website/docs/user-guide/configuration.md +++ b/website/docs/user-guide/configuration.md @@ -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