fix: a declined liveness abort must not cancel a pending compression
Closes the #99758 review P1 (andrexibiza): with a generation claim in play, `interrupt()` called `_admit_hard_cancel()` BEFORE the claim was validated at the final mutation edge, and the production `CompressionCommitFence.cancel_before_commit()` irreversibly sets `_cancelled = True` whenever no commit has started. So a watchdog abort that ultimately DECLINED (real progress landed in the window, claim went stale) had already killed the recovered turn's legitimate pending compression commit — `begin_commit()` refuses a cancelled fence forever. Generation authority covered interrupt publication but not the compression-fence mutation that preceded it. Split hard-cancel admission into two halves: - `_wait_for_compression_commit()` runs pre-claim and is NON-mutating: it only blocks when `commit_in_flight` is true (the started-commit branch of the production fence waits for `finish_commit` without cancelling), so the interrupt still publishes only after an in-flight SessionDB mutation has finished — exactly as before. - `_cancel_pending_compression_commit()` runs AFTER `_consume_claim_and_publish_first_state()` survives, so the destructive pending-commit cancellation can never outlive a stale claim. If a commit crossed its boundary in between, it is no longer fence-cancellable and completes on its own. Regression coverage (both use the real `CompressionCommitFence`): - `test_declined_abort_does_not_cancel_pending_compression_commit`: parks the interrupt at the claim-reservation release, lands real progress (G+1), lets the interrupt decline, then proves `fence.begin_commit()` still admits. Red on the pre-fix tree (mutation-checked: the fence was left cancelled). - `test_declined_abort_parks_and_leaves_fence_operational`: the in-flight-commit window variant — activity lands while the interrupt waits on a started commit; the interrupt declines and a fresh `begin_commit()` still admits afterwards. - The round-6 witness (`...resumes_inside_interrupt_publication`) now models an in-flight commit (`commit_in_flight = True`) so its park point stays inside the pre-claim wait, matching the new admission shape. Also updates the `interrupt()` docstring for the deferred destructive cancellation.
This commit is contained in:
92
run_agent.py
92
run_agent.py
@@ -3419,7 +3419,12 @@ class AIAgent:
|
||||
If provided, the agent will include this in its response context.
|
||||
hard_cancel: Mark this as an explicit stop rather than a redirect or
|
||||
incoming-message interrupt. Compression may honor this
|
||||
atomic signal even while ordinary interrupts are masked.
|
||||
atomic signal even while ordinary interrupts are
|
||||
masked. With a generation claim in play, the
|
||||
destructive compression-fence cancellation is
|
||||
deferred until after the claim survives, so a
|
||||
declined abort never cancels a legitimate pending
|
||||
compression.
|
||||
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).
|
||||
@@ -3471,21 +3476,69 @@ class AIAgent:
|
||||
|
||||
# A hard stop and redirect share one lock so /stop cannot race with an
|
||||
# accepted correction and accidentally turn itself into a retry.
|
||||
def _admit_hard_cancel() -> None:
|
||||
# Cancel any pending compression commit, or wait for an
|
||||
# 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.
|
||||
def _wait_for_compression_commit() -> None:
|
||||
# Pre-claim half of hard-cancel admission (#99758 P1): wait out
|
||||
# a commit that ALREADY crossed its boundary, so the interrupt
|
||||
# is published only after the in-flight SessionDB mutation has
|
||||
# finished — but mutate NOTHING. Cancelling a pending commit is
|
||||
# a destructive, irreversible fence mutation (``begin_commit``
|
||||
# refuses a cancelled fence forever), so it must not run while
|
||||
# a generation claim can still be vetoed: an abort that declines
|
||||
# after the fence was cancelled would have killed the recovered
|
||||
# turn's legitimate pending compression. The destructive half
|
||||
# runs in _cancel_pending_compression_commit(), only after the
|
||||
# claim survived the final mutation edge.
|
||||
fence = vars(self).get("_active_compression_commit_fence")
|
||||
if fence is None:
|
||||
return
|
||||
if not getattr(fence, "commit_in_flight", False):
|
||||
# No commit crossed its boundary — nothing to wait out,
|
||||
# and calling cancel_before_commit here WOULD cancel the
|
||||
# pending commit (the production fence's
|
||||
# cancel_before_commit sets _cancelled whenever no commit
|
||||
# has started). Skip it; the destructive half handles it.
|
||||
return
|
||||
cancel_before_commit = getattr(
|
||||
type(fence), "cancel_before_commit", None
|
||||
)
|
||||
if callable(cancel_before_commit):
|
||||
try:
|
||||
# Marks the fence cancelled (or waits out an already
|
||||
# started commit) without setting the hard-stop Event,
|
||||
# which is published only at the final claim edge.
|
||||
# A commit is in flight (it holds the fence lock
|
||||
# through finish_commit), so this call blocks until
|
||||
# the commit finishes and returns False WITHOUT
|
||||
# setting _cancelled — the started-commit branch of
|
||||
# the production fence never cancels.
|
||||
cancel_before_commit(fence)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"Compression hard-cancel fence wait failed",
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
def _cancel_pending_compression_commit() -> None:
|
||||
# Destructive half of hard-cancel admission (#99758 P1): runs
|
||||
# only AFTER the generation claim survived the final mutation
|
||||
# edge, so an abort that declines can never leave the active
|
||||
# compression fence cancelled. Waiting for an in-flight commit
|
||||
# already happened in _wait_for_compression_commit(); if a
|
||||
# commit crossed its boundary in between, it can no longer be
|
||||
# fence-cancelled (it owns the fence until finish_commit and
|
||||
# completes on its own), so only a still-pending commit is
|
||||
# cancelled here.
|
||||
fence = vars(self).get("_active_compression_commit_fence")
|
||||
if fence is None:
|
||||
return
|
||||
if getattr(fence, "commit_in_flight", False):
|
||||
return
|
||||
cancel_before_commit = getattr(
|
||||
type(fence), "cancel_before_commit", None
|
||||
)
|
||||
if callable(cancel_before_commit):
|
||||
try:
|
||||
# Marks the fence cancelled (or waits out a commit
|
||||
# that started between the wait above and now) without
|
||||
# setting the hard-stop Event, which was already
|
||||
# published at the final claim edge.
|
||||
cancel_before_commit(fence)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
@@ -3548,20 +3601,27 @@ class AIAgent:
|
||||
_redirect_lock = getattr(self, "_pending_redirect_lock", None)
|
||||
if _redirect_lock is not None:
|
||||
with _redirect_lock:
|
||||
# The (potentially blocking) compression fence runs BEFORE
|
||||
# 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.
|
||||
# The (potentially blocking) in-flight-commit wait runs
|
||||
# BEFORE the atomic claim/publication edge; the redirect
|
||||
# lock is still held across it, exactly as before, so /stop
|
||||
# cannot race with an accepted correction. The destructive
|
||||
# pending-commit cancellation runs AFTER the claim survives
|
||||
# (#99758 P1) so a declined abort can never cancel the
|
||||
# recovered turn's legitimate compression.
|
||||
if hard_cancel:
|
||||
_admit_hard_cancel()
|
||||
_wait_for_compression_commit()
|
||||
if not _consume_claim_and_publish_first_state():
|
||||
return False
|
||||
if hard_cancel:
|
||||
_cancel_pending_compression_commit()
|
||||
self._pending_redirect = None
|
||||
else:
|
||||
if hard_cancel:
|
||||
_admit_hard_cancel()
|
||||
_wait_for_compression_commit()
|
||||
if not _consume_claim_and_publish_first_state():
|
||||
return False
|
||||
if hard_cancel:
|
||||
_cancel_pending_compression_commit()
|
||||
self._pending_redirect = None
|
||||
|
||||
# Codex app-server owns its model/tool loop and watches a private
|
||||
|
||||
@@ -85,6 +85,15 @@ class _BlockingCommitFence:
|
||||
self.release = threading.Event()
|
||||
self.calls = 0
|
||||
|
||||
@property
|
||||
def commit_in_flight(self) -> bool:
|
||||
# Mirrors the production fence's lock-free phase marker; this
|
||||
# double models an IN-FLIGHT commit, so the interrupt's pre-claim
|
||||
# wait parks inside cancel_before_commit exactly as against a real
|
||||
# started commit (the production started-commit branch blocks until
|
||||
# finish_commit WITHOUT cancelling).
|
||||
return True
|
||||
|
||||
def cancel_before_commit(self, cancel_event=None):
|
||||
# `cancel_event` is accepted (and ignored) to mirror the production
|
||||
# fence signature; publication happens at the final claim edge, so
|
||||
@@ -829,3 +838,148 @@ def test_interrupt_consumes_claim_and_publishes_first_state_atomically():
|
||||
"interrupt state published after competing activity landed: "
|
||||
f"before={published_before_activity}, after={published_after}"
|
||||
)
|
||||
|
||||
|
||||
def test_declined_abort_does_not_cancel_pending_compression_commit():
|
||||
"""#99758 P1 review: a stale liveness claim must not cancel a legitimate
|
||||
pending compression commit when the abort ultimately declines.
|
||||
|
||||
Schedule under test: the watchdog reserves generation G and parks at
|
||||
the claim-reservation release; real progress lands (G+1, claim
|
||||
invalidated) while the interrupt is parked; the interrupt then runs
|
||||
its (no-op for a pending commit) in-flight wait, declines at the final
|
||||
mutation edge, and the REAL CompressionCommitFence must still admit
|
||||
begin_commit() — the fence must NOT be left cancelled by the stale
|
||||
abort authority. The pre-fix tree cancelled the pending fence BEFORE
|
||||
validating the claim, so begin_commit() refused forever.
|
||||
"""
|
||||
from agent.conversation_compression import CompressionCommitFence
|
||||
|
||||
agent = AIAgent.__new__(AIAgent)
|
||||
agent._turn_liveness_activity_generation = 5
|
||||
agent._turn_liveness_abort_claim = None
|
||||
agent._interrupt_requested = False
|
||||
agent._interrupt_message = None
|
||||
agent._tool_interrupt_reason = None
|
||||
agent._hard_interrupt_requested = threading.Event()
|
||||
agent._execution_thread_id = None
|
||||
agent._active_children_lock = threading.Lock()
|
||||
agent._active_children = set()
|
||||
agent.quiet_mode = True
|
||||
|
||||
# A REAL production fence with a PENDING (not started) commit.
|
||||
fence = CompressionCommitFence()
|
||||
agent._active_compression_commit_fence = fence
|
||||
|
||||
# Park the interrupt right after the claim reservation (release #1 of
|
||||
# the activity lock) so real progress can land in the exact window.
|
||||
parking_lock = _ParkingReleaseLock(threading.Lock())
|
||||
parking_lock.park_on_release = 1
|
||||
agent._turn_liveness_activity_lock = parking_lock
|
||||
|
||||
result = {}
|
||||
|
||||
def interrupt_fn():
|
||||
result["ret"] = AIAgent.interrupt(
|
||||
agent,
|
||||
"watchdog: no progress",
|
||||
hard_cancel=True,
|
||||
require_generation=5,
|
||||
)
|
||||
|
||||
interrupt_thread = threading.Thread(target=interrupt_fn)
|
||||
interrupt_thread.start()
|
||||
assert parking_lock.parked.wait(10.0), (
|
||||
"interrupt never reached the claim-reservation boundary"
|
||||
)
|
||||
|
||||
# Real progress lands while the interrupt is parked between the claim
|
||||
# reservation and the final mutation edge.
|
||||
touched = threading.Event()
|
||||
|
||||
def turn_fn():
|
||||
agent._touch_activity("turn resumed")
|
||||
touched.set()
|
||||
|
||||
turn_thread = threading.Thread(target=turn_fn)
|
||||
turn_thread.start()
|
||||
assert touched.wait(10.0), "competing activity never landed"
|
||||
assert agent._turn_liveness_activity_generation == 6
|
||||
|
||||
parking_lock.release_park.set()
|
||||
interrupt_thread.join(10.0)
|
||||
turn_thread.join(10.0)
|
||||
|
||||
# The abort declined: the claim went stale against generation 6.
|
||||
assert result["ret"] is False, "interrupt should have declined"
|
||||
assert not agent._interrupt_requested
|
||||
assert not agent._hard_interrupt_requested.is_set()
|
||||
# THE P1 INVARIANT: the pending compression commit was NOT cancelled
|
||||
# by the declined abort — begin_commit() still admits.
|
||||
assert fence.begin_commit() is True, (
|
||||
"declined liveness abort left the pending compression fence "
|
||||
"cancelled: begin_commit() refused"
|
||||
)
|
||||
fence.finish_commit()
|
||||
|
||||
|
||||
def test_declined_abort_parks_and_leaves_fence_operational():
|
||||
"""#99758 P1, deterministic window variant: park the interrupt inside the
|
||||
wait phase boundary with a REAL fence whose commit is in flight, resume
|
||||
the turn while parked, and prove both that the interrupt declines AND
|
||||
that the fence can serve a fresh begin_commit afterwards."""
|
||||
from agent.conversation_compression import CompressionCommitFence
|
||||
|
||||
agent = AIAgent.__new__(AIAgent)
|
||||
agent._turn_liveness_activity_generation = 5
|
||||
agent._turn_liveness_abort_claim = None
|
||||
agent._interrupt_requested = False
|
||||
agent._interrupt_message = None
|
||||
agent._tool_interrupt_reason = None
|
||||
agent._hard_interrupt_requested = threading.Event()
|
||||
agent._execution_thread_id = None
|
||||
agent._active_children_lock = threading.Lock()
|
||||
agent._active_children = set()
|
||||
agent.quiet_mode = True
|
||||
|
||||
fence = CompressionCommitFence()
|
||||
agent._active_compression_commit_fence = fence
|
||||
|
||||
# Put the fence in the in-flight state so the wait phase blocks in
|
||||
# cancel_before_commit (the started-commit branch waits for
|
||||
# finish_commit without cancelling).
|
||||
assert fence.begin_commit() is True
|
||||
entered = threading.Event()
|
||||
resumed = threading.Event()
|
||||
|
||||
result = {}
|
||||
|
||||
def interrupt_fn():
|
||||
result["ret"] = AIAgent.interrupt(
|
||||
agent,
|
||||
"watchdog: no progress",
|
||||
hard_cancel=True,
|
||||
require_generation=5,
|
||||
)
|
||||
|
||||
interrupt_thread = threading.Thread(target=interrupt_fn)
|
||||
interrupt_thread.start()
|
||||
# Let the interrupt reach the fence wait (blocking on the held lock).
|
||||
time.sleep(0.2)
|
||||
# Real progress lands while the interrupt waits on the in-flight commit.
|
||||
agent._touch_activity("turn resumed mid-wait")
|
||||
resumed.set()
|
||||
# Release the in-flight commit; the interrupt's wait completes, then
|
||||
# the claim check runs and declines.
|
||||
fence.finish_commit()
|
||||
interrupt_thread.join(10.0)
|
||||
|
||||
assert result["ret"] is False, "interrupt should decline after G+1"
|
||||
assert not agent._interrupt_requested
|
||||
assert not agent._hard_interrupt_requested.is_set()
|
||||
# The declined abort must not have cancelled the fence for FUTURE
|
||||
# commits: a fresh begin_commit still admits.
|
||||
assert fence.begin_commit() is True, (
|
||||
"declined liveness abort left the compression fence cancelled"
|
||||
)
|
||||
fence.finish_commit()
|
||||
|
||||
Reference in New Issue
Block a user