From d63e5d8a10ecff549ce0c68503a42d98bad56e0e Mon Sep 17 00:00:00 2001 From: joaomarcos Date: Sun, 30 Aug 2026 21:17:37 -0300 Subject: [PATCH] fix(cache): source-qualify the peer identity and gate the declared bind Both blockers from @andrexibiza's review of 28a2d7f0ee. 1. The generation lookup was not in the same identity domain as recovery. latest_conversation_boundary() selected on session_key alone, while _declared_conversation_session() is qualified by (source, session_key). X-Hermes-Session-Key accepts any authenticated caller-supplied string, so an API conversation may legally carry the same key as a Telegram row in one database -- a /new over there rotated this conversation's gwk_ generation while recovery correctly refused to cross the same line, moving the affinity identity out from under a physical identity that had not moved. The boundary read now takes (session_key, source), and the carrier is 'source|key|generation' rather than 'key|generation' -- keying on the string alone would also collapse two same-key conversations from different sources onto one routing key, since this value leaves the process verbatim as OpenRouter's sticky session_id and xAI's x-grok-conv-id. The source comes from the agent's own session row, falling back to the platform the row will be created with before it lands. 2. The declared key's stated lower precedence did not survive settlement. Both handlers let stored_session_id / an explicit body session_id win, then called _bind_declared_conversation() unconditionally. record_gateway_session_peer() does SET session_key = ? across compression ancestors, so a request carrying conversation A's chain plus header key B silently rebound A to B: A could no longer be recovered by its own key, and B recovered A's session. Recording is now gated on the declared key having actually selected or minted the session, on both paths. Behind that gate the bind itself refuses to overwrite a row already bound to a different key, so a future caller cannot reintroduce the same defect by opting in wrongly. test_declaration_outranks_the_lineage_root asserted the pre-qualification contract by comparing a DB-backed agent against a DB-less one; it now makes the stronger statement it was written for -- one declared conversation reached through two different physical ids on the same peer. Refs #96811 Found in review by @andrexibiza, whose analysis located each of these defects and specified what a correct fix had to prove. Co-Authored-By: Andrex Ibiza, MBA --- agent/prompt_cache_scope.py | 35 +++++++-- gateway/platforms/api_server.py | 34 +++++++-- gateway/platforms/api_server_runs.py | 17 +++-- hermes_state.py | 19 +++-- .../agent/test_declared_conversation_scope.py | 57 ++++++++++++++- .../test_api_server_declared_conversation.py | 71 +++++++++++++++++++ 6 files changed, 212 insertions(+), 21 deletions(-) diff --git a/agent/prompt_cache_scope.py b/agent/prompt_cache_scope.py index 53ec690819..d3d464cea7 100644 --- a/agent/prompt_cache_scope.py +++ b/agent/prompt_cache_scope.py @@ -97,7 +97,30 @@ def _lineage_root(session_id: str, session_db: Any) -> Optional[str]: return None -def _conversation_generation(session_key: str, session_db: Any) -> str: +def _agent_source(agent: Any, session_id: str, session_db: Any) -> str: + """The ``sessions.source`` this agent's conversation is recorded under. + + Read from the agent's own row when it exists, because that is the value + the peer queries below match on. Before the row lands (the first turn + resolves a scope ahead of ``_ensure_db_session``) it falls back to the + platform, which is what the row will be created with — the two can differ + only when ``HERMES_SESSION_SOURCE`` overrides the platform, and the cost of + that is one cold bucket on the first turn, never a crossed identity. + """ + if session_id and session_db is not None: + try: + row = session_db.get_session(session_id) + except Exception: + logger.debug("declared-scope source lookup failed", exc_info=True) + row = None + if row: + source = str(row.get("source") or "").strip() + if source: + return source + return str(getattr(agent, "platform", "") or "").strip() + + +def _conversation_generation(session_key: str, source: str, session_db: Any) -> str: """Return the generation marker for *session_key*'s current conversation. The declared key is a per-CHAT identifier and deliberately outlives any @@ -134,7 +157,7 @@ def _conversation_generation(session_key: str, session_db: Any) -> str: reader = getattr(session_db, "latest_conversation_boundary", None) if not callable(reader): return "" - boundary = reader(session_key) + boundary = reader(session_key, source) if boundary is None: return "" # (crossings, ended_at). Fixed-point on the timestamp so the carrier is @@ -181,15 +204,19 @@ def declared_conversation_scope(agent: Any) -> Optional[str]: # fork onto its parent's key on a transient DB failure. logger.debug("declared-scope fork check failed", exc_info=True) return None + source = _agent_source(agent, sid, db) if db is not None: try: - generation = _conversation_generation(key, db) + generation = _conversation_generation(key, source, db) except Exception: # Same fail-closed rule as the fork check: an unqualified key # spans /new, so degrade to the physical-id scope instead. logger.debug("declared-scope generation read failed", exc_info=True) return None - carrier = f"{key}|{generation}" if generation else key + # The carrier is the SAME identity tuple the peer queries use: two hosts + # may legally declare the same key string under different sources, and the + # scope leaves this process as a routing key, so it must not collapse them. + carrier = f"{source}|{key}|{generation}" digest = hashlib.sha256(carrier.encode("utf-8", errors="replace")).hexdigest()[:24] return f"{_DECLARED_SCOPE_PREFIX}{digest}" diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index 11b72344b6..3ece6b96c9 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -2356,6 +2356,21 @@ class APIServerAdapter(BasePlatformAdapter): if db is None: return try: + # Defence in depth behind the callers' precedence gate: never + # rewrite a row that already belongs to a different conversation. + # record_gateway_session_peer does SET session_key = ?, so a + # mistaken bind would strand the original conversation and hand its + # session to this request's key. + existing = db.get_session(sid) or {} + current = str(existing.get("session_key") or "").strip() + if current and current != key: + logger.debug( + "[%s] refusing to rebind session %s from a different " + "declared conversation", + self.name, + sid, + ) + return db.record_gateway_session_peer( sid, source=self._SESSION_SOURCE, @@ -6423,6 +6438,12 @@ class APIServerAdapter(BasePlatformAdapter): # the conversation it declared via ``X-Hermes-Session-Key`` before # minting a throwaway id — otherwise every reply is a new conversation # to every affinity surface (#96811). + # The response chain still outranks the declared key. Recording is + # gated on that same precedence: binding a session the chain selected + # would rewrite ITS routing key to this request's header + # (record_gateway_session_peer does SET session_key = ?), stranding the + # original conversation and letting the header key recover it instead. + _declared_selected = not stored_session_id and bool(gateway_session_key) session_id = ( stored_session_id or self._declared_conversation_session(gateway_session_key) @@ -6498,7 +6519,7 @@ class APIServerAdapter(BasePlatformAdapter): tool_complete_callback=_on_tool_complete, agent_ref=agent_ref, gateway_session_key=gateway_session_key, - bind_declared_conversation=True, + bind_declared_conversation=_declared_selected, **agent_overrides, route=route, )) @@ -6534,7 +6555,7 @@ class APIServerAdapter(BasePlatformAdapter): ephemeral_system_prompt=instructions, session_id=session_id, gateway_session_key=gateway_session_key, - bind_declared_conversation=True, + bind_declared_conversation=_declared_selected, **agent_overrides, route=route, ) @@ -7599,10 +7620,11 @@ class APIServerAdapter(BasePlatformAdapter): # (/v1/responses, /v1/runs) record one, so no other # caller's rows change shape. if bind_declared_conversation: - self._bind_declared_conversation( - getattr(agent, "session_id", None) or session_id, - gateway_session_key, - ) + if _declared_selected: + self._bind_declared_conversation( + getattr(agent, "session_id", None) or session_id, + gateway_session_key, + ) clear_session_vars(tokens) self._activate_admitted_request() diff --git a/gateway/platforms/api_server_runs.py b/gateway/platforms/api_server_runs.py index 121eafef20..04b21fefd3 100644 --- a/gateway/platforms/api_server_runs.py +++ b/gateway/platforms/api_server_runs.py @@ -609,6 +609,10 @@ async def _handle_runs( # ``X-Hermes-Session-Key``. Falling straight through to ``run_id`` # made the run id the conversation identity, so a declared channel # re-keyed every affinity surface once per run (#96811). + # Same precedence gate as /v1/responses: an explicit body session_id + # or a chained session owns its own routing key and must not be + # rebound to this request's header key. + _declared_selected = not session_id and bool(gateway_session_key) session_id = ( session_id or self._declared_conversation_session(gateway_session_key) @@ -840,11 +844,14 @@ async def _handle_runs( _clear_turn_process_ownership(agent) # /v1/runs owns its agent lifecycle, so it records # the declared conversation itself rather than - # through _run_agent's bind_declared_conversation. - self._bind_declared_conversation( - getattr(agent, "session_id", None) or session_id, - gateway_session_key, - ) + # through _run_agent's bind_declared_conversation + # -- carrying the same precedence gate, which an + # explicit body session_id turns off. + if _declared_selected: + self._bind_declared_conversation( + getattr(agent, "session_id", None) or session_id, + gateway_session_key, + ) try: unregister_gateway_notify(approval_session_key) finally: diff --git a/hermes_state.py b/hermes_state.py index b10270b94c..7314ca3d85 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -14135,17 +14135,25 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) return bool(session and self._is_explicit_fork_child_row(session)) def latest_conversation_boundary( - self, session_key: str + self, session_key: str, source: str ) -> Optional[Tuple[int, float]]: - """How many conversation boundaries this key has crossed, and when. + """How many conversation boundaries this peer has crossed, and when. - A boundary is a row this key ended at an intentional conversation + A boundary is a row this peer ended at an intentional conversation break — the ``_RESET_END_REASONS`` set (``/new``, ``/switch``, idle, daily, suspended, resume_pending_expired). That is the same fence :meth:`find_latest_gateway_session_for_peer` refuses to reach behind, so the two agree on where one conversation stops and the next begins and cannot drift. + The peer is ``(session_key, source)``, the SAME identity tuple recovery + uses — never the key alone. ``X-Hermes-Session-Key`` accepts any + authenticated caller-supplied string, so an API conversation may + legally carry the same key as a Telegram row in one database; keying + on the string alone would let a ``/new`` on that unrelated row rotate + this conversation's affinity identity while recovery correctly refuses + to cross the same line. + Returns ``(count, latest_ended_at)``, or ``None`` when the key has never been reset. BOTH halves are reported because each one alone has a narrow way to repeat a previous generation: @@ -14163,7 +14171,7 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) ``agent/prompt_cache_scope.py`` uses it to keep a host-declared conversation key from outliving the conversation it names. """ - if not session_key: + if not session_key or not source: return None with self._read_ctx() as conn: row = conn.execute( @@ -14171,10 +14179,11 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) SELECT COUNT(*) AS crossings, MAX(ended_at) AS boundary FROM sessions WHERE session_key = ? + AND source = ? AND ended_at IS NOT NULL AND end_reason IN ({_RESET_END_REASONS_SQL}) """, - (session_key,), + (session_key, source), ).fetchone() boundary = row["boundary"] if row is not None else None if boundary is None: diff --git a/tests/agent/test_declared_conversation_scope.py b/tests/agent/test_declared_conversation_scope.py index 66ed0fcae3..53326a6c73 100644 --- a/tests/agent/test_declared_conversation_scope.py +++ b/tests/agent/test_declared_conversation_scope.py @@ -130,8 +130,12 @@ class TestDeclaredConversationScope: db.create_session("rotated-1", source="webui", parent_session_id="root-sess") scope = resolve_prompt_cache_scope(_agent("rotated-1", db, CHAT_KEY)) - assert scope == declared_conversation_scope(_agent("x", None, CHAT_KEY)) + # Same declared conversation reached through a different physical id + # on the same peer — the property the lineage walk cannot provide. + db.create_session("rotated-2", source="webui") + assert scope == resolve_prompt_cache_scope(_agent("rotated-2", db, CHAT_KEY)) assert scope != "root-sess" + assert scope.startswith("gwk_") def test_branch_child_ignores_the_shared_chat_key(self, db): """/branch keys off session_id, not the chat key — #79161 isolation.""" @@ -534,3 +538,54 @@ class TestConversationGenerationRotates: third = self._keyed(db, "sess-clock-3") scope_third = resolve_prompt_cache_scope(third) assert len({scope_first, scope_second, scope_third}) == 3 + + +class TestPeerIdentityIsSourceQualified: + """The generation and the carrier use the same identity tuple as recovery. + + ``X-Hermes-Session-Key`` accepts any authenticated caller-supplied string, + so an API conversation may legally carry the same key as a Telegram row in + one database. Keying on the string alone let a ``/new`` on that unrelated + row rotate this conversation's affinity identity, while + ``find_latest_gateway_session_for_peer`` correctly refused to cross the + same line — the physical identity stayed put while the affinity identity + moved under it (@andrexibiza on #98811). + """ + + KEY = "shared-key-string" + + def _row(self, db, sid, source): + db.create_session(session_id=sid, source=source, session_key=self.KEY) + return SimpleNamespace( + session_id=sid, _session_db=db, _gateway_session_key=self.KEY, + platform=source, + ) + + def test_a_foreign_sources_reset_does_not_rotate_this_conversation(self, db): + mine = self._row(db, "api-1", "api_server") + before = resolve_prompt_cache_scope(mine) + + # Same key string, different platform, reset only over there. + self._row(db, "tg-1", "telegram") + db.end_session("tg-1", "session_reset") + + assert resolve_prompt_cache_scope(self._row(db, "api-2", "api_server")) == before + + def test_our_own_reset_still_rotates(self, db): + mine = self._row(db, "api-1", "api_server") + before = resolve_prompt_cache_scope(mine) + db.end_session("api-1", "session_reset") + assert resolve_prompt_cache_scope(self._row(db, "api-2", "api_server")) != before + + def test_equal_keys_under_different_sources_never_share_a_scope(self, db): + mine = resolve_prompt_cache_scope(self._row(db, "api-1", "api_server")) + theirs = resolve_prompt_cache_scope(self._row(db, "tg-1", "telegram")) + assert mine != theirs + assert mine.startswith("gwk_") and theirs.startswith("gwk_") + + def test_the_boundary_read_is_peer_scoped(self, db): + db.create_session(session_id="tg-1", source="telegram", session_key=self.KEY) + db.end_session("tg-1", "session_reset") + assert db.latest_conversation_boundary(self.KEY, "telegram") is not None + assert db.latest_conversation_boundary(self.KEY, "api_server") is None + assert db.latest_conversation_boundary(self.KEY, "") is None diff --git a/tests/gateway/test_api_server_declared_conversation.py b/tests/gateway/test_api_server_declared_conversation.py index 0cbcbc27be..e901f9ad61 100644 --- a/tests/gateway/test_api_server_declared_conversation.py +++ b/tests/gateway/test_api_server_declared_conversation.py @@ -293,3 +293,74 @@ class TestRunAgentOptIn: getattr(agent, "session_id", None) or "sess-initial", KEY ) assert calls == [("sess-initial", KEY)] + + +class TestBindFollowsPrecedence: + """Recording is gated on the declared key actually selecting the session. + + `record_gateway_session_peer` performs `SET session_key = ?`, so binding a + session that the response chain or an explicit body id selected would + rewrite THAT conversation's routing key to this request's header: the + original conversation could no longer be recovered by its own key, and the + header key would recover it instead (@andrexibiza on #98811). + """ + + @staticmethod + def _responses_gate(*, stored, key): + # gateway/platforms/api_server.py::_handle_responses + return not stored and bool(key) + + @staticmethod + def _runs_gate(*, body_id, key): + # gateway/platforms/api_server.py::_handle_runs + return not body_id and bool(key) + + def test_chained_session_does_not_record_the_header_key(self, adapter): + assert self._responses_gate(stored="sess-chained", key=KEY) is False + + def test_explicit_body_session_does_not_record_the_header_key(self, adapter): + assert self._runs_gate(body_id="explicit", key=KEY) is False + + def test_declared_or_minted_session_records(self, adapter): + assert self._responses_gate(stored=None, key=KEY) is True + assert self._runs_gate(body_id=None, key=KEY) is True + + def test_undeclared_request_records_nothing(self, adapter): + assert self._responses_gate(stored=None, key=None) is False + assert self._runs_gate(body_id=None, key=None) is False + + def test_a_foreign_binding_is_never_overwritten(self, adapter): + """Defence in depth behind the gate, at the DB layer.""" + a, db = adapter + _seed(db, "sess-A", key=KEY) + + a._bind_declared_conversation("sess-A", OTHER_KEY) + + assert db.get_session("sess-A")["session_key"] == KEY + assert a._declared_conversation_session(KEY) == "sess-A" + assert a._declared_conversation_session(OTHER_KEY) is None + + def test_rebinding_the_same_key_is_idempotent(self, adapter): + a, db = adapter + _seed(db, "sess-A", key=KEY) + a._bind_declared_conversation("sess-A", KEY) + assert db.get_session("sess-A")["session_key"] == KEY + assert a._declared_conversation_session(KEY) == "sess-A" + + def test_an_unbound_row_still_binds(self, adapter): + a, db = adapter + db.create_session(session_id="sess-new", source=SOURCE, model="m") + a._bind_declared_conversation("sess-new", KEY) + assert a._declared_conversation_session(KEY) == "sess-new" + + def test_the_original_conversation_stays_recoverable(self, adapter): + """The end-to-end shape of the defect: A must survive a B-keyed turn.""" + a, db = adapter + _seed(db, "sess-A", key=KEY) + # A request carrying A's chain plus header key B: the gate refuses to + # record, and the DB guard refuses even if something else tried. + assert self._responses_gate(stored="sess-A", key=OTHER_KEY) is False + a._bind_declared_conversation("sess-A", OTHER_KEY) + + assert a._declared_conversation_session(KEY) == "sess-A" + assert a._declared_conversation_session(OTHER_KEY) is None