diff --git a/hermes_state_repair.py b/hermes_state_repair.py index 39039327c2..a4cf6c4f13 100644 --- a/hermes_state_repair.py +++ b/hermes_state_repair.py @@ -1,8 +1,7 @@ """state.db repair, backup and writability preflight (split from hermes_state). -Every name is re-imported into ``hermes_state``; intra-module calls to -patchable helpers go through a lazy ``from hermes_state import ...`` at call -time so monkeypatches there still intercept. +Every name is re-imported into ``hermes_state``; intra-module calls to patchable helpers go through a lazy +``from hermes_state import ...`` at call time so monkeypatches there still intercept. """ from __future__ import annotations @@ -10,6 +9,7 @@ from __future__ import annotations import contextlib import datetime import hashlib +import itertools import json import logging import os @@ -31,27 +31,23 @@ from hermes_state_common import ( logger = logging.getLogger("hermes_state") _REPAIR_LOCK_POLL_SECONDS = 0.1 -# Snapshot copies are data transfer, not locking: bounded separately at 10 MiB/s -# with the historical two-minute floor. +# Snapshot copies are data transfer, not locking: bounded separately at 10 MiB/s (historical two-minute floor). _REPAIR_SNAPSHOT_MIN_THROUGHPUT_BYTES_PER_SECOND = 10 * 1024 * 1024 _MAX_PERSISTENT_REPAIR_ATTEMPTS = 3 _MAX_MALFORMED_BACKUPS = 3 -# Sidecars copied with a damaged DB and pruned with it. ``-journal``: DELETE mode -# (the NFS/SMB/FUSE/ZFS and WAL-reset-bug fallback) leaves a hot journal whenever -# a transaction was open; without it the forensic copy cannot be rolled back. +# Sidecars copied with a damaged DB and pruned with it. ``-journal``: DELETE mode (the NFS/SMB/FUSE/ZFS and +# WAL-reset-bug fallback) leaves a hot journal whenever a transaction was open; without it the copy cannot roll back. _DB_SIDECAR_SUFFIXES = ("-wal", "-shm", "-journal") -# Head/tail bytes sampled by ``_db_fingerprint``: changes on any genuine -# repair/truncation/restore while staying O(1) on a multi-GB file. +# Head/tail bytes sampled by ``_db_fingerprint``: changes on any real repair/truncation/restore, O(1) on a multi-GB file. _FINGERPRINT_SAMPLE_BYTES = 65536 -# Header ranges that move on ordinary commits, not on repair — file change counter -# (24-27), version-valid-for (92-95) — masked out of the sample. A malformed-SCHEMA -# DB still accepts writes and DELETE mode writes the main file directly, so without -# the mask any live write re-keys the ledger and the repair budget resets to 1 +# Header ranges that move on ordinary commits, not on repair — file change counter (24-27), version-valid-for +# (92-95) — masked out of the sample. A malformed-SCHEMA DB still accepts writes and DELETE mode writes the main +# file directly, so without the mask any live write re-keys the ledger and the repair budget resets to 1 # forever. The page-1 sqlite_master b-tree (repair identity) sits after byte 100. _FINGERPRINT_VOLATILE_HEADER_RANGES = ((24, 28), (92, 96)) -# Headroom for the forensic backup (a full raw copy; a repair loop on a large -# state.db is a disk amplifier). Proportional, not a flat multi-GB floor: a refused -# backup is a HARD STOP, and a big reserve would make repair never run on small volumes. +# Headroom for the forensic backup (a full raw copy; a repair loop on a large state.db is a disk amplifier). +# Proportional, not a flat multi-GB floor: a refused backup is a HARD STOP, and a big reserve would make repair +# never run on small volumes. _REPAIR_BACKUP_MIN_FREE_BYTES = 256 * 1024 * 1024 # 256 MiB absolute floor _REPAIR_BACKUP_FREE_FRACTION = 0.02 # plus 2% of the volume _FTS_TABLES = ("messages_fts", "messages_fts_trigram", "messages_fts_cjk") @@ -77,12 +73,10 @@ def _unlink_quiet(path: Path) -> None: def _read_offline(db_path: Path, what: str, reader) -> Optional[str]: """Run *reader()* under ``hermes_cli.sqlite_safe_read.offline_file_access``. - ``close()`` on ANY raw descriptor cancels every POSIX advisory lock this - process holds on the file, including a peer connection's RESERVED lock - (``sqlite_safe_read`` rule 1), so a raw read is only safe with no live - connection; ``None`` when that makes it unsafe or the file is unreadable. - Scaffold/embed installs without hermes_cli have no tracked connections. - """ + ``close()`` on ANY raw descriptor cancels every POSIX advisory lock this process holds on the file, + including a peer connection's RESERVED lock (``sqlite_safe_read`` rule 1), so a raw read is only safe with + no live connection; ``None`` when that makes it unsafe or the file is unreadable. Scaffold/embed installs + without hermes_cli have no tracked connections.""" try: from hermes_cli.sqlite_safe_read import LiveConnectionError, offline_file_access except ImportError: @@ -95,15 +89,13 @@ def _read_offline(db_path: Path, what: str, reader) -> Optional[str]: def _claim_repair_attempt(db_path: Path) -> bool: - """Claim the one-shot per-process repair attempt for *db_path*: True for the - first caller, False afterwards (bounds the repair/reopen loop and stops - concurrent callers racing surgery on one file).""" + """Claim the one-shot per-process repair attempt for *db_path*: True for the first caller, False + afterwards (bounds the repair/reopen loop and stops concurrent callers racing surgery on one file).""" from hermes_state import _repair_attempt_lock, _repair_attempted_paths - key = str(db_path) with _repair_attempt_lock: - if key in _repair_attempted_paths: + if str(db_path) in _repair_attempted_paths: return False - _repair_attempted_paths.add(key) + _repair_attempted_paths.add(str(db_path)) return True @@ -125,6 +117,17 @@ def _msvcrt_lock(handle, flag_name: str) -> None: msvcrt.locking(handle.fileno(), getattr(msvcrt, flag_name), 1) # type: ignore[attr-defined] +def _try_lock_nonblocking(handle) -> None: + """Take the advisory lock on *handle* without waiting (raises on contention).""" + from hermes_state import _IS_WINDOWS + if _IS_WINDOWS: + _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 @@ -148,14 +151,12 @@ def _acquire_repair_lock_windows(lock_path: Path, handle, timeout: float): deadline = time.monotonic() + timeout while True: try: - _msvcrt_lock(handle, "LK_NBLCK") + _try_lock_nonblocking(handle) return True except (BlockingIOError, OSError) as exc: if not is_advisory_lock_contention(exc): - logger.warning( - "Could not acquire state.db repair lock %s (%s) — " - "skipping schema surgery on a non-contention error.", lock_path, exc, - ) + logger.warning("Could not acquire state.db repair lock %s (%s) — skipping schema surgery on a " + "non-contention error.", lock_path, exc) return None if time.monotonic() >= deadline: return False @@ -166,16 +167,13 @@ def _acquire_repair_lock_windows(lock_path: Path, handle, timeout: float): def _cross_process_repair_lock(db_path: Path): """Serialize state.db schema surgery across processes. - Yields True when this process holds the repair lock, False when the bounded - acquire timed out or the lock file could not be opened; on False the caller - must NOT do surgery (unlocked surgery IS the interleaving this prevents). - ``flock``: the kernel drops it when the holder dies (a pidfile would wedge - every future repair); a forked child inheriting the fd is the exception, so - the acquire records pid + start time and breaks a provably dead holder's - lock. Bounded because a live repairer can sit in ``VACUUM`` for minutes. An - unopenable lock file (no space/inodes/descriptors) fails closed too: a - sibling that opened ITS handle before the disk filled may be inside surgery. - """ + Yields True when this process holds the repair lock, False when the bounded acquire timed out or the lock + file could not be opened; on False the caller must NOT do surgery (unlocked surgery IS the interleaving + this prevents). ``flock``: the kernel drops it when the holder dies (a pidfile would wedge every future + repair); a forked child inheriting the fd is the exception, so the acquire records pid + start time and + breaks a provably dead holder's lock. Bounded because a live repairer can sit in ``VACUUM`` for minutes. + An unopenable lock file (no space/inodes/descriptors) fails closed too: a sibling that opened ITS handle + 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.", @@ -196,11 +194,9 @@ def _cross_process_repair_lock(db_path: Path): 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), - ) + 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)) yield acquired finally: if acquired: @@ -210,24 +206,16 @@ def _cross_process_repair_lock(db_path: Path): def _try_acquire_auto_maintenance_lock(db_path: Path) -> Optional[Any]: - """Non-blocking cross-process lock for one auto-maintenance pass (None = skip - the pass: otherwise two startups both pass the interval check and the second - prunes a row the first has only just closed recoverably).""" - from hermes_state import _IS_WINDOWS + """Non-blocking cross-process lock for one auto-maintenance pass (None = skip the pass: otherwise two startups + both pass the interval check and the second prunes a row the first has only just closed recoverably).""" _lock_path, handle = _open_lock_file(db_path, ".auto-maintenance.lock", "auto-maintenance", "skipping automatic maintenance.") - if handle is None: - return None try: - if _IS_WINDOWS: - _msvcrt_lock(handle, "LK_NBLCK") - else: - import fcntl - - fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + if handle is not None: + _try_lock_nonblocking(handle) + return handle except (BlockingIOError, OSError): handle.close() return None - return handle _release_auto_maintenance_lock = _release_lock_handle # release a _try_acquire_auto_maintenance_lock handle @@ -236,11 +224,9 @@ _release_auto_maintenance_lock = _release_lock_handle # release a _try_acquire_ def _bump_schema_cookie(conn: sqlite3.Connection) -> None: """Increment the schema cookie after direct ``sqlite_master`` surgery. - Ordinary DDL bumps it and peers compare it before running a prepared - statement (how they discard a cached schema); ``writable_schema=ON`` edits - do NOT, so live connections would keep firing triggers into ``messages_fts*`` - shadow tables that no longer exist. Best-effort, never raises. - """ + Ordinary DDL bumps it and peers compare it before running a prepared statement (how they discard a cached + schema); ``writable_schema=ON`` edits do NOT, so live connections would keep firing triggers into + ``messages_fts*`` shadow tables that no longer exist. Best-effort, never raises.""" try: current = conn.execute("PRAGMA schema_version").fetchone()[0] # Wrap within SQLite's 32-bit signed range; peers compare for equality. @@ -272,10 +258,7 @@ def _disk_budget(db_path: Path, refusal: str): need = _bundle_bytes(db_path) usage = shutil.disk_usage(db_path.parent) except OSError as exc: - return ( - f"could not determine free space on {db_path.parent} ({exc}); " - f"refusing the {refusal} rather than risk filling the volume" - ) + 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) @@ -289,17 +272,14 @@ def _repair_scratch_space_error(db_path: Path) -> Optional[str]: # 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 " - f"{snapshot_bytes / 1e9:.2f}GB, VACUUM may need another {(2 * snapshot_bytes) / 1e9:.2f}GB, and " - f"{headroom / 1e9:.2f}GB must remain as headroom. Free disk space, then retry." - ) + return (f"only {free / 1e9:.2f}GB free on {db_path.parent}; the repair snapshot needs up to " + f"{snapshot_bytes / 1e9:.2f}GB, VACUUM may need another {(2 * snapshot_bytes) / 1e9:.2f}GB, and " + f"{headroom / 1e9:.2f}GB must remain as headroom. Free disk space, then retry.") 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.""" + """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): @@ -307,20 +287,17 @@ def _backup_free_space_error(db_path: Path) -> Optional[str]: need, free, headroom = budget 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 " - f"leave {headroom / 1e9:.2f}GB headroom. {hint}" - ) + return (f"only {free / 1e9:.2f}GB free on {db_path.parent}; copying the damaged DB needs {need / 1e9:.2f}GB and must " + f"leave {headroom / 1e9:.2f}GB headroom. {hint}") def _repair_snapshot_timeout_seconds(source_path: Path) -> float: - """Bound one SQLite snapshot by source size incl. sidecars (a WAL can hold - committed rows not yet in the main file), so a healthy large-database copy - is not cut off by the repair-lock timeout.""" + """Bound one SQLite snapshot by source size incl. sidecars (a WAL can hold committed rows not yet in the + main file), so a healthy large-database copy is not cut off by the repair-lock timeout.""" from hermes_state import _REPAIR_LOCK_TIMEOUT_SECONDS, _REPAIR_SNAPSHOT_MIN_THROUGHPUT_BYTES_PER_SECOND source_bytes = 0 for candidate in (source_path, *_sidecars(source_path)): - with contextlib.suppress(FileNotFoundError): + with contextlib.suppress(FileNotFoundError): # a sidecar may vanish mid-walk source_bytes += candidate.stat().st_size return max(_REPAIR_LOCK_TIMEOUT_SECONDS, source_bytes / _REPAIR_SNAPSHOT_MIN_THROUGHPUT_BYTES_PER_SECOND) @@ -328,10 +305,8 @@ def _repair_snapshot_timeout_seconds(source_path: Path) -> float: def _repair_failure_consumes_attempt(exc: BaseException) -> bool: """Whether a pre-strategy SQLite failure proves deterministic corruption. - Lock contention, timeouts, disk-full and I/O failures are environmental — a - retry may succeed, so they must not burn the repair ledger. Only SQLite's - corruption/image result codes prove deterministic damage. - """ + Lock contention, timeouts, disk-full and I/O failures are environmental — a retry may succeed, so they + 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) @@ -351,17 +326,14 @@ def _repair_ledger_path(db_path: Path) -> Path: def _db_fingerprint(db_path: Path) -> "Optional[str]": """Cheap identity for a damaged DB file: size + a bounded content sample. - EXCLUDES mtime: a malformed-schema DB still accepts writes, so live writers, - checkpoints and the strategies move mtime between passes; keyed on mtime - every pass looked like a NEW file and the attempt counter reset forever. - Hashing a multi-GB file per open is the cost this ledger avoids, so sample - the head/tail slices any real repair, truncation or restore must change. + EXCLUDES mtime: a malformed-schema DB still accepts writes, so live writers, checkpoints and the + strategies move mtime between passes; keyed on mtime every pass looked like a NEW file and the attempt + counter reset forever. Hashing a multi-GB file per open is the cost this ledger avoids, so sample the + head/tail slices any real repair, truncation or restore must change. - ``None`` = identity unavailable (:func:`_read_offline`; a live peer is expected - here since this runs BEFORE ``_backup_db_file``'s live-connection guard). - Callers MUST NOT substitute a differently-shaped key: the ledger compares for - equality, so alternating shapes never match and the unbounded loop returns. - """ + ``None`` = identity unavailable (:func:`_read_offline`; a live peer is expected here since this runs + BEFORE ``_backup_db_file``'s live-connection guard). Callers MUST NOT substitute a differently-shaped key: + the ledger compares for equality, so alternating shapes never match and the unbounded loop returns.""" def _sample() -> str: st = db_path.stat() with open(db_path, "rb") as fh: @@ -378,22 +350,18 @@ def _db_fingerprint(db_path: Path) -> "Optional[str]": def _backup_content_identity(db_path: Path) -> "Optional[str]": """Recovery-image identity for forensic-backup dedupe: whole file + sidecars. - A DIFFERENT relation from :func:`_db_fingerprint` (never conflate them): - the fingerprint answers "same repair epoch?" from head+tail only, and a live - writer can commit into an *interior* page while preserving size and both - 64 KiB slices — reusing a backup on that basis hands the operator a snapshot - predating real user data. A forensic copy must claim byte identity, so this - digests the ENTIRE main file plus every present sidecar (the WAL can hold - uncheckpointed frames). ``None`` when a live connection makes the read unsafe - (:func:`_read_offline`) — the caller then takes a fresh backup, never a false reuse. - """ + A DIFFERENT relation from :func:`_db_fingerprint` (never conflate them): the fingerprint answers "same + repair epoch?" from head+tail only, and a live writer can commit into an *interior* page while preserving + size and both 64 KiB slices — reusing a backup on that basis hands the operator a snapshot predating real + user data. A forensic copy must claim byte identity, so this digests the ENTIRE main file plus every + present sidecar (the WAL can hold uncheckpointed frames). ``None`` when a live connection makes the read + unsafe (:func:`_read_offline`) — the caller then takes a fresh backup, never a false reuse.""" def _digest() -> str: hasher = hashlib.sha256() members = [("main", db_path), *((sfx, p) for sfx, p in zip(_DB_SIDECAR_SUFFIXES, _sidecars(db_path)) if p.exists())] for label, path in members: - # Length-delimit every member (main file included) so the concatenation - # is prefix-free; otherwise a main-file tail could coincide with a - # main+sidecar split and dedupe two images together. + # Length-delimit every member (main file included) so the concatenation is prefix-free; otherwise a + # main-file tail could coincide with a main+sidecar split and dedupe two images together. hasher.update(f"\0{label}:{path.stat().st_size}\0".encode()) with open(path, "rb") as fh: for chunk in iter(lambda: fh.read(1024 * 1024), b""): @@ -406,56 +374,42 @@ def _backup_content_identity(db_path: Path) -> "Optional[str]": def _read_repair_ledger(db_path: Path) -> "Dict[str, Any]": with contextlib.suppress(OSError, ValueError): raw = json.loads(_repair_ledger_path(db_path).read_text(encoding="utf-8")) - if isinstance(raw, dict): - return raw + return raw if isinstance(raw, dict) else {} return {} def _persistent_repair_attempts_exhausted(db_path: Path) -> bool: """Whether *db_path* has already burned its cross-restart repair budget. - True only when the ledger records ``_MAX_PERSISTENT_REPAIR_ATTEMPTS`` - failures against the CURRENT fingerprint. Never raises; a missing/corrupt - ledger or unstatable DB reads as "not exhausted" (the in-process claim and - cross-process lock still bound one run). Fingerprint unavailable (live - connection) -> compare the SIZE the ledger recorded, otherwise a peer - connection hides an exhausted budget on every pass. - """ + True only when the ledger records ``_MAX_PERSISTENT_REPAIR_ATTEMPTS`` failures against the CURRENT + fingerprint. Never raises; a missing/corrupt ledger or unstatable DB reads as "not exhausted" (the + in-process claim and cross-process lock still bound one run). Fingerprint unavailable (live connection) -> + compare the SIZE the ledger recorded, otherwise a peer connection hides an exhausted budget on every pass.""" ledger = _read_repair_ledger(db_path) recorded = ledger.get("fingerprint") fp = _db_fingerprint(db_path) - if fp is None: - # Size is the one key component that needs no raw read. - try: - size_prefix = f"{db_path.stat().st_size}:" - except OSError: - return False - if not isinstance(recorded, str) or not recorded.startswith(size_prefix): - return False - elif recorded != fp: + try: # size is the one key component that needs no raw read + same = recorded == fp if fp is not None else ( + isinstance(recorded, str) and recorded.startswith(f"{db_path.stat().st_size}:")) + except OSError: return False - return int(ledger.get("failed_attempts", 0)) >= _MAX_PERSISTENT_REPAIR_ATTEMPTS + return same and int(ledger.get("failed_attempts", 0)) >= _MAX_PERSISTENT_REPAIR_ATTEMPTS def _persistent_repair_exhausted_error(db_path: Path) -> str: """The stable operator-facing diagnostic for an exhausted repair budget.""" - return ( - f"automatic repair has already failed {_MAX_PERSISTENT_REPAIR_ATTEMPTS} times on this exact file — the " - f"corruption is beyond the schema/FTS repair strategies (likely b-tree page damage). Manual recovery " - f"required: restore " - f"a backup, or salvage with `sqlite3 {db_path} \".recover\"`. " - f"Delete {_repair_ledger_path(db_path).name} to force another automatic attempt." - ) + return (f"automatic repair has already failed {_MAX_PERSISTENT_REPAIR_ATTEMPTS} times on this exact file — the " + f"corruption is beyond the schema/FTS repair strategies (likely b-tree page damage). Manual recovery " + f"required: restore a backup, or salvage with `sqlite3 {db_path} \".recover\"`. " + f"Delete {_repair_ledger_path(db_path).name} to force another automatic attempt.") def _record_repair_outcome(db_path: Path, *, repaired: bool, fingerprint: "Optional[str]" = None) -> None: """Update the persistent attempt ledger after a repair pass. Never raises. - Keys on the post-attempt fingerprint (what the NEXT exhaustion probe sees). - If a live connection makes it unavailable, keep the recorded key and still - increment — dropping the pass lets a peer reset the budget every time. - Never write a differently shaped key. - """ + Keys on the post-attempt fingerprint (what the NEXT exhaustion probe sees). If a live connection makes it + unavailable, keep the recorded key and still increment — dropping the pass lets a peer reset the budget + every time. Never write a differently shaped key.""" ledger_path = _repair_ledger_path(db_path) try: if repaired: @@ -472,9 +426,7 @@ def _record_repair_outcome(db_path: Path, *, repaired: bool, fingerprint: "Optio fp = recorded attempts = int(ledger.get("failed_attempts", 0)) + 1 if recorded == fp else 1 stamp = datetime.datetime.now().isoformat(timespec="seconds") - ledger_path.write_text( - json.dumps({"fingerprint": fp, "failed_attempts": attempts, "last_attempt": stamp}), encoding="utf-8", - ) + ledger_path.write_text(json.dumps({"fingerprint": fp, "failed_attempts": attempts, "last_attempt": stamp}), encoding="utf-8") except Exception as exc: # pragma: no cover - best effort logger.warning("Could not update state.db repair ledger: %s", exc) @@ -502,26 +454,24 @@ def _prune_malformed_backups(db_path: Path, keep: int = _MAX_MALFORMED_BACKUPS) def _publish_backup_bundle(db_path: Path, staging: Path, backup_path: Path) -> None: """Copy DB + sidecars to *staging* names, then rename each into place. - 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. - """ - pairs: "List[Tuple[Path, Path, Path]]" = [ + 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() - ] + for suffix, sidecar in zip(_DB_SIDECAR_SUFFIXES, _sidecars(db_path)) if sidecar.exists() + ] + [main] published: "List[Path]" = [] try: - shutil.copy2(db_path, staging) - for sidecar, staged, _dst in pairs: - shutil.copy2(sidecar, staged) - for _src, staged, dst in (*pairs, (db_path, staging, backup_path)): + for src, staged, _dst in (main, *triples[:-1]): + shutil.copy2(src, staged) + for _src, staged, dst in triples: os.replace(staged, dst) published.append(dst) except Exception: - for victim in (staging, *(s for _src, s, _d in pairs), *published): + for victim in (*(staged for _s, staged, _d in triples), *published): _unlink_quiet(victim) raise @@ -529,46 +479,37 @@ def _publish_backup_bundle(db_path: Path, staging: Path, backup_path: Path) -> N def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": """Raw-copy a (possibly malformed) DB plus sidecars to a timestamped backup. - Raw bytes on purpose: the DB won't open cleanly, so preserve them exactly - for forensics. Returns ``(backup_path, None)`` or ``(None, reason)``; repair - treats a refused backup as a HARD STOP because the bundle is the recovery - path when every strategy fails. Refuses while a connection to this DB is - live in the process: the raw read would ``close()`` a descriptor and cancel - that connection's POSIX advisory locks (``hermes_cli.sqlite_safe_read``) — + Raw bytes on purpose: the DB won't open cleanly, so preserve them exactly for forensics. Returns ``(backup_path, + None)`` or ``(None, reason)``; repair treats a refused backup as a HARD STOP because the bundle is the recovery + path when every strategy fails. Refuses while a connection to this DB is live in the process: the raw read would + ``close()`` a descriptor and cancel that connection's POSIX advisory locks (``hermes_cli.sqlite_safe_read``) — real case: one SessionDB enters repair while the gateway holds others. - Dedupe: reuse the newest backup when byte-identical to the current recovery - image (``_backup_content_identity`` — 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. - """ + Dedupe: reuse the newest backup when byte-identical to the current recovery image (``_backup_content_identity`` + — 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: 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." - ) + 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}") - # Same-second collision must not overwrite the earlier forensic copy. - seq = 1 - while backup_path.exists(): + for seq in itertools.count(1): # same-second collision must not overwrite the earlier forensic copy + if not backup_path.exists(): + break backup_path = db_path.with_name(f"{db_path.name}.malformed-backup-{stamp}_{seq}") - seq += 1 try: - # Sweep staging debris from an interrupted pass BEFORE the dedupe (it is - # byte-identical to the damaged DB, so dedupe would hand it back as a - # backup). Also the old ``.incomplete`` spelling, which prefix-matches as - # a backup, sorts NEWEST and would otherwise survive prune forever. + # Sweep staging debris from an interrupted pass BEFORE the dedupe (it is byte-identical to the damaged + # DB, so dedupe would hand it back as a backup). Also the old ``.incomplete`` spelling, which + # 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) @@ -597,84 +538,61 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": def preflight_db_writability(db_path: Path, *, db_label: str = "state.db") -> None: """Refuse-or-repair read-only DB files BEFORE the first connection opens. - A stray read-only ``state.db`` / ``-wal`` / ``-shm`` (sudo run, restored - backup, copied dotfiles) otherwise surfaces as an opaque "attempt to write a - readonly database" inside ``_init_schema``, and the obvious wrong "fix" - (deleting the ``-wal``) loses committed transactions. ``chmod u+rw`` repair - only inside the Hermes home 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``. - """ + A stray read-only ``state.db`` / ``-wal`` / ``-shm`` (sudo run, restored backup, copied dotfiles) otherwise + surfaces as an opaque "attempt to write a readonly database" inside ``_init_schema``, and the obvious wrong + "fix" (deleting the ``-wal``) loses committed transactions. ``chmod u+rw`` repair only inside the Hermes home + 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 - raw = str(db_path) - if raw == ":memory:" or raw.startswith("file:"): + if str(db_path) == ":memory:" or str(db_path).startswith("file:"): return - try: - home: Optional[Path] = Path(get_hermes_home()).resolve() - except Exception: # pragma: no cover - defensive - home = None + home: Optional[Path] = None + with contextlib.suppress(Exception): # pragma: no cover - defensive + home = Path(get_hermes_home()).resolve() - def _ensure_writable(p: Path, *, is_dir: bool = False) -> None: - if os.access(p, os.R_OK | os.W_OK): - return + # SQLite needs a writable directory in every journal mode (WAL/SHM sidecars, + # or the rollback journal in DELETE mode). + 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: + 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 "" in_scope = False with contextlib.suppress(OSError, ValueError): in_scope = home is not None and p.resolve().is_relative_to(home) if in_scope: - add = _stat.S_IRUSR | _stat.S_IWUSR | (_stat.S_IXUSR if is_dir else 0) - os.chmod(p, p.stat().st_mode | add) + 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" if is_dir else "") - return - kind = "directory" if is_dir else "file" - wal_note = ( - " Do NOT delete the -wal file — it contains committed data that " - "will be merged into the database once it is writable." - if p.name.endswith("-wal") else "" - ) + logger.info("%s preflight: repaired read-only %s (chmod u+rw%s)", db_label, p, x) + continue + wal_note = (" Do NOT delete the -wal file — it contains committed data that " + "will be merged into the database once it is writable." if p.name.endswith("-wal") else "") raise sqlite3.OperationalError( - f"{db_label} is not writable: {kind} {p} is read-only for this user. Hermes needs read-write access to " - f"open the database. Fix with: chmod u+rw{'x' if is_dir else ''} '{p}' (files owned by another user may " - f"need sudo/chown).{wal_note}" + 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}" ) - # SQLite needs a writable directory in every journal mode (WAL/SHM - # sidecars, or the rollback journal in DELETE mode). - if db_path.parent.is_dir(): - _ensure_writable(db_path.parent, is_dir=True) - for p in (db_path, db_path.with_name(db_path.name + "-wal"), db_path.with_name(db_path.name + "-shm")): - if p.is_file(): - _ensure_writable(p) - def _connect_repair_durable(db_path: Path, *, timeout: float = 5.0) -> sqlite3.Connection: """``sqlite3.connect`` for the repair/probe paths, with macOS write barriers. - These paths bypass ``SessionDB``/:func:`apply_wal_with_fallback`, so they - inherited ``synchronous=NORMAL`` and no ``checkpoint_fullfsync`` — on Darwin - an interrupted ``REINDEX``/``VACUUM``/``writable_schema`` rewrite leaves - half-written b-tree pages. Autocommit (``isolation_level=None``): DDL and - ``VACUUM`` are illegal inside an implicit transaction. Barriers are - best-effort: on a malformed schema even ``PRAGMA synchronous=FULL`` raises, - so whole-file rewrites call :func:`_reapply_durability_barriers` once the - schema parses again. - """ + These paths bypass ``SessionDB``/:func:`apply_wal_with_fallback`, so they inherited ``synchronous=NORMAL`` and + no ``checkpoint_fullfsync`` — on Darwin an interrupted ``REINDEX``/``VACUUM``/``writable_schema`` rewrite leaves + half-written b-tree pages. Autocommit (``isolation_level=None``): DDL and ``VACUUM`` are illegal inside an + implicit transaction. Barriers are best-effort: on a malformed schema even ``PRAGMA synchronous=FULL`` raises, + so whole-file rewrites call :func:`_reapply_durability_barriers` once the schema parses again.""" conn = sqlite3.connect(str(db_path), timeout=timeout, isolation_level=None) _reapply_durability_barriers(conn) return conn -@contextmanager def _repair_conn(db_path: Path, *, timeout: float = 5.0): - """A :func:`_connect_repair_durable` connection, closed on exit.""" - conn = _connect_repair_durable(db_path, timeout=timeout) - try: - yield conn - finally: - conn.close() + """A :func:`_connect_repair_durable` connection as a context manager, closed on exit.""" + return contextlib.closing(_connect_repair_durable(db_path, timeout=timeout)) def _reapply_durability_barriers(conn: sqlite3.Connection) -> bool: @@ -691,10 +609,9 @@ def _reapply_durability_barriers(conn: sqlite3.Connection) -> bool: def apply_durability_barriers(conn: sqlite3.Connection) -> bool: - """Durability barriers for guest users of ``state.db`` that must inherit its - owner's journal mode. Also applies the configured ``database.synchronous`` - level, a per-connection pragma that otherwise only rides on the journal-mode - setup path guests must not run.""" + """Durability barriers for guest users of ``state.db`` that must inherit its owner's journal mode. Also + applies the configured ``database.synchronous`` level, a per-connection pragma that otherwise only + rides on the journal-mode setup path guests must not run.""" from hermes_state import _apply_synchronous_pragma ok = _reapply_durability_barriers(conn) with contextlib.suppress(Exception): @@ -729,19 +646,15 @@ def _open_exclusive(db_path: Path, begin: str) -> sqlite3.Connection: @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. + """Yield ``(conn, None)`` — one live connection that excludes writers for repair surgery — or ``(None, + exc)`` when exclusion could not be taken. - ``locking_mode=EXCLUSIVE`` retains file-level exclusion after the short - ``BEGIN EXCLUSIVE`` is rolled back; the rollback is essential because - ``Connection.backup`` uses this connection as *source* and later as the - promotion *destination*, both transaction-free. It stays open across the - snapshot -> strategies -> promotion window so no writer can commit a change - promotion would overwrite. Existing readers make acquisition fail rather - than being disturbed (fail closed). Timeout 0: the cross-process lock already - serializes repairers, and a partial repair is less safe than "stop the - gateway and retry". - """ + ``locking_mode=EXCLUSIVE`` retains file-level exclusion after the short ``BEGIN EXCLUSIVE`` is rolled + back; the rollback is essential because ``Connection.backup`` uses this connection as *source* and later + as the promotion *destination*, both transaction-free. It stays open across the snapshot -> strategies -> + promotion window so no writer can commit a change promotion would overwrite. Existing readers make + acquisition fail rather than being disturbed (fail closed). Timeout 0: the cross-process lock already + serializes repairers, and a partial repair is less safe than "stop the gateway and retry".""" try: guard = _open_exclusive(db_path, "BEGIN EXCLUSIVE") except (sqlite3.Error, OSError) as exc: @@ -759,44 +672,39 @@ 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.""" + """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.""" # Deadline first: a sidecar vanishing mid-stat must not leak a just-opened descriptor. deadline_seconds = _repair_snapshot_timeout_seconds(source_path) deadline = time.monotonic() + deadline_seconds - source = source_connection or _connect_repair_durable(source_path) - destination = destination_connection def _check_deadline(_status: int, _remaining: int, _total: int) -> None: if time.monotonic() >= deadline: raise TimeoutError(f"timed out copying SQLite repair snapshot after {deadline_seconds:.0f}s") - try: - if destination is None: - destination = _connect_repair_durable(destination_path) - elif destination.in_transaction: - # sqlite3_backup needs a transaction-free destination (the guard holds - # exclusion via locking_mode, not a transaction). + with contextlib.ExitStack() as owned: # closes only connections opened here, destination first + source = source_connection or owned.enter_context(_repair_conn(source_path)) + if destination_connection is not None and destination_connection.in_transaction: + # sqlite3_backup needs a transaction-free destination (the guard holds exclusion via locking_mode). raise sqlite3.ProgrammingError("SQLite repair backup destination has an active transaction") + destination = destination_connection or owned.enter_context(_repair_conn(destination_path)) source.backup(destination, pages=256, progress=_check_deadline, sleep=_REPAIR_LOCK_POLL_SECONDS) - finally: - if destination_connection is None and destination is not None: - destination.close() - if source_connection is None: - source.close() + + +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 def _db_opens_cleanly(db_path: Path) -> Optional[str]: """Probe a DB on a fresh connection. Returns None if healthy, else a reason. - Runs the first statement that trips the malformed-schema parse (``PRAGMA - journal_mode``), ``integrity_check``, a ``sessions`` read, FTS5 MATCH probes - and a rolled-back ``messages`` write — so FTS5 index corruption (reads and - ``integrity_check`` pass, every ``INSERT INTO messages`` fails through the - FTS triggers) is reported as unhealthy. - """ + Runs the first statement that trips the malformed-schema parse (``PRAGMA journal_mode``), + ``integrity_check``, a ``sessions`` read, FTS5 MATCH probes and a rolled-back ``messages`` write — so FTS5 + index corruption (reads and ``integrity_check`` pass, every ``INSERT INTO messages`` fails through the FTS + triggers) is reported as unhealthy.""" from hermes_state import SessionDB, load_fts5_cjk_extension conn = _connect_repair_durable(db_path) try: @@ -817,38 +725,29 @@ def _db_opens_cleanly(db_path: Path) -> Optional[str]: 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" = FTS5 not built yet. - msg = str(exc).lower() - if isinstance(exc, sqlite3.OperationalError) and ( - SessionDB._is_fts5_unavailable_error(exc) or "no such table" in msg or "no such column" in msg - ): - continue - return f"fts5 read probe failed on {fts_table}: {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("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("ROLLBACK") - msg = str(exc).lower() # 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 "no such table" in msg or "no such column" in msg or "no such tokenizer: cjk_unicode61" in msg: + if _schema_not_built(exc) or "no such tokenizer: cjk_unicode61" in str(exc).lower(): return None return str(exc) return None @@ -861,23 +760,19 @@ def _db_opens_cleanly(db_path: Path) -> Optional[str]: def _live_writer_holds_db(db_path: Path) -> bool: """True when a connection outside this call still holds ``db_path`` open. - Asks SQLite for what a repair needs and a live holder cannot grant: - ``locking_mode=EXCLUSIVE`` then ``BEGIN IMMEDIATE`` — in WAL mode that needs - exclusive WAL-index locks, so any other open connection fails it with - SQLITE_BUSY; neither statement parses the schema, so it works on malformed - DBs. Fails **open** (False) on anything but a positive busy/locked signal: - refusing to repair a DB nobody holds would strand the self-heal path. In - ``journal_mode=DELETE`` a held reader takes only SHARED and this returns - False; repair is then serialised only by the cross-process repairer lock. - """ + Asks SQLite for what a repair needs and a live holder cannot grant: ``locking_mode=EXCLUSIVE`` then + ``BEGIN IMMEDIATE`` — in WAL mode that needs exclusive WAL-index locks, so any other open connection fails + it with SQLITE_BUSY; neither statement parses the schema, so it works on malformed DBs. Fails **open** + (False) on anything but a positive busy/locked signal: refusing to repair a DB nobody holds would strand + the self-heal path. In ``journal_mode=DELETE`` a held reader takes only SHARED and this returns False; + repair is then serialised only by the cross-process repairer lock.""" try: probe = _open_exclusive(db_path, "BEGIN IMMEDIATE") with contextlib.suppress(Exception): _close_unpinned(probe) return False except sqlite3.OperationalError as exc: - lowered = str(exc).lower() - return "locked" in lowered or "busy" in lowered + return "locked" in str(exc).lower() or "busy" in str(exc).lower() except Exception: # malformed/unreadable: no evidence of a live holder either way return False @@ -889,30 +784,25 @@ def _repair_skip(report: Dict[str, Any], verb: str, error: str, exc: Optional[Ba if exc is not None and _repair_failure_consumes_attempt(exc): report["_repair_attempted"] = True report["error"] = error - logger.error(f"state.db repair {verb}: %s", report["error"]) + logger.error(f"state.db repair {verb}: %s", error) return report def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, Any]: - """Repair a state.db whose ``sqlite_master`` is malformed or whose FTS - indexes reject writes. + """Repair a state.db whose ``sqlite_master`` is malformed or whose FTS indexes reject writes. - Two corruption classes: malformed schema / "duplicate object definition" - (even ``PRAGMA`` fails), and FTS write-corruption (reads and - ``integrity_check`` pass, writes fail through ``messages_fts*`` triggers). - ``_REPAIR_STRATEGIES`` run least-destructive first on a complete snapshot; - a success is copied back transactionally, so canonical rows are never - modified by a failed attempt. A raw backup is taken first 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}``. - """ + Two corruption classes: malformed schema / "duplicate object definition" (even ``PRAGMA`` fails), and FTS + write-corruption (reads and ``integrity_check`` pass, writes fail through ``messages_fts*`` triggers). + ``_REPAIR_STRATEGIES`` run least-destructive first on a complete snapshot; a success is copied back + transactionally, so canonical rows are never modified by a failed attempt. A raw backup is taken first + 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 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. + # 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) @@ -932,10 +822,8 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A if _db_opens_cleanly(db_path) is None: report["repaired"], report["strategy"] = True, "repaired_by_other_process" else: - report["error"] = ( - "could not obtain the state.db repair lock (held by another process, or the lock file was " - "unopenable); skipped schema surgery to avoid racing a concurrent repairer" - ) + report["error"] = ("could not obtain the state.db repair lock (held by another process, or the lock " + "file was unopenable); skipped schema surgery to avoid racing a concurrent repairer") return report result = report @@ -943,9 +831,8 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A # 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. + # 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", @@ -953,19 +840,16 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, 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. + # 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. before_mode = _probe_journal_mode_for_repair(db_path) result = _repair_state_db_schema_locked(db_path, backup=backup, report=report) if result.get("repaired"): result["journal_mode_before"] = before_mode _restore_journal_mode_after_repair(db_path, before_mode) - # 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. + # 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"): _record_repair_outcome(db_path, repaired=bool(result.get("repaired"))) @@ -973,10 +857,9 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A def _probe_journal_mode_for_repair(db_path: Path) -> Optional[str]: - """Best-effort journal-mode probe: ``wal``/``delete``, or ``None`` when the - file cannot be opened or probed (malformed header, concurrent opener's locks - — both expected on the repair path); callers then fall back to - ``database.journal_mode``.""" + """Best-effort journal-mode probe: ``wal``/``delete``, or ``None`` when the file cannot be opened or + probed (malformed header, concurrent opener's locks — both expected on the repair path); callers then + fall back to ``database.journal_mode``.""" from hermes_state import _on_disk_journal_mode try: with _repair_conn(db_path) as conn: @@ -988,54 +871,39 @@ def _probe_journal_mode_for_repair(db_path: Path) -> Optional[str]: def _restore_journal_mode_after_repair(db_path: Path, before_mode: Optional[str]) -> None: """Re-apply the journal mode after schema surgery. - A rebuilt file comes back in the default (delete) mode; without this a - corruption event silently moves a WAL store out of WAL (the open-time - WAL-reset gate never sees a flip made inside repair). Routed through - :func:`apply_wal_with_fallback`, not a direct pragma, so it inherits the - WAL-reset gate (a vulnerable runtime deliberately keeps DELETE; the - journal_mode-changed WARNING is expected there), the macOS-NFS silent-refusal - handling and the WAL companions. ``before_mode`` is only for the log - comparison; the target is ``database.journal_mode``. Best-effort: the repair - already succeeded, so failures log at WARNING. - """ + A rebuilt file comes back in the default (delete) mode; without this a corruption event silently moves a + WAL store out of WAL (the open-time WAL-reset gate never sees a flip made inside repair). Routed through + :func:`apply_wal_with_fallback`, not a direct pragma, so it inherits the WAL-reset gate (a vulnerable + runtime deliberately keeps DELETE; the journal_mode-changed WARNING is expected there), the macOS-NFS + silent-refusal handling and the WAL companions. ``before_mode`` is only for the log comparison; the target + is ``database.journal_mode``. Best-effort: the repair already succeeded, so failures log at WARNING.""" from hermes_state import apply_wal_with_fallback try: with _repair_conn(db_path) as conn: after = apply_wal_with_fallback(conn, db_label=db_path.name) if before_mode and after != before_mode: - logger.warning( - "state.db repair changed journal_mode %r -> %r (pre-surgery probe %r; restore resolved through " - "apply_wal_with_fallback per database.journal_mode and the WAL-reset gate)", - before_mode, after, before_mode, - ) + logger.warning("state.db repair changed journal_mode %r -> %r (pre-surgery probe %r; restore resolved " + "through apply_wal_with_fallback per database.journal_mode and the WAL-reset gate)", + before_mode, after, before_mode) except (sqlite3.Error, OSError) as exc: - logger.warning( - "state.db repair at %s: post-surgery journal-mode restore " - "failed (%s); verify with PRAGMA journal_mode on the next open", db_path, exc, - ) + logger.warning("state.db repair at %s: post-surgery journal-mode restore failed (%s); verify with PRAGMA " + "journal_mode on the next open", db_path, exc) def _repair_state_db_schema_locked(db_path: Path, *, backup: bool, report: Dict[str, Any]) -> Dict[str, Any]: - """Repair strategies for :func:`repair_state_db_schema`; caller holds the - cross-process repair lock. + """Repair strategies for :func:`repair_state_db_schema`; caller holds the cross-process repair lock. - Strategies run on a SCRATCH COPY, copied back through SQLite's transactional - backup API only once proven to open cleanly, so a failed repair cannot - modify or lose committed data (a WAL checkpoint of committed frames on guard - release is not a repair mutation). WHY not in place: the final strategy - ends in ``VACUUM``, which rebuilds the 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. - """ + Strategies run on a SCRATCH COPY, copied back through SQLite's transactional backup API only once proven to open + cleanly, so a failed repair cannot modify or lose committed data (a WAL checkpoint of committed frames on guard + release is not a repair mutation). WHY not in place: the final strategy ends in ``VACUUM``, which rebuilds the + 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 scratch = db_path.with_name(f"{db_path.name}.repair-scratch") cleanup_error = _unlink_db_triple(scratch) if cleanup_error is not None: - return _repair_skip( - report, "aborted", f"could not remove a stale repair snapshot before probing state.db: {cleanup_error}", - ) + 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). @@ -1046,23 +914,18 @@ def _repair_state_db_schema_locked(db_path: Path, *, backup: bool, report: Dict[ 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}", - ) + 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. 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 schema surgery to avoid overwriting a concurrent " - f"writer. Stop the gateway and retry: {guard_error}", exc=guard_error, - ) + 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: @@ -1074,9 +937,7 @@ def _repair_state_db_schema_locked(db_path: Path, *, backup: bool, report: Dict[ _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, - ) + 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 @@ -1088,11 +949,9 @@ def _repair_state_db_schema_locked(db_path: Path, *, backup: bool, report: Dict[ 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. - 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"], - ) + 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. @@ -1106,8 +965,7 @@ def _promote_repaired_snapshot(scratch: Path, db_path: Path, live_guard: sqlite3 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. - """ + 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) @@ -1128,8 +986,7 @@ def _unlink_db_triple(path: Path) -> Optional[str]: try: victim.unlink(missing_ok=True) except OSError as exc: - # Windows may retain a just-closed SQLite handle for a few scheduler - # ticks; bounded retry (a later open still fails safely if it stays live). + # Windows may retain a just-closed SQLite handle for a few scheduler ticks; bounded retry. if isinstance(exc, PermissionError) and _IS_WINDOWS and attempt < 9: time.sleep(0.05) continue @@ -1172,19 +1029,19 @@ def _strategy_reindex(conn: sqlite3.Connection) -> None: conn.commit() -def _strategy_dedup_schema(conn: sqlite3.Connection) -> None: - """De-duplicate sqlite_master (lowest rowid per type/name), 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) +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) - _edit_sqlite_master(conn, _dedup) + +def _strategy_dedup_schema(conn: sqlite3.Connection) -> None: + """De-duplicate sqlite_master, keeping FTS.""" + _edit_sqlite_master(conn, lambda: _dedup_sqlite_master(conn)) def _strategy_drop_fts_vacuum(conn: sqlite3.Connection) -> None: @@ -1216,9 +1073,8 @@ _REPAIR_STRATEGIES = ( def _run_repair_strategies(db_path: Path, report: Dict[str, Any]) -> Dict[str, Any]: - """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.""" + """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: @@ -1230,7 +1086,6 @@ def _run_repair_strategies(db_path: Path, report: Dict[str, Any]) -> Dict[str, A reason = str(exc) if failure_msg is not None: logger.warning(failure_msg, exc) - continue if reason is None: report["repaired"], report["strategy"] = True, name logger.warning(success_msg, db_path) diff --git a/hermes_state_wal.py b/hermes_state_wal.py index cfda3f2592..c94aea5e34 100644 --- a/hermes_state_wal.py +++ b/hermes_state_wal.py @@ -1,8 +1,7 @@ """SQLite journal-mode and PRAGMA policy for state.db (split from hermes_state). -Every name is re-imported into ``hermes_state``; intra-module calls to -patchable helpers go through a lazy ``from hermes_state import ...`` at call -time so monkeypatches there still intercept. +Every name is re-imported into ``hermes_state``; intra-module calls to patchable helpers go through a lazy +``from hermes_state import ...`` at call time so monkeypatches there still intercept. """ from __future__ import annotations @@ -22,22 +21,16 @@ from hermes_cli.sqlite_runtime import is_sqlite_wal_reset_vulnerable as _is_sqli logger = logging.getLogger("hermes_state") -# WAL needs mmap shared memory + fcntl byte-range locks. Network filesystems (NFS, -# SMB/CIFS, some FUSE, WSL1) raise ``locking protocol``; ZFS corrupts the -shm file -# under concurrent bursts (COW + mmap) -> ``disk I/O error``. Either would silently -# break everything on state.db/kanban.db, so fall back to DELETE (readers block on writes). -_WAL_INCOMPAT_MARKERS = ( - "locking protocol", # SQLITE_PROTOCOL on NFS/SMB - "not authorized", # Some FUSE mounts block WAL pragma outright - "disk i/o error", # ZFS SHM corruption under concurrent connections -) - +# WAL needs mmap shared memory + fcntl byte-range locks. Network filesystems (NFS, SMB/CIFS, some FUSE, WSL1) raise +# ``locking protocol``; ZFS corrupts the -shm file under concurrent bursts (COW + mmap) -> ``disk I/O error``. +# Either would silently break everything on state.db/kanban.db, so fall back to DELETE (readers block on writes). +# "not authorized": some FUSE mounts block the WAL pragma outright. +_WAL_INCOMPAT_MARKERS = ("locking protocol", "not authorized", "disk i/o error") # SQLite's default journal_size_limit is -1 (unlimited); see _apply_wal_size_limit. _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.``; ``_warn_once`` resolves them there at call time. +# 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.``; ``_warn_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() @@ -47,10 +40,9 @@ _delete_overridden_warned_lock = threading.Lock() _journal_upgrade_warned_paths: set = set() _journal_upgrade_warned_lock = threading.Lock() -_CANNOT_VERIFY_DELETE_MSG = ( - "could not verify journal mode before applying configured journal_mode=delete (database is locked — possible " - "concurrent openers); refusing to downgrade a database this process does not exclusively own" -) +_CANNOT_VERIFY_DELETE_MSG = ("could not verify journal mode before applying configured journal_mode=delete (database " + "is locked — possible concurrent openers); refusing to downgrade a database this process " + "does not exclusively own") def _warn_once(lock: threading.Lock, set_name: str, key: str) -> bool: @@ -70,39 +62,33 @@ def _mode_from_row(row) -> str: def _on_disk_journal_mode(conn: sqlite3.Connection) -> Optional[str]: - """Read the journal mode from the DB header; ``None`` if undeterminable - (new DB, or PRAGMA failed) -> callers take their fail-closed "refuse to - downgrade" branch. ``disk i/o error`` can be transient on virtualized block - devices (XFS on cloud hosts), so it is retried a few times first.""" - last_exc: Optional[Exception] = None + """Read the journal mode from the DB header; ``None`` if undeterminable (new DB, or PRAGMA failed) -> + callers take their fail-closed "refuse to downgrade" branch. ``disk i/o error`` can be transient on + virtualized block devices (XFS on cloud hosts), so it is retried a few times first.""" for _ in range(4): try: row = conn.execute("PRAGMA journal_mode").fetchone() except sqlite3.OperationalError as exc: - last_exc = exc if "disk i/o error" not in str(exc).lower(): return None + last_exc = exc time.sleep(0.05) continue - if row is None: - return None - mode = row[0] + mode = row[0] if row else None if isinstance(mode, bytes): # defensive: sqlite3 occasionally returns bytes try: mode = mode.decode("ascii") except UnicodeDecodeError: return None return str(mode).strip().lower() if mode is not None else None - if last_exc is not None: - logger.debug("_on_disk_journal_mode: retries exhausted on disk read (%s)", last_exc) + logger.debug("_on_disk_journal_mode: retries exhausted on disk read (%s)", last_exc) return None def _apply_wal_size_limit(conn: sqlite3.Connection) -> None: - """Bound the WAL so it returns space after big transactions. With the default - (-1) a checkpointed WAL is reused in place, never truncated, so ``state.db-wal`` - keeps the high-water mark of the largest transaction ever (a 3 GB optimize left - a 3 GB WAL). Best-effort: failure only costs disk slack.""" + """Bound the WAL so it returns space after big transactions. With the default (-1) a checkpointed WAL is + reused in place, never truncated, so ``state.db-wal`` keeps the high-water mark of the largest + transaction ever (a 3 GB optimize left a 3 GB WAL). Best-effort: failure only costs disk slack.""" try: conn.execute(f"PRAGMA journal_size_limit={_WAL_SIZE_LIMIT_BYTES}") except sqlite3.OperationalError as exc: # pragma: no cover - defensive @@ -118,10 +104,9 @@ def _darwin_pragma(conn: sqlite3.Connection, pragma: str) -> None: def _apply_macos_checkpoint_barrier(conn: sqlite3.Connection) -> None: - """Enable ``PRAGMA checkpoint_fullfsync`` on macOS. Apple's ``fsync(2)`` - guarantees neither data-on-platter nor ordering, so without ``F_FULLFSYNC`` a - launchd shutdown can turn a "durable" checkpoint into a malformed ``state.db``. - Checkpoint boundaries only (~+0.1 ms/commit vs ~+4 ms for ``fullfsync=1``).""" + """Enable ``PRAGMA checkpoint_fullfsync`` on macOS. Apple's ``fsync(2)`` guarantees neither data-on-platter nor + ordering, so without ``F_FULLFSYNC`` a launchd shutdown can turn a "durable" checkpoint into a malformed + ``state.db``. Checkpoint boundaries only (~+0.1 ms/commit vs ~+4 ms for ``fullfsync=1``).""" _darwin_pragma(conn, "PRAGMA checkpoint_fullfsync=1") @@ -140,10 +125,8 @@ 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 - """ + """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) @@ -151,20 +134,16 @@ def is_sqlite_wal_reset_vulnerable(version_info: Optional[tuple] = None) -> bool def sqlite_source_id() -> str: """Return ``sqlite_source_id()``, or an empty string when unavailable.""" try: - conn = sqlite3.connect(":memory:") - try: + with contextlib.closing(sqlite3.connect(":memory:")) as conn: row = conn.execute("SELECT sqlite_source_id()").fetchone() - finally: - conn.close() except sqlite3.Error: return "" return str(row[0]) if row and row[0] is not None else "" def _database_has_content(conn: sqlite3.Connection) -> bool: - """Whether the file already holds pages (existing vs brand-new DB); lock-free - header read. Fail-quiet False: the only caller gates a warning on this and an - unknown-answer warning would fire on every fresh database.""" + """Whether the file already holds pages (existing vs brand-new DB); lock-free header read. Fail-quiet False: the + only caller gates a warning on this and an unknown-answer warning would fire on every fresh database.""" try: row = conn.execute("PRAGMA page_count").fetchone() return bool(row) and row[0] is not None and int(row[0]) > 0 @@ -173,9 +152,8 @@ def _database_has_content(conn: sqlite3.Connection) -> bool: def resolve_journal_mode() -> str: - """The configured ``database.journal_mode`` (``wal`` default; ``delete`` for - filesystems without WAL-safe durability: macOS virtiofs, NFS, SMB). Invalid - values fail safe to ``wal``.""" + """The configured ``database.journal_mode`` (``wal`` default; ``delete`` for filesystems without WAL-safe + durability: macOS virtiofs, NFS, SMB). Invalid values fail safe to ``wal``.""" try: from hermes_cli.config import load_config_readonly @@ -196,29 +174,23 @@ class WalUnsupportedError(sqlite3.OperationalError): def _verify_configured_delete(actual: str) -> str: """Raise unless SQLite reported ``delete`` for an explicit operator request.""" if actual != "delete": - raise sqlite3.OperationalError( - f"could not set configured journal_mode=delete (got {actual or 'no result'})" - ) + raise sqlite3.OperationalError(f"could not set configured journal_mode=delete (got {actual or 'no result'})") return actual def apply_wal_with_fallback(conn: sqlite3.Connection, *, db_label: str = "state.db", require_wal: bool = False) -> str: """Set ``journal_mode=WAL`` on ``conn``, falling back to DELETE on failure. - Returns the mode actually set. Shared by :class:`SessionDB` and - ``hermes_cli.kanban_db.connect``. WAL-incompatible filesystems either raise - ``OperationalError`` ("locking protocol" / "disk I/O error") or — macOS NFS / - SMB / AgentFS — silently refuse and stay in DELETE; either way log ERROR once - per process per ``db_label`` and fall back. ``require_wal=True`` raises - :class:`WalUnsupportedError` instead. WAL-reset-bug builds - (https://sqlite.org/wal.html#walresetbug) never enable WAL on non-WAL files; - an already-WAL DB keeps WAL with a warning. Gate deliberately RETAINED: - re-measured on the bundled 3.50.4 there is no evidence WAL is safer. + Returns the mode actually set. Shared by :class:`SessionDB` and ``hermes_cli.kanban_db.connect``. + WAL-incompatible filesystems either raise ``OperationalError`` ("locking protocol" / "disk I/O error") or — + macOS NFS / SMB / AgentFS — silently refuse and stay in DELETE; either way log ERROR once per process per + ``db_label`` and fall back. ``require_wal=True`` raises :class:`WalUnsupportedError` instead. WAL-reset-bug + builds (https://sqlite.org/wal.html#walresetbug) never enable WAL on non-WAL files; an already-WAL DB keeps WAL + with a warning. Gate deliberately RETAINED: re-measured on the bundled 3.50.4 there is no evidence WAL is safer. - Invariant on every path: never downgrade to DELETE if the on-disk header - reports WAL or cannot be read — other gateway/cron/worker connections may - hold the DB open, and a live downgrade destroys their uncheckpointed commits. - """ + Invariant on every path: never downgrade to DELETE if the on-disk header reports WAL or cannot be read — other + gateway/cron/worker connections may hold the DB open, and a live downgrade destroys their uncheckpointed + commits.""" from hermes_state import is_sqlite_wal_reset_vulnerable, resolve_journal_mode configured = resolve_journal_mode() @@ -260,9 +232,8 @@ def _enable_wal(conn: sqlite3.Connection, db_label: str, require_wal: bool, curr return "wal" try: - # ``PRAGMA journal_mode=WAL`` RETURNS the resulting mode: macOS NFS, SMB/CIFS - # and the AgentFS overlay refuse WITHOUT raising. Trust the row, not the - # absence of an exception. + # ``PRAGMA journal_mode=WAL`` RETURNS the resulting mode: macOS NFS, SMB/CIFS and the AgentFS overlay + # refuse WITHOUT raising. Trust the row, not the absence of an exception. mode = _mode_from_row(conn.execute("PRAGMA journal_mode=WAL").fetchone()) if mode == "wal": return _wal_activated() @@ -271,10 +242,9 @@ def _enable_wal(conn: sqlite3.Connection, db_label: str, require_wal: bool, curr raise silent_exc _log_wal_fallback_once(db_label, silent_exc) return mode or "delete" + except WalUnsupportedError: + raise # the require_wal silent-refusal raise above — propagate unchanged except sqlite3.OperationalError as exc: - # The require_wal silent-refusal raise above lands here — propagate unchanged. - if isinstance(exc, WalUnsupportedError): - raise msg = str(exc).lower() if not any(marker in msg for marker in _WAL_INCOMPAT_MARKERS): raise # unrelated OperationalError — don't silently swallow @@ -284,8 +254,7 @@ def _enable_wal(conn: sqlite3.Connection, db_label: str, require_wal: bool, curr 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. - existing = _on_disk_journal_mode(conn) - if existing == "wal" or existing is None: + if _on_disk_journal_mode(conn) in ("wal", None): raise if require_wal: raise WalUnsupportedError(str(exc)) from exc @@ -295,11 +264,10 @@ def _enable_wal(conn: sqlite3.Connection, db_label: str, require_wal: bool, curr 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.""" + """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: @@ -316,15 +284,11 @@ def _retry_wal_after_eio(conn: sqlite3.Connection, exc: sqlite3.OperationalError def _set_journal_mode_no_wait(conn: sqlite3.Connection, mode: str) -> str: """Execute ``PRAGMA journal_mode=`` without waiting on other openers. - The ONLY place a non-WAL journal-mode switch may be issued. ``busy_timeout=0`` - turns SQLite's exclusivity requirement into a concurrent-opener detector: - leaving WAL needs exclusive access, so if ANY other connection holds the DB - the pragma fails immediately with ``database is locked`` instead of sneaking - the flip between a writer's transactions (how uncheckpointed WAL commits die). - Callers must treat a raised ``OperationalError`` as "not exclusively owned: - leave the mode alone", never as retryable. Returns the reported mode, ``""`` - if no row. - """ + The ONLY place a non-WAL journal-mode switch may be issued. ``busy_timeout=0`` turns SQLite's exclusivity + requirement into a concurrent-opener detector: leaving WAL needs exclusive access, so if ANY other connection + holds the DB the pragma fails immediately with ``database is locked`` instead of sneaking the flip between a + writer's transactions (how uncheckpointed WAL commits die). Callers must treat a raised ``OperationalError`` as + "not exclusively owned: leave the mode alone", never as retryable. Returns the reported mode, ``""`` if no row.""" try: row = conn.execute("PRAGMA busy_timeout").fetchone() previous_timeout = int(row[0]) if row and row[0] is not None else 0 @@ -341,13 +305,10 @@ def _set_journal_mode_no_wait(conn: sqlite3.Connection, mode: str) -> str: def _apply_delete_for_wal_reset_bug(conn: sqlite3.Connection, *, db_label: str, require_delete: bool = False) -> str: """Avoid enabling WAL when the linked SQLite has the WAL-reset bug. - Already-WAL on disk: keep WAL (no live downgrade) and warn. Mode unreadable - (probe blocked by a concurrent opener): not provably exclusive — leave it and - warn; treating "could not read" as "not WAL" once flipped a live WAL state.db - to DELETE under a writer, destroying its uncheckpointed commits. Otherwise set - DELETE without waiting out openers and warn; an explicit operator request - additionally verifies SQLite accepted DELETE. - """ + Already-WAL on disk: keep WAL (no live downgrade) and warn. Mode unreadable (probe blocked by a concurrent + opener): not provably exclusive — leave it and warn; treating "could not read" as "not WAL" once flipped a live + WAL state.db to DELETE under a writer, destroying its uncheckpointed commits. Otherwise set DELETE without + waiting out openers and warn; an explicit operator request additionally verifies SQLite accepted DELETE.""" current = _on_disk_journal_mode(conn) if current == "wal": _log_wal_reset_bug_once(db_label, kept_wal=True) @@ -367,8 +328,7 @@ def _apply_delete_for_wal_reset_bug(conn: sqlite3.Connection, *, db_label: str, except sqlite3.OperationalError as exc: if require_delete: raise - lowered = str(exc).lower() - if "locked" in lowered or "busy" in lowered: + if "locked" in str(exc).lower() or "busy" in str(exc).lower(): # A concurrent opener appeared between probe and flip: leave the mode as is. _log_wal_reset_bug_once(db_label, kept_wal=True, indeterminate=True) return current or "delete" @@ -388,9 +348,7 @@ def _wal_reset_repair_hint() -> str: cmd = recommended_update_command_for_method(method) if method in {"git", "unknown"}: return f"Hermes-managed installs can repair the embedded runtime with `{cmd}`" - if method == "docker": - return f"update the container image with `{cmd}`" - return cmd # nix/nixos + return f"update the container image with `{cmd}`" if method == "docker" else cmd # else nix/nixos except Exception: return "install a Python build bundled with SQLite 3.51.3+ (or backports 3.50.7 / 3.44.6) and restart Hermes" @@ -399,13 +357,10 @@ def _wal_reset_repair_hint() -> str: # 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). _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 openers)" - ), - "kept_wal": ( - "is already in WAL mode — leaving WAL in place (no live downgrade under concurrent openers)" - ), + "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 " + "openers)"), + "kept_wal": "is already in WAL mode — leaving WAL in place (no live downgrade under concurrent openers)", "delete": "using journal_mode=DELETE instead of enabling WAL", } _ONCE_LOGS = { @@ -419,10 +374,9 @@ _ONCE_LOGS = { ), "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. + # 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: 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 " @@ -500,26 +454,19 @@ def resolve_synchronous_level(raw_value: Any) -> Optional[int]: def _apply_synchronous_pragma(conn: sqlite3.Connection, raw_value: Any, *, db_label: str) -> None: """Set ``PRAGMA synchronous`` from config, never below FULL on macOS. - Kept out of the integer loop in :func:`apply_database_pragmas`: this PRAGMA - decides whether a commit is on 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. - """ + Kept out of the integer loop in :func:`apply_database_pragmas`: this PRAGMA decides whether a commit is on + 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: - logger.warning( - "%s: ignoring unrecognized database.synchronous=%r (expected OFF, NORMAL, FULL, EXTRA, or 0-3)", - db_label, raw_value, - ) + logger.warning("%s: ignoring unrecognized database.synchronous=%r (expected OFF, NORMAL, FULL, EXTRA, or 0-3)", + db_label, raw_value) return if sys.platform == "darwin" and level < _SYNCHRONOUS_FULL: - logger.warning( - "%s: refusing database.synchronous=%s on macOS; keeping FULL. Darwin's fsync() does not guarantee write " - "ordering, so a lower level readmits the half-written btree pages FULL exists to prevent.", - db_label, _SYNCHRONOUS_NAMES[level], - ) + logger.warning("%s: refusing database.synchronous=%s on macOS; keeping FULL. Darwin's fsync() does not " + "guarantee write ordering, so a lower level readmits the half-written btree pages FULL exists " + "to prevent.", db_label, _SYNCHRONOUS_NAMES[level]) return with contextlib.suppress(sqlite3.OperationalError): conn.execute(f"PRAGMA synchronous={level}") @@ -528,14 +475,11 @@ def _apply_synchronous_pragma(conn: sqlite3.Connection, raw_value: Any, *, db_la def apply_database_pragmas(conn: sqlite3.Connection, *, db_label: str = "state.db") -> None: """Apply optional performance and WAL-sizing PRAGMAs from ``config.yaml``. - Journal mode is NOT handled here (owned by :func:`apply_wal_with_fallback`). - ``database:`` keys: ``cache_size`` (negative = KiB, positive = pages), - ``mmap_size`` (bytes, 0 = off), ``temp_store`` (0-3), ``wal_autocheckpoint`` - (pages), ``journal_size_limit`` (bytes), ``synchronous`` (unset leaves the - compile-time default, which differs 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. - """ + Journal mode is NOT handled here (owned by :func:`apply_wal_with_fallback`). ``database:`` keys: ``cache_size`` + (negative = KiB, positive = pages), ``mmap_size`` (bytes, 0 = off), ``temp_store`` (0-3), ``wal_autocheckpoint`` + (pages), ``journal_size_limit`` (bytes), ``synchronous`` (unset leaves the compile-time default, which differs + 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