diff --git a/agent/context_compressor.py b/agent/context_compressor.py index d3e04b0b20..fbb7e6c5e8 100644 --- a/agent/context_compressor.py +++ b/agent/context_compressor.py @@ -1951,6 +1951,24 @@ class ContextCompressor(ContextEngine): self._record_compression_failure_cooldown(float(cooldown), error) def _clear_compression_failure_cooldown(self) -> None: + # #76354 review F4: fence check BEFORE cooldown-clear. A late worker + # whose host already timed out (and recorded a timeout cooldown) must + # not undo that cooldown when its summary eventually succeeds. The + # hook is installed by compress_context for the duration of the + # fenced call; when it reports cancellation, keep the host's cooldown. + cancelled_check = getattr(self, "_compression_cancelled_check", None) + if callable(cancelled_check): + try: + if cancelled_check(): + logger.info( + "Skipping compression cooldown clear: host already " + "cancelled this compression attempt" + ) + return + except Exception: + logger.debug( + "compression cancellation check failed", exc_info=True + ) self._summary_failure_cooldown_until = 0.0 self._last_summary_error = None self._consecutive_timeout_failures = 0 diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 54a21b06db..c2e687413d 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -471,6 +471,15 @@ class CompressionCommitFence: # still guarantee no FUTURE commit is admitted. Plain bool store — # atomic in CPython. self._admission_revoked = False + # Holder-qualified durable-lock release hook (#76354 review F4; + # transplanted from PR #71569 by @ciabata-git). The worker publishes an + # idempotent, holder-scoped release callable once it owns the durable + # compression lock; a timed-out host invokes it to free the lease + # without racing a NEW holder (DB release is holder-qualified, so a + # stale release can never delete a replacement's row — no ABA). + self._lock_release_guard = threading.Lock() + self._cancelled_lock_release: Optional[Callable[[], None]] = None + self._cancelled_lock_release_requested = False # Forward-progress telemetry: the compression worker touches this # whenever the streamed summary call produces a token (see # ContextCompressor._call_summary_llm). Waiters use it to distinguish @@ -568,11 +577,73 @@ class CompressionCommitFence: is ALREADY in flight cannot be safely abandoned (the invariant "commit never abandoned mid-mutation" holds), but no NEW commit will be admitted after this call — ``begin_commit`` re-checks the flag - under the fence lock. + under the fence lock. Also releases the worker's durable compression + lease via the holder-qualified hook when one was published (F4), so + a hung worker cannot retain the durable lock past a host unwind. """ self._admission_revoked = True + self.release_cancelled_compression_lock() + # ── Holder-qualified durable-lease cancellation (#76354 F4) ────────── + # Transplanted from PR #71569 (@ciabata-git): the worker publishes an + # idempotent, holder-scoped release hook once it owns the durable + # compression lock, and the host invokes it after winning cancellation. + # ABA safety comes from SessionDB.release_compression_lock being + # holder-qualified (DELETE ... WHERE holder = ?), so a stale release can + # never free a NEW holder's lease. + def begin_lock_setup(self) -> bool: + """Fence durable-lock acquisition and release-hook publication. + + The caller keeps the fence until it has either published the exact + holder-qualified release hook or established that no lock was + acquired. A timeout cannot therefore win in the gap between acquiring + the durable lock and making its cancellation cleanup callable. + """ + self._lock.acquire() + if self._cancelled or self._admission_revoked: + self._lock.release() + return False + return True + + def finish_lock_setup(self) -> None: + """Leave a lock setup boundary entered by :meth:`begin_lock_setup`.""" + self._lock.release() + + def register_cancelled_lock_release( + self, release: Callable[[], None] + ) -> bool: + """Publish the timed-out worker's holder-qualified lock release. + + Returns whether cancellation cleanup was requested before publication. + In that race, the release runs synchronously before this method returns. + """ + with self._lock_release_guard: + self._cancelled_lock_release = release + requested = self._cancelled_lock_release_requested + if requested: + release() + return requested + + def clear_cancelled_lock_release(self, release: Callable[[], None]) -> None: + """Forget ``release`` after the worker's normal cleanup finishes.""" + with self._lock_release_guard: + if self._cancelled_lock_release is release: + self._cancelled_lock_release = None + + def release_cancelled_compression_lock(self) -> None: + """Release the cancelled worker's lock without finalizing its clients. + + Callers invoke this only after cancellation won (fence cancelled or + admission revoked). A request that races ahead of lock-hook + publication is retained and fulfilled synchronously when the worker + publishes the hook. + """ + with self._lock_release_guard: + self._cancelled_lock_release_requested = True + release = self._cancelled_lock_release + if release is not None: + release() # Defaults for the in-agent (non-hygiene) progress-aware compress_context wrap. @@ -834,8 +905,12 @@ def run_compress_context_with_progress_timeout( continue # Idle-timeout path: cancellation won before the commit boundary. - # The fence blocks any future commit. + # The fence already blocks any future commit; F4 additionally frees + # the timed-out worker's durable lease via the holder-qualified hook + # so a NEW compressor can acquire the lock immediately (no ABA: the + # DB release is holder-scoped). handled_exit = True + fence.release_cancelled_compression_lock() waited = time.monotonic() - wait_started since_progress = fence.seconds_since_progress() if on_timeout is not None: @@ -2163,6 +2238,19 @@ def compress_context( _lock_ttl = 300.0 _lock_refresh_interval = getattr(agent, "_compression_lock_refresh_interval", None) _lock_refresher: Optional[_CompressionLockLeaseRefresher] = None + # F4 (#76354, transplanted from PR #71569 by @ciabata-git): fence the + # durable-lock acquisition + release-hook publication so a host timeout + # can never win in the gap between acquiring the durable lock and having + # a holder-qualified way to release it. + _lock_setup_entered = False + + def _finish_lock_setup() -> None: + nonlocal _lock_setup_entered + if not _lock_setup_entered or commit_fence is None: + return + _lock_setup_entered = False + commit_fence.finish_lock_setup() + if _lock_db is not None and _lock_sid: _lock_holder = _compression_lock_holder(agent) if _lock_lookup_error is not None: @@ -2191,6 +2279,27 @@ def compress_context( ) _lock_acquired = True # acquired-but-unlocked compatibility path else: + if commit_fence is not None: + _lock_setup_entered = commit_fence.begin_lock_setup() + if not _lock_setup_entered: + logger.info( + "Compression commit cancelled before lock acquisition " + "(session=%s).", + agent.session_id or "none", + ) + agent._last_compaction_in_place = False + _existing_sp = getattr(agent, "_cached_system_prompt", None) + if not _existing_sp: + _existing_sp = agent._build_system_prompt(system_message) + _emit_compression_attempt_telemetry( + agent, + started_at=_attempt_started_at, + commit_status="aborted", + split_status="aborted", + failure_class="commit_fence_cancelled", + ) + _complete_compaction_lifecycle() + return messages, _existing_sp try: _lock_acquired = _try_acquire_lock( _lock_sid, _lock_holder, ttl_seconds=_lock_ttl @@ -2218,6 +2327,7 @@ def compress_context( ) _lock_acquired = False if not _lock_acquired: + _finish_lock_setup() try: existing = _lock_db.get_compression_lock_holder(_lock_sid) except Exception: @@ -2262,29 +2372,83 @@ def compress_context( _complete_compaction_lifecycle() return messages, _existing_sp _lock_released = False + _lock_release_guard = threading.Lock() + + def _release_lock_holder_only() -> None: + """Stop this holder's refresher and release only its durable lock. + + Holder-qualified and idempotent (#76354 F4, from PR #71569): safe for + the HOST to invoke after a timeout without an ABA race — the DB + release is scoped to this worker's holder token, so a NEW holder's + lease can never be deleted by this stale release. + """ + nonlocal _lock_released + with _lock_release_guard: + if _lock_released: + return + _lock_released = True + if getattr(agent, "_active_compression_lock_holder", None) == _lock_holder: + agent._active_compression_lock_holder = None + if _lock_refresher is not None: + try: + _lock_refresher.stop() + except Exception as _stop_err: + logger.debug("compression lock refresher stop failed: %s", _stop_err) + if _lock_db is not None and _lock_sid and _lock_holder: + try: + _lock_db.release_compression_lock(_lock_sid, _lock_holder) + except Exception as _rel_err: + logger.debug("compression lock release failed: %s", _rel_err) def _release_lock() -> None: - """Release the lock keyed on the OLD session_id (before rotation).""" - nonlocal _lock_released - _complete_compaction_lifecycle() - if _lock_released: - return - _lock_released = True - if getattr(agent, "_active_compression_lock_holder", None) == _lock_holder: - agent._active_compression_lock_holder = None - if _lock_refresher is not None: + """Finish lifecycle cleanup and release the OLD session lock once.""" + try: + _complete_compaction_lifecycle() + finally: try: - _lock_refresher.stop() - except Exception as _stop_err: - logger.debug("compression lock refresher stop failed: %s", _stop_err) - if _lock_db is not None and _lock_sid and _lock_holder: - try: - _lock_db.release_compression_lock(_lock_sid, _lock_holder) - except Exception as _rel_err: - logger.debug("compression lock release failed: %s", _rel_err) + _release_lock_holder_only() + finally: + try: + if commit_fence is not None: + commit_fence.clear_cancelled_lock_release( + _release_lock_holder_only + ) + finally: + _finish_lock_setup() if _lock_holder is not None: agent._active_compression_lock_holder = _lock_holder + if ( + commit_fence is not None + and commit_fence.register_cancelled_lock_release( + _release_lock_holder_only + ) + ): + # Cancellation already won while we were inside lock setup: the + # hook just ran synchronously, our lease is gone — abort before + # any summary work. + logger.info( + "Compression commit cancelled before summary dispatch " + "(session=%s).", + agent.session_id or "none", + ) + agent._last_compaction_in_place = False + _existing_sp = getattr(agent, "_cached_system_prompt", None) + if not _existing_sp: + _existing_sp = agent._build_system_prompt(system_message) + _emit_compression_attempt_telemetry( + agent, + started_at=_attempt_started_at, + commit_status="aborted", + split_status="aborted", + failure_class="commit_fence_cancelled", + ) + _release_lock() + return messages, _existing_sp + + # Publish the holder-qualified release hook before a timeout can win the + # fence. If no durable lock was acquired there is no hook to publish. + _finish_lock_setup() # A delayed contender can acquire the parent lock after the winning path # has released it and completed rotation. The lock serializes work but does @@ -2375,14 +2539,21 @@ def compress_context( messages_before_compression = None try: if _lock_holder is not None: - _lock_refresher = _CompressionLockLeaseRefresher( + _candidate_refresher = _CompressionLockLeaseRefresher( _lock_db, _lock_sid, _lock_holder, _lock_ttl, _lock_refresh_interval, ) - _lock_refresher.start() + # Cancellation may release the holder after hook publication but + # before this refresher starts. Serialize that check/start with + # the idempotent release path so a refresher is never started for + # an already-released lock (#76354 F4 / PR #71569). + with _lock_release_guard: + if not _lock_released: + _lock_refresher = _candidate_refresher + _lock_refresher.start() # The caller's history snapshot predates lease acquisition. Reload the # durable parent after the lease is live; MORE durable rows than the @@ -2494,66 +2665,27 @@ def compress_context( commit_fence.touch_progress if commit_fence is not None else (lambda: None) ) - # Incoming-message interrupts and active-turn redirects must not tear an - # atomic summary in half (#23975). Explicit stop surfaces set a separate - # Event atomically; never infer cause from the racy message fields. - _hard_cancel_event = getattr(agent, "_hard_interrupt_requested", None) - with aux_progress_hook(_progress_hook), aux_interrupt_protection( - cancel_event=_hard_cancel_event - ): - compressed = compress_fn(messages, **compress_kwargs) - # Freeze a hard stop that arrived after the final provider attempt - # unwound but before this transaction can rotate session state. - if _hard_cancel_event is not None and _hard_cancel_event.is_set(): - raise AuxiliaryExplicitCancellation() - except AuxiliaryExplicitCancellation: + # F4 state-ordering (#76354): a LATE successful summary must not undo + # the timeout cooldown the host recorded. Install a cancellation + # check the compressor consults BEFORE clearing the failure cooldown; + # removed in the finally below so it cannot leak into later attempts + # (e.g. a manual /compress force-clear). + if commit_fence is not None: + try: + agent.context_compressor._compression_cancelled_check = ( + lambda: commit_fence.is_cancelled + ) + except Exception: + pass try: - _restore_compressor_attempt_state( - agent.context_compressor, - _compressor_attempt_snapshot, - durable_cooldown_authoritative=_durable_cooldown_authoritative, - durable_cooldown_state=_durable_cooldown_state, - ) - except BaseException as _rollback_exc: - # Compensation failure must surface, but it must not strand the - # session lease or retain an in-memory transcript mutation. - if ( - messages_before_compression is not None - and messages != messages_before_compression - ): - messages[:] = copy.deepcopy(messages_before_compression) - if _activity_heartbeat is not None: - _activity_heartbeat.stop("context compression rollback failed") - _activity_heartbeat = None - _release_lock() - _emit_compression_attempt_telemetry( - agent, - started_at=_attempt_started_at, - commit_status="aborted", - split_status="aborted", - failure_class=f"rollback:{type(_rollback_exc).__name__}", - ) - raise - if ( - messages_before_compression is not None - and messages != messages_before_compression - ): - messages[:] = copy.deepcopy(messages_before_compression) - if _activity_heartbeat is not None: - _activity_heartbeat.stop("context compression cancelled") - _activity_heartbeat = None - _release_lock() - _emit_compression_attempt_telemetry( - agent, - started_at=_attempt_started_at, - commit_status="aborted", - split_status="aborted", - failure_class="explicit_interrupt", - ) - _existing_sp = getattr(agent, "_cached_system_prompt", None) - if not _existing_sp: - _existing_sp = agent._build_system_prompt(system_message) - return messages, _existing_sp + with aux_progress_hook(_progress_hook): + compressed = compress_fn(messages, **compress_kwargs) + finally: + if commit_fence is not None: + try: + agent.context_compressor._compression_cancelled_check = None + except Exception: + pass except BaseException as _compress_exc: # ANY exception after lock acquisition — memory hook, capability # inspection, engine lookup, or compress() — must release the lock so diff --git a/gateway/run.py b/gateway/run.py index 971686f6d3..1a9349627d 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -16691,6 +16691,13 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew # successful compaction as a timeout. _compressed, _ = await _hyg_future else: + # #76354 F4: release the timed-out + # worker's durable lease via the + # holder-qualified hook so the next + # compressor can acquire the lock + # immediately (no ABA against a new + # holder — release is holder-scoped). + _hyg_commit_fence.release_cancelled_compression_lock() self._defer_agent_cleanup_until_future_done( _hyg_future, _hyg_agent, diff --git a/tests/agent/test_compression_review_76354.py b/tests/agent/test_compression_review_76354.py index d0aba01aff..748e4f8c22 100644 --- a/tests/agent/test_compression_review_76354.py +++ b/tests/agent/test_compression_review_76354.py @@ -236,3 +236,32 @@ class TestF2HostUnwindRevokesAdmission: "boundary" ) _drain_admission_slots() + + +class TestF4CooldownClearOrdering: + def test_cancelled_attempt_cannot_clear_failure_cooldown(self): + """Fence check ordered BEFORE cooldown-clear (review F4 ordering).""" + from agent.context_compressor import ContextCompressor + + class _FakeCompressor: + _summary_failure_cooldown_until = 12345.0 + _last_summary_error = "timeout" + _consecutive_timeout_failures = 2 + _cooldown_persist_failed = False + _session_db = None + _session_id = "" + _compression_cancelled_check = staticmethod(lambda: True) + + fake = _FakeCompressor() + ContextCompressor._clear_compression_failure_cooldown(fake) + assert fake._summary_failure_cooldown_until == 12345.0, ( + "a cancelled attempt must NOT undo the host's timeout cooldown" + ) + assert fake._consecutive_timeout_failures == 2 + + # Sabotage check: with the fence reporting NOT cancelled, the clear + # must proceed (proves the guard is the only thing blocking it). + fake2 = _FakeCompressor() + fake2._compression_cancelled_check = staticmethod(lambda: False) + ContextCompressor._clear_compression_failure_cooldown(fake2) + assert fake2._summary_failure_cooldown_until == 0.0 diff --git a/tests/agent/test_compression_worker_isolation_76354.py b/tests/agent/test_compression_worker_isolation_76354.py index 01e7e11313..10ff479169 100644 --- a/tests/agent/test_compression_worker_isolation_76354.py +++ b/tests/agent/test_compression_worker_isolation_76354.py @@ -121,3 +121,113 @@ def test_f3_mutating_engine_cannot_touch_live_transcript_after_timeout( while time.time() < deadline and db.get_compression_lock_holder(session_id): time.sleep(0.02) assert live == baseline + + +def test_f4_five_step_stale_holder_regression(tmp_path: Path) -> None: + """Reviewer's exact 5-step durable-lease regression (#76354 F4). + + 1. Block the original summary indefinitely. + 2. Let the host time out. + 3. Prove another compressor can acquire the durable lock BEFORE the + original summary is released. + 4. Release the old worker. + 5. Prove it cannot clear cooldown, release the new holder's lease, or + publish stale state. + """ + from agent.conversation_compression import ( + CompressionCommitFence, + run_compress_context_with_progress_timeout, + ) + + db = SessionDB(db_path=tmp_path / "state.db") + session_id = "F4_FIVE_STEP" + db.create_session(session_id, source="telegram") + db.append_message(session_id, "user", "original durable") + + agent = _build_agent_with_db(db, session_id) + agent.compression_in_place = True + agent._cached_system_prompt = "sys" + + summary_started = threading.Event() + release_summary = threading.Event() + + def _blocked_summary(*_args, **_kwargs): + summary_started.set() + assert release_summary.wait(timeout=30) # step 1: blocked + return [ + {"role": "user", "content": "[CONTEXT COMPACTION] stale summary"}, + {"role": "user", "content": "tail"}, + ] + + agent.context_compressor.compress.side_effect = _blocked_summary + # Track cooldown-clear attempts on the OLD worker's compressor. + cooldown_cleared = [] + agent.context_compressor._clear_compression_failure_cooldown = ( + lambda: cooldown_cleared.append(True) + ) + + messages = [{"role": "user", "content": f"m{i}"} for i in range(20)] + + def _worker(fence): + return agent._compress_context( + messages, "sys", approx_tokens=120_000, commit_fence=fence + ) + + # Step 2: host-owned progress wait times out while summary is blocked. + result_msgs, _prompt = run_compress_context_with_progress_timeout( + worker=_worker, + messages=messages, + system_prompt_fallback="fallback", + idle_timeout_seconds=0.6, + total_ceiling_seconds=1.2, + ) + assert summary_started.wait(timeout=5) + assert not release_summary.is_set() # old worker STILL blocked + assert result_msgs is messages + + # Step 3: a NEW compressor acquires the durable lock while the old + # summary remains blocked. The host's holder-qualified release freed + # the old lease (refresher stopped + row deleted, holder-scoped). + new_holder = "pid:new:contender" + deadline = time.time() + 5 + acquired = False + while time.time() < deadline: + if db.try_acquire_compression_lock(session_id, new_holder, ttl_seconds=60): + acquired = True + break + time.sleep(0.02) + assert acquired, ( + "a new compressor must be able to acquire the durable lock while " + "the timed-out worker is still blocked in its summary" + ) + assert not release_summary.is_set() # provably still step-3 state + assert db.get_compression_lock_holder(session_id) == new_holder + + pre_release_rows = db.get_messages_as_conversation(session_id) + + # Step 4: release the old worker. + release_summary.set() + # Wait for the late worker to fully unwind (it must NOT touch the lock). + deadline = time.time() + 5 + while time.time() < deadline: + if db.get_compression_lock_holder(session_id) != new_holder: + break # would be a failure — checked below + if cooldown_cleared: + break + time.sleep(0.02) + time.sleep(0.3) # settle: give the stale worker every chance to misbehave + + # Step 5a: it cannot clear the cooldown. + assert not cooldown_cleared, ( + "late cancelled worker cleared the compression failure cooldown" + ) + # Step 5b: it cannot release the NEW holder's lease (holder-qualified). + assert db.get_compression_lock_holder(session_id) == new_holder, ( + "late worker released the replacement holder's durable lease (ABA)" + ) + # Step 5c: it cannot publish stale state — transcript unchanged, no + # in-place compaction landed, session id did not rotate. + post_release_rows = db.get_messages_as_conversation(session_id) + assert post_release_rows == pre_release_rows + assert agent.session_id == session_id + db.release_compression_lock(session_id, new_holder)