fix(state): retry a transient SQLITE_IOERR on the pooled read path before surfacing it
Since 0.21.0 reads go through mode=ro pooled connections. A read-only OPEN already rides out the millisecond WAL transition window (checkpoint / WAL reset / frame flush by a sibling process; the ro reader cannot rewrite the -shm index) with a bounded retry (#100436), but a WARM pooled reader hitting the same window while its SELECT executes propagated `disk I/O error` straight out of get_session(): 37 identical tracebacks on a multi-process WSL2 ext4-on-vhdx install, each followed by "compression session recovery failed", with quick_check=ok (#100871). The reporter's A/B shows the operator workaround (journal_mode=delete) collapses read throughput ~30000x, so the flake has to be absorbed on the read path. _read_one/_read_all now replay the idempotent statement within the existing read-only IOERR budget (3 x 50 ms) on the SAME connection -- close+reopen would cancel this process's POSIX locks for every sibling connection -- and a persistent IOERR still propagates. No quarantine: EIO on a read is busy, not broken. Every SELECT in the SessionDB siblings (63 call sites) reaches the pool through these two helpers, so the class is covered without a wrapper type. Same-connection retry per #100882's analysis (@fangliquanflq); #100883 (@Sahilvishnaliya) diagnosed the missing recovery in the 0.21.0 read pool. Fixes #100871. Co-authored-by: fangliquanflq <fangliquan@qq.com> Co-authored-by: Sahilvishnaliya <222165401+Sahilvishnaliya@users.noreply.github.com>
This commit is contained in:
@@ -931,13 +931,30 @@ class SessionDB(
|
||||
|
||||
def _read_one(self, sql: str, params: Any = ()) -> Optional[sqlite3.Row]:
|
||||
"""``fetchone()`` of one read-only statement via ``_read_ctx``."""
|
||||
with self._read_ctx() as conn:
|
||||
return conn.execute(sql, params).fetchone()
|
||||
return self._read_retrying_ioerr(lambda conn: conn.execute(sql, params).fetchone())
|
||||
|
||||
def _read_all(self, sql: str, params: Any = ()) -> List[sqlite3.Row]:
|
||||
"""``fetchall()`` of one read-only statement via ``_read_ctx``."""
|
||||
with self._read_ctx() as conn:
|
||||
return conn.execute(sql, params).fetchall()
|
||||
return self._read_retrying_ioerr(lambda conn: conn.execute(sql, params).fetchall())
|
||||
|
||||
def _read_retrying_ioerr(self, fn: Callable[[sqlite3.Connection], T]) -> T:
|
||||
"""Run an idempotent SELECT through ``_read_ctx``, retrying a transient SQLITE_IOERR.
|
||||
|
||||
A warm ``mode=ro`` pooled reader can hit the same millisecond-wide WAL transition window as a
|
||||
read-only OPEN (#100436) when its statement executes or steps: a sibling process's checkpoint /
|
||||
WAL reset / frame flush surfaces ``disk I/O error`` because a read-only connection cannot rewrite
|
||||
the -shm index (#100871, WSL2 ext4-on-vhdx, multi-process). The window closes on its own, so
|
||||
the statement is replayed on the SAME connection within the read-only IOERR budget -- never
|
||||
closed and reopened (close() cancels this process's POSIX locks for every sibling connection),
|
||||
never quarantined (busy is not broken). A persistent IOERR exhausts the budget and propagates."""
|
||||
for attempt in range(_READ_ONLY_IOERR_RETRY_ATTEMPTS + 1):
|
||||
try:
|
||||
with self._read_ctx() as conn:
|
||||
return fn(conn)
|
||||
except sqlite3.OperationalError as exc:
|
||||
if attempt >= _READ_ONLY_IOERR_RETRY_ATTEMPTS or _DISK_IO_ERROR_MARKER not in str(exc).lower():
|
||||
raise
|
||||
time.sleep(_READ_ONLY_IOERR_RETRY_BACKOFF_S)
|
||||
|
||||
def _ensure_db_file_generation(self) -> None:
|
||||
"""Mint a once-per-file generation stamp (state_meta + application_id). First opener wins (INSERT
|
||||
|
||||
62
tests/hermes_state/test_read_path_transient_ioerr.py
Normal file
62
tests/hermes_state/test_read_path_transient_ioerr.py
Normal file
@@ -0,0 +1,62 @@
|
||||
"""#100871: a transient SQLITE_IOERR on a warm pooled read connection is retried on the same
|
||||
connection, then surfaced -- never quarantined, never a close+reopen."""
|
||||
import sqlite3
|
||||
|
||||
import pytest
|
||||
|
||||
import hermes_state
|
||||
from hermes_state import SessionDB
|
||||
|
||||
|
||||
_STATE = {"failures_left": 0, "attempts": 0} # module-level: the tracking factory subclasses _FlakyReads
|
||||
|
||||
|
||||
class _FlakyReads(sqlite3.Connection):
|
||||
"""Real SQLite connection whose first N SELECTs fail the way a mid-checkpoint mode=ro reader does."""
|
||||
|
||||
def execute(self, sql, *args, **kwargs): # type: ignore[override]
|
||||
if str(sql).lstrip().upper().startswith("SELECT"):
|
||||
_STATE["attempts"] += 1
|
||||
if _STATE["failures_left"] > 0:
|
||||
_STATE["failures_left"] -= 1
|
||||
raise sqlite3.OperationalError("disk I/O error")
|
||||
return super().execute(sql, *args, **kwargs)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def db(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr(hermes_state, "_READ_ONLY_IOERR_RETRY_BACKOFF_S", 0.0)
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
db.create_session("s", "cli")
|
||||
if not db._wal_active:
|
||||
db.close()
|
||||
pytest.skip("read pool needs WAL")
|
||||
real_connect = hermes_state._connect_tracked_db
|
||||
|
||||
def flaky_connect(path, *args, **kwargs):
|
||||
if str(path).startswith("file:"):
|
||||
kwargs["factory"] = _FlakyReads
|
||||
return real_connect(path, *args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(hermes_state, "_connect_tracked_db", flaky_connect)
|
||||
while db._evict_one_idle_read_conn(): # the next read opens through the flaky factory
|
||||
pass
|
||||
_STATE.update(failures_left=0, attempts=0)
|
||||
yield db
|
||||
db.close()
|
||||
|
||||
|
||||
def test_transient_ioerr_on_pooled_read_is_retried_on_the_same_connection(db):
|
||||
_STATE["failures_left"] = 1
|
||||
row = db.get_session("s")
|
||||
assert row is not None and row["id"] == "s"
|
||||
assert _STATE["attempts"] == 2 # one failure, one replay
|
||||
assert db._db_corrupt is False and db._db_wal_generation_lost is False
|
||||
|
||||
|
||||
def test_persistent_ioerr_propagates_after_the_budget(db):
|
||||
_STATE["failures_left"] = 10 ** 6
|
||||
with pytest.raises(sqlite3.OperationalError, match="disk I/O error"):
|
||||
db.get_session("s")
|
||||
assert _STATE["attempts"] == hermes_state._READ_ONLY_IOERR_RETRY_ATTEMPTS + 1
|
||||
assert db._db_corrupt is False # busy/EIO is not corruption: no quarantine
|
||||
Reference in New Issue
Block a user