diff --git a/gateway/run.py b/gateway/run.py index 9b9e7eb12e..3f2b91d0c7 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -32621,6 +32621,7 @@ def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop AUTO_ARCHIVE_EVERY = 60 # ticks — poll hourly (state_meta gate owns the real cadence) MEMORY_TRIM_EVERY = 1 # shared helper cooldown bounds actual allocator work MISFIRE_SWEEP_EVERY = 5 # ticks — every 5 minutes (grace window gates real work) + FTS_STALE_RETRY_EVERY = 1 # SessionDB rate-limits the real work (_FTS_STALE_RETRY_SECONDS) # Every platform media cache prunes on the same hourly cadence — one loop # over (name, cleanup_fn), not a copy-pasted try/except per cache. @@ -32755,6 +32756,29 @@ def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop except Exception as e: logger.debug("Auto-archive tick error: %s", e) + # Deferred stale-FTS rebuild retry (#100108). A SessionDB that opened + # while another process held state.db / the rebuild lock fails closed + # and leaves search on the LIKE fallback; a short-lived CLI clears + # that on its next open, but the gateway opens once and stays up for + # days. Retry here, on the existing tick, against the shared + # instances this process already holds: non-blocking admission, no + # new thread, rate-limited inside SessionDB. No-op when nothing is + # stale (one attribute read per instance). + if tick_count % FTS_STALE_RETRY_EVERY == 0: + try: + from hermes_state_registry import live_shared_session_dbs + + for _sdb in live_shared_session_dbs(): + _retry = getattr(_sdb, "retry_deferred_fts_recovery", None) + if callable(_retry) and _retry(): + logger.info( + "Deferred state.db FTS rebuild completed in-process " + "for %s; full-text search restored.", + getattr(_sdb, "db_path", "state.db"), + ) + except Exception as exc: + logger.debug("Deferred FTS retry tick error: %s", exc) + # This is the long-lived messaging-gateway counterpart to the TUI idle # reaper. The helper is config-gated and rate-limited, so calling it on # the 60s housekeeping cadence does not create a trim storm. diff --git a/hermes_state.py b/hermes_state.py index a17da8893f..3cd6f152af 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -100,6 +100,7 @@ from hermes_state_common import ( # noqa: F401 (re-exported for back-compat) _clear_lock_holder_record, _describe_lock_holder, _read_lock_holder_record, + is_advisory_lock_contention, ) from hermes_state_portability import SessionPortabilityMixin from hermes_state_schema import SessionSchemaMixin @@ -1776,13 +1777,14 @@ def _log_wal_reset_bug_once( # for git/pip/system Python installs (#75153). repair_hint = _wal_reset_repair_hint() logger.warning( - "%s: linked SQLite %s is vulnerable to the WAL-reset corruption " - "bug (https://sqlite.org/wal.html#walresetbug) — %s. " + "%s: linked SQLite %s (interpreter %s) is vulnerable to the WAL-reset " + "corruption bug (https://sqlite.org/wal.html#walresetbug) — %s. " "Upgrade to SQLite 3.51.3+ (or backports 3.50.7 / 3.44.6); " "%s. See `hermes doctor`. This warning fires once per " "process per database.", db_label, sqlite3.sqlite_version, + sys.executable, action, repair_hint, ) @@ -2388,7 +2390,15 @@ def _cross_process_repair_lock(db_path: Path): msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) acquired = True break - except (BlockingIOError, OSError): + except (BlockingIOError, OSError) as exc: + if not is_advisory_lock_contention(exc): + logger.warning( + "Could not acquire state.db repair lock %s (%s) — " + "skipping schema surgery on a non-contention error.", + lock_path, exc, + ) + acquired = None + break if time.monotonic() >= deadline: break time.sleep(_REPAIR_LOCK_POLL_SECONDS) @@ -2400,7 +2410,10 @@ def _cross_process_repair_lock(db_path: Path): _REPAIR_LOCK_POLL_SECONDS, "state.db repair lock", ) - if not acquired: + if acquired is None: + # Non-contention failure already logged with its errno. + acquired = False + 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 " diff --git a/hermes_state_common.py b/hermes_state_common.py index fc94d5475f..ae9e63af06 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -7,6 +7,7 @@ hermes_state re-imports every name here for backward compatibility. """ import contextlib +import errno import json import logging import os @@ -925,6 +926,32 @@ _IS_WINDOWS = sys.platform == "win32" # short bounded wait suffices — never re-enter the full timeout. _LOCK_BREAK_REACQUIRE_SECONDS = 5.0 +# errno set for "another process holds this advisory lock". flock() reports +# contention as EWOULDBLOCK/EAGAIN; msvcrt.locking() as EACCES (and EDEADLK +# when its internal retry gives up). Anything else — ESTALE on a dropped NFS +# handle, ENOTSUP/ENOLCK on a filesystem without advisory locks, EIO — is a +# persistent environment failure that no amount of polling turns into an +# acquire. Treating every OSError as contention made such a failure look +# like a live holder and burned the full 120s admission timeout on every +# attempt (#100108, PR #100130). +_LOCK_CONTENTION_ERRNOS = {errno.EAGAIN, errno.EACCES, errno.EWOULDBLOCK} +if hasattr(errno, "EDEADLK"): + _LOCK_CONTENTION_ERRNOS.add(errno.EDEADLK) + + +def is_advisory_lock_contention(exc: BaseException) -> bool: + """True when *exc* means another process holds the advisory lock. + + False for every other ``OSError`` (ESTALE, ENOTSUP, ENOLCK, EIO, ...): + callers must fail closed IMMEDIATELY rather than poll to the deadline, + because retrying cannot succeed and the wait only stalls the caller. + """ + if isinstance(exc, BlockingIOError): + return True + if not isinstance(exc, OSError): + return False + return exc.errno in _LOCK_CONTENTION_ERRNOS + def _proc_start_ticks(pid: int): """Kernel start time of *pid* in clock ticks, or None when unknowable. @@ -1033,7 +1060,11 @@ def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, descript """Bounded POSIX flock acquire with orphaned-holder staleness break. Returns ``(acquired, handle)``; *handle* may have been re-opened (the - caller owns closing whichever handle comes back). + caller owns closing whichever handle comes back). *acquired* is True on + success, False when a holder kept the lock past the deadline, and None + when a non-contention ``OSError`` (ESTALE/ENOTSUP/EIO) made acquisition + impossible — already logged here; callers treat None as "not acquired" + without emitting the held-by-another-process warning. Why breaking exists at all (issue #100108): ``flock`` belongs to the open file DESCRIPTION, which ``fork()`` duplicates into every child. A holder @@ -1056,7 +1087,21 @@ def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, descript while True: try: fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) - except (BlockingIOError, OSError): + except (BlockingIOError, OSError) as exc: + if not is_advisory_lock_contention(exc): + # ESTALE / ENOTSUP / EIO: not a holder, and polling cannot + # fix it. Defer NOW instead of pretending a live process + # held the lock for the whole timeout (#100108). + logger.warning( + "Could not acquire %s %s (%s) — deferring rather than " + "waiting out the %.0fs holder timeout on a " + "non-contention error.", + description, + lock_path, + exc, + timeout_seconds, + ) + return None, handle if time.monotonic() < deadline: time.sleep(poll_seconds) continue @@ -1129,7 +1174,7 @@ def _describe_lock_holder(record) -> str: @contextlib.contextmanager -def fts_rebuild_admission(db_path): +def fts_rebuild_admission(db_path, *, timeout_seconds=None): """Serialize full structural FTS rebuilds on *db_path* across processes. Yields True when this process holds the rebuild authority, False when the @@ -1142,10 +1187,20 @@ def fts_rebuild_admission(db_path): ``db_path`` may be a str or Path; None (in-memory DB / tests without a file path) yields True — a private in-memory DB has no cross-process surface. + + *timeout_seconds* defaults to ``_FTS_REBUILD_LOCK_TIMEOUT_SECONDS``. + Opportunistic in-process retries (``retry_deferred_fts_recovery``) pass + ``0`` so a live holder never stalls a long-lived writer for two minutes; + the orphaned-holder break still applies on the single attempt. """ if db_path is None: yield True return + timeout = ( + _FTS_REBUILD_LOCK_TIMEOUT_SECONDS + if timeout_seconds is None + else max(float(timeout_seconds), 0.0) + ) lock_path = f"{db_path}.fts_rebuild.lock" try: handle = open(lock_path, "a+b") @@ -1171,7 +1226,7 @@ def fts_rebuild_admission(db_path): acquired = False try: if _IS_WINDOWS: - deadline = time.monotonic() + _FTS_REBUILD_LOCK_TIMEOUT_SECONDS + deadline = time.monotonic() + timeout while True: try: import msvcrt @@ -1180,7 +1235,15 @@ def fts_rebuild_admission(db_path): msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) acquired = True break - except (BlockingIOError, OSError): + except (BlockingIOError, OSError) as exc: + if not is_advisory_lock_contention(exc): + logger.warning( + "Could not acquire FTS rebuild lock %s (%s) — " + "deferring on a non-contention error.", + lock_path, exc, + ) + acquired = None + break if time.monotonic() >= deadline: break time.sleep(_FTS_REBUILD_LOCK_POLL_SECONDS) @@ -1188,20 +1251,35 @@ def fts_rebuild_admission(db_path): acquired, handle = _acquire_db_flock( lock_path, handle, - _FTS_REBUILD_LOCK_TIMEOUT_SECONDS, + timeout, _FTS_REBUILD_LOCK_POLL_SECONDS, "FTS rebuild lock", ) - if not acquired: + if acquired is None: + # Non-contention failure: already logged with the real errno; + # a "held by another process" line here would be a lie. + acquired = False + elif not acquired: record = None if _IS_WINDOWS else _read_lock_holder_record(handle) - logger.warning( - "FTS rebuild lock %s held by another process for more than " - "%.0fs — deferring this rebuild to avoid racing the holder " - "(the stale-FTS breadcrumb keeps it retryable). " - "Recorded holder: %s.", - lock_path, _FTS_REBUILD_LOCK_TIMEOUT_SECONDS, - _describe_lock_holder(record), - ) + if timeout <= 0: + # Non-blocking probe from an in-process retry: a busy lock + # is expected and will be tried again, so keep it quiet. + logger.info( + "FTS rebuild lock %s is busy — deferring this retry " + "(the stale-FTS breadcrumb keeps it retryable). " + "Recorded holder: %s.", + lock_path, + _describe_lock_holder(record), + ) + else: + logger.warning( + "FTS rebuild lock %s held by another process for more than " + "%.0fs — deferring this rebuild to avoid racing the holder " + "(the stale-FTS breadcrumb keeps it retryable). " + "Recorded holder: %s.", + lock_path, timeout, + _describe_lock_holder(record), + ) yield acquired finally: try: diff --git a/hermes_state_registry.py b/hermes_state_registry.py index 381e2de539..0f3bcedc20 100644 --- a/hermes_state_registry.py +++ b/hermes_state_registry.py @@ -42,7 +42,7 @@ from __future__ import annotations import logging import threading from pathlib import Path -from typing import TYPE_CHECKING, Dict, Optional, Tuple +from typing import TYPE_CHECKING, Dict, List, Optional, Tuple if TYPE_CHECKING: # pragma: no cover - import cycle guard, typed only from hermes_state import SessionDB @@ -289,6 +289,18 @@ def close_all() -> int: return closed +def live_shared_session_dbs() -> List["SessionDB"]: + """Snapshot of every live (non-retired) shared SessionDB in this process. + + For periodic in-process maintenance (the gateway housekeeping tick's + deferred-FTS retry). Refcounts are NOT touched: the caller only invokes + a method on an instance that some holder already keeps alive; a + concurrent final release closes it and the callee sees ``_conn is None``. + """ + with _lock: + return [g.db for g in _generations.values() if not g.retired] + + def stats() -> Dict[str, int]: """Registry census for tests and diagnostics (no locks held long).""" with _lock: diff --git a/hermes_state_schema.py b/hermes_state_schema.py index 9813b12785..f0f8a9197c 100644 --- a/hermes_state_schema.py +++ b/hermes_state_schema.py @@ -43,6 +43,15 @@ logger = logging.getLogger("hermes_state") _FTS_HOLDER_ESCALATE_ATTEMPTS = 3 _FTS_HOLDER_ESCALATE_SECONDS = 60.0 +# Minimum spacing between in-process retries of a deferred stale-FTS rebuild +# (``retry_deferred_fts_recovery``). The startup open already paid the full +# admission wait once; later retries are non-blocking probes on this cadence +# so a live holder never stalls a long-lived writer. +_FTS_STALE_RETRY_SECONDS = 60.0 +# Each failed retry doubles the spacing up to this cap, so a holder that never +# goes away (a second long-lived writer) costs one deferral warning per hour, +# not one per minute. A successful rebuild clears the stale state entirely. +_FTS_STALE_RETRY_MAX_SECONDS = 3600.0 # Cache for schema_read_probe_statements() — parsing SCHEMA_SQL spins up an # in-memory SQLite database, so derive the statements once per process. @@ -422,8 +431,14 @@ class SessionSchemaMixin: ) return None - def _recover_stale_fts(self, cursor: sqlite3.Cursor, *, legacy: bool) -> bool: - """Atomically rebuild stale base/trigram indexes and resume syncing.""" + def _recover_stale_fts( + self, cursor: sqlite3.Cursor, *, legacy: bool, timeout_seconds=None + ) -> bool: + """Atomically rebuild stale base/trigram indexes and resume syncing. + + *timeout_seconds* bounds the cross-process admission wait; None uses + the full startup budget, ``0`` is the non-blocking in-process retry. + """ foreign_holders = self._foreign_state_db_holders() if foreign_holders: now = time.time() @@ -502,7 +517,9 @@ class SessionSchemaMixin: # authority (fail closed). Losing the race means another process is # already performing this exact recovery; the stale breadcrumb stays # set, so this process simply keeps FTS detached and retries later. - with fts_rebuild_admission(getattr(self, "db_path", None)) as admitted: + with fts_rebuild_admission( + getattr(self, "db_path", None), timeout_seconds=timeout_seconds + ) as admitted: if not admitted: logger.warning( "Deferred stale state.db FTS rebuild: another process " @@ -512,6 +529,65 @@ class SessionSchemaMixin: return False return self._recover_stale_fts_locked(cursor, legacy=legacy) + def retry_deferred_fts_recovery(self) -> bool: + """Retry a deferred stale-FTS rebuild on this open SessionDB. + + ``_recover_stale_fts`` runs at open and fails closed when foreign + holders or the rebuild lock are busy, leaving ``_fts_stale`` set and + search on the LIKE fallback. Live write/search paths must never start + a full rebuild (#97940), so on a short-lived CLI that deferral is + cleared by the next process open — but a gateway opens state.db + once and stays up for days, so "next open" never came (#100108). + This is the in-process retry: bounded backoff from + ``_FTS_STALE_RETRY_SECONDS`` doubling to ``_FTS_STALE_RETRY_MAX_SECONDS``, + non-blocking admission (``timeout=0``) so a live holder is skipped and + tried again later, no new thread — the caller is an existing periodic + tick (gateway housekeeping). + + Returns True only when the index was rebuilt and sync triggers + restored. Never raises. + """ + if not getattr(self, "_fts_stale", False): + return False + if getattr(self, "read_only", False) or getattr(self, "_conn", None) is None: + return False + now = time.monotonic() + if now < getattr(self, "_fts_stale_retry_after", 0.0): + return False + interval = float( + getattr(self, "_fts_stale_retry_interval", 0.0) + ) or _FTS_STALE_RETRY_SECONDS + self._fts_stale_retry_after = now + interval + self._fts_stale_retry_interval = min( + interval * 2.0, _FTS_STALE_RETRY_MAX_SECONDS + ) + try: + with self._lock: + if self._conn is None or not self._fts_stale: + return False + cursor = self._conn.cursor() + legacy = self._db_has_legacy_inline_fts(cursor) + recovered = self._recover_stale_fts( + cursor, legacy=legacy, timeout_seconds=0.0 + ) + if recovered: + # CJK was detached alongside the base indexes; its own + # ensure path decides when it comes back online. + self._ensure_fts_cjk_schema(cursor) + self._fts_stale_retry_interval = 0.0 + try: + self._conn.commit() + except sqlite3.Error: + pass + return recovered + except Exception: # noqa: BLE001 - background retry must never raise + logger.warning( + "In-process retry of the deferred stale state.db FTS rebuild " + "failed; will retry later.", + exc_info=True, + ) + return False + def _recover_stale_fts_locked( self, cursor: sqlite3.Cursor, *, legacy: bool ) -> bool: @@ -1510,7 +1586,8 @@ class SessionSchemaMixin: breadcrumb is persisted, mirroring ``_enter_fts_fail_open``'s ordering contract: triggers must never be live over an index with an unrebuilt gap. FTS stays detached for this instance; the winner's - rebuild — or ``_recover_stale_fts`` at the next startup — restores + rebuild — or ``retry_deferred_fts_recovery`` from the gateway + housekeeping tick, or ``_recover_stale_fts`` at the next startup — restores the index and triggers atomically. """ with fts_rebuild_admission(getattr(self, "db_path", None)) as admitted: diff --git a/hermes_state_search.py b/hermes_state_search.py index 40fddb70fd..3dafeebc4a 100644 --- a/hermes_state_search.py +++ b/hermes_state_search.py @@ -2382,7 +2382,9 @@ class SessionSearchMixin: FAILS CLOSED: if another process holds the rebuild lock beyond the bounded wait, this call defers (returns 0) rather than racing it. Callers already treat 0 as "rebuild made no progress" and fall back - to the stale-FTS breadcrumb path, which retries at next startup. + to the stale-FTS breadcrumb path, which retries in-process from the + gateway housekeeping tick (``retry_deferred_fts_recovery``) and at + next startup. Safe to call when FTS tables don't exist (skips them). Returns the number of FTS indexes that were rebuilt.