From d46ea7bf20811f200ded1abfb231e4076f55232e Mon Sep 17 00:00:00 2001 From: Totoro-qaq Date: Sat, 26 Sep 2026 14:03:14 +0800 Subject: [PATCH] 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. --- gateway/run.py | 8 ++-- gateway/session_db_recovery.py | 25 ++++++++++- .../gateway/test_session_db_handle_sharing.py | 45 +++++++++++++++++++ 3 files changed, 73 insertions(+), 5 deletions(-) diff --git a/gateway/run.py b/gateway/run.py index 40cee7b550..ded6ba3a98 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -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: diff --git a/gateway/session_db_recovery.py b/gateway/session_db_recovery.py index 46e315a03b..9360e1b6e5 100644 --- a/gateway/session_db_recovery.py +++ b/gateway/session_db_recovery.py @@ -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): diff --git a/tests/gateway/test_session_db_handle_sharing.py b/tests/gateway/test_session_db_handle_sharing.py index d874c882cf..59c5cb1b19 100644 --- a/tests/gateway/test_session_db_handle_sharing.py +++ b/tests/gateway/test_session_db_handle_sharing.py @@ -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