fix(agent): holder-qualified durable lease cancellation, cooldown ordering (review F4)

A host timeout previously left the timed-out worker holding the durable
per-session compression lock AND refreshing its lease indefinitely, so a
truly hung summary blocked every later compression attempt; and a LATE
successful summary could clear the failure cooldown the host had just
recorded.

Transplant the lease-cancellation invariants from PR #71569
(@ciabata-git): the worker publishes an idempotent, holder-scoped release
hook on the fence once it owns the durable lock (begin_lock_setup /
register_cancelled_lock_release close the acquire→publish race), the
refresher start is serialized against the release path, and the host
invokes the hook on idle timeout, hygiene timeout, and every unwind
(revoke_commit_admission now also releases). ABA safety: the SessionDB
release is holder-qualified (DELETE ... WHERE holder = ?), so a stale
release can never free a replacement holder's lease.

State ordering: the compressor consults a fence-cancellation check BEFORE
clearing the failure cooldown, so a late worker cannot undo the host's
timeout cooldown; the check is installed only for the fenced call and
removed in a finally.

Regression implements the reviewer's exact 5-step scenario: summary
blocked indefinitely → host timeout → a NEW compressor acquires the
durable lock while the old summary is STILL blocked → old worker released
→ it cannot clear cooldown, release the new holder's lease, or publish
stale state.

PR #76354 review, blocking finding 4 / merge gates 4 + 5.

Co-authored-by: ciabata-git <ciabata-git@users.noreply.github.com>
This commit is contained in:
Teknium
2026-08-01 15:21:04 -07:00
parent 971d81f892
commit fdeb09a596
5 changed files with 376 additions and 80 deletions

View File

@@ -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

View File

@@ -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

View File

@@ -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,

View File

@@ -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

View File

@@ -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)