fix(compression): refresh lease in-transaction before publish; arm cooldown on split failure

Two narrow repairs for #97948 symptom B (large-session rotation aborts with
'Compression lease lost before publication' / session_split_failed, then the
next turn re-runs the identical doomed compression):

1. publish_compression_child gains require_lease_refresh: the lease is
   extended inside the same transaction as the expiry check (same conn, no
   TOCTOU), giving a worker whose refresher thread died from transient DB
   failures one final chance to keep its completed work.

2. A failed compression split now records a 60s failure cooldown, so the
   next turn cannot immediately re-trigger the same compression.

Salvaged from #98137 (author: vsd2807). The timeout-reconciliation half of
that PR is NOT carried: it has a blocking review (runtime sid vs persisted
session_key, one-shot check cannot observe a 6-minute commit, no identity
projection) and needs a redesign.
This commit is contained in:
VVV
2026-08-31 11:31:55 +05:30
committed by kshitij
parent 19dd5bcee6
commit 087cc49a26
4 changed files with 276 additions and 0 deletions
+5
View File
@@ -4714,6 +4714,8 @@ def compress_context(
profile_name=_profile_for_child,
compression_lock_holder=_lock_holder,
require_compression_lease=_lock_holder is not None,
require_lease_refresh=_lock_holder is not None,
lease_ttl_seconds=_lock_ttl,
watermark=(
_commit_watermark
if _foreign_tail_ceiling is not None
@@ -4992,6 +4994,9 @@ def compress_context(
)
else:
logger.warning("Session DB compression split failed — new session will NOT be indexed: %s", e)
agent.context_compressor._record_compression_failure_cooldown(
60, f"session_split_failed: {e}",
)
# Compaction-boundary bookkeeping, computed once. `old_session_id` is only
# bound in the rotation branch; in-place leaves it unset. `_boundary_parent`
+15
View File
@@ -7045,6 +7045,8 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
profile_name: str = None,
compression_lock_holder: str = None,
require_compression_lease: bool = True,
require_lease_refresh: bool = False,
lease_ttl_seconds: float = 300.0,
watermark: Optional[int] = None,
watermark_ceiling: Optional[int] = None,
) -> None:
@@ -7069,8 +7071,21 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
The caller captures ``MAX(id)`` immediately BEFORE that flush; only
rows in ``(watermark, watermark_ceiling]`` are foreign concurrent
tail. ``None`` = unbounded (no internal flush happened).
When *require_lease_refresh* is True, the lease is refreshed inside
the same transaction before the expiry check. This gives a refresher
that stopped due to transient DB failures one final chance to extend
the lease, preventing wasted compression work. The refresh uses the
same ``conn`` as the publication, so there is no TOCTOU window.
"""
def _do(conn):
if require_lease_refresh and compression_lock_holder:
conn.execute(
"UPDATE compression_locks SET expires_at = ? "
"WHERE session_id = ? AND holder = ?",
(time.time() + lease_ttl_seconds, parent_session_id,
compression_lock_holder),
)
lock_row = conn.execute(
"SELECT holder, expires_at FROM compression_locks WHERE session_id = ?",
(parent_session_id,),
@@ -0,0 +1,82 @@
"""Regression tests: a failed compression split must arm the failure cooldown.
Extracted verbatim from PR #98137's test file (author: vsd2807); the sibling
timeout-reconciliation tests were not carried because that production path was
not salvaged (blocking review on #98137).
Issue #97948 symptom B: without the cooldown, the turn after a
session_split_failed abort immediately re-runs the identical doomed
compression.
"""
import time
from unittest.mock import MagicMock
def test_split_failure_records_cooldown():
"""After a session split failure, compression failure cooldown is recorded."""
from agent.context_compressor import ContextCompressor
compressor = ContextCompressor.__new__(ContextCompressor)
compressor._summary_failure_cooldown_until = 0.0
compressor._last_summary_error = None
compressor._cooldown_persist_failed = False
compressor._session_db = None
compressor._session_id = ""
compressor._record_compression_failure_cooldown(60, "session_split_failed: test")
assert compressor._summary_failure_cooldown_until > time.monotonic()
assert compressor._last_summary_error == "session_split_failed: test"
def test_cooldown_blocks_automatic_compression():
from agent.context_compressor import ContextCompressor
compressor = ContextCompressor.__new__(ContextCompressor)
compressor._summary_failure_cooldown_until = time.monotonic() + 60.0
compressor._last_summary_error = "session_split_failed"
compressor._cooldown_persist_failed = False
compressor.quiet_mode = True
assert compressor._automatic_compression_blocked_locally() is True
def test_manual_compress_bypasses_cooldown():
"""Manual /compress (force=True) bypasses cooldown — existing behavior preserved."""
from agent.conversation_compression import compress_context
compressor_mock = MagicMock()
compressor_mock._summary_failure_cooldown_until = time.monotonic() + 60.0
compressor_mock._last_summary_error = "session_split_failed"
compressor_mock._cooldown_persist_failed = False
compressor_mock._anti_thrash_recovery_deadline = 0.0
compressor_mock._consecutive_ineffective_compressions = 0
compressor_mock._summary_failure_streak = 0
compressor_mock._last_compaction_boundary_tokens = None
compressor_mock._verify_compaction_cleared_threshold = False
compressor_mock._proactive_prune_rearm_tokens = None
compressor_mock.compression_count = 0
compressor_mock.threshold_tokens = 10000
compressor_mock.context_length = 128000
compressor_mock.tail_token_budget = 25000
compressor_mock.summary_target_ratio = 0.1
compressor_mock._last_cooldown_refresh_was_authoritative = None
agent_mock = MagicMock()
agent_mock.context_compressor = compressor_mock
agent_mock.compression_enabled = True
agent_mock.session_id = "test-session"
agent_mock._session_db = MagicMock()
agent_mock._session_db.try_acquire_compression_lock.return_value = False
messages = [{"role": "user", "content": "hello"}]
result = compress_context(
agent_mock,
messages,
None,
force=True,
task_id="test",
)
assert result is not None
@@ -0,0 +1,174 @@
"""Tests for RC2: pre-publication lease refresh in publish_compression_child.
When the lease refresher stopped due to transient DB failures, the final
pre-publication refresh inside the same transaction gives one last chance
to extend the lease before the expiry check.
"""
import sqlite3
import threading
import time
from unittest.mock import patch
import pytest
from hermes_state import SessionDB, CompressionSessionBusyError
def _setup_db(tmp_path):
db = SessionDB(tmp_path / "state.db")
return db
def _seed_lock(conn, session_id, holder, expired=False):
now = time.time()
conn.execute(
"INSERT INTO compression_locks (session_id, holder, acquired_at, expires_at) VALUES (?, ?, ?, ?)",
(session_id, holder, now, (now - 10.0) if expired else (now + 300.0)),
)
class TestLeaseRefreshBeforePublish:
def test_refresher_stopped_final_refresh_succeeds(self, tmp_path):
db = _setup_db(tmp_path)
db.create_session("parent-1", source="test")
_seed_lock(db._conn, "parent-1", "holder-1", expired=True)
with patch.object(db, "_execute_write", side_effect=lambda fn: fn(db._conn)):
db.publish_compression_child(
parent_session_id="parent-1",
child_session_id="child-1",
source="test",
messages=[{"role": "user", "content": "hello"}],
compression_lock_holder="holder-1",
require_compression_lease=True,
require_lease_refresh=True,
lease_ttl_seconds=300.0,
)
lock = db._conn.execute(
"SELECT expires_at FROM compression_locks WHERE session_id = ?",
("parent-1",),
).fetchone()
assert lock is not None
assert lock[0] > time.time()
parent = db._conn.execute(
"SELECT ended_at FROM sessions WHERE id = ?",
("parent-1",),
).fetchone()
assert parent is not None
assert parent[0] is not None
def test_refresher_stopped_final_refresh_fails_wrong_holder(self, tmp_path):
db = _setup_db(tmp_path)
_seed_lock(db._conn, "parent-1", "other-holder", expired=True)
with patch.object(db, "_execute_write", side_effect=lambda fn: fn(db._conn)):
with pytest.raises(CompressionSessionBusyError, match="lease lost"):
db.publish_compression_child(
parent_session_id="parent-1",
child_session_id="child-1",
source="test",
messages=[{"role": "user", "content": "hello"}],
compression_lock_holder="holder-1",
require_compression_lease=True,
require_lease_refresh=True,
lease_ttl_seconds=300.0,
)
def test_refresher_healthy_no_duplicate_behavior(self, tmp_path):
db = _setup_db(tmp_path)
db.create_session("parent-1", source="test")
now = time.time()
future = now + 300.0
conn = db._conn
conn.execute(
"INSERT INTO compression_locks (session_id, holder, acquired_at, expires_at) VALUES (?, ?, ?, ?)",
("parent-1", "holder-1", now, future),
)
with patch.object(db, "_execute_write", side_effect=lambda fn: fn(db._conn)):
db.publish_compression_child(
parent_session_id="parent-1",
child_session_id="child-1",
source="test",
messages=[{"role": "user", "content": "hello"}],
compression_lock_holder="holder-1",
require_compression_lease=True,
require_lease_refresh=True,
lease_ttl_seconds=300.0,
)
lock = conn.execute(
"SELECT expires_at FROM compression_locks WHERE session_id = ?",
("parent-1",),
).fetchone()
assert lock is not None
assert lock[0] >= future
def test_stale_holder_cannot_refresh_and_publish(self, tmp_path):
db = _setup_db(tmp_path)
_seed_lock(db._conn, "parent-1", "new-holder", expired=False)
with patch.object(db, "_execute_write", side_effect=lambda fn: fn(db._conn)):
with pytest.raises(CompressionSessionBusyError, match="lease lost"):
db.publish_compression_child(
parent_session_id="parent-1",
child_session_id="child-1",
source="test",
messages=[{"role": "user", "content": "hello"}],
compression_lock_holder="old-holder",
require_compression_lease=True,
require_lease_refresh=True,
lease_ttl_seconds=300.0,
)
def test_no_refresh_when_require_lease_refresh_false(self, tmp_path):
db = _setup_db(tmp_path)
_seed_lock(db._conn, "parent-1", "holder-1", expired=True)
with patch.object(db, "_execute_write", side_effect=lambda fn: fn(db._conn)):
with pytest.raises(CompressionSessionBusyError, match="lease lost"):
db.publish_compression_child(
parent_session_id="parent-1",
child_session_id="child-1",
source="test",
messages=[{"role": "user", "content": "hello"}],
compression_lock_holder="holder-1",
require_compression_lease=True,
require_lease_refresh=False,
lease_ttl_seconds=300.0,
)
def test_refresh_and_lease_check_are_atomic(self, tmp_path):
db = _setup_db(tmp_path)
db.create_session("parent-1", source="test")
_seed_lock(db._conn, "parent-1", "holder-1", expired=True)
real_execute_write = SessionDB._execute_write
def intercepted_execute_write(self, fn, patience_s=None):
original_fn = fn
def wrapper(conn):
result = original_fn(conn)
lock = conn.execute(
"SELECT expires_at FROM compression_locks WHERE session_id = ?",
("parent-1",),
).fetchone()
assert lock is not None
assert lock[0] > time.time()
return result
return real_execute_write(self, wrapper, patience_s)
with patch.object(SessionDB, "_execute_write", intercepted_execute_write):
db.publish_compression_child(
parent_session_id="parent-1",
child_session_id="child-1",
source="test",
messages=[{"role": "user", "content": "hello"}],
compression_lock_holder="holder-1",
require_compression_lease=True,
require_lease_refresh=True,
lease_ttl_seconds=300.0,
)