diff --git a/hermes_state.py b/hermes_state.py index debb49dfe8..c19283f94b 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -1078,7 +1078,7 @@ class SessionDB( 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): diff --git a/hermes_state_fts.py b/hermes_state_fts.py index 2160efd1d3..40a3045722 100644 --- a/hermes_state_fts.py +++ b/hermes_state_fts.py @@ -5,6 +5,7 @@ FTS-scoped corruption detection and the atomic fail-open trigger detach.""" import logging import os import sqlite3 +import time from pathlib import Path from typing import Sequence @@ -350,48 +351,67 @@ class SessionFtsSetupMixin: 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: + 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: + msg = str(detach_exc).lower() + if ( + isinstance(detach_exc, sqlite3.OperationalError) and ("locked" in msg or "busy" in msg) + 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 834d5e49ba..59bb95ae5f 100644 --- a/tests/hermes_state/test_fts_index_fail_open.py +++ b/tests/hermes_state/test_fts_index_fail_open.py @@ -11,9 +11,11 @@ user-facing guidance: ``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; """ import sqlite3 +import threading from types import SimpleNamespace import pytest @@ -115,7 +117,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 +133,48 @@ 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()