Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 10 additions & 10 deletions hermes_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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
Expand Down
91 changes: 56 additions & 35 deletions hermes_state_fts.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
Expand Down
112 changes: 110 additions & 2 deletions tests/hermes_state/test_fts_index_fail_open.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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(
Expand All @@ -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)
Loading