fix(gateway): resolve session storage from the key's profile, not ambient scope

#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/<name>/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 <noreply@anthropic.com>
This commit is contained in:
caya8205-2
2026-08-29 00:03:22 +07:00
committed by Teknium
parent f20bbfa40d
commit 5ffaed6e45
3 changed files with 295 additions and 47 deletions

View File

@@ -0,0 +1 @@
caya8205-2

View File

@@ -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/<name>/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,

View File

@@ -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/<name>/``. 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:<profile>:...``)."""
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/<name>/``. 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