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