diff --git a/agent/context_compressor.py b/agent/context_compressor.py index 53cfad134c..a3daa73b67 100644 --- a/agent/context_compressor.py +++ b/agent/context_compressor.py @@ -2656,6 +2656,7 @@ class ContextCompressor(ContextEngine): self.get_active_compression_failure_cooldown() self._load_fallback_compression_streak() self._load_ineffective_compression_count() + self._load_anti_thrash_recovery_deadline() self._load_proactive_prune_rearm_tokens() def on_session_start(self, session_id: str, **kwargs) -> None: @@ -2821,6 +2822,45 @@ class ContextCompressor(ContextEngine): except Exception as exc: logger.debug("compression ineffective count persist failed (non-sqlite): %s", exc) + def _load_anti_thrash_recovery_deadline(self) -> None: + """Restore the durable recovery deadline (wall-clock epoch, #100185). + + Missing/absent storage leaves the in-memory clock disarmed, so the + next blocked evaluation arms a full fresh window (#54923). + """ + session_db = getattr(self, "_session_db", None) + session_id = getattr(self, "_session_id", "") + getter = getattr(session_db, "get_compression_recovery_deadline", None) + if not session_id or not callable(getter): + return + try: + stored = getter(session_id) + self._anti_thrash_recovery_deadline = max( + 0.0, + float(stored) if isinstance(stored, (int, float, str)) else 0.0, + ) + except (TypeError, ValueError, sqlite3.Error) as exc: + logger.debug("compression recovery deadline lookup failed: %s", exc) + except Exception as exc: + logger.debug("compression recovery deadline lookup failed (non-sqlite): %s", exc) + + def _set_anti_thrash_recovery_deadline(self, deadline: float) -> None: + """Set the recovery deadline, persisting on change only (0 = disarmed).""" + if deadline == self._anti_thrash_recovery_deadline: + return + self._anti_thrash_recovery_deadline = deadline + session_db = getattr(self, "_session_db", None) + session_id = getattr(self, "_session_id", "") + setter = getattr(session_db, "set_compression_recovery_deadline", None) + if not session_id or not callable(setter): + return + try: + setter(session_id, deadline) + except sqlite3.Error as exc: + logger.debug("compression recovery deadline persist failed: %s", exc) + except Exception as exc: + logger.debug("compression recovery deadline persist failed (non-sqlite): %s", exc) + def _record_ineffective_compression_verdict(self, count: int) -> None: """Set the anti-thrash strike counter, keeping the durable copy in sync. @@ -3978,21 +4018,34 @@ class ContextCompressor(ContextEngine): # the worst case in the truly-incompressible state is one compaction # attempt per recovery window — bounded, not thrash. # - # The clock is armed lazily on the first BLOCKED evaluation rather - # than persisted at trip time: a fresh process that loads a durable - # tripped counter (#69872) therefore starts a full window blocked, - # preserving the restart-must-not-disarm contract (#54923). + # The clock is armed lazily on the first BLOCKED evaluation and + # persisted on the session row (#100185): a fresh process/compressor + # that loads a durable tripped counter (#69872) with no stored + # deadline starts a full window blocked, preserving the + # restart-must-not-disarm contract (#54923) — but one that loads an + # already-armed deadline resumes that window instead of restarting it. if ( self._ineffective_compression_count >= 2 or self._fallback_compression_streak >= 2 ): - _now = time.monotonic() - if self._anti_thrash_recovery_deadline <= 0.0: - self._anti_thrash_recovery_deadline = ( + # Wall clock, not monotonic: the deadline is persisted on the + # session row (#100185) so a fresh compressor bound to the same + # session — the gateway rebuilds the AIAgent on every cache + # eviction — resumes the SAME window instead of restarting it. + # Without that, a blocked messaging session never earned its + # probe and stayed blocked forever. + _now = time.time() + if self._anti_thrash_recovery_deadline <= 0.0 or ( + # Clock jumped backwards past a full window: never wait + # longer than one window from now. + self._anti_thrash_recovery_deadline - _now + > self._ANTI_THRASH_RECOVERY_SECONDS + ): + self._set_anti_thrash_recovery_deadline( _now + self._ANTI_THRASH_RECOVERY_SECONDS ) elif _now >= self._anti_thrash_recovery_deadline: - self._anti_thrash_recovery_deadline = 0.0 + self._set_anti_thrash_recovery_deadline(0.0) if self._ineffective_compression_count >= 2: self._record_ineffective_compression_verdict(1) if self._fallback_compression_streak >= 2: @@ -4023,7 +4076,7 @@ class ContextCompressor(ContextEngine): # Guard not tripped (counters were cleared by an effective compaction # or a fitting real-usage reading) — disarm any pending recovery clock # so a LATER trip starts its own full window. - self._anti_thrash_recovery_deadline = 0.0 + self._set_anti_thrash_recovery_deadline(0.0) return False # ------------------------------------------------------------------ diff --git a/hermes_state.py b/hermes_state.py index 7e308704e7..a17da8893f 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -8746,6 +8746,54 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) self._execute_write(_do) + def get_compression_recovery_deadline(self, session_id: str) -> float: + """Return the persisted anti-thrash recovery deadline (wall-clock epoch). + + ``0.0`` means "not armed". The deadline is the durable half of the + #14694 recovery clock: the gateway rebuilds the compressor on every + turn / cache eviction, so a process-local deadline restarted the + wait on each rebuild and a tripped session never earned its probe + (#100185). + """ + if not session_id: + return 0.0 + with self._read_ctx() as conn: + if conn is None: + return 0.0 + row = conn.execute( + "SELECT compression_recovery_deadline FROM sessions WHERE id = ?", + (session_id,), + ).fetchone() + if row is None: + return 0.0 + value = ( + row["compression_recovery_deadline"] + if isinstance(row, sqlite3.Row) + else row[0] + ) + try: + return max(0.0, float(value or 0.0)) + except (TypeError, ValueError): + return 0.0 + + def set_compression_recovery_deadline(self, session_id: str, deadline: float) -> None: + """Persist the anti-thrash recovery deadline; ``0`` / ``None`` disarms it.""" + if not session_id: + return + try: + normalized = max(0.0, float(deadline or 0.0)) + except (TypeError, ValueError): + normalized = 0.0 + stored = normalized if normalized > 0.0 else None + + def _do(conn): + conn.execute( + "UPDATE sessions SET compression_recovery_deadline = ? WHERE id = ?", + (stored, session_id), + ) + + self._execute_write(_do) + # ────────────────────────────────────────────────────────────────────── # Compression locks # ────────────────────────────────────────────────────────────────────── diff --git a/hermes_state_common.py b/hermes_state_common.py index ca081a603a..fc94d5475f 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -353,7 +353,7 @@ def _sql_session_last_active_by_id(session_id_expr: str) -> str: ) -SCHEMA_VERSION = 26 +SCHEMA_VERSION = 27 # FTS storage-layout version, tracked INDEPENDENTLY of SCHEMA_VERSION in the @@ -444,6 +444,7 @@ CREATE TABLE IF NOT EXISTS sessions ( compression_failure_error TEXT, compression_fallback_streak INTEGER NOT NULL DEFAULT 0, compression_ineffective_count INTEGER NOT NULL DEFAULT 0, + compression_recovery_deadline REAL, profile_name TEXT, rewind_count INTEGER NOT NULL DEFAULT 0, archived INTEGER NOT NULL DEFAULT 0, diff --git a/tests/agent/test_compression_anti_thrash_recovery.py b/tests/agent/test_compression_anti_thrash_recovery.py index 109f23c18e..cf245ac9a5 100644 --- a/tests/agent/test_compression_anti_thrash_recovery.py +++ b/tests/agent/test_compression_anti_thrash_recovery.py @@ -17,10 +17,13 @@ The recovery contract pinned here: next recovery waits a FULL fresh window (no immediate re-probe loop). * An effective probe (or any fitting real-usage reading) fully clears the counters through the existing ``update_from_response`` path. -* The recovery clock is armed lazily on the first blocked evaluation and is - NOT durable: a process restart that loads a durable tripped counter - (#69872) starts a full fresh window blocked — a restart must never disarm - or shorten the guard (#54923). +* The recovery clock is armed lazily on the first blocked evaluation and + persisted on the session row as a wall-clock deadline (#100185): a fresh + compressor that loads a durable tripped counter (#69872) with NO stored + deadline starts a full window blocked — a restart must never disarm or + shorten the guard (#54923) — while one that loads an armed deadline + resumes that window instead of restarting it, so gateway agent rebuilds + cannot block a session forever. * The protection itself is preserved: inside the window the gate stays blocked exactly as before. """ @@ -57,10 +60,10 @@ class TestRecoveryWindow: cc = _compressor() _trip(cc) base = 1000.0 - with patch("agent.context_compressor.time.monotonic", return_value=base): + with patch("agent.context_compressor.time.time", return_value=base): assert cc.should_compress(cc.threshold_tokens + 1) is False with patch( - "agent.context_compressor.time.monotonic", + "agent.context_compressor.time.time", return_value=base + cc._ANTI_THRASH_RECOVERY_SECONDS + 1, ): assert cc.should_compress(cc.threshold_tokens + 1) is True @@ -73,10 +76,10 @@ class TestRecoveryWindow: cc = _compressor() cc._fallback_compression_streak = 2 base = 1000.0 - with patch("agent.context_compressor.time.monotonic", return_value=base): + with patch("agent.context_compressor.time.time", return_value=base): assert cc.should_compress(cc.threshold_tokens + 1) is False with patch( - "agent.context_compressor.time.monotonic", + "agent.context_compressor.time.time", return_value=base + cc._ANTI_THRASH_RECOVERY_SECONDS + 1, ): assert cc.should_compress(cc.threshold_tokens + 1) is True @@ -95,13 +98,13 @@ class TestRestartSemantics: cc = _compressor() cc.bind_session_state(session_db=db, session_id="sess-1") assert cc._ineffective_compression_count == 2 - # The recovery clock is process-local and must come up disarmed. + # No stored deadline yet -> the clock comes up disarmed. assert cc._anti_thrash_recovery_deadline == 0.0 base = 5000.0 - with patch("agent.context_compressor.time.monotonic", return_value=base): + with patch("agent.context_compressor.time.time", return_value=base): assert cc.should_compress(cc.threshold_tokens + 1) is False with patch( - "agent.context_compressor.time.monotonic", + "agent.context_compressor.time.time", return_value=base + cc._ANTI_THRASH_RECOVERY_SECONDS + 1, ): assert cc.should_compress(cc.threshold_tokens + 1) is True @@ -113,9 +116,92 @@ class TestRestartSemantics: cc = _compressor() _trip(cc) base = 1000.0 - with patch("agent.context_compressor.time.monotonic", return_value=base): + with patch("agent.context_compressor.time.time", return_value=base): assert cc.should_compress(cc.threshold_tokens + 1) is False assert cc._anti_thrash_recovery_deadline > 0.0 cc.on_session_reset() assert cc._anti_thrash_recovery_deadline == 0.0 assert cc._ineffective_compression_count == 0 + + +class TestDurableDeadline: + """#100185: the gateway rebuilds the compressor on every cache eviction.""" + + def _bound(self, db, session_id="sess-1"): + cc = _compressor() + cc.bind_session_state(session_db=db, session_id=session_id) + return cc + + def test_fresh_compressors_resume_the_same_window(self, tmp_path): + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(session_id="sess-1", source="telegram") + db.set_compression_ineffective_count("sess-1", 2) + base = 5000.0 + first = self._bound(db) + with patch("agent.context_compressor.time.time", return_value=base): + assert first.should_compress(first.threshold_tokens + 1) is False + # Deadline is durable, as a wall-clock epoch. + assert db.get_compression_recovery_deadline("sess-1") == ( + base + first._ANTI_THRASH_RECOVERY_SECONDS + ) + # Fresh compressor (gateway rebuilt the agent) well past the window: + # before the fix it re-armed a new window and stayed blocked forever. + second = self._bound(db) + assert second._anti_thrash_recovery_deadline == ( + base + first._ANTI_THRASH_RECOVERY_SECONDS + ) + with patch( + "agent.context_compressor.time.time", + return_value=base + first._ANTI_THRASH_RECOVERY_SECONDS + 1, + ): + assert second.should_compress(second.threshold_tokens + 1) is True + assert db.get_compression_ineffective_count("sess-1") == 1 + assert db.get_compression_recovery_deadline("sess-1") == 0.0 + + def test_fresh_compressor_inside_window_stays_blocked(self, tmp_path): + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(session_id="sess-1", source="telegram") + db.set_compression_ineffective_count("sess-1", 2) + base = 5000.0 + first = self._bound(db) + with patch("agent.context_compressor.time.time", return_value=base): + assert first.should_compress(first.threshold_tokens + 1) is False + second = self._bound(db) + with patch("agent.context_compressor.time.time", return_value=base + 10): + assert second.should_compress(second.threshold_tokens + 1) is False + assert db.get_compression_ineffective_count("sess-1") == 2 + + def test_backward_clock_jump_is_bounded_to_one_window(self, tmp_path): + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(session_id="sess-1", source="telegram") + db.set_compression_ineffective_count("sess-1", 2) + window = ContextCompressor._ANTI_THRASH_RECOVERY_SECONDS + db.set_compression_recovery_deadline("sess-1", 1_000_000.0) + cc = self._bound(db) + # Wall clock now far BEFORE the stored deadline (clock stepped back). + with patch("agent.context_compressor.time.time", return_value=100.0): + assert cc.should_compress(cc.threshold_tokens + 1) is False + assert db.get_compression_recovery_deadline("sess-1") == 100.0 + window + + def test_clearing_the_guard_disarms_the_durable_deadline(self, tmp_path): + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(session_id="sess-1", source="telegram") + db.set_compression_ineffective_count("sess-1", 2) + cc = self._bound(db) + with patch("agent.context_compressor.time.time", return_value=5000.0): + assert cc.should_compress(cc.threshold_tokens + 1) is False + assert db.get_compression_recovery_deadline("sess-1") > 0.0 + cc._record_ineffective_compression_verdict(0) + with patch("agent.context_compressor.time.time", return_value=5001.0): + assert cc.should_compress(cc.threshold_tokens + 1) is True + assert db.get_compression_recovery_deadline("sess-1") == 0.0 + + def test_session_db_round_trip(self, tmp_path): + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(session_id="sess-1", source="cli") + assert db.get_compression_recovery_deadline("sess-1") == 0.0 + db.set_compression_recovery_deadline("sess-1", 1234.5) + assert db.get_compression_recovery_deadline("sess-1") == 1234.5 + db.set_compression_recovery_deadline("sess-1", 0.0) + assert db.get_compression_recovery_deadline("sess-1") == 0.0 + assert db.get_compression_recovery_deadline("missing") == 0.0 diff --git a/tests/state/test_session_git_metadata_generation.py b/tests/state/test_session_git_metadata_generation.py index 2e1cb30988..b97727725c 100644 --- a/tests/state/test_session_git_metadata_generation.py +++ b/tests/state/test_session_git_metadata_generation.py @@ -237,7 +237,7 @@ def test_legacy_sessions_table_reconciles_generation_column(tmp_path): assert "git_metadata_generation" in columns assert reopened._conn.execute( "SELECT version FROM schema_version" - ).fetchone()[0] == SCHEMA_VERSION == 26 + ).fetchone()[0] == SCHEMA_VERSION reopened.create_session("session", "desktop", cwd="/repo") assert reopened.update_session_cwd("session", "/repo") == 1 finally: