Skip to content
Open
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
58 changes: 57 additions & 1 deletion hermes_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -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().
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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:
Expand Down
77 changes: 77 additions & 0 deletions tests/test_hermes_state_wal_fallback.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading