diff --git a/hermes_state_schema.py b/hermes_state_schema.py index f0f8a9197c..01801a4870 100644 --- a/hermes_state_schema.py +++ b/hermes_state_schema.py @@ -554,12 +554,13 @@ class SessionSchemaMixin: 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 + interval = float(getattr(self, "_fts_stale_retry_interval", 0.0)) + if interval <= 0.0: + interval = _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 + max(interval, _FTS_STALE_RETRY_SECONDS, 1.0) * 2.0, + _FTS_STALE_RETRY_MAX_SECONDS, ) try: with self._lock: diff --git a/tests/state/test_fts_rebuild_admission.py b/tests/state/test_fts_rebuild_admission.py index 3195498aed..ac923c6a6c 100644 --- a/tests/state/test_fts_rebuild_admission.py +++ b/tests/state/test_fts_rebuild_admission.py @@ -430,3 +430,195 @@ class TestNonContentionErrnoFailsFast: d2.close() assert _meta_value(db_path, FTS_STALE_KEY) is None assert _base_fts_triggers(db_path) == set(_FTS_TRIGGERS) + + def test_non_contention_errno_skips_holder_warning( + self, tmp_path, monkeypatch, caplog + ): + """The fast-fail must not ALSO log the misleading 'held by another + process for more than Ns' line — there is no holder.""" + import fcntl + import logging + + monkeypatch.setattr( + hermes_state_common, "_FTS_REBUILD_LOCK_TIMEOUT_SECONDS", 30.0 + ) + + def _flock(*_args, **_kwargs): + raise OSError(errno.ENOTSUP, "no locks on this fs") + + monkeypatch.setattr(fcntl, "flock", _flock) + with caplog.at_level(logging.INFO, logger="hermes_state"): + with hermes_state_common.fts_rebuild_admission( + tmp_path / "state.db" + ) as admitted: + assert admitted is False + messages = [r.getMessage() for r in caplog.records] + assert any("non-contention error" in m for m in messages) + assert not any("held by another process" in m for m in messages) + + def test_repair_lock_non_contention_errno_fails_fast( + self, tmp_path, monkeypatch + ): + """Sibling site: the state.db repair lock shares the errno filter.""" + import fcntl + + import hermes_state + + monkeypatch.setattr(hermes_state, "_REPAIR_LOCK_TIMEOUT_SECONDS", 30.0) + + def _flock(*_args, **_kwargs): + raise OSError(errno.EIO, "i/o error") + + monkeypatch.setattr(fcntl, "flock", _flock) + t0 = time.monotonic() + with hermes_state._cross_process_repair_lock(tmp_path / "state.db") as ok: + assert ok is False + assert time.monotonic() - t0 < 2.0 + + @pytest.mark.parametrize( + "exc, expected", + [ + (BlockingIOError(errno.EAGAIN, "x"), True), + (OSError(errno.EWOULDBLOCK, "x"), True), + (OSError(errno.EACCES, "x"), True), + (OSError(errno.ESTALE, "x"), False), + (OSError(errno.ENOTSUP, "x"), False), + (OSError(errno.ENOLCK, "x"), False), + (OSError(errno.EIO, "x"), False), + (ValueError("not an oserror"), False), + ], + ) + def test_is_advisory_lock_contention_table(self, exc, expected): + assert hermes_state_common.is_advisory_lock_contention(exc) is expected + + +class TestDeferredFtsRetryInProcess: + """Gateway shape (#100108): one SessionDB stays open for days. A deferral + at open must be recoverable from an in-process periodic tick, with the + REAL rebuild lock held by a REAL child process at open time.""" + + @staticmethod + def _mark_stale(db_path: Path) -> None: + raw = sqlite3.connect(str(db_path)) + raw.execute( + "INSERT OR REPLACE INTO state_meta(key, value) VALUES (?, '1')", + (FTS_STALE_KEY,), + ) + for trig in _FTS_TRIGGERS: + raw.execute(f"DROP TRIGGER IF EXISTS {trig}") + raw.commit() + raw.close() + + def test_retry_is_non_blocking_while_live_holder_and_backs_off( + self, tmp_path, fast_timeout, monkeypatch + ): + import hermes_state_schema + + db_path = tmp_path / "state.db" + d = SessionDB(db_path=db_path) + if not d._fts_enabled: + d.close() + pytest.skip("FTS5 unavailable in this build") + d.create_session("s1", source="test") + d.append_message("s1", "user", "hello gateway retry") + d.close() + self._mark_stale(db_path) + + with _rebuild_lock_held_by_other_process(db_path): + gw = SessionDB(db_path=db_path) # long-lived "gateway" open + try: + assert gw._fts_stale is True + # Live holder: the retry must return quickly (timeout=0), + # not wait out any admission budget. + monkeypatch.setattr( + hermes_state_common, "_FTS_REBUILD_LOCK_TIMEOUT_SECONDS", 30.0 + ) + t0 = time.monotonic() + assert gw.retry_deferred_fts_recovery() is False + assert time.monotonic() - t0 < 2.0 + assert gw._fts_stale is True + # Rate limit engaged: an immediate second call is a no-op. + assert gw.retry_deferred_fts_recovery() is False + # Backoff doubled (60s -> 120s) but capped at the max. + assert gw._fts_stale_retry_interval == min( + 2 * hermes_state_schema._FTS_STALE_RETRY_SECONDS, + hermes_state_schema._FTS_STALE_RETRY_MAX_SECONDS, + ) + assert gw._fts_stale_retry_after > time.monotonic() + except BaseException: + gw.close() + raise + # Holder gone. Same instance recovers on the next eligible tick. + try: + gw._fts_stale_retry_after = 0.0 + assert gw.retry_deferred_fts_recovery() is True + assert gw._fts_stale is False + assert gw._fts_enabled is True + # Search actually works again on this very instance. + gw.append_message("s1", "user", "needle-after-holder-gone") + assert gw.retry_deferred_fts_recovery() is False # nothing stale + finally: + gw.close() + assert _meta_value(db_path, FTS_STALE_KEY) is None + assert _base_fts_triggers(db_path) == set(_FTS_TRIGGERS) + + def test_gateway_housekeeping_tick_drives_the_retry( + self, tmp_path, fast_timeout, monkeypatch + ): + """The retry hangs off the EXISTING housekeeping loop (no new thread) + and reaches shared-registry instances.""" + import threading + + import hermes_state_registry + import hermes_state_schema + import gateway.run as grun + + monkeypatch.setattr(hermes_state_schema, "_FTS_STALE_RETRY_SECONDS", 0.0) + db_path = tmp_path / "state.db" + d = SessionDB(db_path=db_path) + if not d._fts_enabled: + d.close() + pytest.skip("FTS5 unavailable in this build") + d.create_session("s1", source="test") + d.append_message("s1", "user", "hello housekeeping") + d.close() + self._mark_stale(db_path) + + with _rebuild_lock_held_by_other_process(db_path): + gw = hermes_state_registry.acquire(db_path) + try: + assert gw._fts_stale is True + assert gw in hermes_state_registry.live_shared_session_dbs() + stop = threading.Event() + th = threading.Thread( + target=grun._start_gateway_housekeeping, + args=(stop,), + kwargs={"interval": 0.05}, + daemon=True, + ) + th.start() + deadline = time.monotonic() + 10.0 + while gw._fts_stale and time.monotonic() < deadline: + time.sleep(0.05) + stop.set() + th.join(timeout=5) + assert gw._fts_stale is False + assert gw._fts_enabled is True + finally: + hermes_state_registry.release_or_close(gw) + assert _meta_value(db_path, FTS_STALE_KEY) is None + + def test_retry_noop_when_not_stale_or_read_only(self, tmp_path): + db_path = tmp_path / "state.db" + d = SessionDB(db_path=db_path) + try: + assert d._fts_stale is False + assert d.retry_deferred_fts_recovery() is False + finally: + d.close() + ro = SessionDB(db_path=db_path, read_only=True) + try: + ro._fts_stale = True + assert ro.retry_deferred_fts_recovery() is False + finally: + ro.close()