fix(state): quarantine SessionDB handle after structural corruption

A bare SQLITE_CORRUPT/NOTADB on a live write (not FTS-scoped, not a
replaced file) now sets a sticky per-instance flag: later writes fail
fast with StateDbCorruptError, the handle never reopens after close(),
and close() skips its explicit PASSIVE WAL checkpoint. Gateway and agent
flush paths divert pending transcripts to JSONL/spool like the replaced
case instead of retrying forever.

Field evidence: a handle that kept writing for ~50 minutes after the
first structural error checkpointed 15 pages under the wrong page
numbers on shutdown (page 1 <- messages_fts_trigram_data leaf), turning
"malformed" into "file is not a database".

Refs #90837, #90950, #97940, #89332, #45383

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CNX8rNYHqA5pT4tAGSzXtb
This commit is contained in:
leomcamilo
2026-09-02 05:22:41 -03:00
committed by kshitij
parent d8616f1c88
commit bcc2e65818
10 changed files with 565 additions and 11 deletions
+38
View File
@@ -23,6 +23,44 @@ cross-process admission lock and foreign-holder guard. If that guarded rebuild
cannot run, FTS remains detached, canonical writes stay available, and cannot run, FTS remains detached, canonical writes stay available, and
`hermes doctor` reports the explicit repair command. `hermes doctor` reports the explicit repair command.
## Live behavior when the file itself is corrupt
If a live write reports bare `SQLITE_CORRUPT` / `SQLITE_NOTADB` (`database
disk image is malformed`, `file is not a database`) with no FTS provenance,
the damage is in a canonical B-tree, the schema, or the freelist. `SessionDB`
then quarantines that handle (`StateDbCorruptError`):
1. the failing write propagates the typed error and nothing is retried;
2. later writes on the handle fail immediately without touching the file;
3. the handle never reopens its connection after `close()`; and
4. `close()` skips its explicit WAL checkpoint.
Stopping the writes is the protection. In the field, a handle that kept
writing for ~50 minutes after the first structural error checkpointed 15
pages under the wrong page numbers on shutdown (page 1 received a
`messages_fts_trigram_data` leaf) and turned a damaged-but-readable file into
one that no longer opened at all. Skipping the explicit checkpoint is the
second line of defence; SQLite may still run its own last-connection
checkpoint when the connection closes, so copy `state.db`, `state.db-wal` and
`state.db-shm` together before restarting anything.
The gateway and the agent flush path treat the quarantine like a replaced
file: pending transcripts go to `sessions/<id>.jsonl` and the gateway
`pending_messages/` spool instead of the retry queue, and the FTS one-shot
rebuild never runs on the damaged file. The quarantine is per process — the
shared handle stays poisoned for every holder until the process restarts on a
repaired or restored file. Do not run `hermes doctor --fix` while the gateway
is still up. Next steps:
```bash
hermes gateway stop
HERMES_HOME="$HOME/.hermes" hermes sessions recover --source "$HOME/.hermes/state.db" --inspect-only
# if recoverable:
HERMES_HOME="$HOME/.hermes" hermes sessions recover --source "$HOME/.hermes/state.db" --output "$HOME/recovered-state.db"
```
or restore the newest snapshot from `state-snapshots/`.
## Explicit repair ## Explicit repair
Stop every process that can open the profile database before repairing it. Stop every process that can open the profile database before repairing it.
+14 -5
View File
@@ -3981,14 +3981,23 @@ class SessionStore:
try: try:
self._append_transcript_message(session_id, msg) self._append_transcript_message(session_id, msg)
except Exception as exc: except Exception as exc:
from hermes_state import CompressionSessionClosedError, StateDbReplacedError from hermes_state import (
CompressionSessionClosedError,
StateDbCorruptError,
StateDbReplacedError,
)
if isinstance(exc, StateDbReplacedError): if isinstance(exc, (StateDbReplacedError, StateDbCorruptError)):
# Both classes mean "this handle must not touch the file
# again": replaced generation (#89332) or structural
# corruption (quarantine). Retrying cannot succeed, and
# the FTS one-shot rebuild below must never run on a
# damaged file. Divert instead.
logger.error( logger.error(
"Session DB was replaced underneath the gateway for %s; " "Session DB refused further writes on this handle for "
"stopping SQLite writes and diverting pending " "%s (%s); stopping SQLite writes and diverting pending "
"transcripts to the on-disk fallback: %s", "transcripts to the on-disk fallback: %s",
session_id, exc, session_id, type(exc).__name__, exc,
) )
with self._transcript_retry_lock: with self._transcript_retry_lock:
remaining = list(self._dirty_transcripts.get(queue_session_id, [])) remaining = list(self._dirty_transcripts.get(queue_session_id, []))
+122 -1
View File
@@ -2265,6 +2265,8 @@ def classify_persistence_error(exc_or_str) -> str:
return "compression" return "compression"
if isinstance(exc_or_str, StateDbReplacedError): if isinstance(exc_or_str, StateDbReplacedError):
return "replaced" return "replaced"
if isinstance(exc_or_str, StateDbCorruptError):
return "corrupt"
text = str(exc_or_str).lower() text = str(exc_or_str).lower()
if "turn lease" in text: if "turn lease" in text:
return "turn_lease" return "turn_lease"
@@ -4361,6 +4363,40 @@ _STATE_DB_REPLACED_MSG = (
) )
class StateDbCorruptError(sqlite3.DatabaseError):
"""A live SessionDB observed structural (non-FTS) corruption and is quarantined.
Raised once a write on this handle reports bare ``SQLITE_CORRUPT`` /
``SQLITE_NOTADB`` that is neither FTS-scoped (``_is_fts_write_corruption_error``)
nor a replaced-file case (``StateDbReplacedError``). Subclasses
``sqlite3.DatabaseError`` so every existing ``except sqlite3.Error``
degrade path keeps working; ``sqlite_errorcode``/``sqlite_errorname``
are copied from the originating error.
The quarantine is sticky for the life of the handle: later writes fail
fast, the handle never reopens after ``close()``, and ``close()`` skips
its own WAL checkpoint. Field evidence (the #90837 lost/reordered-page
signature, the #90950 page-1 clobber): a handle that kept writing for ~50
minutes after the first structural error checkpointed 15 pages under the
wrong page numbers on shutdown, turning a still-readable file into
``file is not a database``. Stopping the writes is what prevents that;
skipping the explicit checkpoint is the second line of defence (SQLite
may still run its own last-connection checkpoint on close — Python's
``sqlite3`` does not expose ``SQLITE_DBCONFIG_NO_CKPT_ON_CLOSE``). The
recovery boundary is a process restart on a repaired or restored file.
"""
_STATE_DB_CORRUPT_MSG = (
"FATAL: state.db reported structural corruption (database disk image is "
"malformed outside the FTS shadow tables) on a live handle; refusing further "
"writes, automatic reopen, and the close-time WAL checkpoint on this file. "
"Stop the gateway, then run `hermes sessions recover --source <state.db> "
"--inspect-only` or restore a snapshot. Unwritten transcripts are diverted to "
"sessions/<id>.jsonl (and the gateway pending_messages spool)."
)
def divert_session_transcript_jsonl(session_id: str, messages) -> "Optional[Path]": def divert_session_transcript_jsonl(session_id: str, messages) -> "Optional[Path]":
"""Append pending messages as JSON lines under HERMES_HOME/sessions. """Append pending messages as JSON lines under HERMES_HOME/sessions.
@@ -5242,6 +5278,12 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
self._db_file_application_id: int = 0 self._db_file_application_id: int = 0
self._db_file_generation_token: str = "" self._db_file_generation_token: str = ""
self._db_replaced = False self._db_replaced = False
# Sticky: set once a write on THIS handle reports bare SQLITE_CORRUPT /
# NOTADB that is not FTS-scoped and not a replaced-file case. Never
# cleared; the recovery boundary is a process restart on a repaired or
# restored file (see StateDbCorruptError).
self._db_corrupt = False
self._db_corrupt_reason = ""
# One-shot guard for the usermerge-floor config write on the # One-shot guard for the usermerge-floor config write on the
# incremental FTS merge cadence (see _merge_fts_incrementally). # incremental FTS merge cadence (see _merge_fts_incrementally).
self._fts_usermerge_floor_applied = False self._fts_usermerge_floor_applied = False
@@ -5783,6 +5825,15 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
# through stale WAL/shm assumptions (#89332). Refuse instead. # through stale WAL/shm assumptions (#89332). Refuse instead.
if self._db_replaced or self._db_file_was_replaced(): if self._db_replaced or self._db_file_was_replaced():
self._halt_db_replaced() self._halt_db_replaced()
# A quarantined handle must never come back: reopening would hand a
# fresh connection (and its own close-time checkpoint) to a file we
# already know is structurally damaged.
if self._db_corrupt:
raise self._corrupt_error(
f"state.db connection for {self.db_path} is quarantined after "
f"structural corruption; refusing to reopen for a {context} "
"after close(). "
)
logger.warning( logger.warning(
"state.db connection for %s was closed while a %s was still in " "state.db connection for %s was closed while a %s was still in "
"flight — reopening (teardown/worker race, #94736)", "flight — reopening (teardown/worker race, #94736)",
@@ -6127,6 +6178,7 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
return "no more rows available" in str(exc).lower() return "no more rows available" in str(exc).lower()
while True: while True:
self._raise_if_db_corrupt()
self._raise_if_db_replaced() self._raise_if_db_replaced()
fn_started = False fn_started = False
try: try:
@@ -6229,6 +6281,11 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
# explicit repair paths retain rebuild ownership. # explicit repair paths retain rebuild ownership.
if self._enter_fts_fail_open(exc): if self._enter_fts_fail_open(exc):
continue continue
# Bare SQLITE_CORRUPT / NOTADB that survived the replaced-file
# check and the FTS-scoped fail-open is structural damage:
# quarantine the handle (see StateDbCorruptError).
if self._is_structural_corruption_error(exc):
self._halt_db_corrupt(exc)
raise raise
except sqlite3.Error as exc: except sqlite3.Error as exc:
# Catch-all for builds that surface 'no more rows available' # Catch-all for builds that surface 'no more rows available'
@@ -6324,6 +6381,54 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
if self._db_file_was_replaced(): if self._db_file_was_replaced():
self._halt_db_replaced() self._halt_db_replaced()
@classmethod
def _is_structural_corruption_error(cls, exc: BaseException) -> bool:
"""Bare SQLITE_CORRUPT/NOTADB with no FTS provenance.
``_is_fts_write_corruption_error`` is the positive FTS classifier;
everything else in the ``corrupt`` bucket of
``classify_persistence_error`` is damage to a canonical B-tree, the
schema, or the freelist — never repairable from the live write path.
"""
if not isinstance(exc, sqlite3.DatabaseError):
return False
if isinstance(exc, StateDbCorruptError):
return False
if cls._is_fts_write_corruption_error(exc):
return False
return classify_persistence_error(exc) == "corrupt"
def _corrupt_error(self, prefix: str = "") -> "StateDbCorruptError":
"""Build the quarantine error for this handle (message assembled once)."""
return StateDbCorruptError(
f"{prefix}{_STATE_DB_CORRUPT_MSG} (cause: {self._db_corrupt_reason})"
)
def _halt_db_corrupt(self, exc: BaseException) -> None:
"""Quarantine this handle and raise; never run in-file repair here."""
self._db_corrupt = True
self._db_corrupt_reason = str(exc)
logger.error(
"state.db %s reported structural corruption outside the FTS "
"indexes (%s); quarantining this handle: no further writes, no "
"automatic reopen, no explicit WAL checkpoint at close. Stop the "
"gateway and run `hermes sessions recover --source %s "
"--inspect-only`.",
self.db_path,
exc,
self.db_path,
)
err = self._corrupt_error()
for attr in ("sqlite_errorcode", "sqlite_errorname"):
value = getattr(exc, attr, None)
if value is not None:
setattr(err, attr, value)
raise err from exc
def _raise_if_db_corrupt(self) -> None:
if self._db_corrupt:
raise self._corrupt_error()
def _sleep_before_write_retry( def _sleep_before_write_retry(
self, deadline: float, patience_s: float self, deadline: float, patience_s: float
) -> bool: ) -> bool:
@@ -6545,6 +6650,7 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
""" """
if not self._fts_enabled or not self._is_fts_write_corruption_error(exc): if not self._fts_enabled or not self._is_fts_write_corruption_error(exc):
return False return False
self._raise_if_db_corrupt()
if self._db_replaced or self._db_file_was_replaced(): if self._db_replaced or self._db_file_was_replaced():
self._halt_db_replaced() self._halt_db_replaced()
@@ -6612,6 +6718,8 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
databases (65K+ pages) due to the exclusive-lock I/O pressure databases (65K+ pages) due to the exclusive-lock I/O pressure
from checkpointing thousands of frames at once (issue #45383). from checkpointing thousands of frames at once (issue #45383).
""" """
if self._db_corrupt:
return # quarantined: never checkpoint over a damaged image
try: try:
with self._lock: with self._lock:
result = self._conn.execute( result = self._conn.execute(
@@ -6703,7 +6811,20 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
self._close_read_conn(conn) self._close_read_conn(conn)
with self._lock: with self._lock:
if self._conn: if self._conn:
if not self.read_only: if self._db_corrupt:
# Quarantined handle (see StateDbCorruptError): no explicit
# checkpoint over a damaged page image.
logger.warning(
"Skipping the close-time WAL checkpoint for %s: this "
"handle observed structural corruption (%s). Take a "
"snapshot of state.db, -wal and -shm before restarting, "
"then run `hermes sessions recover --source %s "
"--inspect-only`.",
self.db_path,
self._db_corrupt_reason,
self.db_path,
)
elif not self.read_only:
# PASSIVE, not TRUNCATE. Every cron run_agent opens+closes a # PASSIVE, not TRUNCATE. Every cron run_agent opens+closes a
# transient SessionDB, so a TRUNCATE here fires a full WAL # transient SessionDB, so a TRUNCATE here fires a full WAL
# reset many times/hour, racing the gateway's long-lived # reset many times/hour, racing the gateway's long-lived
+7 -2
View File
@@ -2670,13 +2670,17 @@ class AIAgent:
# ("storage was busy, send it again") from disk-full/read-only. # ("storage was busy, send it again") from disk-full/read-only.
from hermes_state import ( from hermes_state import (
CompressionSessionClosedError, CompressionSessionClosedError,
StateDbCorruptError,
StateDbReplacedError, StateDbReplacedError,
classify_persistence_error, classify_persistence_error,
divert_session_transcript_jsonl, divert_session_transcript_jsonl,
) )
self._last_persistence_error_cause = classify_persistence_error(e) self._last_persistence_error_cause = classify_persistence_error(e)
if isinstance(e, StateDbReplacedError): if isinstance(e, (StateDbReplacedError, StateDbCorruptError)):
# Replaced generation or quarantined (structurally corrupt)
# handle: SQLite will not take this batch again, so keep it
# on disk instead of only in RAM.
try: try:
divert_session_transcript_jsonl( divert_session_transcript_jsonl(
getattr(self, "session_id", "") or "", getattr(self, "session_id", "") or "",
@@ -2684,7 +2688,8 @@ class AIAgent:
) )
except Exception: except Exception:
logger.warning( logger.warning(
"JSONL divert failed after state.db replace for %s", "JSONL divert failed after state.db %s for %s",
self._last_persistence_error_cause,
getattr(self, "session_id", None), getattr(self, "session_id", None),
exc_info=True, exc_info=True,
) )
@@ -0,0 +1,70 @@
"""Gateway SessionStore must divert, not retry forever, after structural corruption.
Mirrors ``test_session_db_replaced_fallback.py``: once the SessionDB handle
is quarantined (``StateDbCorruptError``) the pending transcript goes to the
JSONL/spool fallback and no FTS surgery runs on the damaged file.
"""
import json
import sqlite3
from gateway.config import GatewayConfig
from gateway.session import SessionStore
class _MalformedConn:
def __init__(self, real_conn):
self._real = real_conn
def execute(self, *args, **kwargs):
raise sqlite3.DatabaseError("database disk image is malformed")
def __getattr__(self, name):
return getattr(self._real, name)
def _assert_diverted(tmp_path, sid, needle):
pending = list((tmp_path / "pending_messages").glob("pending-*.json"))
assert pending, "expected pending_messages/pending-*.json spool"
spooled = False
for path in pending:
payload = json.loads(path.read_text(encoding="utf-8"))
message = (payload.get("data") or {}).get("message") or {}
if needle in str(message.get("content", "")):
spooled = True
break
assert spooled, f"{needle!r} missing from pending spool"
jsonl = tmp_path / "sessions" / f"{sid}.jsonl"
assert jsonl.is_file()
assert needle in jsonl.read_text(encoding="utf-8")
def test_corrupt_state_db_diverts_pending_without_fts_rebuild(tmp_path, monkeypatch):
import hermes_state
live = tmp_path / "state.db"
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
monkeypatch.setattr(hermes_state, "DEFAULT_DB_PATH", live)
store = SessionStore(sessions_dir=tmp_path, config=GatewayConfig())
sid = "gw-corrupt"
store._db.create_session(session_id=sid, source="cli")
store.append_to_transcript(
sid, {"role": "user", "content": "before", "timestamp": 1.0}
)
real_conn = store._db._conn
store._db._conn = _MalformedConn(real_conn)
try:
store.append_to_transcript(
sid, {"role": "user", "content": "after-corrupt", "timestamp": 2.0}
)
assert store._db._db_corrupt is True
# No FTS surgery ran on either layer.
assert store._db._fts_enabled is True
assert store._db._fts_stale is False
assert store._fts_rebuild_attempted is False
assert sid not in store._dirty_transcripts
_assert_diverted(tmp_path, sid, "after-corrupt")
finally:
store._db._conn = real_conn
store.close_all_db_handles()
@@ -0,0 +1,212 @@
"""Quarantine of a live SessionDB handle after structural (non-FTS) corruption.
Field evidence (the #90837 lost/reordered-page-write class): a gateway kept
retrying writes for ~50 minutes after ``gateway_routing`` reported
``database disk image is malformed``; on SIGTERM the close-time
``PRAGMA wal_checkpoint(PASSIVE)`` then wrote 15 pages to the wrong page
numbers (page 1 received a ``messages_fts_trigram_data`` leaf) and the file
stopped opening at all. Once structural corruption is observed on a handle
the only safe policy is to stop touching the file.
"""
import sqlite3
import pytest
from hermes_state import SessionDB, StateDbCorruptError
class _MalformedConn:
"""Connection proxy whose every execute reports bare SQLITE_CORRUPT."""
def __init__(self, real_conn):
self._real = real_conn
def execute(self, *args, **kwargs):
raise sqlite3.DatabaseError("database disk image is malformed")
def __getattr__(self, name):
return getattr(self._real, name)
class TestQuarantineAfterStructuralCorruption:
def test_structural_corruption_sets_sticky_flag_and_raises_typed(self, tmp_path):
db = SessionDB(db_path=tmp_path / "state.db")
real_conn = db._conn
try:
db.create_session(session_id="s1", source="cli", model="test")
db._conn = _MalformedConn(real_conn)
with pytest.raises(StateDbCorruptError, match="malformed") as excinfo:
db.create_session(session_id="s2", source="cli", model="test")
assert isinstance(excinfo.value.__cause__, sqlite3.DatabaseError)
assert db._db_corrupt is True
# Structural damage must never be mistaken for FTS-scoped damage.
assert db._fts_stale is False
finally:
db._conn = real_conn
db.close()
class _RecordingConn:
"""Connection proxy that records every SQL text and delegates."""
def __init__(self, real_conn):
self._real = real_conn
self.recorded = []
def execute(self, sql, *args, **kwargs):
self.recorded.append(str(sql))
return self._real.execute(sql, *args, **kwargs)
def __getattr__(self, name):
return getattr(self._real, name)
def _quarantined_db(tmp_path):
"""A SessionDB whose first corrupt write already tripped the quarantine."""
db = SessionDB(db_path=tmp_path / "state.db")
real_conn = db._conn
db.create_session(session_id="s1", source="cli", model="test")
db._conn = _MalformedConn(real_conn)
with pytest.raises(StateDbCorruptError):
db.create_session(session_id="s2", source="cli", model="test")
db._conn = real_conn
assert db._db_corrupt is True
return db, real_conn
class TestQuarantinedHandleStopsTouchingTheFile:
def test_subsequent_writes_fail_fast_without_touching_connection(self, tmp_path):
db, real_conn = _quarantined_db(tmp_path)
recorder = _RecordingConn(real_conn)
db._conn = recorder
try:
with pytest.raises(StateDbCorruptError):
db.create_session(session_id="s3", source="cli", model="test")
assert recorder.recorded == []
finally:
db._conn = real_conn
db.close()
def test_close_skips_wal_checkpoint_when_quarantined(self, tmp_path, caplog):
db, real_conn = _quarantined_db(tmp_path)
recorder = _RecordingConn(real_conn)
db._conn = recorder
with caplog.at_level("WARNING", logger="hermes_state"):
db.close()
assert not any("wal_checkpoint" in sql for sql in recorder.recorded)
assert db._conn is None
assert any(
"Skipping the close-time WAL checkpoint" in rec.getMessage()
and "hermes sessions recover" in rec.getMessage()
for rec in caplog.records
)
def test_reopen_after_close_refused_when_quarantined(self, tmp_path, monkeypatch):
from unittest.mock import MagicMock
db, real_conn = _quarantined_db(tmp_path)
db.close()
reopen = MagicMock()
monkeypatch.setattr("hermes_state._connect_tracked_db", reopen)
with pytest.raises(StateDbCorruptError, match="structural corruption"):
db.create_session(session_id="s4", source="cli", model="test")
reopen.assert_not_called()
# The read fallback after close() goes through the same reopen path.
with pytest.raises(StateDbCorruptError, match="refusing to reopen"):
db.get_session("s1")
reopen.assert_not_called()
class TestQuarantineScope:
def test_fts_scoped_corruption_does_not_trip_flag(self, tmp_path):
"""Corrupt FTS shadow tables keep the existing fail-open detach path."""
path = tmp_path / "state.db"
db = SessionDB(db_path=path)
db.create_session(session_id="s1", source="cli", model="test")
db.append_message("s1", role="user", content="hello world")
raw = sqlite3.connect(str(path))
raw.execute(
"UPDATE messages_fts_data SET block = X'DEADBEEFDEADBEEFDEADBEEFDEADBEEF'"
)
raw.commit()
raw.close()
try:
db.append_message("s1", role="user", content="healed append")
assert db._db_corrupt is False
assert db._fts_stale is True
assert db._fts_enabled is False
finally:
db.close()
def test_replaced_file_takes_precedence_over_corrupt(self, tmp_path):
import os
from hermes_state import StateDbReplacedError
live = tmp_path / "state.db"
other = tmp_path / "other.db"
db = SessionDB(db_path=live)
real_conn = db._conn
try:
db.create_session(session_id="s1", source="cli", model="test")
if db._db_file_identity is None:
pytest.skip("filesystem does not expose st_dev/st_ino")
alt = SessionDB(db_path=other)
alt.create_session("other", "cli")
alt.close()
os.replace(other, live)
db._conn = _MalformedConn(real_conn)
with pytest.raises(StateDbReplacedError):
db.create_session(session_id="s2", source="cli", model="test")
assert db._db_replaced is True
assert db._db_corrupt is False
finally:
db._conn = real_conn
db.close()
def test_classify_persistence_error_maps_quarantine_to_corrupt(self):
from hermes_state import _STATE_DB_CORRUPT_MSG, classify_persistence_error
assert classify_persistence_error(StateDbCorruptError("x")) == "corrupt"
# The stringified form (RPC boundaries) must classify the same way.
assert classify_persistence_error(_STATE_DB_CORRUPT_MSG) == "corrupt"
@pytest.fixture
def _clean_registry():
import hermes_state_registry as registry
registry.close_all()
registry._generations.clear()
registry._retired.clear()
yield registry
registry.close_all()
registry._generations.clear()
registry._retired.clear()
class TestSharedRegistry:
def test_holders_share_quarantine_and_close_all_skips_checkpoint(
self, tmp_path, _clean_registry
):
registry = _clean_registry
path = tmp_path / "state.db"
holder_a = registry.acquire(path)
holder_b = registry.acquire(path)
assert holder_a is holder_b
real_conn = holder_a._conn
holder_a.create_session(session_id="s1", source="cli", model="test")
holder_a._conn = _MalformedConn(real_conn)
with pytest.raises(StateDbCorruptError):
holder_a.create_session(session_id="s2", source="cli", model="test")
recorder = _RecordingConn(real_conn)
holder_b._conn = recorder
with pytest.raises(StateDbCorruptError):
holder_b.create_session(session_id="s3", source="cli", model="test")
registry.close_all()
assert not any("wal_checkpoint" in sql for sql in recorder.recorded)
assert holder_a._conn is None
@@ -0,0 +1,70 @@
"""Agent flush path: a quarantined (structurally corrupt) SessionDB diverts to JSONL.
Mirrors the replaced-file contract: the batch that SQLite will never take
again is kept on disk under ``sessions/<id>.jsonl`` instead of only in RAM,
the flush fails closed (no retry loop), and the turn-end explanation gets the
``corrupt`` cause.
"""
from __future__ import annotations
from pathlib import Path
from types import SimpleNamespace
from hermes_state import SessionDB, StateDbCorruptError
from run_agent import AIAgent
def _flush_agent(db, session_id):
agent = SimpleNamespace(
_session_db=db,
_session_db_created=True,
_persist_disabled=False,
session_id=session_id,
_session_persist_lock=None,
_flushed_db_message_ids=set(),
_flushed_db_message_session_id=None,
_last_flushed_db_idx=0,
_db_flush_scan_prefix=None,
_persist_user_message_idx=None,
_persist_user_message_override=None,
_persist_user_message_timestamp=None,
_pending_cli_user_message=None,
_active_session_turn_lease_holder=None,
_last_persistence_error_cause=None,
_compression_adoption_failed=False,
)
agent._ensure_db_session = lambda: None
agent._flush_messages_to_session_db = (
AIAgent._flush_messages_to_session_db.__get__(agent, AIAgent)
)
agent._flush_messages_to_session_db_unlocked = (
AIAgent._flush_messages_to_session_db_unlocked.__get__(agent, AIAgent)
)
return agent
def test_flush_diverts_batch_to_jsonl_when_handle_is_quarantined(
tmp_path: Path, monkeypatch
) -> None:
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
db = SessionDB(db_path=tmp_path / "state.db")
try:
db.create_session("live", source="cli")
agent = _flush_agent(db, "live")
def _quarantined(self, *, session_id, messages, **kwargs):
raise StateDbCorruptError("database disk image is malformed (quarantined)")
monkeypatch.setattr(SessionDB, "append_messages_batch", _quarantined)
messages = [{"role": "user", "content": "kept-on-disk-after-corruption"}]
result = agent._flush_messages_to_session_db(messages, [])
assert result is False
assert agent._last_persistence_error_cause == "corrupt"
jsonl = tmp_path / "sessions" / "live.jsonl"
assert jsonl.is_file()
assert "kept-on-disk-after-corruption" in jsonl.read_text(encoding="utf-8")
finally:
db.close()
@@ -438,3 +438,8 @@ def test_run_conversation_partial_stream_recovery_surfaces_explanation():
assert result["response_previewed"] is False assert result["response_previewed"] is False
def test_classify_persistence_error_quarantined_handle_is_corrupt() -> None:
"""A quarantined SessionDB raises the typed error; it stays in the corrupt bucket."""
from hermes_state import StateDbCorruptError, classify_persistence_error
assert classify_persistence_error(StateDbCorruptError("quarantined")) == "corrupt"
+22 -1
View File
@@ -332,7 +332,15 @@ class TestRuntimeFtsRebuild:
with pytest.raises(sqlite3.DatabaseError) as caught: with pytest.raises(sqlite3.DatabaseError) as caught:
db._execute_write(lambda _conn: (_ for _ in ()).throw(structural)) db._execute_write(lambda _conn: (_ for _ in ()).throw(structural))
assert caught.value is structural # Structural corruption quarantines the handle: the typed error wraps
# the original (cause preserved, SQLite result code copied) and the
# sticky flag is set, so later writes fail fast.
from hermes_state import StateDbCorruptError
assert isinstance(caught.value, StateDbCorruptError)
assert caught.value.__cause__ is structural
assert caught.value.sqlite_errorcode == sqlite3.SQLITE_CORRUPT
assert db._db_corrupt is True
assert rebuild_called is False assert rebuild_called is False
assert db._fts_stale is False assert db._fts_stale is False
assert _meta_value(tmp_path / "state.db", FTS_STALE_KEY) is None assert _meta_value(tmp_path / "state.db", FTS_STALE_KEY) is None
@@ -831,6 +839,19 @@ class TestPhysicalCorruptionAcceptance:
# The misdiagnosis message from the field incident must be gone. # The misdiagnosis message from the field incident must be gone.
assert "canonical message rows are preserved" not in caplog.text assert "canonical message rows are preserved" not in caplog.text
assert "attempting one-shot in-place FTS rebuild" not in caplog.text assert "attempting one-shot in-place FTS rebuild" not in caplog.text
# Structural damage quarantines the handle: typed error, sticky
# flag, later writes fail fast, and close() must not checkpoint
# the WAL over a damaged page image (the #90950 page-1 clobber).
from hermes_state import StateDbCorruptError
assert isinstance(caught.value, StateDbCorruptError)
assert db._db_corrupt is True
with pytest.raises(StateDbCorruptError):
db.append_message("s1", "user", "second write after corruption")
caplog.clear()
with caplog.at_level("WARNING", logger="hermes_state"):
db.close()
assert "Skipping the close-time WAL checkpoint" in caplog.text
finally: finally:
db.close() db.close()
+5 -2
View File
@@ -14,7 +14,7 @@ from unittest.mock import MagicMock
import pytest import pytest
from hermes_state import SessionDB, _on_disk_journal_mode from hermes_state import SessionDB, StateDbCorruptError, _on_disk_journal_mode
class _NotADbOnce: class _NotADbOnce:
@@ -42,9 +42,12 @@ class TestFailClosedAfterNotADb:
reopen = MagicMock() reopen = MagicMock()
monkeypatch.setattr("hermes_state._connect_tracked_db", reopen) monkeypatch.setattr("hermes_state._connect_tracked_db", reopen)
db._conn = _NotADbOnce(real_conn) db._conn = _NotADbOnce(real_conn)
with pytest.raises(sqlite3.DatabaseError, match="not a database"): with pytest.raises(sqlite3.DatabaseError, match="not a database") as excinfo:
db.create_session(session_id="s2", source="cli", model="test") db.create_session(session_id="s2", source="cli", model="test")
reopen.assert_not_called() reopen.assert_not_called()
# NOTADB on a live write is structural: the handle is quarantined.
assert isinstance(excinfo.value, StateDbCorruptError)
assert db._db_corrupt is True
finally: finally:
db._conn = real_conn db._conn = real_conn
db.close() db.close()