diff --git a/run_agent.py b/run_agent.py index 36c90e2f78..3db3e7ff1f 100644 --- a/run_agent.py +++ b/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 diff --git a/tests/run_agent/test_turn_liveness_watchdog.py b/tests/run_agent/test_turn_liveness_watchdog.py index cfba7b7076..fd4d9a5376 100644 --- a/tests/run_agent/test_turn_liveness_watchdog.py +++ b/tests/run_agent/test_turn_liveness_watchdog.py @@ -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()