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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
shali10
136 changes: 125 additions & 11 deletions hermes_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -724,20 +724,40 @@ def _on_disk_journal_mode(conn: sqlite3.Connection) -> Optional[str]:

Returns the mode string (e.g. ``"wal"``, ``"delete"``), or ``None``
if the value cannot be determined (new DB, or PRAGMA read failed).

A PRAGMA read can fail transiently with ``disk i/o error`` on
virtualized block devices (XFS on cloud hosts). Treating that as
"mode unknown" pushes callers onto their fail-closed unknown-mode
branch even though the on-disk mode is perfectly readable a few
milliseconds later. Retry the read a few times before giving up:
transient EIO clears, deterministic unsupported-filesystem errors do
not. ``None`` is still returned on final failure so the caller's
existing "unknown → refuse to downgrade" logic applies.
"""
try:
row = conn.execute("PRAGMA journal_mode").fetchone()
except sqlite3.OperationalError:
return None
if row is None:
return None
mode = row[0]
if isinstance(mode, bytes): # defensive: sqlite3 occasionally returns bytes
last_exc: Optional[Exception] = None
for _ in range(4):
try:
mode = mode.decode("ascii")
except UnicodeDecodeError:
row = conn.execute("PRAGMA journal_mode").fetchone()
except sqlite3.OperationalError as exc:
last_exc = exc
if "disk i/o error" not in str(exc).lower():
return None
time.sleep(0.05)
continue
if row is None:
return None
return str(mode).strip().lower() if mode is not None else None
mode = row[0]
if isinstance(mode, bytes): # defensive: sqlite3 occasionally returns bytes
try:
mode = mode.decode("ascii")
except UnicodeDecodeError:
return None
return str(mode).strip().lower() if mode is not None else None
if last_exc is not None:
logger.debug(
"_on_disk_journal_mode: retries exhausted on disk read (%s)", last_exc
)
return None


def _apply_macos_checkpoint_barrier(conn: sqlite3.Connection) -> None:
Expand Down Expand Up @@ -1362,6 +1382,21 @@ def is_malformed_db_error(exc: BaseException) -> bool:
return any(marker in str(exc).lower() for marker in _MALFORMED_SCHEMA_MARKERS)


def _is_not_a_database_error(exc: BaseException) -> bool:
"""True if *exc* is SQLite's 'file is not a database' error.

Raised when a connection's backing file is not a SQLite database — the
runtime connection-corruption class: a sibling process (forked curator
agent, external repair pass) replaced/truncated the file out from under
the live connection. The file on disk may be perfectly healthy; the
CONNECTION is broken. Distinct from the malformed-schema class: the fix
is a reconnect, not schema surgery.
"""
if not isinstance(exc, sqlite3.DatabaseError):
return False
return "file is not a database" in str(exc).lower()


# Markers that mean the host filesystem cannot accept another write. Kept as
# plain substrings so OSError, sqlite3.OperationalError, and wrapped RPC
# error strings all match the same helper.
Expand Down Expand Up @@ -2600,6 +2635,13 @@ def __init__(self, db_path: Path = None, read_only: bool = False):
# in place at most once per SessionDB instance so a genuinely
# unrecoverable database can't put writers into a rebuild loop.
self._fts_runtime_rebuild_attempted = False
# One-shot guard for the runtime connection-reopen recovery on the
# write path. A connection whose backing file was replaced/truncated
# by a sibling process surfaces as "file is not a database" on every
# write; we close and reopen the connection at most once per
# SessionDB instance so a genuinely unrecoverable database can't put
# writers into a reconnect loop.
self._notadb_reconnect_attempted = False
# One-shot guard for the usermerge-floor config write on the
# incremental FTS merge cadence (see _merge_fts_incrementally).
self._fts_usermerge_floor_applied = False
Expand Down Expand Up @@ -3369,6 +3411,20 @@ def _is_no_more_rows(exc: sqlite3.Error) -> bool:
except sqlite3.DatabaseError as exc:
if _is_no_more_rows(exc) and self._sleep_before_write_retry(deadline, patience_s):
continue
# Runtime connection-corruption self-heal: a connection whose
# backing file was replaced/truncated by a sibling process
# (e.g. a forked curator agent inheriting and closing the
# write fd, or an external repair pass) surfaces as "file is
# not a database" on EVERY subsequent write. Without a
# reconnect branch the gateway wedges permanently: every
# transcript/routing write raises, messages stay in memory,
# and swap grows without bound until the process is killed.
# Close the broken connection, reopen the DB file, and retry
# the write once.
if _is_not_a_database_error(exc):
if not self._reconnect_after_notadb():
raise
continue
# Corrupt FTS shadow tables make every write raise the
# malformed/corrupt error class through the FTS sync triggers
# while the canonical messages table is intact. Recover here,
Expand Down Expand Up @@ -3419,6 +3475,64 @@ def _sleep_before_write_retry(
time.sleep(min(jitter, max(deadline - now, 0.001)))
return True

def _reconnect_after_notadb(self) -> bool:
"""Close the corrupted write connection and reopen state.db.

Returns True when the connection was successfully replaced and the
failed write should be retried. Mirrors the constructor's
``_connect_and_init`` so WAL/schema reconciliation runs on the fresh
connection. Never raises — logs and returns False on failure so the
original error propagates.

