diff --git a/hermes_cli/session_recovery.py b/hermes_cli/session_recovery.py index fc679015a02d..ac880218e91c 100644 --- a/hermes_cli/session_recovery.py +++ b/hermes_cli/session_recovery.py @@ -239,34 +239,33 @@ def _disk_space_preflight( def _copy_source_bundle(source: Path, snapshot_dir: Path) -> tuple[Path, list[str]]: """Copy the source DB bundle aside so SQLite never opens the original. - Refuses when a connection to *source* is live in this process. Copying a - database file is an ``open()``/``close()`` on it, and ``close()`` cancels - every POSIX advisory lock the process holds on that file -- including a - running VACUUM's EXCLUSIVE lock (see ``hermes_cli.sqlite_safe_read``). - Recovery normally runs as its own short-lived CLI process against an - offline/quarantined file, so this should never fire; the check keeps this - path consistent with ``hermes_state._backup_db_file``, which refuses the - same situation, rather than leaving two policies for one hazard. - """ - from hermes_cli.sqlite_safe_read import has_live_connection + The whole copy runs inside ``offline_file_access``, which holds the + connection-lifecycle lock for its duration. Checking for a live connection + and *then* copying would be a check/use race: a connection could open in + that window, and the copy's ``close()`` would cancel its POSIX advisory + locks -- the failure class ``hermes_cli.sqlite_safe_read`` exists to + prevent (see #71724). Holding the lock means no connection can appear + mid-copy, across the main file and every sidecar. - if has_live_connection(source): - raise SessionRecoverySafetyError( - f"Refusing to snapshot {source}: a connection to it is still open " - "in this process, and copying the file would cancel that " - "connection's POSIX locks. Close all database handles (stop the " - "gateway/dashboard) and re-run." - ) + Recovery normally runs as its own short-lived CLI process against an + offline/quarantined file, so the refusal should never fire; the guard + keeps this path consistent with ``hermes_state._backup_db_file``. + """ + from hermes_cli.sqlite_safe_read import LiveConnectionError, offline_file_access snapshot_source = snapshot_dir / source.name copied: list[str] = [] - for suffix in _SIDECAR_SUFFIXES: - source_part = _sidecar_path(source, suffix) - if not source_part.exists(): - continue - destination_part = _sidecar_path(snapshot_source, suffix) - shutil.copy2(source_part, destination_part) - copied.append(destination_part.name) + try: + with offline_file_access(source, what="snapshot"): + for suffix in _SIDECAR_SUFFIXES: + source_part = _sidecar_path(source, suffix) + if not source_part.exists(): + continue + destination_part = _sidecar_path(snapshot_source, suffix) + shutil.copy2(source_part, destination_part) + copied.append(destination_part.name) + except LiveConnectionError as exc: + raise SessionRecoverySafetyError(str(exc)) from exc return snapshot_source, copied @@ -751,7 +750,39 @@ def _copy_state_meta_salvage( progress_cb: Optional[ProgressCallback], source_rows: Optional[int], ) -> dict[str, Any]: - """Salvage readable user metadata while regenerating derived FTS state.""" + """Salvage readable user metadata while regenerating derived FTS state. + + Requires both ``key`` and ``value``, matching the non-partial + :func:`_copy_state_meta`. A damaged ``state_meta`` can retain one column + and lose the other; without this check a missing ``key`` raised + ``ValueError`` from ``columns.index("key")`` and aborted the entire + partial recovery, and a missing ``value`` would have copied key-only rows + while reporting the table complete. ``state_meta`` is optional metadata — + recording it as unusable lets ``--allow-partial`` surface the loss as a + warning and carry on recovering sessions and messages. + """ + source_columns = _table_columns(source, "state_meta") + destination_columns = _table_columns(destination, "state_meta") + if not {"key", "value"}.issubset(source_columns): + return { + "mode": "rowid_range_salvage", + "source_meta_rows": source_rows, + "copied_rows": 0, + "columns": ["key", "value"], + "excluded_keys": sorted(_GENERATED_META_KEYS), + "status": "missing", + "error": "source state_meta is missing the key/value columns", + } + if not {"key", "value"}.issubset(destination_columns): + return { + "mode": "rowid_range_salvage", + "source_meta_rows": source_rows, + "copied_rows": 0, + "columns": ["key", "value"], + "excluded_keys": sorted(_GENERATED_META_KEYS), + "status": "failed", + "error": "destination state_meta schema is incomplete", + } def keep_user_meta( row: tuple[Any, ...], diff --git a/hermes_cli/sqlite_safe_read.py b/hermes_cli/sqlite_safe_read.py index e373a91c850f..352228ed522d 100644 --- a/hermes_cli/sqlite_safe_read.py +++ b/hermes_cli/sqlite_safe_read.py @@ -61,6 +61,7 @@ later probe of the real ``Path`` can ever match. from __future__ import annotations +import contextlib import logging import os import sqlite3 @@ -372,3 +373,37 @@ def read_header_bytes_preopen( return handle.read(length) except OSError: return None + + +class LiveConnectionError(RuntimeError): + """A raw file operation was attempted on a database with live connections.""" + + +@contextlib.contextmanager +def offline_file_access(path: Path | str, *, what: str = "read"): + """Hold the connection-lifecycle lock across a raw read of a database file. + + Checking :func:`has_live_connection` and *then* doing the raw I/O is a + check/use race: a connection can be opened in the window between the two, + and the raw ``close()`` will cancel its POSIX advisory locks — the exact + failure class the registry exists to prevent. Any multi-step raw access + (copying a database plus its ``-wal``/``-shm``/``-journal`` sidecars, + hashing a file, moving a bundle aside) must therefore run *inside* this + context manager rather than after a bare check. + + While held, :func:`connect_tracked` blocks, so no new connection can + appear mid-copy. Raises :class:`LiveConnectionError` if a connection is + already live when the guard is entered. + + The lock is only held for the duration of the raw I/O; it never spans + caller work on an open connection, so it does not serialise database use. + """ + with _live_lock: + if _key(path) in _live_connections: + raise LiveConnectionError( + f"Refusing to {what} {path}: a connection to it is still open " + "in this process, and raw file access would cancel that " + "connection's POSIX advisory locks. Close all database " + "handles (stop the gateway/dashboard) and retry." + ) + yield diff --git a/tests/hermes_cli/test_session_recovery.py b/tests/hermes_cli/test_session_recovery.py index a16fccf4d761..d58c1e89495b 100644 --- a/tests/hermes_cli/test_session_recovery.py +++ b/tests/hermes_cli/test_session_recovery.py @@ -288,6 +288,106 @@ def _corrupt_middle_table_leaf( return leaf_page +def test_snapshot_blocks_connections_opened_during_the_copy( + tmp_path: Path, +) -> None: + """No connection may open mid-copy — a bare pre-check is a check/use race. + + Checking has_live_connection() and then copying leaves a window: a + connection can open between the two, and the copy's close() cancels its + POSIX advisory locks. The guard must hold the lifecycle lock across the + whole bundle copy, so connect_tracked() blocks until every raw descriptor + is closed. + """ + import threading + import time + + from hermes_cli import session_recovery as recovery_module + from hermes_cli.sqlite_safe_read import connect_tracked + + source = tmp_path / "racy-state.db" + snapshot_dir = tmp_path / "snapshot" + snapshot_dir.mkdir() + _make_source(source) + + inside_copy = threading.Event() + release_copy = threading.Event() + observed: dict[str, object] = {} + real_copy2 = recovery_module.shutil.copy2 + + def slow_copy2(src, dst, *args, **kwargs): + result = real_copy2(src, dst, *args, **kwargs) + if str(src).endswith("racy-state.db"): + inside_copy.set() + release_copy.wait(10) + return result + + def racer(): + inside_copy.wait(10) + started = time.monotonic() + conn = connect_tracked(source, isolation_level=None, timeout=10.0) + observed["blocked_seconds"] = time.monotonic() - started + observed["copy_finished_first"] = release_copy.is_set() + conn.close() + + thread = threading.Thread(target=racer, daemon=True) + thread.start() + + recovery_module.shutil.copy2 = slow_copy2 + try: + _, copied = recovery_module._copy_source_bundle(source, snapshot_dir) + # Let the racer proceed only once the copy has fully returned. + release_copy.set() + finally: + recovery_module.shutil.copy2 = real_copy2 + release_copy.set() + thread.join(20) + + assert copied, "snapshot should have copied at least the main database" + assert observed.get("copy_finished_first") is True, ( + "a tracked connection opened while raw descriptors on the database " + "were still open — the check/use race is present" + ) + + +def test_partial_recovery_survives_state_meta_missing_key_column( + tmp_path: Path, +) -> None: + """A damaged optional state_meta must degrade, not abort the recovery. + + state_meta can lose ``key`` while ``value`` stays readable. Indexing + ``key`` unconditionally raised ValueError and killed the whole partial + recovery; the table must instead be recorded as unusable so sessions and + messages still come back. + """ + source = tmp_path / "meta-damaged.db" + output = tmp_path / "meta-recovered.db" + expected = _make_source(source) + + conn = sqlite3.connect(str(source), isolation_level=None) + try: + conn.execute("DROP TABLE IF EXISTS state_meta") + conn.execute("CREATE TABLE state_meta(value TEXT)") + conn.execute("INSERT INTO state_meta(value) VALUES ('orphaned')") + finally: + conn.close() + + report = recover_session_database( + source, + output, + work_dir=tmp_path, + chunk_size=4, + allow_partial=True, + ) + + assert report["copy"]["state_meta"]["status"] in {"missing", "failed"} + # The canonical data still recovers. + assert report["verification"]["table_counts"]["sessions"] == expected["sessions"] + assert report["verification"]["table_counts"]["messages"] == expected["messages"] + assert report["verification"]["integrity_check"] == ["ok"] + assert report["installed"] is False + + def test_recovery_refuses_to_snapshot_a_live_database(tmp_path: Path) -> None: """Snapshotting must refuse while a connection to the source is live.