diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 47ff195ead..074cd3255a 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -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` diff --git a/hermes_state.py b/hermes_state.py index 37b7003adc..b4e4db3d0c 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -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,), diff --git a/tests/agent/test_compression_split_failure_cooldown.py b/tests/agent/test_compression_split_failure_cooldown.py new file mode 100644 index 0000000000..1fb200a8a1 --- /dev/null +++ b/tests/agent/test_compression_split_failure_cooldown.py @@ -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 diff --git a/tests/state/test_compression_lease_refresh_before_publish.py b/tests/state/test_compression_lease_refresh_before_publish.py new file mode 100644 index 0000000000..d99ef33a5d --- /dev/null +++ b/tests/state/test_compression_lease_refresh_before_publish.py @@ -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, + )