fix(state): fail fast on non-contention flock errors and retry deferred FTS rebuilds in-process (salvage #100130)

Two pieces of PR #100130 (@HexLab98) re-applied on top of the orphaned-flock
break (894fc35337) and fail-closed admission (#100895) that landed since:

* `is_advisory_lock_contention` (hermes_state_common): only EAGAIN /
  EWOULDBLOCK / EACCES / EDEADLK mean "another process holds the lock".
  ESTALE / ENOTSUP / ENOLCK / EIO from flock or msvcrt.locking are
  environment failures that polling cannot fix — `_acquire_db_flock` and
  both Windows msvcrt loops (FTS rebuild admission, state.db repair lock)
  now defer immediately with the real errno instead of burning the full
  120s / holder timeout and then logging a fake "held by another process".

* `retry_deferred_fts_recovery` (hermes_state_schema): a SessionDB whose
  open-time `_recover_stale_fts` deferred (foreign holders or busy rebuild
  lock) stayed `_fts_stale` — LIKE-only search — until the process
  reopened state.db. Short-lived CLIs reopen every run; the gateway opens
  once and stays up for days, so the deferral was effectively permanent
  (#100108). The retry runs from the EXISTING gateway housekeeping tick
  (`_start_gateway_housekeeping`, 60s) against the shared SessionDB
  instances via `hermes_state_registry.live_shared_session_dbs()`:
  non-blocking admission (`fts_rebuild_admission(timeout_seconds=0)`),
  bounded backoff 60s -> 1h, no new thread, still fails closed on live
  holders. `fts_rebuild_admission` gains the `timeout_seconds` kwarg.

* WAL-reset warning names `sys.executable` so a "linked SQLite 3.45.1"
  line can be matched to the interpreter that actually linked it
  (#100108 point 3).

Deliberately NOT carried from #100130: the "leftover lock file = holder"
premise (a 0-byte lock file never blocked flock; the real cause was the
fork-inherited fd, fixed in 894fc35337) and the `_rebuild_fts_once`
one-shot rework.

Co-authored-by: HexLab98 <liruixinch@outlook.com>
This commit is contained in:
teknium1
2026-09-02 03:57:36 -07:00
committed by Teknium
parent 238b6c1ab9
commit fd05029430
6 changed files with 231 additions and 25 deletions
+24
View File
@@ -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.
+17 -4
View File
@@ -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 "
+93 -15
View File
@@ -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:
+13 -1
View File
@@ -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:
+81 -4
View File
@@ -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:
+3 -1
View File
@@ -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.