refactor(state): repair/wal — inline single-use helpers (_bundle_bytes, _unlink_quiet, _promote_repaired_snapshot, _dedup_sqlite_master, _warn_once, _retry_wal_after_eio), contextlib.closing/suppress, pack calls
This commit is contained in:
+159
-236
@@ -15,8 +15,8 @@ import logging
|
||||
import os
|
||||
import shutil
|
||||
import sqlite3
|
||||
import stat
|
||||
import time
|
||||
from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
@@ -59,17 +59,6 @@ def _sidecars(db_path: Path):
|
||||
return (db_path.with_name(db_path.name + suffix) for suffix in _DB_SIDECAR_SUFFIXES)
|
||||
|
||||
|
||||
def _bundle_bytes(db_path: Path) -> int:
|
||||
"""Size of the main file plus every PRESENT sidecar (OSError propagates)."""
|
||||
return db_path.stat().st_size + sum(p.stat().st_size for p in _sidecars(db_path) if p.exists())
|
||||
|
||||
|
||||
def _unlink_quiet(path: Path) -> None:
|
||||
"""Best-effort unlink; a failure here is never the caller's error."""
|
||||
with contextlib.suppress(OSError):
|
||||
path.unlink(missing_ok=True)
|
||||
|
||||
|
||||
def _read_offline(db_path: Path, what: str, reader) -> Optional[str]:
|
||||
"""Run *reader()* under ``hermes_cli.sqlite_safe_read.offline_file_access``.
|
||||
|
||||
@@ -112,7 +101,6 @@ def _open_lock_file(db_path: Path, suffix: str, what: str, tail: str):
|
||||
|
||||
def _msvcrt_lock(handle, flag_name: str) -> None:
|
||||
import msvcrt
|
||||
|
||||
handle.seek(0)
|
||||
msvcrt.locking(handle.fileno(), getattr(msvcrt, flag_name), 1) # type: ignore[attr-defined]
|
||||
|
||||
@@ -124,26 +112,20 @@ def _try_lock_nonblocking(handle) -> None:
|
||||
_msvcrt_lock(handle, "LK_NBLCK")
|
||||
else:
|
||||
import fcntl
|
||||
|
||||
fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
|
||||
|
||||
def _release_lock_handle(handle, *, clear_record: bool = False) -> None:
|
||||
"""Drop the advisory lock on *handle* (best effort) and close it."""
|
||||
from hermes_state import _IS_WINDOWS
|
||||
try:
|
||||
with contextlib.closing(handle), contextlib.suppress(OSError): # best-effort release; always close
|
||||
if _IS_WINDOWS:
|
||||
_msvcrt_lock(handle, "LK_UNLCK")
|
||||
else:
|
||||
import fcntl
|
||||
|
||||
if clear_record:
|
||||
_clear_lock_holder_record(handle)
|
||||
fcntl.flock(handle.fileno(), fcntl.LOCK_UN)
|
||||
except OSError: # pragma: no cover - best effort release
|
||||
pass
|
||||
finally:
|
||||
handle.close()
|
||||
|
||||
|
||||
def _acquire_repair_lock_windows(lock_path: Path, handle, timeout: float):
|
||||
@@ -176,27 +158,24 @@ def _cross_process_repair_lock(db_path: Path):
|
||||
before the disk filled may be inside surgery."""
|
||||
from hermes_state import _IS_WINDOWS, _REPAIR_LOCK_TIMEOUT_SECONDS
|
||||
lock_path, handle = _open_lock_file(
|
||||
db_path, ".repair.lock", "repair", "skipping schema surgery rather than running it without cross-process authority.",
|
||||
)
|
||||
db_path, ".repair.lock", "repair", "skipping schema surgery rather than running it without cross-process authority.")
|
||||
if handle is None:
|
||||
yield False
|
||||
return
|
||||
|
||||
acquired = False
|
||||
try:
|
||||
if _IS_WINDOWS:
|
||||
acquired = _acquire_repair_lock_windows(lock_path, handle, _REPAIR_LOCK_TIMEOUT_SECONDS)
|
||||
else:
|
||||
acquired, handle = _acquire_db_flock(
|
||||
str(lock_path), handle, _REPAIR_LOCK_TIMEOUT_SECONDS, _REPAIR_LOCK_POLL_SECONDS, "state.db repair lock",
|
||||
)
|
||||
acquired, handle = _acquire_db_flock(str(lock_path), handle, _REPAIR_LOCK_TIMEOUT_SECONDS,
|
||||
_REPAIR_LOCK_POLL_SECONDS, "state.db repair lock")
|
||||
if acquired is None:
|
||||
acquired = False # non-contention failure already logged with its errno
|
||||
elif not acquired:
|
||||
record = None if _IS_WINDOWS else _read_lock_holder_record(handle)
|
||||
logger.warning("state.db repair lock %s held by another process for more than %.0fs — skipping schema "
|
||||
"surgery in this process to avoid racing the repairer. Recorded holder: %s.",
|
||||
lock_path, _REPAIR_LOCK_TIMEOUT_SECONDS, _describe_lock_holder(record))
|
||||
lock_path, _REPAIR_LOCK_TIMEOUT_SECONDS,
|
||||
_describe_lock_holder(None if _IS_WINDOWS else _read_lock_holder_record(handle)))
|
||||
yield acquired
|
||||
finally:
|
||||
if acquired:
|
||||
@@ -250,26 +229,26 @@ def _repair_backup_headroom_bytes(total_bytes: int) -> int:
|
||||
return max(_REPAIR_BACKUP_MIN_FREE_BYTES, int(total_bytes * _REPAIR_BACKUP_FREE_FRACTION))
|
||||
|
||||
|
||||
def _disk_budget(db_path: Path, refusal: str):
|
||||
"""``(bundle_bytes, free_bytes, headroom_bytes)`` for *db_path*'s volume, or an
|
||||
error string. Fails CLOSED on stat()/disk_usage() errors: the nearly-full
|
||||
def _disk_budget(db_path: Path, refusal: str) -> "Tuple[Optional[str], int, int, int]":
|
||||
"""``(error, bundle_bytes, free_bytes, headroom_bytes)`` for *db_path*'s volume (main file plus every PRESENT
|
||||
sidecar); *error* is set (and the sizes zero) on stat()/disk_usage() failure. Fails CLOSED: the nearly-full
|
||||
volume these guards exist for is exactly where they are most likely to fail."""
|
||||
try:
|
||||
need = _bundle_bytes(db_path)
|
||||
need = db_path.stat().st_size + sum(p.stat().st_size for p in _sidecars(db_path) if p.exists())
|
||||
usage = shutil.disk_usage(db_path.parent)
|
||||
except OSError as exc:
|
||||
return f"could not determine free space on {db_path.parent} ({exc}); refusing the {refusal} rather than risk filling the volume"
|
||||
return need, usage.free, _repair_backup_headroom_bytes(usage.total)
|
||||
return (f"could not determine free space on {db_path.parent} ({exc}); refusing the {refusal} rather than risk "
|
||||
"filling the volume"), 0, 0, 0
|
||||
return None, need, usage.free, _repair_backup_headroom_bytes(usage.total)
|
||||
|
||||
|
||||
def _repair_scratch_space_error(db_path: Path) -> Optional[str]:
|
||||
"""Return an error unless snapshot, VACUUM and promotion can fit safely."""
|
||||
budget = _disk_budget(db_path, "repair snapshot")
|
||||
if isinstance(budget, str):
|
||||
return budget
|
||||
snapshot_bytes, free, headroom = budget
|
||||
# VACUUM on the staged DB may need up to 2x the database size (SQLite docs);
|
||||
# the same reserve then covers transactional promotion into the live DB.
|
||||
error, snapshot_bytes, free, headroom = _disk_budget(db_path, "repair snapshot")
|
||||
if error is not None:
|
||||
return error
|
||||
# VACUUM on the staged DB may need up to 2x the database size (SQLite docs); the same reserve then covers
|
||||
# transactional promotion into the live DB.
|
||||
if free >= snapshot_bytes + (2 * snapshot_bytes) + headroom:
|
||||
return None
|
||||
return (f"only {free / 1e9:.2f}GB free on {db_path.parent}; the repair snapshot needs up to "
|
||||
@@ -281,10 +260,9 @@ def _backup_free_space_error(db_path: Path) -> Optional[str]:
|
||||
"""Disk guard for the forensic copy: reason to refuse, or None. A full raw copy on a nearly-full volume (which a
|
||||
preceding repair loop may itself have caused) can finish off the disk and every process on the machine."""
|
||||
hint = _MANUAL_RECOVER_HINT.format(db_path=db_path)
|
||||
budget = _disk_budget(db_path, "forensic copy")
|
||||
if isinstance(budget, str):
|
||||
return f"{budget}. {hint}"
|
||||
need, free, headroom = budget
|
||||
error, need, free, headroom = _disk_budget(db_path, "forensic copy")
|
||||
if error is not None:
|
||||
return f"{error}. {hint}"
|
||||
if free - need >= headroom:
|
||||
return None
|
||||
return (f"only {free / 1e9:.2f}GB free on {db_path.parent}; copying the damaged DB needs {need / 1e9:.2f}GB and must "
|
||||
@@ -309,14 +287,12 @@ def _repair_failure_consumes_attempt(exc: BaseException) -> bool:
|
||||
must not burn the repair ledger. Only SQLite's corruption/image result codes prove deterministic damage."""
|
||||
if not isinstance(exc, sqlite3.DatabaseError):
|
||||
return False
|
||||
error_code = getattr(exc, "sqlite_errorcode", None)
|
||||
if isinstance(error_code, int):
|
||||
if isinstance(error_code := getattr(exc, "sqlite_errorcode", None), int):
|
||||
# Extended result codes keep the primary code in the low byte.
|
||||
return (error_code & 0xFF) in (sqlite3.SQLITE_CORRUPT, sqlite3.SQLITE_NOTADB)
|
||||
# Older sqlite3 without result-code attributes: narrow message match only,
|
||||
# never turning generic "disk is full"/"readonly" into permanent failures.
|
||||
message = str(exc).lower()
|
||||
return "file is not a database" in message or "database disk image is malformed" in message
|
||||
# Older sqlite3 without result-code attributes: narrow message match only, never turning generic
|
||||
# "disk is full"/"readonly" into permanent failures.
|
||||
return any(m in str(exc).lower() for m in ("file is not a database", "database disk image is malformed"))
|
||||
|
||||
|
||||
def _repair_ledger_path(db_path: Path) -> Path:
|
||||
@@ -338,12 +314,9 @@ def _db_fingerprint(db_path: Path) -> "Optional[str]":
|
||||
st = db_path.stat()
|
||||
with open(db_path, "rb") as fh:
|
||||
head = fh.read(_FINGERPRINT_SAMPLE_BYTES)
|
||||
tail = b""
|
||||
if st.st_size > _FINGERPRINT_SAMPLE_BYTES:
|
||||
fh.seek(max(0, st.st_size - _FINGERPRINT_SAMPLE_BYTES))
|
||||
tail = fh.read(_FINGERPRINT_SAMPLE_BYTES)
|
||||
fh.seek(max(0, st.st_size - _FINGERPRINT_SAMPLE_BYTES))
|
||||
tail = fh.read(_FINGERPRINT_SAMPLE_BYTES) if st.st_size > _FINGERPRINT_SAMPLE_BYTES else b""
|
||||
return f"{st.st_size}:{hashlib.sha256(_mask_volatile_header(head) + tail).hexdigest()[:32]}"
|
||||
|
||||
return _read_offline(db_path, "fingerprint", _sample)
|
||||
|
||||
|
||||
@@ -367,7 +340,6 @@ def _backup_content_identity(db_path: Path) -> "Optional[str]":
|
||||
for chunk in iter(lambda: fh.read(1024 * 1024), b""):
|
||||
hasher.update(chunk)
|
||||
return hasher.hexdigest()
|
||||
|
||||
return _read_offline(db_path, "backup-identity", _digest)
|
||||
|
||||
|
||||
@@ -420,8 +392,8 @@ def _record_repair_outcome(db_path: Path, *, repaired: bool, fingerprint: "Optio
|
||||
fp = fingerprint if fingerprint is not None else _db_fingerprint(db_path)
|
||||
if fp is None:
|
||||
if not isinstance(recorded, str):
|
||||
# No prior key to extend and no safe way to mint one; the
|
||||
# in-process claim and cross-process lock still bound this run.
|
||||
# No prior key to extend and no safe way to mint one; the in-process claim and cross-process lock
|
||||
# still bound this run.
|
||||
return
|
||||
fp = recorded
|
||||
attempts = int(ledger.get("failed_attempts", 0)) + 1 if recorded == fp else 1
|
||||
@@ -457,22 +429,20 @@ def _publish_backup_bundle(db_path: Path, staging: Path, backup_path: Path) -> N
|
||||
ORDER MATTERS: the main DB name is the bundle's commit marker (what ``_existing_malformed_backups``
|
||||
counts), so sidecars publish FIRST and the main DB LAST — a failure partway never leaves a countable
|
||||
backup over a missing sidecar. On failure, staging files AND anything promoted are removed."""
|
||||
# (source, staged, destination); main DB copied first, published LAST.
|
||||
main = (db_path, staging, backup_path)
|
||||
triples: "List[Tuple[Path, Path, Path]]" = [
|
||||
(sidecar, staging.with_name(staging.name + suffix), backup_path.with_name(backup_path.name + suffix))
|
||||
for suffix, sidecar in zip(_DB_SIDECAR_SUFFIXES, _sidecars(db_path)) if sidecar.exists()
|
||||
] + [main]
|
||||
sidecars = [(sidecar, staging.with_name(staging.name + suffix), backup_path.with_name(backup_path.name + suffix))
|
||||
for suffix, sidecar in zip(_DB_SIDECAR_SUFFIXES, _sidecars(db_path)) if sidecar.exists()]
|
||||
published: "List[Path]" = []
|
||||
try:
|
||||
for src, staged, _dst in (main, *triples[:-1]):
|
||||
for src, staged, _dst in (main, *sidecars):
|
||||
shutil.copy2(src, staged)
|
||||
for _src, staged, dst in triples:
|
||||
for _src, staged, dst in (*sidecars, main):
|
||||
os.replace(staged, dst)
|
||||
published.append(dst)
|
||||
except Exception:
|
||||
for victim in (*(staged for _s, staged, _d in triples), *published):
|
||||
_unlink_quiet(victim)
|
||||
for victim in (*(staged for _s, staged, _d in (*sidecars, main)), *published):
|
||||
with contextlib.suppress(OSError): # best effort; a failure here is never the caller's error
|
||||
victim.unlink(missing_ok=True)
|
||||
raise
|
||||
|
||||
|
||||
@@ -489,17 +459,13 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]":
|
||||
— NOT mtime, NOT ``_db_fingerprint``); a repair loop once re-copied the same bytes on every restart. Staging
|
||||
names live OUTSIDE the ``.malformed-backup-`` prefix: inside it they count as a backup, sort NEWEST (prune kept
|
||||
partials, deleted intact copies) and dedupe could return one with no real forensic copy on disk."""
|
||||
try:
|
||||
with contextlib.suppress(ImportError): # scaffold/embed installs without hermes_cli track no connections
|
||||
from hermes_cli.sqlite_safe_read import has_live_connection
|
||||
live = has_live_connection(db_path)
|
||||
except ImportError:
|
||||
live = False
|
||||
if live:
|
||||
reason = (f"a connection to {db_path} is still open in this process; raw-copying it would cancel that "
|
||||
"connection's POSIX advisory locks. Close all SessionDB handles first.")
|
||||
logger.error("Refusing to raw-copy %s for backup: %s", db_path, reason)
|
||||
return None, reason
|
||||
|
||||
if has_live_connection(db_path):
|
||||
reason = (f"a connection to {db_path} is still open in this process; raw-copying it would cancel that "
|
||||
"connection's POSIX advisory locks. Close all SessionDB handles first.")
|
||||
logger.error("Refusing to raw-copy %s for backup: %s", db_path, reason)
|
||||
return None, reason
|
||||
stamp = datetime.datetime.now().strftime("%Y%m%d_%H%M%S")
|
||||
backup_path = db_path.with_name(f"{db_path.name}.malformed-backup-{stamp}")
|
||||
for seq in itertools.count(1): # same-second collision must not overwrite the earlier forensic copy
|
||||
@@ -512,22 +478,19 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]":
|
||||
# prefix-matches as a backup, sorts NEWEST and would otherwise survive prune forever.
|
||||
for pattern in (f"{db_path.name}.backup-staging-*", f"{db_path.name}.malformed-backup-*.incomplete*"):
|
||||
for old in db_path.parent.glob(pattern):
|
||||
_unlink_quiet(old)
|
||||
with contextlib.suppress(OSError):
|
||||
old.unlink(missing_ok=True)
|
||||
with contextlib.suppress(OSError):
|
||||
# Hash the source only when there is a candidate: hashing a multi-GB
|
||||
# file right before copying it is waste on the first-corruption pass.
|
||||
newest = _existing_malformed_backups(db_path)[:1]
|
||||
if newest:
|
||||
# Hash the source only when there is a candidate: hashing a multi-GB file before copying it is waste.
|
||||
if newest := _existing_malformed_backups(db_path)[:1]:
|
||||
src_id = _backup_content_identity(db_path)
|
||||
if src_id is not None and _backup_content_identity(newest[0]) == src_id:
|
||||
logger.info("Reusing existing forensic backup %s (identical to the damaged DB).", newest[0])
|
||||
return newest[0], None
|
||||
reason = _backup_free_space_error(db_path)
|
||||
if reason is not None:
|
||||
if (reason := _backup_free_space_error(db_path)) is not None:
|
||||
logger.error("Refusing forensic backup of %s: %s", db_path, reason)
|
||||
return None, reason
|
||||
staging = db_path.with_name(f"{db_path.name}.backup-staging-{stamp}")
|
||||
_publish_backup_bundle(db_path, staging, backup_path)
|
||||
_publish_backup_bundle(db_path, db_path.with_name(f"{db_path.name}.backup-staging-{stamp}"), backup_path)
|
||||
_prune_malformed_backups(db_path)
|
||||
return backup_path, None
|
||||
except Exception as exc: # pragma: no cover - best effort
|
||||
@@ -544,19 +507,14 @@ def preflight_db_writability(db_path: Path, *, db_label: str = "state.db") -> No
|
||||
tree (Hermes owns those files; ``chmod`` fails on files the user doesn't own, bounding the repair exactly);
|
||||
otherwise fail fast naming the file and command. Never deletes/truncates a WAL sidecar — once writable, the
|
||||
normal open checkpoints it. ``:memory:``/``file:`` skipped. Shared with ``kanban_db``."""
|
||||
import stat as _stat
|
||||
|
||||
if str(db_path) == ":memory:" or str(db_path).startswith("file:"):
|
||||
return
|
||||
home: Optional[Path] = None
|
||||
with contextlib.suppress(Exception): # pragma: no cover - defensive
|
||||
home = Path(get_hermes_home()).resolve()
|
||||
|
||||
# SQLite needs a writable directory in every journal mode (WAL/SHM sidecars,
|
||||
# or the rollback journal in DELETE mode).
|
||||
# SQLite needs a writable directory in every journal mode (WAL/SHM sidecars, or the DELETE-mode journal).
|
||||
sidecars = (db_path.with_name(db_path.name + "-wal"), db_path.with_name(db_path.name + "-shm"))
|
||||
targets = [(db_path.parent, True)] + [(p, False) for p in (db_path, *sidecars) if p.is_file()]
|
||||
for p, is_dir in targets:
|
||||
for p, is_dir in [(db_path.parent, True), *((p, False) for p in (db_path, *sidecars) if p.is_file())]:
|
||||
if (is_dir and not p.is_dir()) or os.access(p, os.R_OK | os.W_OK):
|
||||
continue
|
||||
x = "x" if is_dir else ""
|
||||
@@ -564,7 +522,7 @@ def preflight_db_writability(db_path: Path, *, db_label: str = "state.db") -> No
|
||||
with contextlib.suppress(OSError, ValueError):
|
||||
in_scope = home is not None and p.resolve().is_relative_to(home)
|
||||
if in_scope:
|
||||
os.chmod(p, p.stat().st_mode | _stat.S_IRUSR | _stat.S_IWUSR | (_stat.S_IXUSR if is_dir else 0))
|
||||
os.chmod(p, p.stat().st_mode | stat.S_IRUSR | stat.S_IWUSR | (stat.S_IXUSR if is_dir else 0))
|
||||
if in_scope and os.access(p, os.R_OK | os.W_OK):
|
||||
logger.info("%s preflight: repaired read-only %s (chmod u+rw%s)", db_label, p, x)
|
||||
continue
|
||||
@@ -573,8 +531,7 @@ def preflight_db_writability(db_path: Path, *, db_label: str = "state.db") -> No
|
||||
raise sqlite3.OperationalError(
|
||||
f"{db_label} is not writable: {'directory' if is_dir else 'file'} {p} is read-only for this user. Hermes "
|
||||
f"needs read-write access to open the database. Fix with: chmod u+rw{x} '{p}' (files owned by another "
|
||||
f"user may need sudo/chown).{wal_note}"
|
||||
)
|
||||
f"user may need sudo/chown).{wal_note}")
|
||||
|
||||
|
||||
def _connect_repair_durable(db_path: Path, *, timeout: float = 5.0) -> sqlite3.Connection:
|
||||
@@ -616,9 +573,7 @@ def apply_durability_barriers(conn: sqlite3.Connection) -> bool:
|
||||
ok = _reapply_durability_barriers(conn)
|
||||
with contextlib.suppress(Exception):
|
||||
from hermes_cli.config import cfg_get, load_config_readonly # local: avoids an import cycle
|
||||
|
||||
raw_synchronous = cfg_get(load_config_readonly(), "database", "synchronous", default=None)
|
||||
if raw_synchronous is not None:
|
||||
if (raw_synchronous := cfg_get(load_config_readonly(), "database", "synchronous", default=None)) is not None:
|
||||
_apply_synchronous_pragma(conn, raw_synchronous, db_label="state.db (guest)")
|
||||
return ok
|
||||
|
||||
@@ -635,16 +590,15 @@ def _open_exclusive(db_path: Path, begin: str) -> sqlite3.Connection:
|
||||
*begin*; closed (unpinned) and re-raised when exclusion cannot be taken."""
|
||||
conn = _connect_repair_durable(db_path, timeout=0.0)
|
||||
try:
|
||||
conn.execute("PRAGMA locking_mode=EXCLUSIVE")
|
||||
conn.execute(begin)
|
||||
conn.execute("ROLLBACK")
|
||||
for statement in ("PRAGMA locking_mode=EXCLUSIVE", begin, "ROLLBACK"):
|
||||
conn.execute(statement)
|
||||
except BaseException:
|
||||
_close_unpinned(conn)
|
||||
raise
|
||||
return conn
|
||||
|
||||
|
||||
@contextmanager
|
||||
@contextlib.contextmanager
|
||||
def _exclusive_repair_db_guard(db_path: Path):
|
||||
"""Yield ``(conn, None)`` — one live connection that excludes writers for repair surgery — or ``(None,
|
||||
exc)`` when exclusion could not be taken.
|
||||
@@ -663,15 +617,14 @@ def _exclusive_repair_db_guard(db_path: Path):
|
||||
try:
|
||||
yield guard, None
|
||||
finally:
|
||||
# Releasing the exclusive locks before close keeps a close-time checkpoint
|
||||
# from being mistaken for a repair write by callers that reopen immediately.
|
||||
# Releasing the exclusive locks before close keeps a close-time checkpoint from being mistaken for a
|
||||
# repair write by callers that reopen immediately.
|
||||
_close_unpinned(guard)
|
||||
|
||||
|
||||
def _copy_database_snapshot(
|
||||
source_path: Path, destination_path: Path, *, source_connection: Optional[sqlite3.Connection] = None,
|
||||
destination_connection: Optional[sqlite3.Connection] = None,
|
||||
) -> None:
|
||||
def _copy_database_snapshot(source_path: Path, destination_path: Path, *,
|
||||
source_connection: Optional[sqlite3.Connection] = None,
|
||||
destination_connection: Optional[sqlite3.Connection] = None) -> None:
|
||||
"""Copy one complete SQLite snapshot without replacing either file inode: the online backup API folds
|
||||
committed WAL frames into the source snapshot and writes the destination in one transaction (rolled
|
||||
back if interrupted), so ``state.db`` is never swapped out from under handles that refer to it."""
|
||||
@@ -694,8 +647,7 @@ def _copy_database_snapshot(
|
||||
|
||||
def _schema_not_built(exc: BaseException) -> bool:
|
||||
"""``no such table/column``: FTS5 / core tables not created yet (brand new file mid-init)."""
|
||||
msg = str(exc).lower()
|
||||
return "no such table" in msg or "no such column" in msg
|
||||
return any(m in str(exc).lower() for m in ("no such table", "no such column"))
|
||||
|
||||
|
||||
def _db_opens_cleanly(db_path: Path) -> Optional[str]:
|
||||
@@ -708,53 +660,51 @@ def _db_opens_cleanly(db_path: Path) -> Optional[str]:
|
||||
from hermes_state import SessionDB, load_fts5_cjk_extension
|
||||
conn = _connect_repair_durable(db_path)
|
||||
try:
|
||||
# Best-effort tokenizer load: messages_fts_cjk needs cjk_unicode61 before any
|
||||
# statement can touch it; tokenizer absence must never classify as corruption.
|
||||
load_fts5_cjk_extension(conn)
|
||||
conn.execute("PRAGMA journal_mode").fetchone()
|
||||
rows = conn.execute("PRAGMA integrity_check").fetchall()
|
||||
problems = [str(r[0]) for r in rows if r and str(r[0]).lower() != "ok"]
|
||||
if problems:
|
||||
return "; ".join(problems[:3])
|
||||
conn.execute("SELECT COUNT(*) FROM sessions").fetchone()
|
||||
|
||||
# FTS5 read probe: partial shadow-table corruption makes MATCH/snippet/rank
|
||||
# raise while integrity_check reports healthy. MATCH '""' (empty phrase)
|
||||
# parses, scans zero rows and exercises the shadow tables; FTS5 rejects MATCH ''.
|
||||
for fts_table in _FTS_TABLES:
|
||||
with contextlib.closing(conn):
|
||||
# Best-effort tokenizer load: messages_fts_cjk needs cjk_unicode61 before any statement can touch it;
|
||||
# tokenizer absence must never classify as corruption.
|
||||
load_fts5_cjk_extension(conn)
|
||||
conn.execute("PRAGMA journal_mode").fetchone()
|
||||
rows = conn.execute("PRAGMA integrity_check").fetchall()
|
||||
problems = [str(r[0]) for r in rows if r and str(r[0]).lower() != "ok"]
|
||||
if problems:
|
||||
return "; ".join(problems[:3])
|
||||
conn.execute("SELECT COUNT(*) FROM sessions").fetchone()
|
||||
# FTS5 read probe: partial shadow-table corruption makes MATCH/snippet/rank raise while integrity_check
|
||||
# reports healthy. MATCH '""' (empty phrase) parses, scans zero rows and exercises the shadow tables;
|
||||
# FTS5 rejects MATCH ''.
|
||||
for fts_table in _FTS_TABLES:
|
||||
try:
|
||||
conn.execute(f"SELECT 1 FROM {fts_table} WHERE {fts_table} MATCH '\"\"' LIMIT 1").fetchone()
|
||||
except sqlite3.DatabaseError as exc:
|
||||
# Builds without fts5/trigram raise "no such module|tokenizer"; calling that corruption would
|
||||
# send the DB into repair, whose final fallback deletes messages_fts%. "no such table/column" =
|
||||
# not built yet.
|
||||
benign = SessionDB._is_fts5_unavailable_error(exc) or _schema_not_built(exc)
|
||||
if not (isinstance(exc, sqlite3.OperationalError) and benign):
|
||||
return f"fts5 read probe failed on {fts_table}: {exc}"
|
||||
# FTS write probe: drive a row through the messages_fts* triggers in a transaction that is always
|
||||
# rolled back.
|
||||
probe_session_id = f"_hermes_fts_health_probe_{time.time_ns()}"
|
||||
try:
|
||||
conn.execute(f"SELECT 1 FROM {fts_table} WHERE {fts_table} MATCH '\"\"' LIMIT 1").fetchone()
|
||||
except sqlite3.DatabaseError as exc:
|
||||
# Builds without fts5/trigram raise "no such module|tokenizer"; calling that corruption would send the
|
||||
# DB into repair, whose final fallback deletes messages_fts%. "no such table/column" = not built yet.
|
||||
benign = SessionDB._is_fts5_unavailable_error(exc) or _schema_not_built(exc)
|
||||
if not (isinstance(exc, sqlite3.OperationalError) and benign):
|
||||
return f"fts5 read probe failed on {fts_table}: {exc}"
|
||||
|
||||
# FTS write probe: drive a row through the messages_fts* triggers in a
|
||||
# transaction that is always rolled back.
|
||||
probe_session_id = f"_hermes_fts_health_probe_{time.time_ns()}"
|
||||
try:
|
||||
conn.execute("BEGIN IMMEDIATE")
|
||||
conn.execute("INSERT INTO sessions (id, source, started_at) VALUES (?, ?, ?)",
|
||||
(probe_session_id, "_health_probe", time.time()))
|
||||
conn.execute("INSERT INTO messages (session_id, role, content, timestamp) VALUES (?, ?, ?, ?)",
|
||||
(probe_session_id, "user", "_fts_health_probe", time.time()))
|
||||
conn.execute("ROLLBACK")
|
||||
except sqlite3.OperationalError as exc:
|
||||
with contextlib.suppress(sqlite3.Error):
|
||||
conn.execute("BEGIN IMMEDIATE")
|
||||
conn.execute("INSERT INTO sessions (id, source, started_at) VALUES (?, ?, ?)",
|
||||
(probe_session_id, "_health_probe", time.time()))
|
||||
conn.execute("INSERT INTO messages (session_id, role, content, timestamp) VALUES (?, ?, ?, ?)",
|
||||
(probe_session_id, "user", "_fts_health_probe", time.time()))
|
||||
conn.execute("ROLLBACK")
|
||||
# Missing messages/sessions tables = brand new file mid-init, not corruption.
|
||||
# "no such tokenizer": this process lacks the cjk extension the DB's index
|
||||
# needs — capability gap; a tokenizer-less SessionDB drops the triggers itself.
|
||||
if _schema_not_built(exc) or "no such tokenizer: cjk_unicode61" in str(exc).lower():
|
||||
return None
|
||||
return str(exc)
|
||||
return None
|
||||
except sqlite3.OperationalError as exc:
|
||||
with contextlib.suppress(sqlite3.Error):
|
||||
conn.execute("ROLLBACK")
|
||||
# Missing messages/sessions tables = brand new file mid-init, not corruption. "no such tokenizer":
|
||||
# this process lacks the cjk extension the DB's index needs — capability gap; a tokenizer-less
|
||||
# SessionDB drops the triggers itself.
|
||||
if _schema_not_built(exc) or "no such tokenizer: cjk_unicode61" in str(exc).lower():
|
||||
return None
|
||||
return str(exc)
|
||||
return None
|
||||
except sqlite3.DatabaseError as exc:
|
||||
return str(exc)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _live_writer_holds_db(db_path: Path) -> bool:
|
||||
@@ -768,7 +718,7 @@ def _live_writer_holds_db(db_path: Path) -> bool:
|
||||
repair is then serialised only by the cross-process repairer lock."""
|
||||
try:
|
||||
probe = _open_exclusive(db_path, "BEGIN IMMEDIATE")
|
||||
with contextlib.suppress(Exception):
|
||||
with contextlib.suppress(Exception): # a close() error is not evidence of a holder
|
||||
_close_unpinned(probe)
|
||||
return False
|
||||
except sqlite3.OperationalError as exc:
|
||||
@@ -798,27 +748,26 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A
|
||||
unless ``backup=False``. Serialised across processes (gateway, Desktop backend and CLI open the same file;
|
||||
concurrent ``writable_schema`` surgery is itself a corruption source). Returns ``{repaired, strategy,
|
||||
backup_path, error}``."""
|
||||
from hermes_state import _cross_process_repair_lock, _db_opens_cleanly, _live_writer_holds_db, _persistent_repair_attempts_exhausted, _probe_journal_mode_for_repair, _record_repair_outcome, _repair_state_db_schema_locked
|
||||
from hermes_state import (_cross_process_repair_lock, _db_opens_cleanly, _live_writer_holds_db,
|
||||
_persistent_repair_attempts_exhausted, _probe_journal_mode_for_repair,
|
||||
_record_repair_outcome, _repair_state_db_schema_locked)
|
||||
report: Dict[str, Any] = {"repaired": False, "strategy": None, "backup_path": None, "error": None}
|
||||
|
||||
# Startup-watchdog lease: repair is I/O-bound (near-zero CPU), which the watchdog's CPU fallback would
|
||||
# misread as a parked deadlock. One lease (clamped to _MAX_LEASE_S=900) beats per-chunk renewal complexity.
|
||||
report_startup_progress(900.0, phase="state_db_repair")
|
||||
|
||||
db_path = Path(db_path)
|
||||
if not db_path.exists():
|
||||
report["error"] = f"{db_path} does not exist"
|
||||
return report
|
||||
|
||||
# Cross-restart cap: the in-memory claim bounds one process, but unhealable
|
||||
# b-tree damage used to re-run surgery + a fresh backup on EVERY restart.
|
||||
# Cross-restart cap: the in-memory claim bounds one process, but unhealable b-tree damage used to re-run
|
||||
# surgery + a fresh backup on EVERY restart.
|
||||
if _persistent_repair_attempts_exhausted(db_path):
|
||||
return _repair_skip(report, "skipped", _persistent_repair_exhausted_error(db_path))
|
||||
|
||||
with _cross_process_repair_lock(db_path) as holding_lock:
|
||||
if not holding_lock:
|
||||
# Another process holds the lock (or the lock file was unopenable);
|
||||
# it may have healed the file already, so re-probe before failing.
|
||||
# Another process holds the lock (or the lock file was unopenable); it may have healed the file
|
||||
# already, so re-probe before failing.
|
||||
if _db_opens_cleanly(db_path) is None:
|
||||
report["repaired"], report["strategy"] = True, "repaired_by_other_process"
|
||||
else:
|
||||
@@ -827,18 +776,15 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A
|
||||
return report
|
||||
|
||||
result = report
|
||||
# Recheck exhaustion after acquisition: a queued repairer can have
|
||||
# recorded the final failure while this process waited.
|
||||
# Recheck exhaustion after acquisition: a queued repairer can have recorded the final failure while
|
||||
# this process waited.
|
||||
if _persistent_repair_attempts_exhausted(db_path):
|
||||
_repair_skip(report, "skipped", _persistent_repair_exhausted_error(db_path))
|
||||
# WAL-holder preflight: fail closed for active readers before a backup is taken. Not the race defence — the
|
||||
# exclusive guard in the locked routine excludes writers through promotion and sees DELETE-mode readers too.
|
||||
elif _live_writer_holds_db(db_path):
|
||||
_repair_skip(
|
||||
report, "skipped",
|
||||
"a live writer still holds state.db; skipped schema surgery to avoid tearing b-tree pages under a "
|
||||
"concurrent writer. Stop the gateway (hermes gateway stop) and retry.",
|
||||
)
|
||||
_repair_skip(report, "skipped", "a live writer still holds state.db; skipped schema surgery to avoid tearing "
|
||||
"b-tree pages under a concurrent writer. Stop the gateway (hermes gateway stop) and retry.")
|
||||
else:
|
||||
# Probe journal mode BEFORE surgery: a rebuilt file comes back in the default (delete) mode and nothing
|
||||
# else records the flip. Unprobeable (damaged file) -> database.journal_mode is the restore target.
|
||||
@@ -850,8 +796,7 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A
|
||||
# Environmental aborts (before a strategy mutates the snapshot) are retriable, not proof of exhaustion;
|
||||
# the private marker stays out of the public report. The ledger update stays under the cross-process
|
||||
# lock so two repairers cannot lose each other's updates; a queued loser must not record at all.
|
||||
attempted = bool(result.pop("_repair_attempted", False))
|
||||
if attempted or result.get("repaired"):
|
||||
if result.pop("_repair_attempted", False) or result.get("repaired"):
|
||||
_record_repair_outcome(db_path, repaired=bool(result.get("repaired")))
|
||||
return result
|
||||
|
||||
@@ -899,84 +844,69 @@ def _repair_state_db_schema_locked(db_path: Path, *, backup: bool, report: Dict[
|
||||
file from the schema SQLite can still parse — when the damage IS in the schema b-tree (the ``malformed database
|
||||
schema ()`` class) every table hanging off the unreadable part is silently dropped, the probe still reports
|
||||
malformed, and repair returned ``repaired=False`` having destroyed what it was asked to save."""
|
||||
from hermes_state import _backup_db_file, _copy_database_snapshot, _db_opens_cleanly, _repair_scratch_space_error, _run_repair_strategies, _unlink_db_triple
|
||||
from hermes_state import (_backup_db_file, _copy_database_snapshot, _db_opens_cleanly,
|
||||
_repair_scratch_space_error, _run_repair_strategies, _unlink_db_triple)
|
||||
scratch = db_path.with_name(f"{db_path.name}.repair-scratch")
|
||||
cleanup_error = _unlink_db_triple(scratch)
|
||||
if cleanup_error is not None:
|
||||
if (cleanup_error := _unlink_db_triple(scratch)) is not None:
|
||||
return _repair_skip(report, "aborted", f"could not remove a stale repair snapshot before probing state.db: {cleanup_error}")
|
||||
|
||||
# Re-probe under the lock: a process we queued behind may have just repaired
|
||||
# the file; redoing surgery would undo it (the repair/re-corrupt cascade).
|
||||
# Re-probe under the lock: a process we queued behind may have just repaired the file; redoing surgery
|
||||
# would undo it (the repair/re-corrupt cascade).
|
||||
if _db_opens_cleanly(db_path) is None:
|
||||
report["repaired"], report["strategy"] = True, "already_healthy"
|
||||
return report
|
||||
|
||||
if backup:
|
||||
bpath, backup_error = _backup_db_file(db_path)
|
||||
report["backup_path"] = str(bpath) if bpath else None
|
||||
if bpath is None: # HARD STOP: the forensic image is the recovery path when every strategy fails.
|
||||
return _repair_skip(report, "aborted", "pre-repair backup refused; aborting schema repair to avoid "
|
||||
f"mutating the only copy of the damaged DB: {backup_error}")
|
||||
|
||||
# The forensic copy precedes this guard on purpose: its live-holder checks
|
||||
# would be poisoned by our own exclusive connection. Everything touching the
|
||||
# repair image or live promotion happens only under writer exclusion.
|
||||
# The forensic copy precedes this guard on purpose: its live-holder checks would be poisoned by our own
|
||||
# exclusive connection. Everything touching the repair image or live promotion happens under writer exclusion.
|
||||
with _exclusive_repair_db_guard(db_path) as (live_guard, guard_error):
|
||||
if live_guard is None:
|
||||
return _repair_skip(report, "skipped", "could not acquire exclusive state.db repair ownership; skipped "
|
||||
f"schema surgery to avoid overwriting a concurrent writer. Stop the gateway and retry: "
|
||||
f"{guard_error}", exc=guard_error)
|
||||
|
||||
space_error = _repair_scratch_space_error(db_path)
|
||||
if space_error is not None:
|
||||
if (space_error := _repair_scratch_space_error(db_path)) is not None:
|
||||
return _repair_skip(report, "aborted", space_error)
|
||||
|
||||
try:
|
||||
# Source = live_guard: it owns the exclusion, and a second connection
|
||||
# could be blocked by our own EXCLUSIVE lock on some SQLite builds.
|
||||
# Source = live_guard: it owns the exclusion, and a second connection could be blocked by our own
|
||||
# EXCLUSIVE lock on some SQLite builds.
|
||||
_copy_database_snapshot(db_path, scratch, source_connection=live_guard)
|
||||
except (OSError, sqlite3.Error, TimeoutError) as exc:
|
||||
_unlink_db_triple(scratch)
|
||||
return _repair_skip(report, "aborted", f"could not stage a complete SQLite repair snapshot of {db_path}: {exc}", exc=exc)
|
||||
|
||||
try:
|
||||
# Private marker for the outer wrapper: a strategy failure consumes the
|
||||
# persistent budget; a promotion failure is classified separately.
|
||||
# Private marker for the outer wrapper: a strategy failure consumes the persistent budget; a
|
||||
# promotion failure is classified separately.
|
||||
report["_repair_attempted"] = True
|
||||
_run_repair_strategies(scratch, report)
|
||||
if report.get("repaired"):
|
||||
_promote_repaired_snapshot(scratch, db_path, live_guard, report)
|
||||
# Never ``os.replace`` the live DB: Windows rejects replacement under open handles and POSIX would
|
||||
# leave those handles on the old inode. The guard keeps writer exclusion throughout.
|
||||
try:
|
||||
_copy_database_snapshot(scratch, db_path, destination_connection=live_guard)
|
||||
except (OSError, sqlite3.Error, TimeoutError) as exc:
|
||||
report.update(repaired=False, strategy=None,
|
||||
_repair_attempted=_repair_failure_consumes_attempt(exc))
|
||||
report["error"] = f"repaired snapshot could not be promoted transactionally: {exc}"
|
||||
logger.error("state.db repair promotion failed: %s", exc)
|
||||
else:
|
||||
logger.warning("state.db repaired via '%s' and promoted transactionally: %s",
|
||||
report.get("strategy"), db_path)
|
||||
if not report.get("repaired"):
|
||||
# Logged HERE, not in the strategies: they see the scratch copy, and
|
||||
# the message a human acts on must name a path that still exists.
|
||||
# Logged HERE, not in the strategies: they see the scratch copy, and the message a human acts on
|
||||
# must name a path that still exists.
|
||||
logger.error("state.db schema repair could not recover %s automatically (no committed canonical data "
|
||||
"was modified or lost; backup: %s); manual restore from backup may be required.",
|
||||
db_path, report["backup_path"])
|
||||
return report
|
||||
finally:
|
||||
# Never leave a half-repaired file beside the DB to be mistaken for the real thing.
|
||||
cleanup_error = _unlink_db_triple(scratch)
|
||||
if cleanup_error is not None:
|
||||
if (cleanup_error := _unlink_db_triple(scratch)) is not None:
|
||||
logger.warning("Could not remove state.db repair snapshot after repair: %s", cleanup_error)
|
||||
|
||||
|
||||
def _promote_repaired_snapshot(scratch: Path, db_path: Path, live_guard: sqlite3.Connection, report: Dict[str, Any]) -> None:
|
||||
"""Copy the repaired *scratch* back into the live DB through *live_guard*.
|
||||
|
||||
Never ``os.replace`` the live DB: Windows rejects replacement under open
|
||||
handles and POSIX would leave those handles on the old inode. The guard
|
||||
keeps writer exclusion throughout. On failure the report reverts to unrepaired."""
|
||||
from hermes_state import _copy_database_snapshot
|
||||
try:
|
||||
_copy_database_snapshot(scratch, db_path, destination_connection=live_guard)
|
||||
except (OSError, sqlite3.Error, TimeoutError) as exc:
|
||||
report.update(repaired=False, strategy=None, _repair_attempted=_repair_failure_consumes_attempt(exc))
|
||||
report["error"] = f"repaired snapshot could not be promoted transactionally: {exc}"
|
||||
logger.error("state.db repair promotion failed: %s", exc)
|
||||
else:
|
||||
logger.warning("state.db repaired via '%s' and promoted transactionally: %s", report.get("strategy"), db_path)
|
||||
|
||||
|
||||
def _unlink_db_triple(path: Path) -> Optional[str]:
|
||||
"""Remove *path* and every SQLite sidecar; return any cleanup failure."""
|
||||
from hermes_state import _IS_WINDOWS
|
||||
@@ -1014,10 +944,8 @@ def _strategy_rebuild_fts(conn: sqlite3.Connection) -> None:
|
||||
# The cjk index can only be rebuilt with its tokenizer loaded (best-effort).
|
||||
load_fts5_cjk_extension(conn)
|
||||
for table_name in _FTS_TABLES:
|
||||
try:
|
||||
with contextlib.suppress(sqlite3.OperationalError): # table absent (FTS disabled / trigram off / cjk absent)
|
||||
conn.execute(f"INSERT INTO {table_name}({table_name}) VALUES('rebuild')")
|
||||
except sqlite3.OperationalError:
|
||||
continue # table absent (FTS disabled / trigram off / cjk not present)
|
||||
|
||||
|
||||
def _strategy_reindex(conn: sqlite3.Connection) -> None:
|
||||
@@ -1029,19 +957,18 @@ def _strategy_reindex(conn: sqlite3.Connection) -> None:
|
||||
conn.commit()
|
||||
|
||||
|
||||
def _dedup_sqlite_master(conn: sqlite3.Connection) -> bool:
|
||||
"""Delete duplicate sqlite_master rows (lowest rowid per type/name wins)."""
|
||||
dupes = conn.execute(
|
||||
"SELECT type, name, COUNT(*) AS c, MIN(rowid) AS keep FROM sqlite_master GROUP BY type, name HAVING c > 1"
|
||||
).fetchall()
|
||||
for type_, name, _count, keep in dupes:
|
||||
conn.execute("DELETE FROM sqlite_master WHERE type IS ? AND name IS ? AND rowid <> ?", (type_, name, keep))
|
||||
return bool(dupes)
|
||||
|
||||
|
||||
def _strategy_dedup_schema(conn: sqlite3.Connection) -> None:
|
||||
"""De-duplicate sqlite_master, keeping FTS."""
|
||||
_edit_sqlite_master(conn, lambda: _dedup_sqlite_master(conn))
|
||||
"""De-duplicate sqlite_master (lowest rowid per type/name wins), keeping FTS."""
|
||||
def _dedup() -> bool:
|
||||
dupes = conn.execute(
|
||||
"SELECT type, name, COUNT(*) AS c, MIN(rowid) AS keep FROM sqlite_master GROUP BY type, name HAVING c > 1"
|
||||
).fetchall()
|
||||
for type_, name, _count, keep in dupes:
|
||||
conn.execute("DELETE FROM sqlite_master WHERE type IS ? AND name IS ? AND rowid <> ?",
|
||||
(type_, name, keep))
|
||||
return bool(dupes)
|
||||
|
||||
_edit_sqlite_master(conn, _dedup)
|
||||
|
||||
|
||||
def _strategy_drop_fts_vacuum(conn: sqlite3.Connection) -> None:
|
||||
@@ -1057,18 +984,15 @@ def _strategy_drop_fts_vacuum(conn: sqlite3.Connection) -> None:
|
||||
# (name, body, success log, failure log) in escalation order. failure log None =
|
||||
# final strategy: its failure lands in report["error"] instead of logged-and-skipped.
|
||||
_REPAIR_STRATEGIES = (
|
||||
("rebuild_fts", _strategy_rebuild_fts,
|
||||
"state.db FTS indexes rebuilt in place (schema preserved): %s",
|
||||
("rebuild_fts", _strategy_rebuild_fts, "state.db FTS indexes rebuilt in place (schema preserved): %s",
|
||||
"state.db FTS in-place rebuild pass failed: %s"),
|
||||
("reindex_btree", _strategy_reindex,
|
||||
"state.db B-tree indexes rebuilt via REINDEX: %s",
|
||||
("reindex_btree", _strategy_reindex, "state.db B-tree indexes rebuilt via REINDEX: %s",
|
||||
"state.db REINDEX pass failed: %s"),
|
||||
("dedup_schema", _strategy_dedup_schema,
|
||||
"state.db schema repaired by de-duplicating sqlite_master (FTS index preserved): %s",
|
||||
"state.db dedup repair pass failed: %s"),
|
||||
("drop_fts_rebuild", _strategy_drop_fts_vacuum,
|
||||
"state.db schema repaired by dropping FTS schema; indexes will rebuild from messages on next open: %s",
|
||||
None),
|
||||
"state.db schema repaired by dropping FTS schema; indexes will rebuild from messages on next open: %s", None),
|
||||
)
|
||||
|
||||
|
||||
@@ -1076,7 +1000,6 @@ def _run_repair_strategies(db_path: Path, report: Dict[str, Any]) -> Dict[str, A
|
||||
"""Escalating repair attempts, applied to *db_path* IN PLACE — only ever a scratch copy nothing else holds open,
|
||||
never the user's database. The "could not recover" log lives in the caller so it names the user's database."""
|
||||
from hermes_state import _db_opens_cleanly
|
||||
|
||||
for name, body, success_msg, failure_msg in _REPAIR_STRATEGIES:
|
||||
try:
|
||||
with _repair_conn(db_path) as conn:
|
||||
|
||||
+55
-93
@@ -30,7 +30,7 @@ _WAL_INCOMPAT_MARKERS = ("locking protocol", "not authorized", "disk i/o error")
|
||||
_WAL_SIZE_LIMIT_BYTES = 64 * 1024 * 1024 # 64 MiB
|
||||
|
||||
# Once-per-process-per-db_label dedup sets (kanban_db.connect() runs on every kanban operation, so an undeduped
|
||||
# line would repeat per connection). Tests clear these via ``hermes_state.<name>``; ``_warn_once`` resolves them there.
|
||||
# line would repeat per connection). Tests clear these via ``hermes_state.<name>``; ``_log_once`` resolves them there.
|
||||
_wal_fallback_warned_paths: set[str] = set()
|
||||
_wal_fallback_warned_lock = threading.Lock()
|
||||
_wal_reset_bug_warned_paths: set[str] = set()
|
||||
@@ -45,17 +45,6 @@ _CANNOT_VERIFY_DELETE_MSG = ("could not verify journal mode before applying conf
|
||||
"does not exclusively own")
|
||||
|
||||
|
||||
def _warn_once(lock: threading.Lock, set_name: str, key: str) -> bool:
|
||||
"""True the first time *key* is seen in ``hermes_state.<set_name>``."""
|
||||
import hermes_state
|
||||
seen = getattr(hermes_state, set_name)
|
||||
with lock:
|
||||
if key in seen:
|
||||
return False
|
||||
seen.add(key)
|
||||
return True
|
||||
|
||||
|
||||
def _mode_from_row(row) -> str:
|
||||
"""Lower-cased mode from a ``PRAGMA journal_mode`` row, ``""`` if no row."""
|
||||
return str(row[0]).strip().lower() if row and row[0] is not None else ""
|
||||
@@ -97,10 +86,9 @@ def _apply_wal_size_limit(conn: sqlite3.Connection) -> None:
|
||||
|
||||
def _darwin_pragma(conn: sqlite3.Connection, pragma: str) -> None:
|
||||
"""Best-effort PRAGMA on macOS only (no-op elsewhere, never raises)."""
|
||||
if sys.platform != "darwin":
|
||||
return
|
||||
with contextlib.suppress(sqlite3.OperationalError):
|
||||
conn.execute(pragma)
|
||||
if sys.platform == "darwin":
|
||||
with contextlib.suppress(sqlite3.OperationalError):
|
||||
conn.execute(pragma)
|
||||
|
||||
|
||||
def _apply_macos_checkpoint_barrier(conn: sqlite3.Connection) -> None:
|
||||
@@ -127,8 +115,7 @@ def _apply_wal_companions(conn: sqlite3.Connection) -> None:
|
||||
def is_sqlite_wal_reset_vulnerable(version_info: Optional[tuple] = None) -> bool:
|
||||
"""True when the linked SQLite has the WAL-reset bug (3.7.0–3.51.2; fixed 3.51.3+, backports 3.50.7 /
|
||||
3.44.6). Pre-WAL libraries are safe. https://sqlite.org/wal.html#walresetbug"""
|
||||
info = version_info if version_info is not None else sqlite3.sqlite_version_info
|
||||
return _is_sqlite_wal_reset_vulnerable(info)
|
||||
return _is_sqlite_wal_reset_vulnerable(sqlite3.sqlite_version_info if version_info is None else version_info)
|
||||
|
||||
|
||||
def sqlite_source_id() -> str:
|
||||
@@ -156,7 +143,6 @@ def resolve_journal_mode() -> str:
|
||||
durability: macOS virtiofs, NFS, SMB). Invalid values fail safe to ``wal``."""
|
||||
try:
|
||||
from hermes_cli.config import load_config_readonly
|
||||
|
||||
database = (load_config_readonly() or {}).get("database", {})
|
||||
raw = database.get("journal_mode", "wal") if isinstance(database, dict) else "wal"
|
||||
except Exception:
|
||||
@@ -194,13 +180,11 @@ def apply_wal_with_fallback(conn: sqlite3.Connection, *, db_label: str = "state.
|
||||
from hermes_state import is_sqlite_wal_reset_vulnerable, resolve_journal_mode
|
||||
configured = resolve_journal_mode()
|
||||
|
||||
# Vulnerable SQLite: never enable WAL on non-WAL files. Configured mode is
|
||||
# resolved first so an explicit DELETE request is still verified.
|
||||
# Vulnerable SQLite: never enable WAL on non-WAL files (configured mode resolved first so an explicit DELETE
|
||||
# request is still verified).
|
||||
if is_sqlite_wal_reset_vulnerable():
|
||||
return _apply_delete_for_wal_reset_bug(conn, db_label=db_label, require_delete=configured == "delete")
|
||||
|
||||
# Read-only probe (no flock/checkpoint/WAL-SHM unlink): WAL-init must not
|
||||
# unlink files other connections hold open.
|
||||
# Read-only probe (no flock/checkpoint/WAL-SHM unlink): WAL-init must not unlink files other connections hold.
|
||||
current_mode = _on_disk_journal_mode(conn)
|
||||
if current_mode == "wal":
|
||||
if configured == "delete":
|
||||
@@ -214,15 +198,13 @@ def apply_wal_with_fallback(conn: sqlite3.Connection, *, db_label: str = "state.
|
||||
# Probe failed (locked/busy): ownership not provably exclusive. Fail loudly.
|
||||
raise sqlite3.OperationalError(_CANNOT_VERIFY_DELETE_MSG)
|
||||
return _verify_configured_delete(_set_journal_mode_no_wait(conn, "DELETE"))
|
||||
|
||||
return _enable_wal(conn, db_label, require_wal, current_mode)
|
||||
|
||||
|
||||
def _enable_wal(conn: sqlite3.Connection, db_label: str, require_wal: bool, current_mode: Optional[str]) -> str:
|
||||
"""Flip a non-WAL, non-vulnerable connection to WAL, or fall back to DELETE."""
|
||||
# Decide BEFORE the flip whether it overwrites a mode somebody chose (probe and
|
||||
# page_count are only readable while the file is untouched). A 0-page DB has
|
||||
# no prior choice, and every caller reaches this before creating schema.
|
||||
# Decide BEFORE the flip whether it overwrites a mode somebody chose (probe and page_count are only readable
|
||||
# while the file is untouched). A 0-page DB has no prior choice, and every caller reaches this before schema.
|
||||
upgrading_existing_db = current_mode is not None and current_mode != "wal" and _database_has_content(conn)
|
||||
|
||||
def _wal_activated() -> str:
|
||||
@@ -249,11 +231,23 @@ def _enable_wal(conn: sqlite3.Connection, db_label: str, require_wal: bool, curr
|
||||
if not any(marker in msg for marker in _WAL_INCOMPAT_MARKERS):
|
||||
raise # unrelated OperationalError — don't silently swallow
|
||||
if "disk i/o error" in msg:
|
||||
activated, exc = _retry_wal_after_eio(conn, exc)
|
||||
if activated:
|
||||
return _wal_activated()
|
||||
# Never downgrade if WAL is on disk or the mode cannot be read (probe blocked
|
||||
# by a concurrent opener) — ownership is not provably exclusive either way.
|
||||
# Retry twice: EIO is either deterministic WAL-incompatibility (ZFS / APFS-CoW) or a one-shot transient,
|
||||
# and treating a transient as a permanent downgrade produced mixed-mode corruption (A downgrades to
|
||||
# DELETE while siblings set WAL). A non-EIO retry error propagates.
|
||||
for _ in range(2):
|
||||
time.sleep(0.05)
|
||||
try:
|
||||
row = conn.execute("PRAGMA journal_mode=WAL").fetchone()
|
||||
except sqlite3.OperationalError as retry_exc:
|
||||
if "disk i/o error" not in str(retry_exc).lower():
|
||||
raise
|
||||
exc = retry_exc
|
||||
continue
|
||||
if _mode_from_row(row) == "wal":
|
||||
return _wal_activated()
|
||||
break
|
||||
# Never downgrade if WAL is on disk or the mode cannot be read (probe blocked by a concurrent opener) —
|
||||
# ownership is not provably exclusive either way.
|
||||
if _on_disk_journal_mode(conn) in ("wal", None):
|
||||
raise
|
||||
if require_wal:
|
||||
@@ -263,24 +257,6 @@ def _enable_wal(conn: sqlite3.Connection, db_label: str, require_wal: bool, curr
|
||||
return "delete"
|
||||
|
||||
|
||||
def _retry_wal_after_eio(conn: sqlite3.Connection, exc: sqlite3.OperationalError):
|
||||
"""Retry ``journal_mode=WAL`` twice after ``disk i/o error``: EIO is either deterministic
|
||||
WAL-incompatibility (ZFS / APFS-CoW) or a one-shot transient, and treating a transient as a permanent
|
||||
downgrade produced mixed-mode corruption (A downgrades to DELETE while siblings set WAL). Returns
|
||||
``(wal_activated, last_exc)``; a non-EIO retry error propagates."""
|
||||
for _ in range(2):
|
||||
time.sleep(0.05)
|
||||
try:
|
||||
row = conn.execute("PRAGMA journal_mode=WAL").fetchone()
|
||||
except sqlite3.OperationalError as retry_exc:
|
||||
if "disk i/o error" not in str(retry_exc).lower():
|
||||
raise
|
||||
exc = retry_exc
|
||||
continue
|
||||
return _mode_from_row(row) == "wal", exc
|
||||
return False, exc
|
||||
|
||||
|
||||
def _set_journal_mode_no_wait(conn: sqlite3.Connection, mode: str) -> str:
|
||||
"""Execute ``PRAGMA journal_mode=<mode>`` without waiting on other openers.
|
||||
|
||||
@@ -343,7 +319,6 @@ def _wal_reset_repair_hint() -> str:
|
||||
"""Repair hint matching what ``hermes update`` can actually do for this install type."""
|
||||
try:
|
||||
from hermes_cli.config import detect_install_method, get_project_root, recommended_update_command_for_method
|
||||
|
||||
method = detect_install_method(get_project_root())
|
||||
cmd = recommended_update_command_for_method(method)
|
||||
if method in {"git", "unknown"}:
|
||||
@@ -353,9 +328,8 @@ def _wal_reset_repair_hint() -> str:
|
||||
return "install a Python build bundled with SQLite 3.51.3+ (or backports 3.50.7 / 3.44.6) and restart Hermes"
|
||||
|
||||
|
||||
# Once-per-(process, db_label) log table. Levels are deliberate: falling back to
|
||||
# DELETE and an ignored ``journal_mode: delete`` are real losses (ERROR); a non-WAL
|
||||
# -> WAL flip is normally desirable and only its invisibility was the problem (WARNING).
|
||||
# Once-per-(process, db_label) log table. Levels are deliberate: falling back to DELETE and an ignored
|
||||
# ``journal_mode: delete`` are real losses (ERROR); a non-WAL -> WAL flip is normally desirable (WARNING).
|
||||
_WAL_RESET_BUG_ACTIONS = {
|
||||
"indeterminate": ("journal mode could not be verified or exclusively switched (database is locked — possible "
|
||||
"concurrent openers); leaving the journal mode untouched (no live downgrade under concurrent "
|
||||
@@ -364,50 +338,45 @@ _WAL_RESET_BUG_ACTIONS = {
|
||||
"delete": "using journal_mode=DELETE instead of enabling WAL",
|
||||
}
|
||||
_ONCE_LOGS = {
|
||||
"wal_reset_bug": (
|
||||
_wal_reset_bug_warned_lock, "_wal_reset_bug_warned_paths", logging.WARNING,
|
||||
# Install-type-aware so the warning never promises a repair path that
|
||||
# doesn't exist for git/pip/system Python installs.
|
||||
"wal_reset_bug": (_wal_reset_bug_warned_lock, "_wal_reset_bug_warned_paths", logging.WARNING,
|
||||
# Install-type-aware so the warning never promises a repair path that doesn't exist for git/pip installs.
|
||||
"%s: linked SQLite %s (interpreter %s) is vulnerable to the WAL-reset corruption bug "
|
||||
"(https://sqlite.org/wal.html#walresetbug) — %s. Upgrade to SQLite 3.51.3+ (or backports 3.50.7 / 3.44.6); "
|
||||
"%s. See `hermes doctor`. This warning fires once per process per database.",
|
||||
),
|
||||
"journal_upgrade": (
|
||||
_journal_upgrade_warned_lock, "_journal_upgrade_warned_paths", logging.WARNING,
|
||||
# journal_mode is a property of the FILE: switching an existing DB to WAL rewrites its header and
|
||||
# outlives the process. Operators set DELETE on the file directly (the documented WAL-reset-bug
|
||||
# mitigation) and nothing told them the next open would silently put WAL back.
|
||||
"%s. See `hermes doctor`. This warning fires once per process per database."),
|
||||
"journal_upgrade": (_journal_upgrade_warned_lock, "_journal_upgrade_warned_paths", logging.WARNING,
|
||||
# journal_mode is a property of the FILE: switching an existing DB to WAL rewrites its header and outlives
|
||||
# the process. Operators set DELETE on the file (the WAL-reset-bug mitigation) and nothing told them the
|
||||
# next open would silently put WAL back.
|
||||
"%s: on-disk journal_mode was %s and has been switched to WAL. This rewrites the database header and "
|
||||
"persists after this process exits. If %s was a deliberate choice (for example the mitigation for the SQLite "
|
||||
"WAL-reset bug, or a WAL-unsafe filesystem), setting it with PRAGMA on the file will not survive -- every "
|
||||
"open re-applies the configured mode. Set `database.journal_mode: delete` in config.yaml to make it stick. "
|
||||
"This message fires once per process per database.",
|
||||
),
|
||||
"wal_fallback": (
|
||||
_wal_fallback_warned_lock, "_wal_fallback_warned_paths", logging.ERROR,
|
||||
"This message fires once per process per database."),
|
||||
"wal_fallback": (_wal_fallback_warned_lock, "_wal_fallback_warned_paths", logging.ERROR,
|
||||
# Under kanban dispatcher + workers a DELETE-mode write blocks readers as SQLITE_BUSY.
|
||||
"%s: WAL journal_mode unsupported on this filesystem (%s) — falling back to journal_mode=DELETE (slower "
|
||||
"rollback-journal mode; reduces concurrency but works on NFS/SMB/FUSE/ZFS). See "
|
||||
"https://www.sqlite.org/wal.html for details. This message fires once per process per database.",
|
||||
),
|
||||
"delete_overridden": (
|
||||
_delete_overridden_warned_lock, "_delete_overridden_warned_paths", logging.ERROR,
|
||||
# Never-live-downgrade keeps WAL; without this the operator never learns
|
||||
# that ``database.journal_mode: delete`` had no effect.
|
||||
"https://www.sqlite.org/wal.html for details. This message fires once per process per database."),
|
||||
"delete_overridden": (_delete_overridden_warned_lock, "_delete_overridden_warned_paths", logging.ERROR,
|
||||
# Never-live-downgrade keeps WAL; without this the operator never learns their delete had no effect.
|
||||
"%s: database.journal_mode=delete is configured but the on-disk database is already WAL; keeping WAL (a live "
|
||||
"downgrade under open connections can corrupt the DB). To apply journal_mode=DELETE, stop all connections to "
|
||||
"this DB and run a one-time offline 'PRAGMA journal_mode=DELETE' on the file. This message fires once per "
|
||||
"process per database.",
|
||||
),
|
||||
"process per database."),
|
||||
}
|
||||
|
||||
|
||||
def _log_once(kind: str, db_label: str, *args: Any) -> None:
|
||||
"""Emit ``_ONCE_LOGS[kind]`` once per (process, db_label). Callable *args* are
|
||||
resolved only after the dedupe check, so install-method probes run once."""
|
||||
import hermes_state
|
||||
lock, set_name, level, message = _ONCE_LOGS[kind]
|
||||
if _warn_once(lock, set_name, db_label):
|
||||
logger.log(level, message, db_label, *(a() if callable(a) else a for a in args))
|
||||
seen = getattr(hermes_state, set_name)
|
||||
with lock:
|
||||
if db_label in seen:
|
||||
return
|
||||
seen.add(db_label)
|
||||
logger.log(level, message, db_label, *(a() if callable(a) else a for a in args))
|
||||
|
||||
|
||||
def _log_wal_reset_bug_once(db_label: str, *, kept_wal: bool, indeterminate: bool = False) -> None:
|
||||
@@ -426,8 +395,7 @@ _log_wal_fallback_once = functools.partial(_log_once, "wal_fallback")
|
||||
_log_configured_delete_overridden_once = functools.partial(_log_once, "delete_overridden")
|
||||
|
||||
|
||||
# Operators write synchronous as a name; mapped here rather than passed through
|
||||
# so a typo becomes a warning instead of a silently different durability level.
|
||||
# Operators write synchronous as a name; mapped so a typo becomes a warning, not a silently different level.
|
||||
_SYNCHRONOUS_LEVELS: Dict[str, int] = {"OFF": 0, "NORMAL": 1, "FULL": 2, "EXTRA": 3}
|
||||
_SYNCHRONOUS_NAMES: Dict[int, str] = {v: k for k, v in _SYNCHRONOUS_LEVELS.items()}
|
||||
_SYNCHRONOUS_FULL = 2
|
||||
@@ -438,8 +406,7 @@ def resolve_synchronous_level(raw_value: Any) -> Optional[int]:
|
||||
case, or ``0``-``3``) to its PRAGMA integer; None for anything else so the
|
||||
caller warns and leaves the level untouched (guessing at durability is worse)."""
|
||||
if isinstance(raw_value, bool):
|
||||
# bool is an int subclass and YAML turns bare `on`/`off` into one.
|
||||
# "off" is a real durability choice; True is meaningless.
|
||||
# bool is an int subclass and YAML turns bare `on`/`off` into one; "off" is a real choice, True meaningless.
|
||||
return 0 if raw_value is False else None
|
||||
if isinstance(raw_value, int):
|
||||
return raw_value if raw_value in _SYNCHRONOUS_NAMES else None
|
||||
@@ -458,8 +425,7 @@ def _apply_synchronous_pragma(conn: sqlite3.Connection, raw_value: Any, *, db_la
|
||||
the platter, so an unrecognised value must not fall through to "SQLite default" the way a bad
|
||||
``cache_size`` can. Darwin floor: this runs after :func:`_enforce_macos_synchronous_full`, so a configured
|
||||
``NORMAL`` would silently undo the btree protection — raising is allowed, lowering is refused out loud."""
|
||||
level = resolve_synchronous_level(raw_value)
|
||||
if level is None:
|
||||
if (level := resolve_synchronous_level(raw_value)) is None:
|
||||
logger.warning("%s: ignoring unrecognized database.synchronous=%r (expected OFF, NORMAL, FULL, EXTRA, or 0-3)",
|
||||
db_label, raw_value)
|
||||
return
|
||||
@@ -481,15 +447,12 @@ def apply_database_pragmas(conn: sqlite3.Connection, *, db_label: str = "state.d
|
||||
between bundled/distro/Homebrew builds). Best-effort: failures are ignored so DB init never breaks on a
|
||||
malformed section. Applied to ALL connection types: writer, read_only, WAL readers."""
|
||||
try:
|
||||
# Local import avoids a circular import with hermes_cli.config.
|
||||
from hermes_cli.config import cfg_get, load_config_readonly
|
||||
|
||||
from hermes_cli.config import cfg_get, load_config_readonly # local: avoids a circular import
|
||||
cfg = load_config_readonly()
|
||||
except Exception:
|
||||
return
|
||||
for pragma_name in ("cache_size", "mmap_size", "temp_store", "wal_autocheckpoint", "journal_size_limit"):
|
||||
raw_value = cfg_get(cfg, "database", pragma_name, default=None)
|
||||
if raw_value is None:
|
||||
if (raw_value := cfg_get(cfg, "database", pragma_name, default=None)) is None:
|
||||
continue
|
||||
try:
|
||||
value = int(str(raw_value).strip())
|
||||
@@ -499,6 +462,5 @@ def apply_database_pragmas(conn: sqlite3.Connection, *, db_label: str = "state.d
|
||||
with contextlib.suppress(sqlite3.OperationalError):
|
||||
conn.execute(f"PRAGMA {pragma_name}={value}")
|
||||
# Last: sizing pragmas cannot change durability (see _apply_synchronous_pragma).
|
||||
raw_synchronous = cfg_get(cfg, "database", "synchronous", default=None)
|
||||
if raw_synchronous is not None:
|
||||
if (raw_synchronous := cfg_get(cfg, "database", "synchronous", default=None)) is not None:
|
||||
_apply_synchronous_pragma(conn, raw_synchronous, db_label=db_label)
|
||||
|
||||
Reference in New Issue
Block a user