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