test(state): non-contention errno table, repair-lock sibling, in-process deferred-FTS retry via housekeeping tick
Regression coverage for the #100130 salvage, all against real SessionDB files and a real child process holding the flock: * errno table for `is_advisory_lock_contention` (EAGAIN/EWOULDBLOCK/EACCES contend; ESTALE/ENOTSUP/ENOLCK/EIO fail fast); no misleading "held by another process" line on the fast-fail path; `_cross_process_repair_lock` shares the filter (sibling site). * `retry_deferred_fts_recovery`: open under a live holder -> stale; retry returns in <2s with a 30s admission budget (timeout=0); rate limit + 60s->120s backoff engaged; holder dies -> same instance recovers, triggers restored, breadcrumb cleared; no-op when not stale / read-only. * `_start_gateway_housekeeping` tick (real loop, 50ms interval) recovers a stale shared-registry SessionDB with no direct call and no extra thread. Backoff floor: a monkeypatched 0s base interval must not zero the doubled interval (min 1s), so the cap math is testable. Sabotage run (source at origin/main, these tests): 16 failed / 35 passed, including 30s timeouts on the fast-fail tests.
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user