From 401eb80b2955f656f1ad48578068036e9fafb39e Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 22:24:18 -0700 Subject: [PATCH] =?UTF-8?q?refactor(state):=20repair/wal=20=E2=80=94=20inl?= =?UTF-8?q?ine=20single-use=20helpers=20(=5Fbundle=5Fbytes,=20=5Funlink=5F?= =?UTF-8?q?quiet,=20=5Fpromote=5Frepaired=5Fsnapshot,=20=5Fdedup=5Fsqlite?= =?UTF-8?q?=5Fmaster,=20=5Fwarn=5Fonce,=20=5Fretry=5Fwal=5Fafter=5Feio),?= =?UTF-8?q?=20contextlib.closing/suppress,=20pack=20calls?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- hermes_state_repair.py | 395 +++++++++++++++++------------------------ hermes_state_wal.py | 148 ++++++--------- 2 files changed, 214 insertions(+), 329 deletions(-) diff --git a/hermes_state_repair.py b/hermes_state_repair.py index a4cf6c4f13..29357f309c 100644 --- a/hermes_state_repair.py +++ b/hermes_state_repair.py @@ -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: diff --git a/hermes_state_wal.py b/hermes_state_wal.py index c94aea5e34..1092ceffff 100644 --- a/hermes_state_wal.py +++ b/hermes_state_wal.py @@ -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.``; ``_warn_once`` resolves them there. +# line would repeat per connection). Tests clear these via ``hermes_state.``; ``_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.``.""" - 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=`` 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)