One-shot per instance: a genuinely unrecoverable database must not
put writers into a reconnect loop that pins CPU on every write.
"""
if self._notadb_reconnect_attempted:
return False
self._notadb_reconnect_attempted = True
logger.warning(
"state.db connection reported 'file is not a database' — closing "
"and reopening the connection to self-heal (one-shot)."
)
try:
with self._lock:
if self._conn is not None:
try:
self._conn.close()
except Exception:
pass
self._conn = None
new_conn = _connect_tracked_db(
str(self.db_path),
tracking_path=self.db_path,
check_same_thread=False,
timeout=1.0,
isolation_level=None,
)
new_conn.row_factory = sqlite3.Row
# Publish BEFORE schema init: _init_schema/_reconcile_columns
# operate on self._conn, not on the local variable.
self._conn = new_conn
self._wal_active = (
apply_wal_with_fallback(new_conn, db_label="state.db")
== "wal"
)
apply_database_pragmas(new_conn, db_label="state.db")
new_conn.execute("PRAGMA foreign_keys=ON")
self._fts_cjk_loaded = load_fts5_cjk_extension(new_conn)
self._init_schema()
except Exception as exc:
logger.error(
"state.db reconnect after 'file is not a database' failed (%s); "
"the database may need the full offline repair path.",
exc,
)
return False
logger.warning(
"state.db connection reopened successfully; retrying the failed write."
)
return True

@staticmethod
def _is_fts_write_corruption_error(exc: sqlite3.DatabaseError) -> bool:
"""True for the error class a corrupt FTS index raises on writes.
Expand Down
125 changes: 125 additions & 0 deletions tests/test_state_db_notadb_selfheal.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
"""Tests for the state.db runtime connection self-heal (PR #82280 remainder).

Covers the two independently-valuable pieces salvaged from the state.db
hardening rollup:

* one-shot reconnect when a live write connection reports
``file is not a database`` (backing file replaced/truncated by a sibling
process — the connection is broken, the on-disk file may be healthy);
* transient ``disk i/o error`` retry in ``_on_disk_journal_mode`` so a
one-shot EIO doesn't push callers onto the fail-closed unknown-mode branch.
"""

import sqlite3
from unittest.mock import MagicMock

import pytest

from hermes_state import SessionDB, _is_not_a_database_error, _on_disk_journal_mode


class _NotADbOnce:
"""Connection proxy that raises 'file is not a database' on execute."""

def __init__(self, real_conn):
self._real = real_conn

def execute(self, *args, **kwargs):
raise sqlite3.DatabaseError("file is not a database")

def __getattr__(self, name):
return getattr(self._real, name)


class TestIsNotADatabaseError:
def test_matches_sqlite_message(self):
assert _is_not_a_database_error(
sqlite3.DatabaseError("file is not a database")
)

def test_rejects_other_database_errors(self):
assert not _is_not_a_database_error(
sqlite3.DatabaseError("database disk image is malformed")
)

def test_rejects_non_sqlite_exceptions(self):
assert not _is_not_a_database_error(ValueError("file is not a database"))


class TestReconnectAfterNotADb:
def test_write_self_heals_when_connection_breaks(self, tmp_path):
"""A broken connection over a healthy file reconnects and retries."""
db = SessionDB(db_path=tmp_path / "state.db")
try:
db.create_session(session_id="s1", source="cli", model="test")
# Simulate the runtime corruption class: the connection starts
# raising 'file is not a database' while the on-disk file is
# perfectly healthy (sibling replaced/truncated the old inode).
db._conn = _NotADbOnce(db._conn)

db.create_session(session_id="s2", source="cli", model="test")

assert db._notadb_reconnect_attempted is True
assert db.get_session("s2") is not None
# The pre-existing row survived (same on-disk file).
assert db.get_session("s1") is not None
finally:
db.close()

def test_reconnect_is_one_shot(self, tmp_path):
"""A second 'file is not a database' propagates instead of looping."""
db = SessionDB(db_path=tmp_path / "state.db")
try:
db._notadb_reconnect_attempted = True
db._conn = _NotADbOnce(db._conn)
with pytest.raises(sqlite3.DatabaseError, match="not a database"):
db.create_session(session_id="s3", source="cli", model="test")
finally:
db._conn = None
db.close()

def test_failed_reconnect_returns_false_and_original_error_propagates(
self, tmp_path, monkeypatch
):
"""If the reopen itself fails, the original write error surfaces."""
db = SessionDB(db_path=tmp_path / "state.db")
try:
monkeypatch.setattr(
"hermes_state._connect_tracked_db",
MagicMock(side_effect=sqlite3.DatabaseError("file is not a database")),
)
db._conn = _NotADbOnce(db._conn)
with pytest.raises(sqlite3.DatabaseError, match="not a database"):
db.create_session(session_id="s4", source="cli", model="test")
assert db._notadb_reconnect_attempted is True
finally:
db._conn = None
db.close()


class TestOnDiskJournalModeEioRetry:
def _conn_raising_then(self, failures, result_rows):
conn = MagicMock()
cursor = MagicMock()
cursor.fetchone.return_value = result_rows
conn.execute.side_effect = list(failures) + [cursor]
return conn

def test_transient_eio_clears_on_retry(self):
conn = self._conn_raising_then(
[sqlite3.OperationalError("disk i/o error")] * 2, ("wal",)
)
assert _on_disk_journal_mode(conn) == "wal"

def test_persistent_eio_returns_none(self):
conn = MagicMock()
conn.execute.side_effect = sqlite3.OperationalError("disk i/o error")
assert _on_disk_journal_mode(conn) is None
# Bounded: retried a handful of times, not forever.
assert conn.execute.call_count == 4

def test_non_eio_operational_error_fails_fast(self):
conn = MagicMock()
conn.execute.side_effect = sqlite3.OperationalError("database is locked")
assert _on_disk_journal_mode(conn) is None
assert conn.execute.call_count == 1
Loading