diff --git a/hermes_state_dbfile.py b/hermes_state_dbfile.py index ba0af9d837..a6ee6c9c5a 100644 --- a/hermes_state_dbfile.py +++ b/hermes_state_dbfile.py @@ -31,31 +31,27 @@ from hermes_state_common import ( # Log-record parity with the origin module (caplog tests pin "hermes_state"). logger = logging.getLogger("hermes_state") - # _read_sqlite_application_id runs on EVERY write via _raise_if_db_replaced, # against the LIVE state.db. A bare open()/read()/close() there is the # howtocorrupt §2.2 bug: close() cancels every POSIX advisory lock this -# process holds on the file — measured on Linux/SQLite 3.53.1, one probe call -# drops the WAL-mode DMS shared lock the writer connection holds on state.db -# (see hermes_cli/sqlite_safe_read.py for the module built around this rule). -# With the DMS lock gone, a fresh opener in another process can treat this -# writer as dead and rerun WAL-index recovery underneath it. +# process holds on the file — one probe call drops the WAL-mode DMS shared +# lock the writer connection holds (see hermes_cli/sqlite_safe_read.py). With +# the DMS lock gone, a fresh opener in another process can treat this writer +# as dead and rerun WAL-index recovery underneath it. # # The probe therefore reads through a per-path fd cached for the life of the # process: opening an fd never cancels locks (only close() does), and # os.pread takes no shared file position. When the path is re-pointed at a # new inode (the very replacement this probe exists to detect), the stale fd # is RETIRED, never closed — closing it would cancel the live connection's -# locks on the old file, the exact bug being avoided. Replacement events are -# rare and halt writes anyway, so the leak is bounded. +# locks on the old file. Replacement events are rare and halt writes anyway, +# so the leak is bounded. _HEADER_PROBE_LOCK = threading.Lock() - - _HEADER_PROBE_FDS: "dict[str, tuple[int, int, int]]" = {} # key -> (fd, dev, ino) - - _RETIRED_HEADER_PROBE_FDS: "list[int]" = [] # intentionally never closed +_FTS_TABLE_NAMES = ("messages_fts", "messages_fts_trigram", "messages_fts_cjk") + def _pread_db_header(db_path: Path, length: int) -> "Optional[bytes]": """Lock-safe raw header read of a possibly-live SQLite database. @@ -93,8 +89,7 @@ def _pread_db_header(db_path: Path, length: int) -> "Optional[bytes]": except OSError: _RETIRED_HEADER_PROBE_FDS.append(fd) return None - cached = (fd, fst.st_dev, fst.st_ino) - _HEADER_PROBE_FDS[key] = cached + cached = _HEADER_PROBE_FDS[key] = (fd, fst.st_dev, fst.st_ino) try: return os.pread(cached[0], length, 0) except OSError: @@ -104,24 +99,15 @@ def _pread_db_header(db_path: Path, length: int) -> "Optional[bytes]": def _read_sqlite_application_id(db_path: Path) -> "Optional[int]": """Read application_id from the SQLite header without opening a connection. - Safe against live databases: routed through :func:`_pread_db_header`, - which never issues a ``close()`` that would cancel this process's POSIX - locks on the file (howtocorrupt §2.2). + Routed through :func:`_pread_db_header`, which never issues a ``close()`` + that would cancel this process's POSIX locks on the file. """ from hermes_state import _STATE_DB_APPLICATION_ID_OFFSET - header = _pread_db_header(db_path, _STATE_DB_APPLICATION_ID_OFFSET + 4) - if header is None: + end = _STATE_DB_APPLICATION_ID_OFFSET + 4 + header = _pread_db_header(db_path, end) + if header is None or len(header) < end or header[:16] != b"SQLite format 3\x00": return None - if len(header) < _STATE_DB_APPLICATION_ID_OFFSET + 4: - return None - if header[:16] != b"SQLite format 3\x00": - return None - return int( - struct.unpack( - ">I", - header[_STATE_DB_APPLICATION_ID_OFFSET:_STATE_DB_APPLICATION_ID_OFFSET + 4], - )[0] - ) + return int(struct.unpack(">I", header[_STATE_DB_APPLICATION_ID_OFFSET:end])[0]) def _stat_sqlite_sidecar_identity(db_path: Path) -> Dict[str, tuple]: @@ -142,10 +128,24 @@ def _canonical_sqlite_path(path: str) -> str: def _watched_sqlite_sidecar_paths(db_path) -> Set[str]: base = os.path.abspath(os.fspath(db_path)) - return { - _canonical_sqlite_path(base + "-wal"), - _canonical_sqlite_path(base + "-shm"), - } + return {_canonical_sqlite_path(base + "-wal"), _canonical_sqlite_path(base + "-shm")} + + +def _iter_proc_fd_targets(): + """Yield ``(pid, readlink target)`` for every readable ``/proc//fd`` entry.""" + for pid_str in os.listdir("/proc"): + if not pid_str.isdigit(): + continue + fd_dir = f"/proc/{pid_str}/fd" + try: + fds = os.listdir(fd_dir) + except OSError: + continue # process gone or not ours + for fd in fds: + try: + yield int(pid_str), os.readlink(f"{fd_dir}/{fd}") + except OSError: + continue def iter_deleted_sqlite_sidecar_holders(db_path) -> List[Tuple[int, str]]: @@ -163,31 +163,14 @@ def iter_deleted_sqlite_sidecar_holders(db_path) -> List[Tuple[int, str]]: """ if not sys.platform.startswith("linux"): return [] - holders: List[Tuple[int, str]] = [] watched = _watched_sqlite_sidecar_paths(db_path) try: - for pid_str in os.listdir("/proc"): - if not pid_str.isdigit(): - continue - pid = int(pid_str) - fd_dir = f"/proc/{pid}/fd" - try: - fds = os.listdir(fd_dir) - except OSError: - continue - for fd in fds: - try: - target = os.readlink(f"{fd_dir}/{fd}") - except OSError: - continue - if " (deleted)" not in target: - continue - if _canonical_sqlite_path(target) in watched: - holders.append((pid, target)) + for pid, target in _iter_proc_fd_targets(): + if " (deleted)" in target and _canonical_sqlite_path(target) in watched: + holders.append((pid, target)) except Exception as exc: logger.debug("deleted-WAL holder scan failed for %s: %s", db_path, exc) - return holders return holders @@ -198,8 +181,7 @@ def refuse_deleted_wal_generation(db_path) -> None: replacement WAL inode while a live writer still holds the orphan. """ from hermes_state import DeletedWalGenerationError, _DELETED_WAL_GENERATION_MSG - holders = iter_deleted_sqlite_sidecar_holders(db_path) - if not holders: + if not iter_deleted_sqlite_sidecar_holders(db_path): return logger.error(_DELETED_WAL_GENERATION_MSG) raise DeletedWalGenerationError(_DELETED_WAL_GENERATION_MSG) @@ -216,8 +198,7 @@ def _connect_tracked_db(path, tracking_path=None, **kwargs): The ONLY tolerated fallback is the helper being absent entirely (scaffold/embed installs that ship hermes_state without hermes_cli). A real connection failure must propagate: silently retrying an *untracked* - connect would disable the guard for the lifetime of that connection, - which is precisely the failure mode this module exists to prevent. + connect would disable the guard for the lifetime of that connection. """ try: from hermes_cli.sqlite_safe_read import connect_tracked @@ -228,22 +209,14 @@ def _connect_tracked_db(path, tracking_path=None, **kwargs): path, ) return sqlite3.connect(str(path), **kwargs) - # Open through THIS module's sqlite3.connect so callers (and tests) that # patch hermes_state.sqlite3.connect keep control of connection creation; # the helper still owns tracking. - return connect_tracked( - path, - tracking_path=tracking_path, - connect_fn=sqlite3.connect, - **kwargs, - ) + return connect_tracked(path, tracking_path=tracking_path, connect_fn=sqlite3.connect, **kwargs) -def is_zeroed_state_db( - path: Path, *, probe_bytes: int = 100, force: bool = False -) -> bool: - """Detect the #68474/#97568 zeroed state.db signature (0-byte or NUL header). +def is_zeroed_state_db(path: Path, *, probe_bytes: int = 100, force: bool = False) -> bool: + """Detect the zeroed state.db signature (0-byte or NUL header). Byte-level probe, so it is only safe BEFORE any connection to *path* exists in this process: ``close()`` cancels every POSIX advisory lock the @@ -277,10 +250,7 @@ def is_zeroed_state_db( if not force and has_live_connection(path): return False - - head = read_header_bytes_preopen( - path, length=max(16, probe_bytes), force=force - ) + head = read_header_bytes_preopen(path, length=max(16, probe_bytes), force=force) if head is None: return False if len(head) == 0: @@ -300,67 +270,56 @@ def quarantine_cross_process_lock(path: Path, timeout: float = 5.0): handle = lock_path.open("a+b") acquired = False try: - deadline = time.monotonic() + timeout if platform.system() == "Windows": import msvcrt - while True: - try: - handle.seek(0) - msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) - acquired = True - break - except OSError: - if time.monotonic() >= deadline: - break - time.sleep(0.020) + def _try_lock(): + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) + + def _unlock(): + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1) else: import fcntl - while True: - try: - fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) - acquired = True + def _try_lock(): + fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + + def _unlock(): + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + deadline = time.monotonic() + timeout + while True: + try: + _try_lock() + acquired = True + break + except OSError: + if time.monotonic() >= deadline: break - except (BlockingIOError, OSError): - if time.monotonic() >= deadline: - break - time.sleep(0.020) + time.sleep(0.020) yield acquired finally: try: if acquired: - if platform.system() == "Windows": - import msvcrt - - handle.seek(0) - msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1) - else: - import fcntl - - fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + _unlock() except (OSError, AttributeError): pass finally: handle.close() -def quarantine_zeroed_state_db( - path: Path, *, already_locked: bool = False -) -> Optional[Path]: +def quarantine_zeroed_state_db(path: Path, *, already_locked: bool = False) -> Optional[Path]: """Move a zeroed state.db aside (preserve bytes) and return quarantine path. - Uses a cross-process lock (``#68805``) so two concurrent startups cannot - race: the first process moves the zeroed file and the second re-checks - under the lock, finding the file already gone (or a fresh DB in its place) - instead of clobbering the quarantine. + Uses a cross-process lock so two concurrent startups cannot race: the first + process moves the zeroed file and the second re-checks under the lock, + finding the file already gone (or a fresh DB in its place) instead of + clobbering the quarantine. """ def _do_quarantine(): if not path.exists(): - logger.info( - "quarantine_zeroed_state_db: %s already moved by another process", - path, - ) + logger.info("quarantine_zeroed_state_db: %s already moved by another process", path) return None if not is_zeroed_state_db(path): logger.info( @@ -369,20 +328,16 @@ def quarantine_zeroed_state_db( path, ) return None - try: ts = time.strftime("%Y%m%d-%H%M%S") except Exception: ts = "unknown" - dest = path.with_name( - f"{path.name}.zeroed-{ts}-{os.getpid()}.bak" - ) + stem = f"{path.name}.zeroed-{ts}-{os.getpid()}" + dest = path.with_name(f"{stem}.bak") n = 0 while dest.exists(): n += 1 - dest = path.with_name( - f"{path.name}.zeroed-{ts}-{os.getpid()}-{n}.bak" - ) + dest = path.with_name(f"{stem}-{n}.bak") try: path.rename(dest) except OSError as exc: @@ -399,7 +354,6 @@ def quarantine_zeroed_state_db( if already_locked: return _do_quarantine() - with quarantine_cross_process_lock(path) as acquired: if not acquired: logger.error( @@ -421,97 +375,69 @@ def collect_state_db_stats(db_path: Path) -> Dict[str, Any]: run against a *live* database held by a gateway without ever taking a write lock or mutating the file. Every field is collected independently: a failed pragma/SELECT yields ``None`` for that field, and the helper - itself never raises. + itself never raises. Deliberately does NOT instantiate :class:`SessionDB` + — its constructor runs schema DDL, which a diagnostics probe must never do. - Deliberately does NOT instantiate :class:`SessionDB` — its constructor - runs schema DDL (migrations, FTS table creation), which is exactly the - kind of write a diagnostics probe must never perform. - - Returned keys (all present, any may be None on failure): - - - ``page_count``, ``page_size``, ``freelist_count`` — PRAGMA values - - ``logical_size_bytes`` — page_count * page_size (post-checkpoint size) - - ``wal_size_bytes`` — stat() of ``-wal`` (0 when absent) - - ``journal_mode`` — PRAGMA journal_mode string - - ``messages`` / ``sessions`` — row counts - - ``fts_tables`` — dict of {table_name: bool} presence for - messages_fts / messages_fts_trigram / messages_fts_cjk - - ``fts_storage_version`` — int from state_meta, None when the marker is - absent (legacy pre-v23 inline layout) - - ``fts_rebuild_pending`` — True when the deferred v23 backfill has not - finished (high_water present and progress < high_water) - - ``fts_rebuild_high_water`` / ``fts_rebuild_progress`` — raw ints - - ``fts_rebuild_deferral`` — durable blocked-repair diagnostic, when present + Returned keys (all present, any may be None on failure): ``page_count``, + ``page_size``, ``freelist_count``, ``logical_size_bytes`` (page_count * + page_size), ``wal_size_bytes`` (stat of ``-wal``, 0 when absent), + ``journal_mode``, ``messages`` / ``sessions`` row counts, ``fts_tables`` + ({name: present}), ``fts_storage_version`` (None = legacy inline layout), + ``fts_rebuild_pending`` (deferred backfill unfinished), + ``fts_rebuild_high_water`` / ``fts_rebuild_progress`` raw ints, and + ``fts_rebuild_deferral`` (durable blocked-repair diagnostic). """ from hermes_state import _connect_tracked_db - stats: Dict[str, Any] = { - "page_count": None, - "page_size": None, - "freelist_count": None, - "logical_size_bytes": None, - "wal_size_bytes": None, - "journal_mode": None, - "messages": None, - "sessions": None, - "fts_tables": None, - "fts_storage_version": None, - "fts_rebuild_pending": None, - "fts_rebuild_high_water": None, - "fts_rebuild_progress": None, - "fts_rebuild_deferral": None, - } - + stats: Dict[str, Any] = dict.fromkeys(( + "page_count", "page_size", "freelist_count", "logical_size_bytes", "wal_size_bytes", + "journal_mode", "messages", "sessions", "fts_tables", "fts_storage_version", + "fts_rebuild_pending", "fts_rebuild_high_water", "fts_rebuild_progress", + "fts_rebuild_deferral", + )) # WAL sidecar size needs no connection at all. try: wal_path = Path(str(db_path) + "-wal") stats["wal_size_bytes"] = wal_path.stat().st_size if wal_path.exists() else 0 except OSError: pass - - conn = None try: - # mode=ro refuses to create the file and refuses every write; a - # short timeout keeps doctor snappy when a writer holds the lock. - # Route through the tracked connect so byte-probe helpers - # (read_header_bytes_preopen) see this connection and refuse raw + # mode=ro refuses to create the file and refuses every write; a short + # timeout keeps doctor snappy when a writer holds the lock. The tracked + # connect lets byte-probe helpers see this connection and refuse raw # opens that could cancel our POSIX locks mid-read. conn = _connect_tracked_db( - f"file:{Path(db_path)}?mode=ro", - tracking_path=Path(db_path), - uri=True, - timeout=2.0, + f"file:{Path(db_path)}?mode=ro", tracking_path=Path(db_path), uri=True, timeout=2.0 ) except Exception as exc: - logger.debug("collect_state_db_stats: cannot open %s read-only: %s", - db_path, exc) + logger.debug("collect_state_db_stats: cannot open %s read-only: %s", db_path, exc) return stats - def _scalar(sql: str) -> Any: + def _scalar(sql: str, params=()) -> Any: try: - row = conn.execute(sql).fetchone() + row = conn.execute(sql, params).fetchone() return row[0] if row else None except Exception: return None + def _int(value) -> Optional[int]: + return int(value) if value is not None else None + + def _meta_int(key: str) -> Optional[int]: + try: + return _int(_scalar("SELECT value FROM state_meta WHERE key = ?", (key,))) + except Exception: + return None + try: - pc = _scalar("PRAGMA page_count") - ps = _scalar("PRAGMA page_size") - stats["page_count"] = int(pc) if pc is not None else None - stats["page_size"] = int(ps) if ps is not None else None + stats["page_count"] = _int(_scalar("PRAGMA page_count")) + stats["page_size"] = _int(_scalar("PRAGMA page_size")) if stats["page_count"] is not None and stats["page_size"] is not None: stats["logical_size_bytes"] = stats["page_count"] * stats["page_size"] - - fl = _scalar("PRAGMA freelist_count") - stats["freelist_count"] = int(fl) if fl is not None else None - + stats["freelist_count"] = _int(_scalar("PRAGMA freelist_count")) jm = _scalar("PRAGMA journal_mode") stats["journal_mode"] = str(jm) if jm is not None else None - - msgs = _scalar("SELECT COUNT(*) FROM messages") - stats["messages"] = int(msgs) if msgs is not None else None - sess = _scalar("SELECT COUNT(*) FROM sessions") - stats["sessions"] = int(sess) if sess is not None else None - + stats["messages"] = _int(_scalar("SELECT COUNT(*) FROM messages")) + stats["sessions"] = _int(_scalar("SELECT COUNT(*) FROM sessions")) # FTS table presence via sqlite_master (never SELECTs from the # virtual tables themselves — a corrupt index must not fail stats). try: @@ -520,39 +446,22 @@ def collect_state_db_stats(db_path: Path) -> Dict[str, Any]: for row in conn.execute( "SELECT name FROM sqlite_master WHERE type = 'table' " "AND name IN (?, ?, ?)", - ("messages_fts", "messages_fts_trigram", "messages_fts_cjk"), + _FTS_TABLE_NAMES, ).fetchall() } - stats["fts_tables"] = { - t: (t in names) - for t in ("messages_fts", "messages_fts_trigram", "messages_fts_cjk") - } + stats["fts_tables"] = {t: (t in names) for t in _FTS_TABLE_NAMES} except Exception: pass - # Raw state_meta reads — cheap, and independent of SessionDB. - def _meta_int(key: str) -> Optional[int]: - try: - row = conn.execute( - "SELECT value FROM state_meta WHERE key = ?", (key,) - ).fetchone() - return int(row[0]) if row and row[0] is not None else None - except Exception: - return None - stats["fts_storage_version"] = _meta_int("fts_storage_version") high_water = _meta_int("fts_rebuild_high_water") progress = _meta_int("fts_rebuild_progress") stats["fts_rebuild_high_water"] = high_water stats["fts_rebuild_progress"] = progress - if high_water is None: - stats["fts_rebuild_pending"] = False - else: - stats["fts_rebuild_pending"] = (progress or 0) < high_water + stats["fts_rebuild_pending"] = False if high_water is None else (progress or 0) < high_water try: row = conn.execute( - "SELECT value FROM state_meta WHERE key = ? LIMIT 1", - (FTS_REBUILD_DEFERRAL_KEY,), + "SELECT value FROM state_meta WHERE key = ? LIMIT 1", (FTS_REBUILD_DEFERRAL_KEY,) ).fetchone() if row: parsed = json.loads(row[0]) @@ -565,7 +474,6 @@ def collect_state_db_stats(db_path: Path) -> Dict[str, Any]: conn.close() except Exception: pass - return stats @@ -582,33 +490,13 @@ def count_db_holders(db_path: Path) -> Optional[int]: if not sys.platform.startswith("linux"): return None target = os.path.realpath(str(db_path)) - holders = 0 - for pid in os.listdir("/proc"): - if not pid.isdigit(): - continue - fd_dir = f"/proc/{pid}/fd" - try: - fds = os.listdir(fd_dir) - except OSError: - continue # process gone or not ours - for fd in fds: - try: - if os.readlink(f"{fd_dir}/{fd}") == target: - holders += 1 - break # one hit per PID - except OSError: - continue - return holders + return len({pid for pid, link in _iter_proc_fd_targets() if link == target}) except Exception: return None def _is_inactive_orphan_desktop_holder( - *, - ppid: int, - age_seconds: float, - min_age_seconds: float, - ephemeral_backend: bool, + *, ppid: int, age_seconds: float, min_age_seconds: float, ephemeral_backend: bool, connection_statuses: List[str], ) -> bool: """Pure safety predicate for the narrow Desktop holder reap.""" @@ -620,35 +508,21 @@ def _is_inactive_orphan_desktop_holder( ) -def _concrete_state_db_holder_pids( - db_path: Path, holders: List[Tuple[int, str]] -) -> List[int]: +def _concrete_state_db_holder_pids(db_path: Path, holders: List[Tuple[int, str]]) -> List[int]: """Return unique PIDs proven to hold this DB or one of its sidecars.""" canonical_db = os.path.normcase(os.path.abspath(os.fspath(db_path))) - watched = { - canonical_db, - canonical_db + "-wal", - canonical_db + "-shm", - } + watched = {canonical_db, canonical_db + "-wal", canonical_db + "-shm"} pids: List[int] = [] - seen = set() for pid, path in holders: - canonical_path = os.path.normcase( - os.path.abspath(path.removesuffix(" (deleted)")) - ) - if pid <= 0 or pid in seen or canonical_path not in watched: + if pid <= 0 or pid in pids or _canonical_sqlite_path(path) not in watched: continue - seen.add(pid) pids.append(pid) return pids def _read_proc_cmdline(pid: int) -> Optional[str]: - """Read /proc//cmdline, world-readable even when fd table is not. - - Returns the cmdline as a space-joined string, or None when unreadable - (process exited, or hidepid mount). - """ + """Read /proc//cmdline (world-readable even when the fd table is not) + as a space-joined string; None when unreadable (exited, hidepid mount).""" try: with open(f"/proc/{pid}/cmdline", "rb") as f: raw = f.read() @@ -664,11 +538,8 @@ _HERMES_CMDLINE_MARKERS = ("hermes_cli.main", "hermes_cli/main", "hermes serve", def _looks_like_hermes(cmdline: str) -> bool: - """Heuristic: does this cmdline look like a Hermes process? - - Used to decide whether an uninspectable process (fd table unreadable - due to different user) should be treated as a potential state.db holder. - We only flag processes that look like Hermes, not every system daemon. - """ + """Heuristic: does this cmdline look like a Hermes process? Decides whether + an uninspectable process (fd table unreadable, different user) is treated + as a potential state.db holder; system daemons are not flagged.""" lower = cmdline.lower() return any(marker in lower for marker in _HERMES_CMDLINE_MARKERS) diff --git a/hermes_state_repair.py b/hermes_state_repair.py index ba10d93376..183e8e46a9 100644 --- a/hermes_state_repair.py +++ b/hermes_state_repair.py @@ -31,6 +31,73 @@ from hermes_state_common import ( # Log-record parity with the origin module (caplog tests pin "hermes_state"). logger = logging.getLogger("hermes_state") +_REPAIR_LOCK_POLL_SECONDS = 0.1 +# Snapshot copies are data transfer, not inter-process locking: bound them +# separately at 10 MiB/s, with the 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 alongside a damaged DB and pruned with it. ``-journal`` +# matters because rollback-journal (DELETE) mode — the fallback on NFS/SMB/ +# FUSE/ZFS and WAL-reset-vulnerable builds — leaves a hot journal whenever a +# transaction was open; without it the forensic copy cannot be rolled 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. +_FINGERPRINT_SAMPLE_BYTES = 65536 +# Header ranges that move on ordinary commits rather than on repair, masked +# out of the content sample: file change counter (24-27) and version-valid-for +# (92-95). In DELETE mode a commit writes the main file directly and a +# malformed-SCHEMA DB still accepts writes, 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 — what repair identity depends on — sits after byte 100. +_FINGERPRINT_VOLATILE_HEADER_RANGES = ((24, 28), (92, 96)) +# Free-space headroom for the pre-repair forensic backup (a full raw copy of +# the damaged DB plus sidecars; 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 large absolute reserve would turn "repair loops" into +# "repair never runs" on small container/VM 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") +_MANUAL_RECOVER_HINT = 'Free disk space, then retry (or recover manually with `sqlite3 {db_path} ".recover"`).' + + +def _sidecars(db_path: Path): + """The three SQLite sidecar paths for *db_path* (present or not).""" + 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.""" + try: + path.unlink(missing_ok=True) + except OSError: + pass + + +def _offline_access_tools(): + """``(offline_file_access, LiveConnectionError)`` from hermes_cli, or inert + stand-ins on scaffold/embed installs (no tracked connections exist there, + so a raw read is safe).""" + try: + from hermes_cli.sqlite_safe_read import LiveConnectionError, offline_file_access + return offline_file_access, LiveConnectionError + except ImportError: + @contextmanager + def offline_file_access(_path, **_kw): + yield + + class LiveConnectionError(Exception): + pass + + return offline_file_access, LiveConnectionError + def _claim_repair_attempt(db_path: Path) -> bool: """Claim the one-shot per-process repair attempt for *db_path*. @@ -47,12 +114,40 @@ def _claim_repair_attempt(db_path: Path) -> bool: return True -_REPAIR_LOCK_POLL_SECONDS = 0.1 +def _open_lock_file(db_path: Path, suffix: str, what: str, tail: str): + """Open ``.`` for locking; on failure warn and return None.""" + lock_path = db_path.with_name(db_path.name + suffix) + try: + lock_path.parent.mkdir(parents=True, exist_ok=True) + return lock_path, lock_path.open("a+b") + except OSError as exc: + logger.warning(f"Could not open state.db {what} lock %s (%s) — {tail}", lock_path, exc) + return lock_path, None -# Snapshot copies are data transfer, not inter-process locking: bound them -# separately at 10 MiB/s, with the historical two-minute floor. -_REPAIR_SNAPSHOT_MIN_THROUGHPUT_BYTES_PER_SECOND = 10 * 1024 * 1024 +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] + + +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: + 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() @contextlib.contextmanager @@ -61,9 +156,8 @@ def _cross_process_repair_lock(db_path: Path): Yields True when this process holds the repair lock for *db_path*, False when the bounded acquire timed out or the lock file could not be opened. - Unlike the kanban init lock (idempotent critical section), running surgery - unlocked IS the unsafe interleaving this prevents: a caller that gets - False must NOT do surgery. + Running surgery unlocked IS the unsafe interleaving this prevents: a caller + that gets False must NOT do surgery. ``flock`` because the kernel drops it when the holder dies (a pidfile would wedge every future repair); a forked child that inherited the fd is @@ -71,23 +165,16 @@ def _cross_process_repair_lock(db_path: Path): breaks the lock when that holder is provably dead (``_acquire_db_flock``). The acquire is bounded because a *live* repairer can sit in ``VACUUM`` for minutes, and an unbounded wait would hang the caller's open silently. + An unopenable lock file (out of space/inodes/descriptors) fails closed + too: a sibling that opened ITS handle before the disk filled may still be + inside surgery. """ from hermes_state import _IS_WINDOWS, _REPAIR_LOCK_TIMEOUT_SECONDS - lock_path = db_path.with_name(db_path.name + ".repair.lock") - try: - lock_path.parent.mkdir(parents=True, exist_ok=True) - handle = lock_path.open("a+b") - except OSError as exc: - # Fail closed, like a timed-out acquire. An unopenable lock file means - # out of space/inodes/descriptors — and a sibling that opened ITS - # handle before the disk filled may still be inside surgery; yielding - # True here once let two processes run surgery on the same live - # state.db. Callers handle False by re-probing. - logger.warning( - "Could not open state.db repair lock %s (%s) — skipping schema " - "surgery rather than running it without cross-process authority.", - lock_path, exc, - ) + lock_path, handle = _open_lock_file( + db_path, ".repair.lock", "repair", + "skipping schema surgery rather than running it without cross-process authority.", + ) + if handle is None: yield False return @@ -97,10 +184,7 @@ def _cross_process_repair_lock(db_path: Path): deadline = time.monotonic() + _REPAIR_LOCK_TIMEOUT_SECONDS while True: try: - import msvcrt - - handle.seek(0) - msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) + _msvcrt_lock(handle, "LK_NBLCK") acquired = True break except (BlockingIOError, OSError) as exc: @@ -117,41 +201,24 @@ def _cross_process_repair_lock(db_path: Path): time.sleep(_REPAIR_LOCK_POLL_SECONDS) else: acquired, handle = _acquire_db_flock( - str(lock_path), - handle, - _REPAIR_LOCK_TIMEOUT_SECONDS, - _REPAIR_LOCK_POLL_SECONDS, - "state.db repair lock", + str(lock_path), handle, _REPAIR_LOCK_TIMEOUT_SECONDS, + _REPAIR_LOCK_POLL_SECONDS, "state.db repair lock", ) if acquired is None: - # Non-contention failure already logged with its errno. - acquired = False + 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(record), ) yield acquired finally: - try: - if acquired: - if _IS_WINDOWS: - import msvcrt - - handle.seek(0) - msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1) - else: - import fcntl - - _clear_lock_holder_record(handle) - fcntl.flock(handle.fileno(), fcntl.LOCK_UN) - except OSError: # pragma: no cover - best effort release - pass - finally: + if acquired: + _release_lock_handle(handle, clear_record=True) + else: handle.close() @@ -163,27 +230,14 @@ def _try_acquire_auto_maintenance_lock(db_path: Path) -> Optional[Any]: check and the second prunes a row the first has only just closed recoverably. """ from hermes_state import _IS_WINDOWS - lock_path = db_path.with_name(db_path.name + ".auto-maintenance.lock") - try: - lock_path.parent.mkdir(parents=True, exist_ok=True) - handle = lock_path.open("a+b") - except OSError as exc: - logger.warning( - "Could not open state.db auto-maintenance lock %s (%s) — skipping " - "automatic maintenance.", - lock_path, - exc, - ) + _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: - import msvcrt - - handle.seek(0) - msvcrt.locking( # type: ignore[attr-defined] - handle.fileno(), msvcrt.LK_NBLCK, 1 # type: ignore[attr-defined] - ) + _msvcrt_lock(handle, "LK_NBLCK") else: import fcntl @@ -196,23 +250,7 @@ def _try_acquire_auto_maintenance_lock(db_path: Path) -> Optional[Any]: def _release_auto_maintenance_lock(handle: Any) -> None: """Release a handle returned by :func:`_try_acquire_auto_maintenance_lock`.""" - from hermes_state import _IS_WINDOWS - try: - if _IS_WINDOWS: - import msvcrt - - handle.seek(0) - msvcrt.locking( # type: ignore[attr-defined] - handle.fileno(), msvcrt.LK_UNLCK, 1 # type: ignore[attr-defined] - ) - else: - import fcntl - - fcntl.flock(handle.fileno(), fcntl.LOCK_UN) - except OSError: # pragma: no cover - best effort release - pass - finally: - handle.close() + _release_lock_handle(handle) def _bump_schema_cookie(conn: sqlite3.Connection) -> None: @@ -233,36 +271,6 @@ def _bump_schema_cookie(conn: sqlite3.Connection) -> None: logger.warning("Could not bump state.db schema cookie: %s", exc) -_MAX_PERSISTENT_REPAIR_ATTEMPTS = 3 - - -_MAX_MALFORMED_BACKUPS = 3 - - -# Sidecars copied alongside a damaged DB and pruned with it. ``-journal`` -# matters because rollback-journal (DELETE) mode — Hermes's fallback on -# NFS/SMB/FUSE/ZFS and WAL-reset-vulnerable SQLite builds — leaves a hot -# journal whenever a transaction was open; without it the forensic copy -# cannot be rolled back to a consistent state by hand. -_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. -_FINGERPRINT_SAMPLE_BYTES = 65536 - - -# Header ranges that move on ordinary commits rather than on repair, masked -# out of the content sample: file change counter (24-27) and version-valid-for -# (92-95). In DELETE mode a commit writes the main file directly and a -# malformed-SCHEMA DB still accepts writes, so without the mask any live write -# re-keys the ledger and the repair budget resets to 1 forever. (WAL mode -# routes commits to the -wal sidecar; masking is harmless there.) The page-1 -# sqlite_master b-tree — what repair identity depends on — sits after byte 100 -# and stays in the sample. -_FINGERPRINT_VOLATILE_HEADER_RANGES = ((24, 28), (92, 96)) - - def _mask_volatile_header(head: bytes) -> bytes: """Zero the commit-counter fields so ordinary writes don't re-key the ledger.""" if len(head) < 96: @@ -273,26 +281,9 @@ def _mask_volatile_header(head: bytes) -> bytes: return bytes(buf) -# Free-space headroom for the pre-repair forensic backup: a full raw copy of -# the damaged DB plus sidecars, so a repair loop on a large state.db is a disk -# amplifier (one incident wrote ~98MB every ~10s until the volume was nearly -# full). Proportional, not a flat floor: an absolute multi-GB reserve would -# refuse backups that fit on small container/VM volumes, and since a refused -# backup is a HARD STOP that would turn "repair loops" into "repair never -# runs" there. Require the copy plus a small slice of the volume, with a -# modest floor. -_REPAIR_BACKUP_MIN_FREE_BYTES = 256 * 1024 * 1024 # 256 MiB absolute floor - - -_REPAIR_BACKUP_FREE_FRACTION = 0.02 # plus 2% of the volume - - def _repair_backup_headroom_bytes(total_bytes: int) -> int: """Free space required *beyond* the copy itself, for a volume of *total_bytes*.""" - return max( - _REPAIR_BACKUP_MIN_FREE_BYTES, - int(total_bytes * _REPAIR_BACKUP_FREE_FRACTION), - ) + return max(_REPAIR_BACKUP_MIN_FREE_BYTES, int(total_bytes * _REPAIR_BACKUP_FREE_FRACTION)) def _repair_scratch_space_error(db_path: Path) -> Optional[str]: @@ -300,19 +291,13 @@ def _repair_scratch_space_error(db_path: Path) -> Optional[str]: import shutil try: - main_bytes = db_path.stat().st_size - snapshot_bytes = main_bytes - for suffix in _DB_SIDECAR_SUFFIXES: - sidecar = db_path.with_name(db_path.name + suffix) - if sidecar.exists(): - snapshot_bytes += sidecar.stat().st_size + snapshot_bytes = _bundle_bytes(db_path) usage = shutil.disk_usage(db_path.parent) headroom = _repair_backup_headroom_bytes(usage.total) # Strategy 2 runs VACUUM on the staged DB, which SQLite documents may # need up to 2x the database size in extra space; the same reserve then # covers transactional promotion into the live DB. - required = snapshot_bytes + (2 * snapshot_bytes) + headroom - if usage.free >= required: + if usage.free >= snapshot_bytes + (2 * snapshot_bytes) + headroom: return None return ( f"only {usage.free / 1e9:.2f}GB free on {db_path.parent}; the " @@ -336,20 +321,12 @@ def _repair_snapshot_timeout_seconds(source_path: Path) -> float: """ from hermes_state import _REPAIR_LOCK_TIMEOUT_SECONDS, _REPAIR_SNAPSHOT_MIN_THROUGHPUT_BYTES_PER_SECOND source_bytes = 0 - for suffix in ("", *_DB_SIDECAR_SUFFIXES): - candidate = ( - source_path - if not suffix - else source_path.with_name(source_path.name + suffix) - ) + for candidate in (source_path, *_sidecars(source_path)): try: source_bytes += candidate.stat().st_size except FileNotFoundError: continue - return max( - _REPAIR_LOCK_TIMEOUT_SECONDS, - source_bytes / _REPAIR_SNAPSHOT_MIN_THROUGHPUT_BYTES_PER_SECOND, - ) + return max(_REPAIR_LOCK_TIMEOUT_SECONDS, source_bytes / _REPAIR_SNAPSHOT_MIN_THROUGHPUT_BYTES_PER_SECOND) def _repair_failure_consumes_attempt(exc: BaseException) -> bool: @@ -366,16 +343,11 @@ def _repair_failure_consumes_attempt(exc: BaseException) -> bool: error_code = getattr(exc, "sqlite_errorcode", None) if isinstance(error_code, int): # Extended result codes keep the primary code in the low byte. - primary_code = error_code & 0xFF - return primary_code in (sqlite3.SQLITE_CORRUPT, sqlite3.SQLITE_NOTADB) - + 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 - ) + return "file is not a database" in message or "database disk image is malformed" in message def _repair_ledger_path(db_path: Path) -> Path: @@ -387,11 +359,10 @@ def _db_fingerprint(db_path: Path) -> "Optional[str]": Deliberately EXCLUDES mtime: the malformed-schema class still accepts writes, so live writers, WAL checkpoints and the strategies themselves move - mtime between passes; keyed on mtime, every pass looked like a NEW file, - the attempt counter reset to 1 forever and each pass wrote another - full-size forensic copy. Hashing a multi-GB file on every open is the cost - this ledger exists to avoid, so sample the head/tail slices any real - repair, truncation or restore necessarily changes. + mtime between passes; keyed on mtime, every pass looked like a NEW file and + the attempt counter reset to 1 forever. Hashing a multi-GB file on every + open is the cost this ledger exists to avoid, so sample the head/tail + slices any real repair, truncation or restore necessarily changes. Runs under ``offline_file_access``: ``close()`` on ANY raw descriptor cancels every POSIX advisory lock this process holds on the file, including @@ -400,35 +371,19 @@ def _db_fingerprint(db_path: Path) -> "Optional[str]": ``_backup_db_file``'s ``has_live_connection`` guard). Returns ``None`` ("identity unavailable") when that makes the read unsafe. Callers MUST NOT substitute a differently-shaped key: the ledger compares keys for equality, - so alternating shapes never matches and the unbounded loop returns; the - ledger helpers keep the recorded key instead. + so alternating shapes never matches and the unbounded loop returns. """ try: st = db_path.stat() - try: - from hermes_cli.sqlite_safe_read import ( - LiveConnectionError, - offline_file_access, - ) - except ImportError: - # Scaffold/embed installs ship hermes_state without hermes_cli; no - # tracked connections exist there, so the raw read is safe. - @contextmanager - def offline_file_access(_path, **_kw): - yield - - class LiveConnectionError(Exception): - pass - + offline_file_access, LiveConnectionError = _offline_access_tools() try: with offline_file_access(db_path, what="fingerprint"): 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) - else: - tail = b"" except LiveConnectionError: return None digest = hashlib.sha256(_mask_volatile_header(head) + tail).hexdigest()[:32] @@ -441,31 +396,18 @@ def _backup_content_identity(db_path: Path) -> "Optional[str]": """Recovery-image identity for forensic-backup dedupe: whole file + sidecars. A DIFFERENT equivalence relation from :func:`_db_fingerprint`; never - conflate them. The fingerprint answers "same repair epoch?" and masks - commit counters / samples only head+tail so an ordinary write does not mint - a fresh repair budget. A live writer can commit rows into an *interior* - page while preserving size and the first/last 64 KiB, so two materially - different recovery images share one fingerprint; 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 committed - frames). The O(n) read is cheaper than the O(n) write it avoids. Runs under - ``offline_file_access`` (same POSIX-lock reason as ``_db_fingerprint``); - ``None`` when a live connection makes the read unsafe — the caller then - takes a fresh backup, never a false reuse. + conflate them. The fingerprint answers "same repair epoch?" and samples + only head+tail; a live writer can commit rows into an *interior* page while + preserving size and both 64 KiB slices, so two materially different + recovery images share one fingerprint, and 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 committed frames). Runs + under ``offline_file_access`` (same POSIX-lock reason); ``None`` when a + live connection makes the read unsafe — the caller then takes a fresh + backup, never a false reuse. """ - try: - from hermes_cli.sqlite_safe_read import ( - LiveConnectionError, - offline_file_access, - ) - except ImportError: - @contextmanager - def offline_file_access(_path, **_kw): - yield - - class LiveConnectionError(Exception): - pass + offline_file_access, LiveConnectionError = _offline_access_tools() def _hash_whole(path: Path, hasher: "Any") -> None: with open(path, "rb") as fh: @@ -480,15 +422,12 @@ def _backup_content_identity(db_path: Path) -> "Optional[str]": # coincide with a main+sidecar split and dedupe two images together. hasher.update(f"\0main:{db_path.stat().st_size}\0".encode()) _hash_whole(db_path, hasher) - for suffix in _DB_SIDECAR_SUFFIXES: - sidecar = db_path.with_name(db_path.name + suffix) + for suffix, sidecar in zip(_DB_SIDECAR_SUFFIXES, _sidecars(db_path)): if sidecar.exists(): hasher.update(f"\0{suffix}:{sidecar.stat().st_size}\0".encode()) _hash_whole(sidecar, hasher) return hasher.hexdigest() - except LiveConnectionError: - return None - except OSError: + except (LiveConnectionError, OSError): return None @@ -541,9 +480,7 @@ def _persistent_repair_exhausted_error(db_path: Path) -> str: ) -def _record_repair_outcome( - db_path: Path, *, repaired: bool, fingerprint: "Optional[str]" = None -) -> None: +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. Defaults to the post-attempt fingerprint (what the NEXT exhaustion probe @@ -565,21 +502,15 @@ def _record_repair_outcome( # 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 - ) + attempts = int(ledger.get("failed_attempts", 0)) + 1 if recorded == fp else 1 import datetime ledger_path.write_text( - json.dumps( - { - "fingerprint": fp, - "failed_attempts": attempts, - "last_attempt": datetime.datetime.now().isoformat( - timespec="seconds" - ), - } - ), + json.dumps({ + "fingerprint": fp, + "failed_attempts": attempts, + "last_attempt": datetime.datetime.now().isoformat(timespec="seconds"), + }), encoding="utf-8", ) except Exception as exc: # pragma: no cover - best effort @@ -591,10 +522,8 @@ def _existing_malformed_backups(db_path: Path) -> "List[Path]": prefix = f"{db_path.name}.malformed-backup-" try: found = [ - p - for p in db_path.parent.iterdir() - if p.name.startswith(prefix) - and not p.name.endswith(_DB_SIDECAR_SUFFIXES) + p for p in db_path.parent.iterdir() + if p.name.startswith(prefix) and not p.name.endswith(_DB_SIDECAR_SUFFIXES) ] except OSError: return [] @@ -604,16 +533,76 @@ def _existing_malformed_backups(db_path: Path) -> "List[Path]": def _prune_malformed_backups(db_path: Path, keep: int = _MAX_MALFORMED_BACKUPS) -> None: """Delete all but the *keep* newest forensic backups (and sidecars).""" for stale in _existing_malformed_backups(db_path)[keep:]: - for victim in ( - stale, - *(stale.with_name(stale.name + suffix) for suffix in _DB_SIDECAR_SUFFIXES), - ): + for victim in (stale, *_sidecars(stale)): try: victim.unlink(missing_ok=True) except OSError as exc: # pragma: no cover - best effort logger.warning("Could not prune stale DB backup %s: %s", victim, exc) +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. Fails CLOSED on stat()/disk_usage() errors: the nearly-full volume + this guard exists for is exactly where they are most likely to fail, and + proceeding would take the copy that finishes off the disk. + """ + import shutil + + hint = _MANUAL_RECOVER_HINT.format(db_path=db_path) + try: + need = _bundle_bytes(db_path) + usage = shutil.disk_usage(db_path.parent) + headroom = _repair_backup_headroom_bytes(usage.total) + if usage.free - need >= headroom: + return None + return ( + f"only {usage.free / 1e9:.2f}GB free on {db_path.parent}; " + f"copying the damaged DB needs {need / 1e9:.2f}GB and must " + f"leave {headroom / 1e9:.2f}GB headroom. {hint}" + ) + except OSError as exc: + return ( + f"could not determine free space on {db_path.parent} ({exc}); " + f"refusing the forensic copy rather than risk filling the volume. {hint}" + ) + + +def _publish_backup_bundle(db_path: Path, staging: Path, backup_path: Path) -> None: + """Copy DB + sidecars to *staging* names, then rename each into place. + + PUBLICATION ORDER MATTERS: the main DB name is the bundle's commit marker + (what ``_existing_malformed_backups`` counts), so sidecars go FIRST and the + main DB LAST; a failure partway then never leaves a countable main backup + over a missing sidecar. On any failure, unpublished staging files AND + anything already promoted are removed, so no official ``backup_path`` + survives a partial bundle. + """ + import shutil + + pairs: "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() + ] + 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)): + os.replace(staged, dst) + published.append(dst) + except Exception: + for staged in (staging, *(s for _src, s, _d in pairs)): + _unlink_quiet(staged) + for dst in published: + _unlink_quiet(dst) + raise + + def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": """Raw-copy a (possibly malformed) DB plus sidecars to a timestamped backup. @@ -623,11 +612,17 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": because the forensic bundle is the recovery path when every strategy fails. Refuses while a connection to this DB is live in the process: reading the file would ``close()`` a descriptor and cancel that connection's POSIX - advisory locks (see ``hermes_cli.sqlite_safe_read``) — a real case, since - one SessionDB can enter repair while the gateway holds others. + advisory locks (see ``hermes_cli.sqlite_safe_read``). + + Dedupe: if the newest existing backup is byte-identical to the current + recovery image (``_backup_content_identity`` — NOT mtime, NOT + ``_db_fingerprint``), reuse it; a repair loop once copied the same damaged + bytes on every restart. The copy lands under a staging name OUTSIDE the + ``.malformed-backup-`` prefix: a staging name inside it counts as a backup, + sorts NEWEST (prune kept partials and deleted intact copies) and dedupe + could return it with no real forensic copy on disk. """ import datetime - import shutil try: from hermes_cli.sqlite_safe_read import has_live_connection @@ -645,13 +640,10 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": 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 (two damaged states within one second) must not - # overwrite the earlier forensic copy. + # Same-second collision must not overwrite the earlier forensic copy. seq = 1 while backup_path.exists(): - backup_path = db_path.with_name( - f"{db_path.name}.malformed-backup-{stamp}_{seq}" - ) + backup_path = db_path.with_name(f"{db_path.name}.malformed-backup-{stamp}_{seq}") seq += 1 try: # Sweep staging debris from an earlier interrupted pass BEFORE the @@ -659,22 +651,9 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": # dedupe would hand it back as a legitimate backup. Also sweeps the # pre-merge ``.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 pattern in (f"{db_path.name}.backup-staging-*", f"{db_path.name}.malformed-backup-*.incomplete*"): for old in db_path.parent.glob(pattern): - try: - old.unlink(missing_ok=True) - except OSError: # pragma: no cover - best effort - pass - # Dedupe: a repair loop used to copy the SAME damaged bytes on every - # restart (~900MB a pass, 89GB over 11 days). If the newest existing - # backup is byte-identical to the current recovery image, reuse it. - # Match on ``_backup_content_identity`` — NOT mtime (the malformed- - # SCHEMA class still accepts writes, so mtime missed every pass) and - # NOT ``_db_fingerprint`` (see its docstring: an interior-page write - # changes the recovery image without changing the fingerprint). + _unlink_quiet(old) try: # Only hash the source when there is a candidate to compare # against; hashing a multi-GB source right before copying it is @@ -691,90 +670,12 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": return existing, None except OSError: pass - # Disk guard: 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. Refuse while there is room to. - try: - need = db_path.stat().st_size - for suffix in _DB_SIDECAR_SUFFIXES: - sidecar = db_path.with_name(db_path.name + suffix) - if sidecar.exists(): - need += sidecar.stat().st_size - usage = shutil.disk_usage(db_path.parent) - headroom = _repair_backup_headroom_bytes(usage.total) - if usage.free - need < headroom: - reason = ( - f"only {usage.free / 1e9:.2f}GB free on {db_path.parent}; " - f"copying the damaged DB needs {need / 1e9:.2f}GB and must " - f"leave {headroom / 1e9:.2f}GB headroom. Free disk space, " - f"then retry (or recover manually with `sqlite3 {db_path} " - '".recover"`).' - ) - logger.error("Refusing forensic backup of %s: %s", db_path, reason) - return None, reason - except OSError as exc: - # Fail CLOSED: the nearly-full volume this guard exists for is - # exactly where stat()/disk_usage() is most likely to fail, and - # proceeding would take the copy that finishes off the disk. Repair - # then waits (HARD STOP) for a human to free space — the safe side. - reason = ( - f"could not determine free space on {db_path.parent} ({exc}); " - "refusing the forensic copy rather than risk filling the " - f"volume. Free disk space, then retry (or recover manually " - f'with `sqlite3 {db_path} ".recover"`).' - ) + reason = _backup_free_space_error(db_path) + if reason is not None: logger.error("Refusing forensic backup of %s: %s", db_path, reason) return None, reason - # Copy to a staging name OUTSIDE the ``.malformed-backup-`` prefix and - # rename into place only once every copy succeeded. A staging name - # inside the prefix (e.g. ``…-.incomplete``) counts as a backup, - # sorts NEWEST (so prune kept partials and deleted intact copies), and - # dedupe could return it as ``backup_path`` — passing the hard stop with - # no real forensic copy on disk. staging = db_path.with_name(f"{db_path.name}.backup-staging-{stamp}") - # (staging_src, final_dst) pairs. PUBLICATION ORDER MATTERS: the main - # DB name is the bundle's commit marker (what - # ``_existing_malformed_backups`` counts), so sidecars go FIRST and the - # main DB LAST; a failure partway then never leaves a countable main - # backup over a missing sidecar that would pass the hard stop and dedupe. - staged_sidecars: "List[Tuple[Path, Path, Path]]" = [] - for suffix in _DB_SIDECAR_SUFFIXES: - sidecar = db_path.with_name(db_path.name + suffix) - if sidecar.exists(): - side_staging = staging.with_name(staging.name + suffix) - side_dst = backup_path.with_name(backup_path.name + suffix) - staged_sidecars.append((sidecar, side_staging, side_dst)) - main_pair = (staging, backup_path) - published: "List[Path]" = [] - all_staging_srcs = [staging] + [s for _src, s, _d in staged_sidecars] - try: - shutil.copy2(db_path, staging) - for sidecar, side_staging, _side_dst in staged_sidecars: - shutil.copy2(sidecar, side_staging) - publish_order = [ - (s, d) for _src, s, d in staged_sidecars - ] + [main_pair] - for src, dst in publish_order: - os.replace(src, dst) - published.append(dst) - except Exception: - # Roll back unpublished staging files AND anything already promoted, - # so a failure after the main os.replace leaves no official backup_path. - for src in all_staging_srcs: - try: - src.unlink(missing_ok=True) - except OSError: - pass - for dst in published: - try: - dst.unlink(missing_ok=True) - except OSError: - pass - try: - staging.unlink(missing_ok=True) - except OSError: - pass - raise + _publish_backup_bundle(db_path, staging, backup_path) _prune_malformed_backups(db_path) return backup_path, None except Exception as exc: # pragma: no cover - best effort @@ -782,11 +683,7 @@ def _backup_db_file(db_path: Path) -> "Tuple[Optional[Path], Optional[str]]": return None, f"backup copy failed: {exc}" -def preflight_db_writability( - db_path: Path, - *, - db_label: str = "state.db", -) -> None: +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 @@ -803,7 +700,6 @@ def preflight_db_writability( raw = str(db_path) if raw == ":memory:" or raw.startswith("file:"): return - try: home: Optional[Path] = Path(get_hermes_home()).resolve() except Exception: # pragma: no cover - defensive @@ -831,17 +727,14 @@ def preflight_db_writability( if 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 "", + 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 "" + if p.name.endswith("-wal") else "" ) raise sqlite3.OperationalError( f"{db_label} is not writable: {kind} {p} is read-only for this " @@ -855,16 +748,13 @@ def preflight_db_writability( # SQLite needs a writable directory in every journal mode (WAL/SHM # sidecars, or the rollback journal in DELETE mode). _ensure_writable(parent, is_dir=True) - for suffix in ("", "-wal", "-shm"): p = db_path.with_name(db_path.name + suffix) if suffix else db_path if p.is_file(): _ensure_writable(p) -def _connect_repair_durable( - db_path: Path, *, timeout: float = 5.0 -) -> sqlite3.Connection: +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 open ``state.db`` directly (not via ``SessionDB`` / @@ -897,10 +787,7 @@ def _reapply_durability_barriers(conn: sqlite3.Connection) -> bool: _apply_macos_checkpoint_barrier(conn) _enforce_macos_synchronous_full(conn) return True - except sqlite3.DatabaseError: - # Schema still unparseable — pragmas cannot be set yet. - return False - except Exception: + except Exception: # schema still unparseable, or anything else — pragmas cannot be set yet return False @@ -921,14 +808,21 @@ def apply_durability_barriers(conn: sqlite3.Connection) -> bool: cfg = load_config_readonly() raw_synchronous = cfg_get(cfg, "database", "synchronous", default=None) if raw_synchronous is not None: - _apply_synchronous_pragma( - conn, raw_synchronous, db_label="state.db (guest)" - ) + _apply_synchronous_pragma(conn, raw_synchronous, db_label="state.db (guest)") except Exception: pass return ok +def _close_unpinned(conn: sqlite3.Connection) -> None: + """Leave EXCLUSIVE locking mode (so the file is never left pinned) and close.""" + try: + conn.execute("PRAGMA locking_mode=NORMAL") + except Exception: + pass + conn.close() + + @contextmanager def _exclusive_repair_db_guard(db_path: Path): """Yield one live connection that excludes writers for repair surgery. @@ -940,44 +834,32 @@ def _exclusive_repair_db_guard(db_path: Path): stays open across the whole snapshot -> strategies -> promotion window, so no other writer can commit a change promotion would overwrite. Existing readers make acquisition fail rather than being disturbed: repair fails - closed unless this process owns the whole window. + closed unless this process owns the whole window. Timeout 0: the + cross-process repair lock already serializes repairers, and a partial + repair is less safe than an explicit "stop the gateway and retry". """ guard: Optional[sqlite3.Connection] = None try: - # The cross-process repair lock already serializes repairers. Do not - # wait behind an ordinary application connection: a partial repair is - # less safe than an explicit "stop the gateway and retry". guard = _connect_repair_durable(db_path, timeout=0.0) guard.execute("PRAGMA locking_mode=EXCLUSIVE") guard.execute("BEGIN EXCLUSIVE") guard.execute("ROLLBACK") except (sqlite3.Error, OSError) as exc: if guard is not None: - try: - guard.execute("PRAGMA locking_mode=NORMAL") - except Exception: - pass - guard.close() + _close_unpinned(guard) yield None, exc return - try: yield guard, None finally: - try: - # Release the exclusive locks before close; also keeps a close-time - # checkpoint from being mistaken for a repair write by callers that - # immediately reopen state.db. - guard.execute("PRAGMA locking_mode=NORMAL") - except Exception: - pass - guard.close() + # Releasing the exclusive locks before close also keeps a close-time + # checkpoint from being mistaken for a repair write by callers that + # immediately reopen state.db. + _close_unpinned(guard) def _copy_database_snapshot( - source_path: Path, - destination_path: Path, - *, + source_path: Path, destination_path: Path, *, source_connection: Optional[sqlite3.Connection] = None, destination_connection: Optional[sqlite3.Connection] = None, ) -> None: @@ -993,15 +875,10 @@ def _copy_database_snapshot( deadline = time.monotonic() + deadline_seconds source = source_connection or _connect_repair_durable(source_path) destination = destination_connection - own_source = source_connection is None - own_destination = destination_connection is None def _check_deadline(_status: int, _remaining: int, _total: int) -> None: if time.monotonic() >= deadline: - raise TimeoutError( - "timed out copying SQLite repair snapshot after " - f"{deadline_seconds:.0f}s" - ) + raise TimeoutError(f"timed out copying SQLite repair snapshot after {deadline_seconds:.0f}s") try: if destination is None: @@ -1009,19 +886,12 @@ def _copy_database_snapshot( elif destination.in_transaction: # sqlite3_backup needs a transaction-free destination; the exclusive # guard retains exclusion via locking_mode, not a transaction. - raise sqlite3.ProgrammingError( - "SQLite repair backup destination has an active transaction" - ) - source.backup( - destination, - pages=256, - progress=_check_deadline, - sleep=_REPAIR_LOCK_POLL_SECONDS, - ) + raise sqlite3.ProgrammingError("SQLite repair backup destination has an active transaction") + source.backup(destination, pages=256, progress=_check_deadline, sleep=_REPAIR_LOCK_POLL_SECONDS) finally: - if own_destination and destination is not None: + if destination_connection is None and destination is not None: destination.close() - if own_source: + if source_connection is None: source.close() @@ -1039,9 +909,7 @@ def _db_opens_cleanly(db_path: Path) -> Optional[str]: try: # Best-effort tokenizer load: messages_fts_cjk needs cjk_unicode61 # before any statement (incl. the trigger-driven write probe) can touch - # it. Without it this probe sees the DB as a tokenizer-less SessionDB - # would (which drops the cjk triggers), so tokenizer absence must never - # classify as corruption. + # 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() @@ -1050,33 +918,24 @@ def _db_opens_cleanly(db_path: Path) -> Optional[str]: return "; ".join(problems[:3]) conn.execute("SELECT COUNT(*) FROM sessions").fetchone() - # FTS5 read probe. The write probe below misses partial shadow-table - # corruption where MATCH / snippet / rank raise DatabaseError("database - # disk image is malformed"), silently breaking session_search and - # /resume title resolution while check-only reports healthy. - for fts_table in ("messages_fts", "messages_fts_trigram", "messages_fts_cjk"): + # FTS5 read probe: partial shadow-table corruption makes MATCH / + # snippet / rank raise while check-only reports healthy. MATCH '""' + # (empty phrase) parses, scans zero rows and exercises the shadow-table + # read path; FTS5 rejects MATCH '' outright. + for fts_table in _FTS_TABLES: try: - # Trigram backs title resolution, so probe it too. MATCH '""' - # (empty phrase) parses, scans zero rows and exercises the - # shadow-table read path; FTS5 rejects MATCH '' outright. - conn.execute( - f"SELECT 1 FROM {fts_table} WHERE {fts_table} MATCH '\"\"' LIMIT 1" - ).fetchone() + conn.execute(f"SELECT 1 FROM {fts_table} WHERE {fts_table} MATCH '\"\"' LIMIT 1").fetchone() except sqlite3.OperationalError as exc: - # Canonical capability classifier: on builds without fts5 a - # legacy messages_fts table may exist and MATCH raises "no such - # module: fts5"; treating that as corruption would send the DB - # into repair, whose final fallback deletes the messages_fts% - # schema. Covers "no such tokenizer: trigram" too. + # Builds without fts5 / trigram raise "no such module|tokenizer"; + # treating that as corruption would send the DB into repair, + # whose final fallback deletes the messages_fts% schema. if SessionDB._is_fts5_unavailable_error(exc): continue msg = str(exc).lower() if "no such table" in msg or "no such column" in msg: - # FTS5 not built yet (brand new file mid-init). - continue + continue # FTS5 not built yet (brand new file mid-init) return f"fts5 read probe failed on {fts_table}: {exc}" except sqlite3.DatabaseError as exc: - # Partial shadow-table damage: MATCH raises though the table parses. return f"fts5 read probe failed on {fts_table}: {exc}" # FTS write probe: drive a row through the messages_fts* triggers in a @@ -1096,7 +955,6 @@ def _db_opens_cleanly(db_path: Path) -> Optional[str]: ) conn.execute("ROLLBACK") except sqlite3.OperationalError as exc: - # Missing tables / FTS disabled — not the corruption class we probe. try: conn.execute("ROLLBACK") except sqlite3.Error: @@ -1128,10 +986,9 @@ def _live_writer_holds_db(db_path: Path) -> bool: busy/locked signal — refusing to repair a DB nobody holds would strand the self-heal path. - Scope: WAL mode only. In ``journal_mode=DELETE`` (Hermes's fallback on - WAL-reset-vulnerable builds and NFS/SMB) a held reader takes only SHARED - and this returns False; repair is then serialised only by the cross-process - repairer lock. Broadening to DELETE mode is a follow-up. + Scope: WAL mode only. 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. """ probe = None try: @@ -1143,25 +1000,23 @@ def _live_writer_holds_db(db_path: Path) -> bool: except sqlite3.OperationalError as exc: lowered = str(exc).lower() return "locked" in lowered or "busy" in lowered - except sqlite3.DatabaseError: - # Malformed/unreadable: no evidence of a live holder either way. - return False - except Exception: + except Exception: # malformed/unreadable: no evidence of a live holder either way return False finally: if probe is not None: try: - # Drop exclusive mode before close so the probe never leaves - # the file pinned. - probe.execute("PRAGMA locking_mode=NORMAL") - except Exception: - pass - try: - probe.close() + _close_unpinned(probe) except Exception: pass +def _repair_skip(report: Dict[str, Any], verb: str, error: str) -> Dict[str, Any]: + """Record *error* on *report* and log it as ``state.db repair ``.""" + report["error"] = error + logger.error(f"state.db repair {verb}: %s", report["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. @@ -1184,12 +1039,7 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A error: str|None}``. """ 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, - } + report: Dict[str, Any] = {"repaired": False, "strategy": None, "backup_path": None, "error": None} # Startup-watchdog progress lease: repair is I/O-bound (near-zero CPU), # which the watchdog's CPU fallback would misread as a parked deadlock. A @@ -1204,12 +1054,9 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A # Cross-restart attempt cap: the in-memory claim bounds one process, but a # class the strategies cannot heal (b-tree page damage) used to re-run the - # whole surgery, with a fresh forensic backup, on EVERY restart. After - # _MAX_PERSISTENT_REPAIR_ATTEMPTS failures on the same file, stop. + # whole surgery, with a fresh forensic backup, on EVERY restart. if _persistent_repair_attempts_exhausted(db_path): - report["error"] = _persistent_repair_exhausted_error(db_path) - logger.error("state.db repair skipped: %s", report["error"]) - return report + return _repair_skip(report, "skipped", _persistent_repair_exhausted_error(db_path)) result = report with _cross_process_repair_lock(db_path) as holding_lock: @@ -1229,32 +1076,27 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A ) else: # Recheck exhaustion after acquisition: a queued repairer can have - # recorded the final failure while this process waited, and this - # process must not start a fourth attempt. + # recorded the final failure while this process waited. if _persistent_repair_attempts_exhausted(db_path): - report["error"] = _persistent_repair_exhausted_error(db_path) - logger.error("state.db repair skipped: %s", report["error"]) + _repair_skip(report, "skipped", _persistent_repair_exhausted_error(db_path)) # WAL-holder preflight: fail-closed for active readers before a # forensic backup is taken. Not the race defence — the exclusive # guard in the locked routine excludes writers through promotion and # rejects DELETE-mode readers this probe cannot see. elif _live_writer_holds_db(db_path): - report["error"] = ( + _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." + "Stop the gateway (hermes gateway stop) and retry.", ) - logger.error("state.db repair skipped: %s", report["error"]) else: # Probe the journal mode BEFORE surgery: a rebuilt file comes # back in the default (delete) mode and nothing else records - # the flip (see _restore_journal_mode_after_repair). The probe - # may fail on a damaged file; then database.journal_mode is - # the restore target. + # the flip. The probe may fail on a damaged file; then + # 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 - ) + 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) @@ -1266,9 +1108,7 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A # 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")) - ) + _record_repair_outcome(db_path, repaired=bool(result.get("repaired"))) return result @@ -1297,12 +1137,10 @@ def _restore_journal_mode_after_repair(db_path: Path, before_mode: Optional[str] 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 vulnerable-SQLite WAL-reset gate (a rebuilt file IS a new - database; on a vulnerable runtime the gate deliberately keeps DELETE, and - "could not reach WAL" is expected there), the macOS-NFS silent-refusal - handling, and the WAL companions (size limit, checkpoint barrier, - synchronous=FULL). ``before_mode`` (None if unprobeable) is only for the - log comparison; the target comes from ``database.journal_mode``. + inherits the vulnerable-SQLite WAL-reset gate (on a vulnerable runtime the + gate deliberately keeps DELETE), the macOS-NFS silent-refusal handling, + and the WAL companions. ``before_mode`` (None if unprobeable) is only for + the log comparison; the target comes from ``database.journal_mode``. Best-effort: the repair already succeeded, so failures log at WARNING. """ from hermes_state import apply_wal_with_fallback @@ -1328,36 +1166,30 @@ def _restore_journal_mode_after_repair(db_path: Path, before_mode: Optional[str] ) -def _repair_state_db_schema_locked( - db_path: Path, *, backup: bool, report: Dict[str, Any] -) -> Dict[str, Any]: +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 must hold the cross-process repair lock for *db_path*. Strategies run on a SCRATCH COPY; the result is copied back through SQLite's transactional backup API only once proven to open cleanly, so a failed - repair cannot modify or lose committed canonical data. (A WAL checkpoint - of already-committed frames on guard release is not a repair mutation.) + repair cannot modify or lose committed canonical data. WHY not in place: Strategy 2 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 handled here — every table hanging off the unreadable part is silently dropped, the probe then correctly reports STILL malformed, and repair returned ``repaired=False`` - having destroyed what it was asked to save. The forensic backup does not - close this (nothing reads it back). Not mutating the original is the - property that holds without a human in the loop. + having destroyed what it was asked to save. Not mutating the original is + the property that holds without a human in the loop. """ 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: - report["error"] = ( - "could not remove a stale repair snapshot before probing state.db: " - f"{cleanup_error}" + return _repair_skip( + report, "aborted", + f"could not remove a stale repair snapshot before probing state.db: {cleanup_error}", ) - logger.error("state.db repair aborted: %s", report["error"]) - return report # Re-probe under the lock: a process we queued behind may have just # repaired the file; redoing surgery would undo its work (the @@ -1373,12 +1205,11 @@ def _repair_state_db_schema_locked( if bpath is None: # HARD STOP: the forensic image is still required when corruption # defeats every strategy, even though strategies run on a snapshot. - report["error"] = ( + 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}" + f"mutating the only copy of the damaged DB: {backup_error}", ) - logger.error("state.db repair aborted: %s", report["error"]) - return report # The forensic copy deliberately precedes this guard: its raw-copy safety # checks inspect real live holders and would be poisoned by our exclusive @@ -1386,38 +1217,31 @@ def _repair_state_db_schema_locked( # happens only after writer exclusion is held. with _exclusive_repair_db_guard(db_path) as (live_guard, guard_error): if live_guard is None: - report["error"] = ( + if guard_error is not None and _repair_failure_consumes_attempt(guard_error): + report["_repair_attempted"] = True + 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}" + f"writer. Stop the gateway and retry: {guard_error}", ) - if guard_error is not None and _repair_failure_consumes_attempt( - guard_error - ): - report["_repair_attempted"] = True - logger.error("state.db repair skipped: %s", report["error"]) - return report space_error = _repair_scratch_space_error(db_path) if space_error is not None: - report["error"] = space_error - logger.error("state.db repair aborted: %s", report["error"]) - return report + return _repair_skip(report, "aborted", space_error) try: # Reuse live_guard rather than a second source connection: the # guard 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 - ) + _copy_database_snapshot(db_path, scratch, source_connection=live_guard) except (OSError, sqlite3.Error, TimeoutError) as exc: - report["error"] = ( - f"could not stage a complete SQLite repair snapshot of {db_path}: {exc}" - ) if _repair_failure_consumes_attempt(exc): report["_repair_attempted"] = True - logger.error("state.db repair aborted: %s", report["error"]) + _repair_skip( + report, "aborted", + f"could not stage a complete SQLite repair snapshot of {db_path}: {exc}", + ) _unlink_db_triple(scratch) return report @@ -1433,27 +1257,17 @@ def _repair_state_db_schema_locked( # under open handles and POSIX would leave those handles on # the old inode. The guard that staged the live image # receives the promotion, keeping writer exclusion throughout. - _copy_database_snapshot( - scratch, - db_path, - destination_connection=live_guard, - ) + _copy_database_snapshot(scratch, db_path, destination_connection=live_guard) except (OSError, sqlite3.Error, TimeoutError) as exc: report["repaired"] = False report["strategy"] = None - report["_repair_attempted"] = _repair_failure_consumes_attempt( - exc - ) - report["error"] = ( - "repaired snapshot could not be promoted transactionally: " - f"{exc}" - ) + report["_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, + report.get("strategy"), db_path, ) if not report.get("repaired"): # Logged HERE, not in the strategies: they see the scratch copy, @@ -1463,8 +1277,7 @@ def _repair_state_db_schema_locked( "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"], + db_path, report["backup_path"], ) return report finally: @@ -1472,18 +1285,14 @@ def _repair_state_db_schema_locked( # or human to mistake for the real thing. cleanup_error = _unlink_db_triple(scratch) if cleanup_error is not None: - logger.warning( - "Could not remove state.db repair snapshot after repair: %s", - cleanup_error, - ) + logger.warning("Could not remove state.db repair snapshot after repair: %s", cleanup_error) def _unlink_db_triple(path: Path) -> Optional[str]: """Remove *path* and every SQLite sidecar; return any cleanup failure.""" from hermes_state import _IS_WINDOWS failures: List[str] = [] - for suffix in ("", *_DB_SIDECAR_SUFFIXES): - victim = path if not suffix else path.with_name(path.name + suffix) + for victim in (path, *_sidecars(path)): for attempt in range(10): try: victim.unlink() @@ -1505,135 +1314,126 @@ def _unlink_db_triple(path: Path) -> Optional[str]: return "; ".join(failures) or None -def _run_repair_strategies( - db_path: Path, report: Dict[str, Any] -) -> Dict[str, Any]: +# ── Repair strategies, least destructive first (each mutates its connection's DB) ── + +def _strategy_rebuild_fts(conn: sqlite3.Connection) -> None: + """FTS5 'rebuild' rewrites each index from the content table: the least- + destructive fix for an index that rejects writes while reads work.""" + from hermes_state import load_fts5_cjk_extension + # The cjk index can only be rebuilt with its tokenizer loaded + # (best-effort; a tokenizer-less host skips it below). + load_fts5_cjk_extension(conn) + for table_name in _FTS_TABLES: + try: + 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: + """integrity_check reports "wrong # of entries in index" when a B-tree + index drifts from its base table; REINDEX rewrites it from canonical rows.""" + # REINDEX rewrites every index b-tree; take the barriers now that the + # schema parses, in case the open-time attempt was refused. + _reapply_durability_barriers(conn) + conn.execute("REINDEX") + conn.commit() + + +def _strategy_dedup_schema(conn: sqlite3.Connection) -> None: + """De-duplicate sqlite_master (lowest rowid per type/name), keeping FTS.""" + conn.execute("PRAGMA writable_schema=ON") + 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), + ) + if dupes: + _bump_schema_cookie(conn) + conn.execute("PRAGMA writable_schema=OFF") + conn.commit() + + +def _strategy_drop_fts_vacuum(conn: sqlite3.Connection) -> None: + """Drop all FTS schema and VACUUM; indexes rebuild on the next open. + + The destructive one, and why the strategies run on a scratch copy: on a + damaged schema b-tree VACUUM silently drops every table hanging off the + unreadable part (see _repair_state_db_schema_locked). + """ + conn.execute("PRAGMA writable_schema=ON") + conn.execute("DELETE FROM sqlite_master WHERE name LIKE 'messages_fts%'") + _bump_schema_cookie(conn) + conn.execute("PRAGMA writable_schema=OFF") + conn.commit() + # The schema parses now, so the barriers can finally stick — and VACUUM + # rewrites the entire file, the worst operation to lose halfway. + _reapply_durability_barriers(conn) + conn.execute("VACUUM") + + +# (strategy name, body, success log, failure log) in escalation order. The +# first three log a warning on failure and fall through; the last records the +# failure in report["error"] for the caller. +_REPAIR_STRATEGIES = ( + ("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", + "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"), +) +_FINAL_STRATEGY = ( + "drop_fts_rebuild", _strategy_drop_fts_vacuum, + "state.db schema repaired by dropping FTS schema; indexes will rebuild from messages on next open: %s", +) + + +def _run_repair_strategies(db_path: Path, report: Dict[str, Any]) -> Dict[str, Any]: """Escalating repair attempts, applied to *db_path* IN PLACE. Every strategy mutates its argument, so this is only ever called by :func:`_repair_state_db_schema_locked` on a scratch copy nothing else - holds open — never on the user's database. + holds open — never on the user's database. The "could not recover" log + lives in the caller: it must name the user's database, not the scratch copy. """ - from hermes_state import _db_opens_cleanly, load_fts5_cjk_extension - # ── Strategy 0: rebuild FTS indexes in place (FTS write-corruption) ── - # FTS5 'rebuild' rewrites the index from the content table: the - # least-destructive fix for an index that rejects writes while reads work. - try: - conn = _connect_repair_durable(db_path) - try: - # The cjk index can only be rebuilt with its tokenizer loaded - # (best-effort; a tokenizer-less host skips it below). - load_fts5_cjk_extension(conn) - for table_name in ( - "messages_fts", "messages_fts_trigram", "messages_fts_cjk" - ): - try: - conn.execute( - f"INSERT INTO {table_name}({table_name}) VALUES('rebuild')" - ) - except sqlite3.OperationalError: - # Table absent (FTS disabled / trigram off / cjk not present - # or tokenizer unavailable). - continue - finally: - conn.close() - if _db_opens_cleanly(db_path) is None: - report["repaired"] = True - report["strategy"] = "rebuild_fts" - logger.warning( - "state.db FTS indexes rebuilt in place (schema preserved): %s", - db_path, - ) - return report - except sqlite3.DatabaseError as exc: - logger.warning("state.db FTS in-place rebuild pass failed: %s", exc) + from hermes_state import _db_opens_cleanly - # ── Strategy 0.5: rebuild stale B-tree indexes ── - # integrity_check reports "wrong # of entries in index" when a B-tree index - # drifts from its base table; REINDEX rewrites it from the canonical rows. - try: + def _apply(body) -> Optional[str]: conn = _connect_repair_durable(db_path) try: - # REINDEX rewrites every index b-tree; take the barriers now that - # the schema parses, in case the open-time attempt was refused. - _reapply_durability_barriers(conn) - conn.execute("REINDEX") - conn.commit() + body(conn) finally: conn.close() - if _db_opens_cleanly(db_path) is None: - report["repaired"] = True - report["strategy"] = "reindex_btree" - logger.warning( - "state.db B-tree indexes rebuilt via REINDEX: %s", db_path - ) - return report - except sqlite3.DatabaseError as exc: - logger.warning("state.db REINDEX pass failed: %s", exc) + return _db_opens_cleanly(db_path) - # ── Strategy 1: de-duplicate sqlite_master (keeps FTS index) ── - try: - conn = _connect_repair_durable(db_path) - try: - conn.execute("PRAGMA writable_schema=ON") - 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), - ) - if dupes: - _bump_schema_cookie(conn) - conn.execute("PRAGMA writable_schema=OFF") - conn.commit() - finally: - conn.close() - if _db_opens_cleanly(db_path) is None: - report["repaired"] = True - report["strategy"] = "dedup_schema" - logger.warning( - "state.db schema repaired by de-duplicating sqlite_master " - "(FTS index preserved): %s", db_path - ) - return report - except sqlite3.DatabaseError as exc: - logger.warning("state.db dedup repair pass failed: %s", exc) + def _succeed(name: str, message: str) -> Dict[str, Any]: + report["repaired"] = True + report["strategy"] = name + logger.warning(message, db_path) + return report - # ── Strategy 2: drop all FTS schema, VACUUM, rebuild on next open ── - # The destructive one, and why this path runs on a scratch copy: on a - # damaged schema b-tree VACUUM silently drops every table hanging off the - # unreadable part (see _repair_state_db_schema_locked). - try: - conn = _connect_repair_durable(db_path) + for name, body, success_msg, failure_msg in _REPAIR_STRATEGIES: try: - conn.execute("PRAGMA writable_schema=ON") - conn.execute("DELETE FROM sqlite_master WHERE name LIKE 'messages_fts%'") - _bump_schema_cookie(conn) - conn.execute("PRAGMA writable_schema=OFF") - conn.commit() - # The schema parses now, so the barriers can finally stick — and - # VACUUM rewrites the entire file, the worst operation to lose halfway. - _reapply_durability_barriers(conn) - conn.execute("VACUUM") - finally: - conn.close() - reason = _db_opens_cleanly(db_path) + if _apply(body) is None: + return _succeed(name, success_msg) + except sqlite3.DatabaseError as exc: + logger.warning(failure_msg, exc) + + name, body, success_msg = _FINAL_STRATEGY + try: + reason = _apply(body) if reason is None: - report["repaired"] = True - report["strategy"] = "drop_fts_rebuild" - logger.warning( - "state.db schema repaired by dropping FTS schema; indexes " - "will rebuild from messages on next open: %s", db_path - ) - return report + return _succeed(name, success_msg) report["error"] = reason except sqlite3.DatabaseError as exc: report["error"] = str(exc) - - # The "could not recover" log lives in the caller: it must name the user's - # database, not the scratch copy. return report