From 69245e65bda72aaf320d18c4e60d143cba7959a8 Mon Sep 17 00:00:00 2001 From: loulanyue <260355617@qq.com> Date: Sun, 30 Aug 2026 08:10:25 +0800 Subject: [PATCH] fix(state): serialize startup across zero-byte check, quarantine, connect, and schema commit (#97568) - Guard against concurrent-opener race where newly created 0-byte state.db was falsely quarantined before first schema write - Wrap startup in quarantine_cross_process_lock when database is uninitialized or zeroed - Guard is_zeroed_sqlite_file and is_zeroed_state_db against active live connections in current process - Add concurrent-opener and live-connection regression tests --- hermes_cli/backup.py | 5 +- hermes_state.py | 176 ++++++++++++++++++++-------------- tests/test_zeroed_state_db.py | 79 +++++++++++++++ 3 files changed, 189 insertions(+), 71 deletions(-) diff --git a/hermes_cli/backup.py b/hermes_cli/backup.py index 7d540991c8..b2a068797b 100644 --- a/hermes_cli/backup.py +++ b/hermes_cli/backup.py @@ -451,7 +451,10 @@ def is_zeroed_sqlite_file( return False if size < 0: return False - from hermes_cli.sqlite_safe_read import read_header_bytes_preopen + from hermes_cli.sqlite_safe_read import has_live_connection, read_header_bytes_preopen + + if not force and has_live_connection(path): + return False head = read_header_bytes_preopen( path, length=max(16, probe_bytes), force=force diff --git a/hermes_state.py b/hermes_state.py index 28946a6b06..c619e440c2 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -3954,7 +3954,10 @@ def is_zeroed_state_db( return False if size < 0: return False - from hermes_cli.sqlite_safe_read import read_header_bytes_preopen + from hermes_cli.sqlite_safe_read import has_live_connection, read_header_bytes_preopen + + if not force and has_live_connection(path): + return False head = read_header_bytes_preopen( path, length=max(16, probe_bytes), force=force @@ -3968,14 +3971,9 @@ def is_zeroed_state_db( return all(byte == 0 for byte in head) -def quarantine_zeroed_state_db(path: Path) -> 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. - """ +@contextlib.contextmanager +def quarantine_cross_process_lock(path: Path, timeout: float = 5.0): + """Acquire the cross-process lock for path.quarantine.lock.""" import platform lock_path = path.with_name(path.name + ".quarantine.lock") @@ -3983,9 +3981,10 @@ def quarantine_zeroed_state_db(path: Path) -> Optional[Path]: handle = lock_path.open("a+b") acquired = False try: - deadline = time.monotonic() + 5.0 + deadline = time.monotonic() + timeout if platform.system() == "Windows": import msvcrt + while True: try: handle.seek(0) @@ -3998,6 +3997,7 @@ def quarantine_zeroed_state_db(path: Path) -> Optional[Path]: time.sleep(0.020) else: import fcntl + while True: try: fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) @@ -4007,21 +4007,36 @@ def quarantine_zeroed_state_db(path: Path) -> Optional[Path]: if time.monotonic() >= deadline: break time.sleep(0.020) - if not acquired: - # Fail closed: do NOT proceed without the lock. A slow or paused - # startup that still owns the lock can overlap this fallback and - # the two processes can act on the same live file (#68805 review). - logger.error( - "quarantine lock for %s not acquired within 5s — refusing to " - "quarantine without the cross-process lock. The zeroed file " - "is left in place. If sessions fail to load, restore from " - "state-snapshots via `hermes snapshot list` / " - "`hermes snapshot restore `.", - path, - ) - return None - # Re-check under the lock: another process may have already quarantined - # the file, leaving a fresh DB (or no file at all) in its place. + 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) + except (OSError, AttributeError): + pass + finally: + handle.close() + + +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. + """ + def _do_quarantine(): if not path.exists(): logger.info( "quarantine_zeroed_state_db: %s already moved by another process", @@ -4040,12 +4055,9 @@ def quarantine_zeroed_state_db(path: Path) -> Optional[Path]: ts = time.strftime("%Y%m%d-%H%M%S") except Exception: ts = "unknown" - # Unique destination with PID suffix to avoid collision across - # concurrent startups that somehow both enter the lock. dest = path.with_name( f"{path.name}.zeroed-{ts}-{os.getpid()}.bak" ) - # Non-clobbering: if dest somehow exists, append a counter. n = 0 while dest.exists(): n += 1 @@ -4057,7 +4069,6 @@ def quarantine_zeroed_state_db(path: Path) -> Optional[Path]: except OSError as exc: logger.error("Failed to quarantine zeroed %s: %s", path, exc) return None - # Also move empty WAL/SHM if present so a fresh open is clean for suffix in ("-wal", "-shm"): side = Path(str(path) + suffix) if side.exists(): @@ -4066,20 +4077,22 @@ def quarantine_zeroed_state_db(path: Path) -> Optional[Path]: except OSError: pass return dest - 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) - except (OSError, AttributeError): - pass - finally: - handle.close() + + if already_locked: + return _do_quarantine() + + with quarantine_cross_process_lock(path) as acquired: + if not acquired: + logger.error( + "quarantine lock for %s not acquired within 5s — refusing to " + "quarantine without the cross-process lock. The zeroed file " + "is left in place. If sessions fail to load, restore from " + "state-snapshots via `hermes snapshot list` / " + "`hermes snapshot restore `.", + path, + ) + return None + return _do_quarantine() # ── Read-only health/stats probes (hermes doctor, dashboards) ────────── @@ -4669,34 +4682,43 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) if not read_only: preflight_db_writability(self.db_path, db_label="state.db") - # #68474: zeroed state.db (size>0, all-NUL header) used to fail as a - # generic "file is not a database" with no recovery path. Quarantine - # the bytes (do not delete) and continue so a fresh DB can open; - # point the operator at pre-update snapshots. - if ( + # #68474 / #97568: Serialize startup across zero-byte check, quarantine, + # connect, and schema commit so concurrent openers don't race on an + # absent-path -> connect -> schema-commit window. + needs_startup_guard = ( not read_only - and self.db_path.exists() - and is_zeroed_state_db(self.db_path) - ): - try: - zsize = self.db_path.stat().st_size - except OSError: - zsize = -1 - qpath = quarantine_zeroed_state_db(self.db_path) - snaps = self.db_path.parent / "state-snapshots" - msg = ( - f"state.db looks ZEROED ({zsize} bytes, no SQLite header). " - f"Preserved at {qpath or '(quarantine failed — file left in place)'}. " - f"Restore from {snaps} via `hermes snapshot list` / " - f"`hermes snapshot restore ` if available. " - "Opening a fresh empty database so the agent can start." + and ( + not self.db_path.exists() + or is_zeroed_state_db(self.db_path) ) - logger.error(msg) - _set_last_init_error(msg) - # If quarantine failed, do not open the zeroed file (would fail - # opaquely or risk further damage). Raise with the clear message. - if qpath is None and self.db_path.exists() and is_zeroed_state_db(self.db_path): - raise sqlite3.DatabaseError(msg) + ) + + def _handle_quarantine_if_zeroed(already_locked: bool = False): + if ( + self.db_path.exists() + and is_zeroed_state_db(self.db_path) + ): + try: + zsize = self.db_path.stat().st_size + except OSError: + zsize = -1 + qpath = quarantine_zeroed_state_db( + self.db_path, already_locked=already_locked + ) + snaps = self.db_path.parent / "state-snapshots" + msg = ( + f"state.db looks ZEROED ({zsize} bytes, no SQLite header). " + f"Preserved at {qpath or '(quarantine failed — file left in place)'}. " + f"Restore from {snaps} via `hermes snapshot list` / " + f"`hermes snapshot restore ` if available. " + "Opening a fresh empty database so the agent can start." + ) + logger.error(msg) + _set_last_init_error(msg) + # If quarantine failed, do not open the zeroed file (would fail + # opaquely or risk further damage). Raise with the clear message. + if qpath is None and self.db_path.exists() and is_zeroed_state_db(self.db_path): + raise sqlite3.DatabaseError(msg) def _connect_and_init(): self._conn = _connect_tracked_db( @@ -4758,8 +4780,22 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) ) ) + def _open_with_optional_startup_guard(): + if needs_startup_guard: + with quarantine_cross_process_lock(self.db_path) as lock_acquired: + if not lock_acquired: + logger.warning( + "startup quarantine lock for %s not acquired within 5s; proceeding", + self.db_path, + ) + _handle_quarantine_if_zeroed(already_locked=lock_acquired) + _connect_and_init_with_lock_patience() + else: + _handle_quarantine_if_zeroed(already_locked=False) + _connect_and_init_with_lock_patience() + try: - _connect_and_init_with_lock_patience() + _open_with_optional_startup_guard() except sqlite3.DatabaseError as exc: # The malformed-schema class (e.g. a duplicate sqlite_master # row for messages_fts) fails on the very first statement — diff --git a/tests/test_zeroed_state_db.py b/tests/test_zeroed_state_db.py index 470dd9dfae..a3d05b2ab2 100644 --- a/tests/test_zeroed_state_db.py +++ b/tests/test_zeroed_state_db.py @@ -200,3 +200,82 @@ def test_quarantine_fails_closed_when_lock_held(tmp_path): # Release the lock so the background thread can exit cleanly release_lock.set() holder.join(timeout=5) + + +def test_concurrent_openers_zero_byte_startup_serialization(tmp_path, monkeypatch): + """#97580: Verify that two concurrent SessionDB openers on a non-existent + database serialize through the startup lock, avoid racing on the initial + 0-byte creation window, and do not falsely quarantine each other's live file. + """ + import hermes_state as hs + import threading + + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + db = tmp_path / "state.db" + + errors = [None, None] + results = [None, None] + + def worker(idx): + try: + sdb = hs.SessionDB(db_path=db) + try: + # Confirm schema is active + row = sdb._conn.execute("SELECT 1").fetchone() + assert row[0] == 1 + results[idx] = "ok" + finally: + sdb.close() + except Exception as exc: + errors[idx] = exc + + t1 = threading.Thread(target=worker, args=(0,)) + t2 = threading.Thread(target=worker, args=(1,)) + t1.start() + t2.start() + t1.join(timeout=10) + t2.join(timeout=10) + + assert errors[0] is None, f"Opener 0 failed: {errors[0]}" + assert errors[1] is None, f"Opener 1 failed: {errors[1]}" + assert results[0] == "ok" + assert results[1] == "ok" + + # No spurious quarantine backups should have been created + backups = list(tmp_path.glob("state.db.zeroed-*.bak")) + assert len(backups) == 0, f"Expected 0 quarantine backups, got: {backups}" + assert db.exists() + assert not hs.is_zeroed_state_db(db) + + +def test_live_connection_0_byte_not_quarantined_in_process(tmp_path, monkeypatch): + """#97580: A live 0-byte connection tracked in this process must not be + quarantined by is_zeroed_state_db / SessionDB. + """ + import hermes_state as hs + from hermes_cli.sqlite_safe_read import connect_tracked + + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + db = tmp_path / "state.db" + + # Create a live tracked 0-byte connection + conn = connect_tracked(str(db)) + try: + assert db.exists() and db.stat().st_size == 0 + # is_zeroed_state_db must recognize the live connection and refuse to declare it zeroed + assert hs.is_zeroed_state_db(db) is False + + # SessionDB open must not quarantine this live file + sdb = hs.SessionDB(db_path=db) + try: + backups = list(tmp_path.glob("state.db.zeroed-*.bak")) + assert len(backups) == 0, f"Spurious quarantine occurred: {backups}" + finally: + sdb.close() + + # The original connection can still write safely + conn.execute("CREATE TABLE live_check (id INTEGER)") + conn.commit() + finally: + conn.close() +