fix(gateway): reopen a session DB handle the registry tore down

Profile stop, restart and delete release the profile's handles through
hermes_state_registry.close_all_under(), which tears the generation
down. The gateway's per-path handle caches (the session store's, and
the runner's async wrapper around it) kept serving that dead SessionDB.
Its close-race self-heal then reopened a writer the registry does not
know about: the agent's next acquire() got a second writer on the same
file, and later close_all_under() calls could not release it. After a
delete and same-name recreate, every call raised StateDbReplacedError,
and the leaked handle kept the old WAL open, so other processes refused
to open the recreated store.

Remember which cached handles the registry owned, and when the registry
has since torn one down, drop it and reopen through the registry.
This commit is contained in:
Totoro-qaq
2026-09-26 14:03:14 +08:00
committed by Teknium
parent 8c0f24b1ce
commit d46ea7bf20
3 changed files with 73 additions and 5 deletions

View File

@@ -3836,10 +3836,10 @@ class GatewayRunner(
# The store owns/sweeps it at shutdown; this cache holds only the async wrapper (close_all).
# Both caches resolve the SAME ``_default_db_path()``, so the process was holding two writer
# connections and two read pools against one state.db — the fd budget doubled for nothing, and
# doubled again per profile on a multiplexed gateway (#98573). A borrowed wrapper cannot go
# stale in practice: the store's cache only drops handles in close_all_db_handles() (shutdown),
# and while the store's own open is failing there is nothing to borrow, so nothing is cached
# here either.
# doubled again per profile on a multiplexed gateway (#98573). A borrowed wrapper goes stale only
# when the registry tears its generation down (profile unserve/delete); both caches then drop
# the dead handle and reopen through the registry. While the store's own open is failing there
# is nothing to borrow, so nothing is cached here either.
store = getattr(self, "session_store", None)
borrowed = getattr(store, "_db", None) if store is not None else None
if borrowed is not None:

View File

@@ -43,6 +43,14 @@ def _publish_health(source: _HealthSource, path: Path, state: str) -> None:
pass # Runtime health is diagnostic only; persistence must not depend on it.
def _registry_owned(handle: Any) -> bool:
"""True while the process-wide registry owns *handle*; an ``AsyncSessionDB`` (the only handle
type with a ``_db``) is judged by the SessionDB it wraps. ``close_all`` / ``close_all_under``
clear the flag when they tear the generation down."""
inner = getattr(handle, "_db", handle)
return getattr(inner, "_shared_registry_owned", None) is True
class RecoverableHandleCache:
"""Cache handles by path while allowing failed opens to heal in-process.
@@ -62,6 +70,9 @@ class RecoverableHandleCache:
self._initial_retry_delay = max(0.0, float(initial_retry_delay))
self._max_retry_delay = max(self._initial_retry_delay, float(max_retry_delay))
self._unavailable: dict[Path, _Unavailable] = {}
# Paths whose cached handle the registry owned when it was cached: once the registry tears
# that generation down (profile unserve/delete), the entry is dead and must be reopened.
self._registry_backed: set[Path] = set()
self._health_source = _HealthSource()
self._generation = 0
self._close_rejected: Callable[[Any], None] | None = None
@@ -81,7 +92,14 @@ class RecoverableHandleCache:
path = Path(path)
with self.lock:
if path in self.handles:
return self.handles[path]
handle = self.handles[path]
if path not in self._registry_backed or _registry_owned(handle):
return handle
# The registry force-closed this generation. Serving it would let its self-heal reopen
# a writer the registry cannot see (a second writer beside the next ``acquire``), or keep
# raising StateDbReplacedError after a delete + recreate: reopen through the registry.
del self.handles[path]
self._registry_backed.discard(path)
unavailable = self._unavailable.setdefault(path, _Unavailable())
if unavailable.in_flight or self._clock() < unavailable.next_retry_at:
return None
@@ -116,6 +134,10 @@ class RecoverableHandleCache:
stale = self._is_stale(path, unavailable, generation)
if not stale:
self.handles[path] = handle
if _registry_owned(handle):
self._registry_backed.add(path)
else:
self._registry_backed.discard(path)
self._unavailable.pop(path, None)
close_rejected = self._close_rejected if stale else None
if stale:
@@ -136,6 +158,7 @@ class RecoverableHandleCache:
handles = list(self.handles.values())
paths = set(self.handles) | set(self._unavailable)
self.handles.clear()
self._registry_backed.clear()
self._unavailable.clear()
for handle in handles:
with contextlib.suppress(Exception):

View File

@@ -156,3 +156,48 @@ def test_unavailable_store_handle_does_not_resurrect_a_second_open(store):
assert runner._session_db_handles == {}, (
"a duplicate handle was cached on the store's failure path"
)
def test_a_handle_the_registry_tore_down_is_reopened_through_the_registry(store, home):
"""Profile unserve and delete force-close the profile's generation (``close_all_under``).
The store's cache kept serving that dead object. Its self-heal then reopened a writer the
registry does not know about, so the agent's ``acquire`` got a second writer on the same file,
and after a delete + recreate every call raised ``StateDbReplacedError`` until restart.
"""
import hermes_state_registry as registry
profile = home / "profiles" / "work"
profile.mkdir(parents=True)
(profile / "config.yaml").write_text("{}\n", encoding="utf-8")
path = profile / "state.db"
first = store._open_session_db_for_active_scope(db_path=path)
first.create_session("before-unserve", source="telegram")
assert registry.close_all_under(profile) == 1
second = store._open_session_db_for_active_scope(db_path=path)
assert second is not first
second.create_session("after-unserve", source="telegram")
assert first._conn is None, "the torn-down handle was revived outside the registry"
acquired = registry.acquire(path)
try:
assert acquired is second, "one file, one writer: the agent must share the store's handle"
finally:
registry.release(acquired)
def test_the_runner_wrapper_follows_the_reopened_store_handle(store, home):
"""The runner caches an async wrapper around the store's handle; a torn-down inner handle must
not keep being served through it."""
import hermes_state_registry as registry
runner = _runner_with(store)
first = runner._open_session_db_for_active_scope()
assert first is not None and first._db is store._db
registry.close_all_under(home)
second = runner._open_session_db_for_active_scope()
assert second is not first
assert second._db is store._db and second._db is not first._db