From 3f7ebed586f398c484afbd230c751849fb03a2e4 Mon Sep 17 00:00:00 2001 From: Kennedy Umege Date: Fri, 28 Aug 2026 23:11:59 +0100 Subject: [PATCH] fix(state): converge concurrent WAL initialization MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A burst of fresh connections to a database that is still in DELETE mode all run `PRAGMA journal_mode=WAL`; only one can win the exclusive header transition. SQLite answers the losers with SQLITE_BUSY *immediately* on this lock-upgrade path — the connection's busy_timeout is not consulted — so `apply_wal_with_fallback` raised "database is locked" and the opener failed, even though a peer was establishing WAL at that very moment. Measured on macOS with SQLite 3.53.1 (16 concurrent fresh openers, 40 trials): 81 of 640 opens through `apply_wal_with_fallback` failed on unpatched main, with a maximum observed wait of 13 ms — i.e. the 5 s busy timeout never engaged. Fix: `_set_wal_with_busy_convergence` retries only SQLITE_BUSY, with jitter, inside one bounded monotonic window (1 s). After each BUSY it re-probes the on-disk journal mode, so a loser converges on the peer's WAL instead of retrying the pragma. Exhaustion re-raises the original SQLITE_BUSY. It never downgrades to DELETE and deliberately does not retry SQLITE_LOCKED, WAL-incompatible storage errors, or other faults, so the existing NFS/SMB/ZFS fallback and the WAL-reset gate are untouched. With the fix the same measurement opens 640 of 640 in WAL, maximum wait 54 ms. Tests cover the four contracts: peer-established WAL is accepted after one BUSY, a bounded retry may establish WAL locally, deadline exhaustion fails closed without ever issuing journal_mode=DELETE, and SQLITE_LOCKED is not treated as the convergence race. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01VEiv7g2yvXE19vwSr3oRLt --- hermes_state.py | 58 ++++++++++++++++++- tests/test_hermes_state_wal_fallback.py | 77 +++++++++++++++++++++++++ 2 files changed, 134 insertions(+), 1 deletion(-) diff --git a/hermes_state.py b/hermes_state.py index 7e308704e7156..0cf4bf21eda5f 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -970,6 +970,15 @@ def _ensure_test_isolation(db_path: Path) -> None: "disk i/o error", # ZFS SHM corruption under concurrent connections ) +# A burst of fresh connections can all probe DELETE before one wins the +# exclusive header transition to WAL. SQLite may return SQLITE_BUSY from +# PRAGMA journal_mode=WAL immediately in that lock-upgrade race even when the +# connection has a longer busy_timeout. Give the winner a short bounded window +# to publish WAL, then let the losing connections converge by re-probing. +_WAL_BUSY_CONVERGENCE_SECONDS = 1.0 +_WAL_BUSY_RETRY_MIN_SECONDS = 0.01 +_WAL_BUSY_RETRY_MAX_SECONDS = 0.05 + # Upper bound for the write-ahead log. SQLite defaults to -1 (unlimited), # which lets state.db-wal keep the high-water mark of the largest-ever # transaction forever. See _apply_wal_size_limit(). @@ -1197,6 +1206,53 @@ def _on_disk_journal_mode(conn: sqlite3.Connection) -> Optional[str]: return None +def _is_sqlite_busy(exc: sqlite3.OperationalError) -> bool: + """Return whether ``exc`` is SQLITE_BUSY (not SQLITE_LOCKED).""" + error_code = getattr(exc, "sqlite_errorcode", None) + if isinstance(error_code, int): + # Extended result codes retain the primary code in the low byte. + return error_code & 0xFF == sqlite3.SQLITE_BUSY + # Compatibility for runtimes/test doubles without sqlite_errorcode. + return str(exc).strip().lower() in {"database is locked", "database is busy"} + + +def _set_wal_with_busy_convergence(conn: sqlite3.Connection): + """Set WAL while absorbing only the fresh-connection SQLITE_BUSY race. + + A losing connection re-probes after BUSY: if a peer established WAL, that + is success; otherwise retry until one monotonic deadline. Exhaustion + re-raises SQLITE_BUSY. This never switches to DELETE and deliberately does + not retry SQLITE_LOCKED, WAL-incompatible storage errors, or other faults. + """ + deadline = time.monotonic() + _WAL_BUSY_CONVERGENCE_SECONDS + last_busy: Optional[sqlite3.OperationalError] = None + while True: + try: + return conn.execute("PRAGMA journal_mode=WAL").fetchone() + except sqlite3.OperationalError as exc: + if not _is_sqlite_busy(exc): + raise + last_busy = exc + + # Re-probe even if SQLite consumed the time budget waiting: a peer may + # have completed the WAL transition while this setter was blocked. + if _on_disk_journal_mode(conn) == "wal": + return ("wal",) + remaining = deadline - time.monotonic() + if remaining <= 0: + assert last_busy is not None + raise last_busy + time.sleep( + min( + remaining, + random.uniform( + _WAL_BUSY_RETRY_MIN_SECONDS, + _WAL_BUSY_RETRY_MAX_SECONDS, + ), + ) + ) + + def _apply_wal_size_limit(conn: sqlite3.Connection) -> None: """Bound the WAL so it returns space to the OS after big transactions. @@ -1524,7 +1580,7 @@ def apply_wal_with_fallback( # returned row, not the mere absence of an exception; otherwise we # report a false ``"wal"`` AND skip the fallback WARNING, leaving the # DB silently in DELETE (reader-blocks-writer) with no signal. - row = conn.execute("PRAGMA journal_mode=WAL").fetchone() + row = _set_wal_with_busy_convergence(conn) mode = str(row[0]).strip().lower() if row and row[0] is not None else "" if mode == "wal": if _upgrading_existing_db: diff --git a/tests/test_hermes_state_wal_fallback.py b/tests/test_hermes_state_wal_fallback.py index b71dd32b61768..c0bbd5b1588bb 100644 --- a/tests/test_hermes_state_wal_fallback.py +++ b/tests/test_hermes_state_wal_fallback.py @@ -114,6 +114,83 @@ def test_succeeds_on_local_fs(self, tmp_path): assert cur.fetchone()[0].lower() == "wal" conn.close() + def test_busy_setter_converges_when_peer_establishes_wal(self, monkeypatch): + """A loser in the fresh-connection WAL race accepts the peer's WAL.""" + attempts = [0] + + class _Conn: + def execute(self, sql): + attempts[0] += 1 + raise sqlite3.OperationalError("database is locked") + + monkeypatch.setattr( + hermes_state, "_on_disk_journal_mode", lambda _conn: "wal" + ) + + assert hermes_state._set_wal_with_busy_convergence(_Conn()) == ("wal",) + assert attempts[0] == 1 + + def test_busy_setter_retries_when_mode_remains_delete(self, monkeypatch): + """If no peer won yet, one bounded retry may establish WAL locally.""" + attempts = [0] + + class _Cursor: + def fetchone(self): + return ("wal",) + + class _Conn: + def execute(self, sql): + attempts[0] += 1 + if attempts[0] == 1: + raise sqlite3.OperationalError("database is locked") + return _Cursor() + + monkeypatch.setattr( + hermes_state, "_on_disk_journal_mode", lambda _conn: "delete" + ) + monkeypatch.setattr(hermes_state.time, "sleep", lambda _seconds: None) + + assert hermes_state._set_wal_with_busy_convergence(_Conn()) == ("wal",) + assert attempts[0] == 2 + + def test_busy_setter_deadline_fails_closed_without_delete(self, monkeypatch): + """Persistent BUSY is bounded and never converted into a downgrade.""" + calls = [] + clock = iter([0.0, 2.0]) + + class _Conn: + def execute(self, sql): + calls.append(sql) + raise sqlite3.OperationalError("database is locked") + + monkeypatch.setattr( + hermes_state, "_on_disk_journal_mode", lambda _conn: "delete" + ) + monkeypatch.setattr(hermes_state.time, "monotonic", lambda: next(clock)) + + with pytest.raises(sqlite3.OperationalError, match="database is locked"): + hermes_state._set_wal_with_busy_convergence(_Conn()) + + assert calls == ["PRAGMA journal_mode=WAL"] + assert not any("journal_mode=DELETE" in sql for sql in calls) + + def test_sqlite_locked_is_not_retried_as_busy(self, monkeypatch): + """SQLITE_LOCKED is not the cross-connection convergence race.""" + attempts = [0] + + class _Conn: + def execute(self, sql): + attempts[0] += 1 + raise sqlite3.OperationalError("database table is locked") + + monkeypatch.setattr( + hermes_state, "_on_disk_journal_mode", lambda _conn: "delete" + ) + + with pytest.raises(sqlite3.OperationalError, match="database table is locked"): + hermes_state._set_wal_with_busy_convergence(_Conn()) + assert attempts[0] == 1 + def test_falls_back_to_delete_on_locking_protocol(self, tmp_path, caplog): """NFS-style ``locking protocol`` error → DELETE mode + one ERROR.""" conn, _ = _open_blocking(tmp_path / "nfs.db", isolation_level=None)