From 5ffaed6e45bfbb31fcded6c692165ee9e11f444e Mon Sep 17 00:00:00 2001 From: caya8205-2 Date: Sat, 29 Aug 2026 00:03:22 +0700 Subject: [PATCH] fix(gateway): resolve session storage from the key's profile, not ambient scope MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #88734 made SessionStore._db follow the ambient HERMES_HOME so a multiplexed profile's rows reach its own state.db. That is correct for the inbound message path, which installs the scope via _profile_runtime_scope. Nothing else does. _session_expiry_watcher (gateway/run.py) walks the single process-wide _entries dict — every profile's keys — and finalizes expired sessions with no scope installed, so _db resolved the ROOT store for rows that live under profiles//state.db. The scoped inbound path and the unscoped background path then maintained two copies of the same logical session whose end_reason drifted apart independently. Once they disagreed, the #54878 stale-routing guard read one copy while the routing index pointed at the other, and a live conversation was dropped and recreated — silently, since that branch only sets was_auto_reset when a reset policy also fired. Field evidence from a live two-profile install: session 20260814_234313 was end_reason=None in the root store but agent_close in the profile store, while 20260822_225807 was inverted. Both directions, which rules out a single mis-scoped writer. The owning profile is already encoded in the session key, so derive the store from it: _profile_home_for_key / _db_for_key, plus _db_for_session_id for the entry points addressed by session id. 40 self._db uses across 14 methods now resolve that way. No signature changed and no existing test was modified. _profile_home_for_key returns None when multiplexing is off, when the key carries the legacy agent:main namespace, or when the profile has no live directory, so single-profile installs resolve exactly where they always did. The explicit-path branch still goes through SessionDB.__init__ -> _ensure_test_isolation, keeping the live-DB guard over per-profile paths. Part of #66887. The routing-index half — _routing_scope() and the sessions.json mirror still pinned to one frozen sessions_dir while the handle moves — is left for a follow-up rather than mixed in here. Co-Authored-By: Claude Opus 5 --- contributors/emails/yaandere200@gmail.com | 1 + gateway/session.py | 192 +++++++++++++----- ...test_multiplex_session_db_profile_scope.py | 149 ++++++++++++++ 3 files changed, 295 insertions(+), 47 deletions(-) create mode 100644 contributors/emails/yaandere200@gmail.com diff --git a/contributors/emails/yaandere200@gmail.com b/contributors/emails/yaandere200@gmail.com new file mode 100644 index 0000000000..c0c6e1bb72 --- /dev/null +++ b/contributors/emails/yaandere200@gmail.com @@ -0,0 +1 @@ +caya8205-2 diff --git a/gateway/session.py b/gateway/session.py index facdb6f4d8..9fa3ed15c7 100644 --- a/gateway/session.py +++ b/gateway/session.py @@ -1319,6 +1319,10 @@ class SessionStore: self._db_pinned = _DB_UNPINNED self._db_handles: Dict[Path, Any] = {} self._db_handles_lock = threading.Lock() + # profile name -> its HERMES_HOME (or None to use the ambient scope). + # Memoized so the per-key store lookup stays a dict hit instead of a + # profile-directory stat on every transcript append. + self._profile_home_cache: Dict[str, Optional[Path]] = {} from gateway.session_db_recovery import RecoverableHandleCache self._db_handle_cache = RecoverableHandleCache( @@ -1327,9 +1331,14 @@ class SessionStore: ) self._open_session_db_for_active_scope() - def _open_session_db_for_active_scope(self): + def _open_session_db_for_active_scope(self, db_path: Optional[Path] = None): """Return the SessionDB for the profile scope active on this task. + ``db_path`` pins the store explicitly instead of consulting the + ambient scope. ``_db_for_key`` uses it so work running outside + ``_profile_runtime_scope`` still reaches the profile that owns the + row it is about to touch. + ``SessionDB(db_path=None)`` resolves ``_default_db_path()`` at call time, and that helper follows the context-local HERMES_HOME override installed by ``_profile_runtime_scope``. Resolving here rather than @@ -1345,10 +1354,10 @@ class SessionStore: """ from hermes_state import SessionDB, _default_db_path - path = Path(_default_db_path()) + path = Path(db_path) if db_path is not None else Path(_default_db_path()) def _open(): try: - return SessionDB() + return SessionDB(db_path=path) if db_path is not None else SessionDB() except RuntimeError as e: if "live-system guard" in str(e): # Test-isolation guard fired: a pytest-context process @@ -1389,6 +1398,86 @@ class SessionStore: def _db(self, value) -> None: self._db_pinned = value + def _profile_home_for_key(self, session_key: Optional[str]) -> Optional[Path]: + """HERMES_HOME of the profile that owns *session_key*, or None. + + None means "resolve exactly the way we always have": multiplexing is + off, the key carries the legacy ``agent:main`` namespace, or the + profile has no live directory. A single-profile gateway therefore + never reaches a different store than before this helper existed. + """ + if not getattr(self.config, "multiplex_profiles", False): + return None + profile = self._profile_from_session_key(session_key) + if not profile or profile == "default": + return None + cache = self._profile_home_cache + if profile in cache: + return cache[profile] + home: Optional[Path] = None + try: + from hermes_cli.profiles import get_profile_dir, profile_exists + + if profile_exists(profile): + home = Path(get_profile_dir(profile)) + except Exception as exc: + logger.debug( + "Could not resolve profile home for %r: %s", session_key, exc + ) + home = None + cache[profile] = home + return home + + def _db_for_key(self, session_key: Optional[str]): + """The SessionDB holding *session_key*'s rows, whatever scope is active. + + ``_db`` follows the ambient HERMES_HOME, and only the inbound message + path installs one (``_profile_runtime_scope``). Background work runs + unscoped while operating on every profile's keys out of the single + process-wide ``_entries`` dict — ``_session_expiry_watcher`` is the + clearest case — so it reads and writes the ROOT store for rows that + actually live under ``profiles//state.db``. The two writers + then drift apart on the same logical session until the routing index + disagrees with the row and the #54878 self-heal drops a live + conversation (#66887). + + The owning profile is already encoded in the key, so deriving the + store from it makes every caller agree on one file per session + without threading scope through each call site. + """ + if self._db_pinned is not _DB_UNPINNED: + return self._db_pinned + home = self._profile_home_for_key(session_key) + if home is None: + return self._db + try: + return self._open_session_db_for_active_scope(db_path=home / "state.db") + except Exception: + # Same contract as ``_db``: a failed open degrades to the JSONL + # fallback rather than taking routing down. + return None + + def _db_for_session_id(self, session_id: Optional[str]): + """The SessionDB holding *session_id*'s row. + + Transcript, compression and rewind entry points are addressed by + session id rather than routing key, so recover the owning profile + from the in-memory index. Deliberately lock-free: several of those + callers already hold ``_lock``, and a miss simply falls back to the + ambient store — the behavior that predates ``_db_for_key``. + """ + if not session_id: + return self._db + key = None + try: + for entry in list(self._entries.values()): + if entry.session_id == session_id: + key = entry.session_key + break + except Exception: + key = None + return self._db_for_key(key) + def close_all_db_handles(self) -> None: """Close every SessionDB handle this store opened, one per resolved path. @@ -2142,9 +2231,9 @@ class SessionStore: another team's session. The caller performs one explicit exact lookup of the old unscoped key instead. """ - if not self._db: + if not self._db_for_key(session_key): return None - finder = getattr(self._db, "find_latest_gateway_session_for_peer", None) + finder = getattr(self._db_for_key(session_key), "find_latest_gateway_session_for_peer", None) if not callable(finder): return None try: @@ -2226,11 +2315,11 @@ class SessionStore: reset_reason = self._should_reset(entry, source) if reset_reason: try: - promote = getattr(self._db, "promote_to_session_reset", None) + promote = getattr(self._db_for_key(session_key), "promote_to_session_reset", None) if callable(promote): promote(entry.session_id, reset_reason) else: - self._db.end_session(entry.session_id, reset_reason) + self._db_for_key(session_key).end_session(entry.session_id, reset_reason) except Exception as exc: logger.debug( "Gateway recovered-session reset promotion failed for %s: %s", @@ -2239,7 +2328,7 @@ class SessionStore: ) return None try: - self._db.reopen_session(entry.session_id) + self._db_for_key(session_key).reopen_session(entry.session_id) except Exception as exc: logger.debug("Gateway session DB reopen failed for %s: %s", session_key, exc) if migrated_legacy: @@ -2317,9 +2406,9 @@ class SessionStore: include_compression_ancestors: bool = False, ) -> None: """Persist the routing peer for an existing gateway session row.""" - if not self._db or not source: + if not self._db_for_key(session_key) or not source: return - recorder = getattr(self._db, "record_gateway_session_peer", None) + recorder = getattr(self._db_for_key(session_key), "record_gateway_session_peer", None) if not callable(recorder): return try: @@ -2377,8 +2466,12 @@ class SessionStore: # rehydrate it after the in-memory override was popped. entry.model_override = None self._save() - if self._db: - setter = getattr(self._db, "set_expiry_finalized", None) + # The expiry watcher calls this from a background task that never + # entered ``_profile_runtime_scope``, so resolve the store from the + # key rather than from the ambient scope (#66887). + _db = self._db_for_key(entry.session_key) + if _db: + setter = getattr(_db, "set_expiry_finalized", None) if callable(setter): try: setter(entry.session_id, True) @@ -2397,7 +2490,7 @@ class SessionStore: # live rows or rows ended with ``agent_close``. Explicit # boundaries (compression, session_reset, new_command, etc.) # are preserved — the first writer wins. - self._db.promote_to_session_reset(entry.session_id) + _db.promote_to_session_reset(entry.session_id) except Exception as exc: logger.debug( "Session DB promote_to_session_reset failed for %s: %s", @@ -2492,8 +2585,13 @@ class SessionStore: live routing key and silently swallow every subsequent message until the next restart (#54878 — the live-gateway variant of #52804/FM9). DB errors are non-fatal — never block routing on a failed lookup. + + The store is resolved from the row's owning profile rather than the + ambient scope: an unscoped background writer keeps its own copy of + the same session, and comparing against that copy reports a live + session as ended (#66887). """ - db = getattr(self, "_db", None) + db = self._db_for_session_id(session_id) if not db or not session_id: return False try: @@ -2557,10 +2655,10 @@ class SessionStore: mapping pointing at the compressed parent. Heal that on read so the next inbound message resumes the child instead of reloading the parent. """ - if not session_id or self._db is None: + if not session_id or self._db_for_session_id(session_id) is None: return session_id try: - return self._db.get_compression_tip(session_id) or session_id + return self._db_for_session_id(session_id).get_compression_tip(session_id) or session_id except Exception: logger.debug( "Compression-tip lookup failed for session %s", @@ -2891,7 +2989,7 @@ class SessionStore: prev_session_id = recovered.session_id else: try: - self._db.reopen_session(recovered.session_id) + self._db_for_key(session_key).reopen_session(recovered.session_id) except Exception as exc: logger.debug( "Gateway session DB reopen failed for %s: %s", @@ -2971,7 +3069,7 @@ class SessionStore: self._save_entries() # SQLite operations outside the lock (unchanged). - if self._db and db_end_session_id: + if self._db_for_key(session_key) and db_end_session_id: # Use the specific reset reason so state.db is auditable (e.g. # "resume_pending_expired" is distinguishable from a normal # "session_reset" caused by idle/daily expiry). @@ -2982,11 +3080,11 @@ class SessionStore: # (agent_close / ws_orphan_reap), which first-reason-wins # end_session would preserve — leaving the reset session # resurrectable by stale-route recovery (#61220, #61993). - _promote = getattr(self._db, "promote_to_session_reset", None) + _promote = getattr(self._db_for_key(session_key), "promote_to_session_reset", None) if callable(_promote): _promote(db_end_session_id, _db_end_reason) else: - self._db.end_session(db_end_session_id, _db_end_reason) + self._db_for_key(session_key).end_session(db_end_session_id, _db_end_reason) except Exception as e: # A failed end-write leaves a zombie open row still holding # this chat's session_key: restart recovery will resolve the @@ -2999,9 +3097,9 @@ class SessionStore: db_end_session_id, session_key, e, ) - if self._db and db_create_kwargs: + if self._db_for_key(session_key) and db_create_kwargs: try: - self._db.create_session(**db_create_kwargs) + self._db_for_key(session_key).create_session(**db_create_kwargs) self._record_gateway_session_peer( session_id, session_key, @@ -3471,17 +3569,17 @@ class SessionStore: "model_config": {"_reset_from": db_end_session_id}, } - if self._db and db_end_session_id: + if self._db_for_key(session_key) and db_end_session_id: try: # Promote (not plain end_session): an accidental # agent_close/ws_orphan_reap end must not survive an explicit # user reset, or recovery resurrects the reset session # (#61993 — the user's /new was silently undone). - _promote = getattr(self._db, "promote_to_session_reset", None) + _promote = getattr(self._db_for_key(session_key), "promote_to_session_reset", None) if callable(_promote): _promote(db_end_session_id, "session_reset") else: - self._db.end_session(db_end_session_id, "session_reset") + self._db_for_key(session_key).end_session(db_end_session_id, "session_reset") except Exception as e: # Zombie hazard — see the get_or_create twin path (#82616). logger.warning( @@ -3491,9 +3589,9 @@ class SessionStore: db_end_session_id, session_key, e, ) - if self._db and db_create_kwargs: + if self._db_for_key(session_key) and db_create_kwargs: try: - self._db.create_session(**db_create_kwargs) + self._db_for_key(session_key).create_session(**db_create_kwargs) self._record_gateway_session_peer( session_id, session_key, @@ -3590,23 +3688,23 @@ class SessionStore: self._entries[session_key] = new_entry self._save() - if self._db and db_end_session_id: + if self._db_for_key(session_key) and db_end_session_id: try: # Promote (not plain end_session): a stale agent_close / # ws_orphan_reap end on the outgoing session must be upgraded # to the explicit switch boundary, or recovery can resurrect # it over the user's /resume choice (#61220 bug class). - _promote = getattr(self._db, "promote_to_session_reset", None) + _promote = getattr(self._db_for_key(session_key), "promote_to_session_reset", None) if callable(_promote): _promote(db_end_session_id, "session_switch") else: - self._db.end_session(db_end_session_id, "session_switch") + self._db_for_key(session_key).end_session(db_end_session_id, "session_switch") except Exception as e: logger.debug("Session DB end_session failed: %s", e) - if self._db: + if self._db_for_key(session_key): try: - self._db.reopen_session(target_session_id) + self._db_for_key(session_key).reopen_session(target_session_id) except Exception as e: logger.debug("Session DB reopen_session failed: %s", e) self._record_gateway_session_peer( @@ -3680,7 +3778,7 @@ class SessionStore: def append_to_transcript(self, session_id: str, message: Dict[str, Any], skip_db: bool = False) -> None: """Serialize transcript draining across queue migration boundaries.""" - if not self._db or skip_db: + if not self._db_for_session_id(session_id) or skip_db: return with self._get_transcript_drain_lock(): reroutes = getattr(self, "_transcript_reroutes", None) @@ -3802,9 +3900,9 @@ class SessionStore: # continuation exists; adopt only a different, still-live # tip, otherwise fail closed as before. child_id = "" - tip = self._db.get_compression_tip(session_id) + tip = self._db_for_session_id(session_id).get_compression_tip(session_id) if tip and tip != session_id: - tip_row = self._db.get_session(tip) + tip_row = self._db_for_session_id(session_id).get_session(tip) if tip_row is not None and tip_row.get("ended_at") is None: child_id = str(tip) if child_id: @@ -3933,7 +4031,7 @@ class SessionStore: def _append_transcript_message(self, session_id: str, message: Dict[str, Any]) -> None: """Write one transcript row. Caller handles retry queuing.""" - self._db.append_message( + self._db_for_session_id(session_id).append_message( session_id=session_id, role=message.get("role", "unknown"), content=message.get("content"), @@ -4048,10 +4146,10 @@ class SessionStore: when no DB is available (in-memory sessions). Used by the gateway's transient-failure dedupe guard (#47237). """ - if not self._db: + if not self._db_for_session_id(session_id): return False try: - return self._db.has_platform_message_id( + return self._db_for_session_id(session_id).has_platform_message_id( session_id, platform_message_id ) except Exception: @@ -4089,11 +4187,11 @@ class SessionStore: own the cross-process turn lease. It leaves internal rewrite policy unchanged for existing callers unless they opt in explicitly. """ - if not self._db: + if not self._db_for_session_id(session_id): return True with self._get_transcript_drain_lock(): try: - self._db.replace_messages( + self._db_for_session_id(session_id).replace_messages( session_id, messages, active_only=active_only, @@ -4119,7 +4217,7 @@ class SessionStore: "vanished" (disk=0) even though every message sat healthy under the child session. """ - if not self._db: + if not self._db_for_session_id(session_id): return [] # Follow the write-side reroute chain (cycle-guarded, same shape as # append_to_transcript). @@ -4131,7 +4229,7 @@ class SessionStore: try: # Durable successor: a compression child published to state.db # survives restart even though the in-memory reroute map doesn't. - tip = self._db.get_compression_tip(session_id) + tip = self._db_for_session_id(session_id).get_compression_tip(session_id) if tip: session_id = tip except Exception: @@ -4141,7 +4239,7 @@ class SessionStore: # user;user wedge (e.g. a turn that persisted no assistant row) # would otherwise re-trigger the pre-request repair on every # request forever — heal it once at the restore boundary. - return self._db.get_messages_as_conversation( + return self._db_for_session_id(session_id).get_messages_as_conversation( session_id, repair_alternation=True ) except Exception as e: @@ -4176,7 +4274,7 @@ class SessionStore: selected current turn must still be a composite carrier, and its live payload must be losslessly replayable as text before anything changes. """ - if not self._db: + if not self._db_for_session_id(session_id): return None with self._get_transcript_drain_lock(): if n < 1: @@ -4188,8 +4286,8 @@ class SessionStore: ) try: - expected_active_ids = self._db.get_active_message_ids(session_id) - durable = self._db.get_messages_as_conversation( + expected_active_ids = self._db_for_session_id(session_id).get_active_message_ids(session_id) + durable = self._db_for_session_id(session_id).get_messages_as_conversation( session_id, include_row_ids=True, ) @@ -4218,7 +4316,7 @@ class SessionStore: # so /retry can explain why the selected carrier is unsafe. target_text = retryable_user_text(target_view.get("content")) try: - result = self._db.rewind_to_message( + result = self._db_for_session_id(session_id).rewind_to_message( session_id, target_id, preserve_compaction_handoff=handoff is not None, diff --git a/tests/gateway/test_multiplex_session_db_profile_scope.py b/tests/gateway/test_multiplex_session_db_profile_scope.py index 7ced892a24..ca6580ef5f 100644 --- a/tests/gateway/test_multiplex_session_db_profile_scope.py +++ b/tests/gateway/test_multiplex_session_db_profile_scope.py @@ -13,6 +13,15 @@ These tests pin the handle to the *active* scope rather than to construction time. ``test_write_under_profile_scope_lands_in_profile_store`` is the one that reproduces the report; it fails against the pre-fix code with the session row sitting in the root store. + +The second group covers #66887: scoping the handle to the *active* scope is +only half an answer, because only the inbound path ever installs one. Every +background caller — the expiry watcher above all — walks the single +process-wide ``_entries`` dict, which holds every profile's keys, with no +scope at all, and so resolved the root store for rows living under +``profiles//``. Those tests resolve the store from the profile encoded +in the key instead, and pin that single-profile installs still resolve exactly +where they always did. """ import asyncio @@ -457,3 +466,143 @@ def test_runner_session_db_follows_the_active_profile_scope(multiplex_homes): assert runner._session_db_handles == {} assert root_db._db._conn is None assert profile_db._db._conn is None + + +# --------------------------------------------------------------------------- +# #66887 — the store must follow the key, not whatever scope happens to be on +# --------------------------------------------------------------------------- + + +def _expiry_finalized_flag(db_path: Path, session_id: str): + """Read one session's expiry_finalized flag, or None when the row is absent.""" + if not db_path.exists(): + return None + conn = sqlite3.connect(str(db_path)) + try: + row = conn.execute( + "SELECT expiry_finalized FROM sessions WHERE id = ?", (session_id,) + ).fetchone() + return None if row is None else row[0] + except sqlite3.OperationalError: + return None + finally: + conn.close() + + +def _multiplex_store(root: Path) -> SessionStore: + """A store whose keys carry the profile namespace (``agent::...``).""" + with patch("gateway.session.SessionStore._ensure_loaded"): + store = SessionStore( + sessions_dir=root / "sessions", + config=GatewayConfig(multiplex_profiles=True), + ) + store._loaded = True + return store + + +def _profile_source() -> SessionSource: + return SessionSource( + platform=Platform.TELEGRAM, chat_id="555", user_id="u1", profile="fitness" + ) + + +def test_scoped_inbound_turn_lands_in_profile_store(multiplex_homes): + """Control for the test below: the scoped path was already correct. + + #88734 fixed the inbound path, which runs inside ``_profile_runtime_scope``. + Pinning it here makes the next test unambiguous — the only difference + between the two is whether a scope is installed. + """ + root, profile = multiplex_homes + store = _multiplex_store(root) + + token = set_hermes_home_override(str(profile)) + try: + entry = store.get_or_create_session(_profile_source()) + finally: + reset_hermes_home_override(token) + + assert entry.session_key.startswith("agent:fitness:") + assert _session_ids(profile / "state.db") == {entry.session_id} + assert _session_ids(root / "state.db") == set() + + +def test_unscoped_background_finalize_reaches_the_key_owner_store(multiplex_homes): + """Background work carries no scope but owns every profile's keys. + + ``_session_expiry_watcher`` walks the process-wide ``_entries`` dict and + finalizes expired sessions without entering ``_profile_runtime_scope``, so + resolving from the ambient home wrote the flag to the ROOT store while the + row lives under ``profiles//``. Two copies of one session then drift + apart until the #54878 guard drops a live conversation. + + Fails before this change with ``expiry_finalized`` still 0 on the profile row. + """ + root, profile = multiplex_homes + store = _multiplex_store(root) + + token = set_hermes_home_override(str(profile)) + try: + entry = store.get_or_create_session(_profile_source()) + finally: + reset_hermes_home_override(token) + + # No scope installed — exactly how the watcher calls this. + store.set_expiry_finalized(entry) + + assert _expiry_finalized_flag(profile / "state.db", entry.session_id) == 1 + assert _session_ids(root / "state.db") == set() + + +def test_unscoped_staleness_check_reads_the_key_owner_store(multiplex_homes): + """The routing guard must consult the row it actually routes to. + + ``_is_session_ended_in_db`` decides whether the #54878 self-heal fires. + Reading the ambient store lets another store's copy answer the question, + which is how a live session gets reported as ended and dropped. + """ + root, profile = multiplex_homes + store = _multiplex_store(root) + + token = set_hermes_home_override(str(profile)) + try: + entry = store.get_or_create_session(_profile_source()) + finally: + reset_hermes_home_override(token) + + # Alive in the profile store, and the root store has never heard of it. + assert store._is_session_ended_in_db(entry.session_id) is False + + token = set_hermes_home_override(str(profile)) + try: + store._db.end_session(entry.session_id, "agent_close") + finally: + reset_hermes_home_override(token) + + assert store._is_session_ended_in_db(entry.session_id) is True + + +def test_default_namespace_keeps_ambient_resolution(multiplex_homes): + """Guardrail: the legacy ``agent:main`` namespace must not change stores. + + Single-profile installs are the overwhelming majority. A key without a + named profile has to resolve exactly where it did before ``_db_for_key`` + existed, or this fix would silently relocate their history. + """ + root, _profile = multiplex_homes + store = _make_store(root) # multiplex off -> agent:main keys + + assert store._profile_home_for_key("agent:main:telegram:dm:1") is None + assert store._db_for_key("agent:main:telegram:dm:1") is store._db + assert store._db_for_key(None) is store._db + + +def test_pinned_handle_still_wins_over_key_resolution(multiplex_homes): + """``store._db = fake`` stays authoritative, as the rest of the suite assumes.""" + root, _profile = multiplex_homes + store = _multiplex_store(root) + + sentinel = object() + store._db = sentinel + assert store._db_for_key("agent:fitness:telegram:dm:1") is sentinel + assert store._db_for_session_id("whatever") is sentinel