diff --git a/hermes_state.py b/hermes_state.py index debb49dfe8a47..8ecc3b5c266a2 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -979,14 +979,7 @@ def _execute_write( # mutations, not just idempotent UPSERTs. ioerr_begin_retried = False while True: - self._raise_if_db_corrupt() - if storage_state(self.db_path) == STORAGE_CORRUPT: - # Another handle in this process already saw structural damage on this file. - # Quarantine this one before it touches SQLite; the error type is the same - # StateDbCorruptError, so every transcript-diversion owner handles it unchanged. - self._halt_db_corrupt(sqlite3.DatabaseError( - "database disk image is malformed (reported earlier in this process: " - f"{storage_corrupt_reason(self.db_path)})")) + self._raise_if_db_corrupt(storage=True) # NOTE: the replaced/generation live probe runs INSIDE the lock below, # not here. close() mutates _conn and _db_sidecar_identity under that # same lock, ending the WAL generation (SQLite unlinks the -wal/-shm @@ -1078,7 +1071,7 @@ def _execute_write( self._raise_if_db_replaced() # Corrupt FTS shadow tables fail every write via the sync triggers while canonical # rows are intact: detach the derived indexes atomically and retry (never rebuild here). - if self._enter_fts_fail_open(exc): + if self._enter_fts_fail_open(exc, deadline=deadline, patience_s=patience_s): continue # What survives both checks is structural damage: quarantine. if self._is_structural_corruption_error(exc): @@ -1397,9 +1390,16 @@ def _settle_lost_generation_locked(self) -> bool: ) return retire_without_close - def _raise_if_db_corrupt(self) -> None: + def _raise_if_db_corrupt(self, *, storage: bool = False) -> None: if self._db_corrupt: raise self._corrupt_error() + if storage and storage_state(self.db_path) == STORAGE_CORRUPT: + # Another handle in this process already saw structural damage on this file. + # Quarantine this one before it touches SQLite; the error type is the same + # StateDbCorruptError, so every transcript-diversion owner handles it unchanged. + self._halt_db_corrupt(sqlite3.DatabaseError( + "database disk image is malformed (reported earlier in this process: " + f"{storage_corrupt_reason(self.db_path)})")) def _sleep_before_write_retry(self, deadline: float, patience_s: float) -> bool: """Sleep one jitter interval if the budget allows; True = retry, False = deadline passed. Small diff --git a/hermes_state_fts.py b/hermes_state_fts.py index 2160efd1d3b9e..3ced10cb2c217 100644 --- a/hermes_state_fts.py +++ b/hermes_state_fts.py @@ -5,13 +5,14 @@ import logging import os import sqlite3 +import time from pathlib import Path from typing import Sequence from hermes_constants import get_hermes_home from hermes_state_common import (FTS_CJK_STALE_KEY, FTS_STALE_KEY, _FTS_CJK_TRIGGERS, _FTS_TRIGGERS, routed_sessions_setting) -from hermes_state_errors import is_fts_scoped_corruption_error +from hermes_state_errors import is_fts_scoped_corruption_error, is_sqlite_lock_error # caplog tests pin the "hermes_state" logger name. logger = logging.getLogger("hermes_state") @@ -350,48 +351,68 @@ def _is_fts_write_corruption_error(exc: sqlite3.DatabaseError) -> bool: gateway transcript retry: see :func:`hermes_state_errors.is_fts_scoped_corruption_error`.""" return is_fts_scoped_corruption_error(exc) - def _enter_fts_fail_open(self, exc: sqlite3.DatabaseError) -> bool: + def _enter_fts_fail_open( + self, exc: sqlite3.DatabaseError, *, deadline: float | None = None, patience_s: float | None = None, + ) -> bool: """Detach corrupt FTS indexes so canonical writes can continue. Breadcrumb + trigger drop commit atomically: once triggers are absent the index has a - gap of unknown extent, so nobody may reinstall them without a full rebuild.""" + gap of unknown extent, so nobody may reinstall them without a full rebuild. + + A busy write lock is waited out on the caller's write budget (default + ``_WRITE_PATIENCE_S``), like ``_execute_write``: the writer connection's busy + timeout is only 1 s, and the usual holder is a sibling writer detaching the + same corrupt index — giving up after 1 s cost that turn's canonical write.""" if not self._fts_enabled or not self._is_fts_write_corruption_error(exc): return False - self._raise_if_db_corrupt() - try: - with self._lock: - self._raise_if_db_replaced() - if self._conn is None: - self._reopen_after_close_locked(context="write") - self._conn.execute("BEGIN IMMEDIATE") - try: - self._conn.execute( - "INSERT INTO state_meta (key, value) VALUES (?, '1') " - "ON CONFLICT(key) DO UPDATE SET value = excluded.value", - (FTS_STALE_KEY,), - ) - cjk_triggers_present = self._conn.execute( - "SELECT 1 FROM sqlite_master WHERE type = 'trigger' " - f"AND name IN ({','.join('?' for _ in _FTS_CJK_TRIGGERS)}) " - "LIMIT 1", - _FTS_CJK_TRIGGERS, - ).fetchone() - if cjk_triggers_present: + if patience_s is None: + patience_s = self._WRITE_PATIENCE_S + if deadline is None: + deadline = time.monotonic() + patience_s + while True: + # Re-checked every attempt: a sibling may quarantine the file while we wait for the + # lock, and nothing may be committed on a quarantined handle. + self._raise_if_db_corrupt(storage=True) + try: + with self._lock: + self._raise_if_db_replaced() + if self._conn is None: + self._reopen_after_close_locked(context="write") + self._conn.execute("BEGIN IMMEDIATE") + try: self._conn.execute( "INSERT INTO state_meta (key, value) VALUES (?, '1') " "ON CONFLICT(key) DO UPDATE SET value = excluded.value", - (FTS_CJK_STALE_KEY,), + (FTS_STALE_KEY,), ) - self._drop_all_fts_triggers(self._conn.cursor()) - self._conn.commit() - except BaseException: - self._conn.rollback() - raise - except sqlite3.Error as detach_exc: - logger.error( - "Could not detach corrupt FTS indexes; canonical write still cannot proceed: %s", - detach_exc, - ) - return False + cjk_triggers_present = self._conn.execute( + "SELECT 1 FROM sqlite_master WHERE type = 'trigger' " + f"AND name IN ({','.join('?' for _ in _FTS_CJK_TRIGGERS)}) " + "LIMIT 1", + _FTS_CJK_TRIGGERS, + ).fetchone() + if cjk_triggers_present: + self._conn.execute( + "INSERT INTO state_meta (key, value) VALUES (?, '1') " + "ON CONFLICT(key) DO UPDATE SET value = excluded.value", + (FTS_CJK_STALE_KEY,), + ) + self._drop_all_fts_triggers(self._conn.cursor()) + self._conn.commit() + except BaseException: + self._conn.rollback() + raise + break + except sqlite3.Error as detach_exc: + if ( + isinstance(detach_exc, sqlite3.OperationalError) and is_sqlite_lock_error(detach_exc) + and self._sleep_before_write_retry(deadline, patience_s) + ): + continue + logger.error( + "Could not detach corrupt FTS indexes; canonical write still cannot proceed: %s", + detach_exc, + ) + return False self._fts_stale = True self._fts_enabled = False self._trigram_available = False diff --git a/tests/hermes_state/test_fts_index_fail_open.py b/tests/hermes_state/test_fts_index_fail_open.py index 834d5e49ba500..496ca0f7a2e28 100644 --- a/tests/hermes_state/test_fts_index_fail_open.py +++ b/tests/hermes_state/test_fts_index_fail_open.py @@ -11,14 +11,18 @@ ``messages`` (the turn proceeds); * an FTS-scoped error that still escapes (detach refused) classifies as ``fts_index`` and never quarantines the handle; +* a sibling process holding the write lock when the detach runs is waited out, not a lost write; +* a quarantine that lands while the detach waits stops it: nothing is committed on the file; """ import sqlite3 +import threading from types import SimpleNamespace import pytest -from hermes_state import SessionDB +from hermes_state import SessionDB, StateDbCorruptError +from hermes_state_health import mark_storage_corrupt, reset_storage_state from run_agent import AIAgent @@ -115,7 +119,7 @@ def test_escaped_fts_only_error_is_index_scoped_not_quarantined(tmp_path, monkey try: _seed(db) _stomp_fts_shadow(db_path) - monkeypatch.setattr(db, "_enter_fts_fail_open", lambda exc: False) + monkeypatch.setattr(db, "_enter_fts_fail_open", lambda exc, **_: False) agent = _flush_agent(db, "s1") ok = agent._flush_messages_to_session_db( @@ -131,3 +135,107 @@ def test_escaped_fts_only_error_is_index_scoped_not_quarantined(tmp_path, monkey assert "refused detach" not in _contents(db_path) finally: db.close() + + +def test_detach_waits_out_a_sibling_holding_the_write_lock(tmp_path): + """Gateway + TUI hit the same corrupt index: one detaches while the other waits. A sibling + that takes the write lock between this writer's corruption error and its detach, and holds + it past the writer connection's 1 s busy timeout, must be waited out on the write budget — + the canonical row lands instead of escaping as 'database disk image is malformed'.""" + db_path = tmp_path / "state.db" + db = SessionDB(db_path=db_path) + try: + _seed(db, rows=5) + _stomp_fts_shadow(db_path) + held, release = threading.Event(), threading.Event() + + def sibling(): + raw = sqlite3.connect(str(db_path), timeout=30, isolation_level=None) + raw.execute("BEGIN IMMEDIATE") + held.set() + release.wait(10) + raw.execute("COMMIT") + raw.close() + + real_check = db._is_fts_write_corruption_error + holder = [] + + def check_then_contend(exc): + hit = real_check(exc) + if hit and not holder: # the writer has rolled back; the sibling grabs the lock now + holder.append(threading.Thread(target=sibling)) + holder[0].start() + assert held.wait(10) + threading.Timer(1.6, release.set).start() + return hit + + db._is_fts_write_corruption_error = check_then_contend + db.append_message("s1", "user", "lands after the sibling lets go") + if not holder: + pytest.skip("this SQLite build defers FTS shadow corruption past the insert trigger") + holder[0].join(10) + + assert _contents(db_path)[-1] == "lands after the sibling lets go" + assert db._fts_stale is True + assert db._db_corrupt is False + finally: + db.close() + + +def test_quarantine_while_detach_waits_commits_nothing(tmp_path): + """The detach may now wait up to the write budget for the lock. A sibling that quarantines + this file meanwhile (structural corruption latched process-wide) must stop it: the retry + drops no triggers, commits no stale breadcrumb, and the corrupt error surfaces.""" + db_path = tmp_path / "state.db" + db = SessionDB(db_path=db_path) + try: + _seed(db, rows=5) + _stomp_fts_shadow(db_path) + held, release = threading.Event(), threading.Event() + + def sibling(): + raw = sqlite3.connect(str(db_path), timeout=30, isolation_level=None) + raw.execute("BEGIN IMMEDIATE") + held.set() + release.wait(10) + raw.execute("COMMIT") + raw.close() + + real_check, real_sleep = db._is_fts_write_corruption_error, db._sleep_before_write_retry + holder = [] + + def check_then_contend(exc): + hit = real_check(exc) + if hit and not holder: + holder.append(threading.Thread(target=sibling)) + holder[0].start() + assert held.wait(10) + return hit + + def quarantine_then_sleep(deadline, patience_s): + mark_storage_corrupt(db_path, "database disk image is malformed (sibling handle)") + release.set() + return real_sleep(deadline, patience_s) + + db._is_fts_write_corruption_error = check_then_contend + db._sleep_before_write_retry = quarantine_then_sleep + with pytest.raises(StateDbCorruptError): + db.append_message("s1", "user", "must not land on a quarantined file") + if not holder: + pytest.skip("this SQLite build defers FTS shadow corruption past the insert trigger") + holder[0].join(10) + + raw = sqlite3.connect(str(db_path)) + try: + triggers = raw.execute( + "SELECT count(*) FROM sqlite_master WHERE type = 'trigger' AND name LIKE 'messages_fts%'" + ).fetchone()[0] + stale = raw.execute("SELECT value FROM state_meta WHERE key LIKE 'fts%stale%'").fetchall() + finally: + raw.close() + assert triggers > 0 + assert stale == [] + assert db._fts_stale is False + finally: + db.close() + reset_storage_state(db_path)