fix(api): resolve the declared conversation instead of minting a session per request
POST /v1/responses and POST /v1/runs parse and authenticate the client's X-Hermes-Session-Key, pass it downstream for memory scoping, and then mint a throwaway physical session id anyway whenever the client manages its own history (no previous_response_id chain to carry one forward). Every conversation-affinity hint Hermes sends is derived from that physical id, so all four re-keyed on every single reply: prompt_cache_key on both OpenAI-wire transports, the OpenRouter and Nous sticky session_id, and xAI's x-grok-conv-id. The conversation never landed back on a warm prefix. Fix the identity rather than the four consumers. The declared key resolves to its live session through find_latest_gateway_session_for_peer -- the same reset-fenced recovery every native gateway platform already uses -- and the turn records the row it ended on through record_gateway_session_peer, which AIAgent._ensure_db_session never did (it knows the key and writes the row unkeyed, so the mapping the next reply needs did not exist). Because the lookup is fenced on sessions.end_reason, the generation that must rotate is already durable: session_reset (/new), session_switch, idle, daily, suspended and resume_pending_expired all return None, so a new conversation gets a new id and a cold affinity scope, and a retired generation can never be resolved again. No counter, no new persisted field, and no new precedence rule in the cache-scope resolver -- /branch, delegate and tool children keep the isolation of #79161/#79017 byte for byte. Precedence is unchanged where it already worked: an explicit body session_id and the previous_response_id chain both still outrank the declared key, and a request that declares nothing keeps its per-request id. Recording is opt-in (bind_declared_conversation), so no other _run_agent caller's rows change. Refs #96811 (cherry picked from commit e7c83dddf36784d1012bf483240ebc7f6b2ef9aa)
This commit is contained in:
@@ -2282,6 +2282,94 @@ class APIServerAdapter(BasePlatformAdapter):
|
||||
# that the sanitized form is safe to pass into Honcho / state.db.
|
||||
_MAX_SESSION_HEADER_LEN = 256
|
||||
|
||||
# Source stamped on every session row this platform owns. Hardwired in
|
||||
# both places that create one — ``_bind_api_server_session`` (session
|
||||
# ContextVars) and ``_create_agent`` (``platform="api_server"``) — so the
|
||||
# peer lookup below can filter on it without guessing.
|
||||
_SESSION_SOURCE = "api_server"
|
||||
|
||||
def _declared_conversation_session(
|
||||
self, gateway_session_key: Optional[str]
|
||||
) -> Optional[str]:
|
||||
"""Resolve the live session a client declared with ``X-Hermes-Session-Key``.
|
||||
|
||||
The key names the *conversation*; ``session_id`` names the transcript
|
||||
that conversation is currently on. A client that manages its own
|
||||
history has no ``previous_response_id`` chain to carry the transcript
|
||||
forward, so the handlers used to mint a fresh id per request — and
|
||||
every conversation-affinity hint Hermes sends off that id
|
||||
(``prompt_cache_key`` on both OpenAI-wire transports, the
|
||||
OpenRouter/Nous sticky ``session_id``, and xAI's ``x-grok-conv-id``)
|
||||
re-keyed on every single reply (#96811).
|
||||
|
||||
This is the same reset-fenced recovery every native gateway platform
|
||||
already uses (``SessionStore._recover_session_for_peer``): rows ended
|
||||
at a conversation boundary — ``session_reset`` (/new),
|
||||
``session_switch``, ``idle``, ``daily``, ``suspended`` — are fenced
|
||||
out, so a new conversation still gets a new id and a cold affinity
|
||||
scope. The generation that must rotate is therefore already durable
|
||||
in ``sessions.end_reason``; nothing here needs a counter.
|
||||
|
||||
Returns ``None`` when nothing was declared, when no live row is
|
||||
recorded for the declared key, or on any DB error — every one of
|
||||
those leaves the caller's per-request id exactly as it is today.
|
||||
"""
|
||||
key = (gateway_session_key or "").strip()
|
||||
if not key:
|
||||
return None
|
||||
db = self._ensure_session_db()
|
||||
if db is None:
|
||||
return None
|
||||
try:
|
||||
row = db.find_latest_gateway_session_for_peer(
|
||||
source=self._SESSION_SOURCE, session_key=key
|
||||
)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"[%s] declared-conversation lookup failed", self.name, exc_info=True
|
||||
)
|
||||
return None
|
||||
return str(row["id"]) if row and row.get("id") else None
|
||||
|
||||
def _bind_declared_conversation(
|
||||
self, session_id: Optional[str], gateway_session_key: Optional[str]
|
||||
) -> None:
|
||||
"""Record the declared conversation key on the session row.
|
||||
|
||||
Counterpart to :meth:`_declared_conversation_session`. Without it the
|
||||
row is written unkeyed by ``AIAgent._ensure_db_session`` (which knows
|
||||
the key but does not persist it), and the reset-fenced lookup can
|
||||
never see it — the mapping the next reply needs would not exist.
|
||||
|
||||
``include_compression_ancestors`` carries the key up a mid-turn
|
||||
compression rotation so the pre- and post-rotation rows of one
|
||||
conversation share it, while that same walk deliberately stops at
|
||||
``/branch``, delegate and tool children (#79161). The statement is an
|
||||
UPDATE, so it is a harmless no-op on a turn that failed before the row
|
||||
was created.
|
||||
"""
|
||||
key = (gateway_session_key or "").strip()
|
||||
sid = str(session_id or "").strip()
|
||||
if not key or not sid:
|
||||
return
|
||||
db = self._ensure_session_db()
|
||||
if db is None:
|
||||
return
|
||||
try:
|
||||
db.record_gateway_session_peer(
|
||||
sid,
|
||||
source=self._SESSION_SOURCE,
|
||||
session_key=key,
|
||||
include_compression_ancestors=True,
|
||||
)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"[%s] declared-conversation bind failed for %s",
|
||||
self.name,
|
||||
sid,
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
def _parse_session_key_header(
|
||||
self, request: "web.Request"
|
||||
) -> tuple[Optional[str], Optional["web.Response"]]:
|
||||
@@ -6330,8 +6418,16 @@ class APIServerAdapter(BasePlatformAdapter):
|
||||
conversation_history = _auto_truncate_response_history(conversation_history)
|
||||
|
||||
# Reuse session from previous_response_id chain so the dashboard
|
||||
# groups the entire conversation under one session entry.
|
||||
session_id = stored_session_id or str(uuid.uuid4())
|
||||
# groups the entire conversation under one session entry. A client
|
||||
# that manages its own history has no chain to reuse, so fall back to
|
||||
# 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).
|
||||
session_id = (
|
||||
stored_session_id
|
||||
or self._declared_conversation_session(gateway_session_key)
|
||||
or str(uuid.uuid4())
|
||||
)
|
||||
|
||||
stream = _coerce_request_bool(body.get("stream"), default=False)
|
||||
route = self._resolve_route(body.get("model"))
|
||||
@@ -6402,6 +6498,7 @@ class APIServerAdapter(BasePlatformAdapter):
|
||||
tool_complete_callback=_on_tool_complete,
|
||||
agent_ref=agent_ref,
|
||||
gateway_session_key=gateway_session_key,
|
||||
bind_declared_conversation=True,
|
||||
**agent_overrides,
|
||||
route=route,
|
||||
))
|
||||
@@ -6437,6 +6534,7 @@ class APIServerAdapter(BasePlatformAdapter):
|
||||
ephemeral_system_prompt=instructions,
|
||||
session_id=session_id,
|
||||
gateway_session_key=gateway_session_key,
|
||||
bind_declared_conversation=True,
|
||||
**agent_overrides,
|
||||
route=route,
|
||||
)
|
||||
@@ -7265,6 +7363,7 @@ class APIServerAdapter(BasePlatformAdapter):
|
||||
requested_runtime: Optional[Dict[str, Any]] = None,
|
||||
route_source: str = "global",
|
||||
confirmed_runtime_lock: bool = False,
|
||||
bind_declared_conversation: bool = False,
|
||||
) -> tuple:
|
||||
"""
|
||||
Create an agent and run a conversation in a thread executor.
|
||||
@@ -7491,6 +7590,19 @@ class APIServerAdapter(BasePlatformAdapter):
|
||||
# shutdown. pop() is a no-op when _create_agent
|
||||
# succeeded but the turn never reached registration.
|
||||
self._shutdown_interruptible_agents.pop(id(agent), None)
|
||||
# Record the declared conversation on the row the turn
|
||||
# actually ended on — ``agent.session_id`` already
|
||||
# carries a mid-turn compression rotation (#16938), so
|
||||
# the next reply resolves the live transcript rather
|
||||
# than its retired parent. Opt-in: only the routes that
|
||||
# resolve their session id from the declared key
|
||||
# (/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,
|
||||
)
|
||||
clear_session_vars(tokens)
|
||||
|
||||
self._activate_admitted_request()
|
||||
|
||||
@@ -604,7 +604,16 @@ async def _handle_runs(
|
||||
|
||||
run_id = f"run_{uuid.uuid4().hex}"
|
||||
self._run_owners[run_id] = self._run_idempotency_scope(request)
|
||||
session_id = session_id or run_id
|
||||
# Same rule as /v1/responses: an explicit body session_id wins, then
|
||||
# the response chain, then the conversation the client declared via
|
||||
# ``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).
|
||||
session_id = (
|
||||
session_id
|
||||
or self._declared_conversation_session(gateway_session_key)
|
||||
or run_id
|
||||
)
|
||||
# Approval queues gate host-side tool execution and must be isolated
|
||||
# per API run. Client-provided session IDs and memory session keys are
|
||||
# conversation/memory scopes, not authorization namespaces: multiple
|
||||
@@ -829,6 +838,13 @@ async def _handle_runs(
|
||||
# run deliberately left running (same race-window
|
||||
# guard as gateway/run.py and _run_agent above).
|
||||
_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,
|
||||
)
|
||||
try:
|
||||
unregister_gateway_notify(approval_session_key)
|
||||
finally:
|
||||
|
||||
@@ -0,0 +1,208 @@
|
||||
"""Four-stage probe for #96811 — per-response session ids churn every
|
||||
conversation-affinity key.
|
||||
|
||||
Run: ``python scratch/repro_96811.py`` (no pytest, no network, temp state.db).
|
||||
|
||||
The probe walks the causal chain end to end on a real ``SessionDB``:
|
||||
|
||||
S1 reproduce the churn on current main's ``/v1/responses`` shape
|
||||
S2 show why the durable key->session mapping cannot help today
|
||||
S3 show the two-line repair restores affinity across replies
|
||||
S4 isolation: /new rotates, distinct keys never collide, no-key unchanged
|
||||
|
||||
Every stage prints PASS/FAIL; the process exits non-zero if any stage fails.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sys
|
||||
import tempfile
|
||||
import time
|
||||
import types
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
|
||||
|
||||
from agent.prompt_cache_scope import resolve_prompt_cache_scope # noqa: E402
|
||||
from hermes_state import SessionDB # noqa: E402
|
||||
|
||||
SOURCE = "api_server"
|
||||
DECLARED_KEY = "agent:main:api_server:room-42:member-7"
|
||||
OTHER_KEY = "agent:main:api_server:room-42:member-8"
|
||||
|
||||
_failures: list[str] = []
|
||||
|
||||
|
||||
def check(stage: str, label: str, ok: bool, detail: str = "") -> None:
|
||||
status = "PASS" if ok else "FAIL"
|
||||
print(f" [{status}] {stage} {label}" + (f" -- {detail}" if detail else ""))
|
||||
if not ok:
|
||||
_failures.append(f"{stage} {label} {detail}".strip())
|
||||
|
||||
|
||||
def agent_for(session_id: str, db: SessionDB) -> types.SimpleNamespace:
|
||||
"""Minimal stand-in for the attributes resolve_prompt_cache_scope reads."""
|
||||
return types.SimpleNamespace(session_id=session_id, _session_db=db)
|
||||
|
||||
|
||||
def responses_turn(
|
||||
db: SessionDB,
|
||||
*,
|
||||
declared_key: str | None,
|
||||
resolve_declared: bool,
|
||||
stamp_peer: bool,
|
||||
) -> str:
|
||||
"""One POST /v1/responses turn with client-managed conversation_history.
|
||||
|
||||
``resolve_declared`` / ``stamp_peer`` toggle the two proposed repairs so
|
||||
the same code path reproduces both the defect and the fix.
|
||||
"""
|
||||
# gateway/platforms/api_server.py::_handle_responses — client supplies its
|
||||
# own history, so the previous_response_id chain yields nothing.
|
||||
stored_session_id = None
|
||||
|
||||
declared_session_id = None
|
||||
if resolve_declared and declared_key:
|
||||
row = db.find_latest_gateway_session_for_peer(
|
||||
source=SOURCE, session_key=declared_key
|
||||
)
|
||||
declared_session_id = row["id"] if row else None
|
||||
|
||||
# main today: `session_id = stored_session_id or str(uuid.uuid4())`
|
||||
session_id = stored_session_id or declared_session_id or str(uuid.uuid4())
|
||||
|
||||
# AIAgent._ensure_db_session() — note: no session_key is written.
|
||||
if db.get_session(session_id) is None:
|
||||
db.create_session(session_id=session_id, source=SOURCE, model="m")
|
||||
|
||||
if stamp_peer and declared_key:
|
||||
db.record_gateway_session_peer(
|
||||
session_id, source=SOURCE, session_key=declared_key
|
||||
)
|
||||
return session_id
|
||||
|
||||
|
||||
def affinity_surfaces(session_id: str, db: SessionDB) -> str:
|
||||
"""The single value all four wire surfaces are derived from.
|
||||
|
||||
prompt_cache_key (codex + chat_completions transports) reads
|
||||
resolve_prompt_cache_scope; the OpenRouter/Nous sticky ``session_id`` and
|
||||
xAI's ``x-grok-conv-id`` read the conversation root, which is the same
|
||||
lineage walk over the same physical id. One value, four consumers.
|
||||
"""
|
||||
return resolve_prompt_cache_scope(agent_for(session_id, db))
|
||||
|
||||
|
||||
def stage_1_reproduce(db: SessionDB) -> tuple[str, str]:
|
||||
print("\nS1 reproduce: two replies on one declared conversation")
|
||||
a = responses_turn(db, declared_key=DECLARED_KEY, resolve_declared=False, stamp_peer=False)
|
||||
time.sleep(0.01)
|
||||
b = responses_turn(db, declared_key=DECLARED_KEY, resolve_declared=False, stamp_peer=False)
|
||||
check("S1", "physical ids differ per reply", a != b, f"{a[:8]} != {b[:8]}")
|
||||
check(
|
||||
"S1",
|
||||
"affinity scope churns per reply",
|
||||
affinity_surfaces(a, db) != affinity_surfaces(b, db),
|
||||
"prompt_cache_key / sticky session_id / x-grok-conv-id all re-key",
|
||||
)
|
||||
return a, b
|
||||
|
||||
|
||||
def stage_2_mapping_unreachable(db: SessionDB, session_id: str) -> None:
|
||||
print("\nS2 why the durable mapping cannot help today")
|
||||
row = db.get_session(session_id)
|
||||
check("S2", "row is created unkeyed", not (row or {}).get("session_key"),
|
||||
f"session_key={(row or {}).get('session_key')!r}")
|
||||
found = db.find_latest_gateway_session_for_peer(
|
||||
source=SOURCE, session_key=DECLARED_KEY
|
||||
)
|
||||
check("S2", "reset-fenced peer lookup finds nothing", found is None,
|
||||
"unkeyed rows are invisible to recovery")
|
||||
|
||||
|
||||
def stage_3_repair(db: SessionDB) -> str:
|
||||
print("\nS3 repair: stamp the routing key, resolve the id from it")
|
||||
a = responses_turn(db, declared_key=DECLARED_KEY, resolve_declared=True, stamp_peer=True)
|
||||
time.sleep(0.01)
|
||||
b = responses_turn(db, declared_key=DECLARED_KEY, resolve_declared=True, stamp_peer=True)
|
||||
c = responses_turn(db, declared_key=DECLARED_KEY, resolve_declared=True, stamp_peer=True)
|
||||
check("S3", "physical id is stable across replies", a == b == c, f"{a[:8]}")
|
||||
scopes = {affinity_surfaces(x, db) for x in (a, b, c)}
|
||||
check("S3", "affinity scope is stable across replies", len(scopes) == 1,
|
||||
next(iter(scopes))[:16])
|
||||
return a
|
||||
|
||||
|
||||
def stage_4_isolation(db: SessionDB, live_id: str) -> None:
|
||||
print("\nS4 isolation: rotation, collision, opt-out")
|
||||
|
||||
# /new — SessionStore.reset_session ends the row with 'session_reset',
|
||||
# which is inside _RESET_END_REASONS, so the recovery fence blocks it.
|
||||
live_scope = affinity_surfaces(live_id, db)
|
||||
db.end_session(live_id, "session_reset")
|
||||
after_reset = responses_turn(
|
||||
db, declared_key=DECLARED_KEY, resolve_declared=True, stamp_peer=True
|
||||
)
|
||||
check("S4", "/new rotates the physical id", after_reset != live_id,
|
||||
f"{live_id[:8]} -> {after_reset[:8]}")
|
||||
check(
|
||||
"S4",
|
||||
"/new rotates the affinity scope",
|
||||
affinity_surfaces(after_reset, db) != live_scope,
|
||||
"no ABA: the boundary is durable in sessions.end_reason",
|
||||
)
|
||||
|
||||
# idle / daily / suspended auto-resets use the same fence set.
|
||||
for reason in ("idle", "daily", "suspended"):
|
||||
prev = responses_turn(
|
||||
db, declared_key=DECLARED_KEY, resolve_declared=True, stamp_peer=True
|
||||
)
|
||||
db.end_session(prev, reason)
|
||||
nxt = responses_turn(
|
||||
db, declared_key=DECLARED_KEY, resolve_declared=True, stamp_peer=True
|
||||
)
|
||||
check("S4", f"'{reason}' auto-reset rotates", nxt != prev,
|
||||
f"{prev[:8]} -> {nxt[:8]}")
|
||||
|
||||
# A second declared channel never lands on the first one's conversation.
|
||||
mine = responses_turn(
|
||||
db, declared_key=DECLARED_KEY, resolve_declared=True, stamp_peer=True
|
||||
)
|
||||
theirs = responses_turn(
|
||||
db, declared_key=OTHER_KEY, resolve_declared=True, stamp_peer=True
|
||||
)
|
||||
check("S4", "distinct declared keys stay isolated", mine != theirs,
|
||||
f"{mine[:8]} != {theirs[:8]}")
|
||||
|
||||
# A client that declares nothing keeps today's per-request identity.
|
||||
n1 = responses_turn(db, declared_key=None, resolve_declared=True, stamp_peer=True)
|
||||
n2 = responses_turn(db, declared_key=None, resolve_declared=True, stamp_peer=True)
|
||||
check("S4", "no declared key -> behavior unchanged", n1 != n2,
|
||||
"undeclared clients keep per-request ids")
|
||||
|
||||
|
||||
def main() -> int:
|
||||
# ignore_cleanup_errors: SessionDB holds the SQLite handle open, and
|
||||
# Windows refuses to unlink a mapped file (WinError 32).
|
||||
with tempfile.TemporaryDirectory(
|
||||
prefix="repro96811_", ignore_cleanup_errors=True
|
||||
) as tmp:
|
||||
db = SessionDB(Path(tmp) / "state.db")
|
||||
churn_a, _ = stage_1_reproduce(db)
|
||||
stage_2_mapping_unreachable(db, churn_a)
|
||||
live = stage_3_repair(db)
|
||||
stage_4_isolation(db, live)
|
||||
|
||||
print("\n" + "=" * 62)
|
||||
if _failures:
|
||||
print(f"{len(_failures)} FAILED:")
|
||||
for f in _failures:
|
||||
print(f" - {f}")
|
||||
return 1
|
||||
print("all four stages PASS")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -0,0 +1,295 @@
|
||||
"""Declared-conversation identity on the API server (#96811).
|
||||
|
||||
A client that manages its own history has no ``previous_response_id`` chain,
|
||||
so ``/v1/responses`` and ``/v1/runs`` used to mint a throwaway physical
|
||||
session id per request even when the request declared its conversation with
|
||||
``X-Hermes-Session-Key``. Every conversation-affinity hint Hermes sends is
|
||||
derived from that physical id — ``prompt_cache_key`` on both OpenAI-wire
|
||||
transports, the OpenRouter/Nous sticky ``session_id``, and xAI's
|
||||
``x-grok-conv-id`` — so all four re-keyed on every single reply.
|
||||
|
||||
These tests pin the identity contract itself rather than the four consumers:
|
||||
the declared key resolves to one live session, the resolution is fenced by
|
||||
the durable conversation boundaries already recorded in
|
||||
``sessions.end_reason``, and nothing that does not declare a key changes.
|
||||
"""
|
||||
|
||||
import types
|
||||
|
||||
import pytest
|
||||
|
||||
from gateway.config import PlatformConfig
|
||||
from gateway.platforms.api_server import APIServerAdapter
|
||||
from hermes_state import SessionDB
|
||||
|
||||
KEY = "agent:main:api_server:room-42:member-7"
|
||||
OTHER_KEY = "agent:main:api_server:room-42:member-8"
|
||||
SOURCE = "api_server"
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def adapter(tmp_path):
|
||||
"""An adapter whose lazy SessionDB is pinned to a scratch state.db."""
|
||||
a = APIServerAdapter(
|
||||
PlatformConfig(enabled=True, extra={"host": "127.0.0.1", "port": 0, "key": "k"})
|
||||
)
|
||||
db = SessionDB(tmp_path / "state.db")
|
||||
a._session_db = db
|
||||
try:
|
||||
yield a, db
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def _seed(db, session_id, *, key=KEY, source=SOURCE):
|
||||
db.create_session(session_id=session_id, source=source, model="m")
|
||||
if key:
|
||||
db.record_gateway_session_peer(session_id, source=source, session_key=key)
|
||||
|
||||
|
||||
class TestDeclaredConversationResolution:
|
||||
def test_declared_key_resolves_the_live_session(self, adapter):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert a._declared_conversation_session(KEY) == "sess-live"
|
||||
|
||||
def test_replies_land_on_one_session(self, adapter):
|
||||
"""The defect in one line: same key, three replies, one identity."""
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert {a._declared_conversation_session(KEY) for _ in range(3)} == {"sess-live"}
|
||||
|
||||
def test_no_declared_key_resolves_nothing(self, adapter):
|
||||
"""Undeclared clients keep today's per-request identity."""
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert a._declared_conversation_session(None) is None
|
||||
assert a._declared_conversation_session("") is None
|
||||
assert a._declared_conversation_session(" ") is None
|
||||
|
||||
def test_unknown_key_resolves_nothing(self, adapter):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert a._declared_conversation_session(OTHER_KEY) is None
|
||||
|
||||
def test_distinct_keys_stay_isolated(self, adapter):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-mine", key=KEY)
|
||||
_seed(db, "sess-theirs", key=OTHER_KEY)
|
||||
assert a._declared_conversation_session(KEY) == "sess-mine"
|
||||
assert a._declared_conversation_session(OTHER_KEY) == "sess-theirs"
|
||||
|
||||
def test_another_platforms_key_is_not_adopted(self, adapter):
|
||||
"""The source filter keeps a telegram chat key out of the API server."""
|
||||
a, db = adapter
|
||||
_seed(db, "sess-telegram", key=KEY, source="telegram")
|
||||
assert a._declared_conversation_session(KEY) is None
|
||||
|
||||
def test_db_failure_degrades_to_a_fresh_id(self, adapter):
|
||||
a, _ = adapter
|
||||
|
||||
class _Boom:
|
||||
def find_latest_gateway_session_for_peer(self, **_kw):
|
||||
raise RuntimeError("db down")
|
||||
|
||||
a._session_db = _Boom()
|
||||
assert a._declared_conversation_session(KEY) is None
|
||||
|
||||
def test_missing_db_degrades_to_a_fresh_id(self, adapter, monkeypatch):
|
||||
a, _ = adapter
|
||||
monkeypatch.setattr(a, "_ensure_session_db", lambda: None)
|
||||
assert a._declared_conversation_session(KEY) is None
|
||||
|
||||
|
||||
class TestConversationBoundariesRotate:
|
||||
"""The generation that must rotate is already durable in end_reason.
|
||||
|
||||
#79017/#86733's contract: an affinity scope stays warm across
|
||||
continuation and compression rotation, and goes cold on a new
|
||||
conversation. ``_RESET_END_REASONS`` is that boundary set, and the
|
||||
recovery fence honours every member of it — including the idle/daily
|
||||
policy resets, not just ``/new``.
|
||||
"""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"reason",
|
||||
["session_reset", "session_switch", "idle", "daily", "suspended",
|
||||
"resume_pending_expired"],
|
||||
)
|
||||
def test_boundary_rotates_the_conversation(self, adapter, reason):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-old")
|
||||
db.end_session("sess-old", reason)
|
||||
assert a._declared_conversation_session(KEY) is None
|
||||
|
||||
def test_boundary_cannot_be_reached_behind(self, adapter):
|
||||
"""A later row wins, and the retired one never comes back (no ABA)."""
|
||||
a, db = adapter
|
||||
_seed(db, "sess-gen1")
|
||||
db.end_session("sess-gen1", "session_reset")
|
||||
_seed(db, "sess-gen2")
|
||||
assert a._declared_conversation_session(KEY) == "sess-gen2"
|
||||
db.end_session("sess-gen2", "session_reset")
|
||||
assert a._declared_conversation_session(KEY) is None
|
||||
|
||||
def test_accidental_end_stays_resumable(self, adapter):
|
||||
"""An accidental close is not a conversation boundary."""
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
db.end_session("sess-live", "agent_close")
|
||||
assert a._declared_conversation_session(KEY) == "sess-live"
|
||||
|
||||
|
||||
class TestBindDeclaredConversation:
|
||||
def test_bind_makes_the_row_resolvable(self, adapter):
|
||||
"""Without the bind the row is written unkeyed and is invisible."""
|
||||
a, db = adapter
|
||||
db.create_session(session_id="sess-new", source=SOURCE, model="m")
|
||||
assert a._declared_conversation_session(KEY) is None
|
||||
|
||||
a._bind_declared_conversation("sess-new", KEY)
|
||||
assert a._declared_conversation_session(KEY) == "sess-new"
|
||||
|
||||
def test_bind_is_a_noop_without_a_key(self, adapter):
|
||||
a, db = adapter
|
||||
db.create_session(session_id="sess-new", source=SOURCE, model="m")
|
||||
a._bind_declared_conversation("sess-new", None)
|
||||
a._bind_declared_conversation("sess-new", " ")
|
||||
assert db.get_session("sess-new").get("session_key") in (None, "")
|
||||
|
||||
def test_bind_is_a_noop_without_a_session(self, adapter):
|
||||
a, _ = adapter
|
||||
a._bind_declared_conversation(None, KEY)
|
||||
a._bind_declared_conversation("", KEY)
|
||||
assert a._declared_conversation_session(KEY) is None
|
||||
|
||||
def test_bind_survives_a_db_failure(self, adapter):
|
||||
a, _ = adapter
|
||||
|
||||
class _Boom:
|
||||
def record_gateway_session_peer(self, *_a, **_kw):
|
||||
raise RuntimeError("db down")
|
||||
|
||||
a._session_db = _Boom()
|
||||
a._bind_declared_conversation("sess-new", KEY) # must not raise
|
||||
|
||||
def test_bind_follows_a_compression_rotation(self, adapter):
|
||||
"""The turn binds the row it ended on; the retired parent follows.
|
||||
|
||||
``include_compression_ancestors`` keys the whole compression lineage,
|
||||
so the next reply resolves the live child rather than its parent.
|
||||
"""
|
||||
a, db = adapter
|
||||
db.create_session(session_id="sess-parent", source=SOURCE, model="m")
|
||||
db.end_session("sess-parent", "compression")
|
||||
db.create_session(
|
||||
session_id="sess-child",
|
||||
source=SOURCE,
|
||||
model="m",
|
||||
parent_session_id="sess-parent",
|
||||
)
|
||||
a._bind_declared_conversation("sess-child", KEY)
|
||||
assert a._declared_conversation_session(KEY) == "sess-child"
|
||||
|
||||
|
||||
class TestOtherStructuresUnaffected:
|
||||
"""The bind uses the same routing-peer record every native platform uses.
|
||||
|
||||
These pin the two places a newly keyed row could leak into: the channel
|
||||
directory (contact lists) and the SessionStore's own routing table.
|
||||
"""
|
||||
|
||||
def test_channel_directory_skips_declared_api_rows(self, adapter):
|
||||
"""Rows carry no chat_id/origin, so the directory has no entry to build."""
|
||||
from gateway import channel_directory
|
||||
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
row = db.get_session("sess-live")
|
||||
origin = {
|
||||
"chat_id": row.get("chat_id"),
|
||||
"thread_id": row.get("thread_id"),
|
||||
"chat_name": row.get("display_name"),
|
||||
}
|
||||
assert channel_directory._session_entry_id(origin) is None
|
||||
|
||||
def test_session_store_routing_table_is_untouched(self, adapter):
|
||||
"""SessionStore loads from gateway_routing, not from sessions.session_key."""
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert db.load_gateway_routing_entries() == {}
|
||||
|
||||
|
||||
class TestHandlerWiring:
|
||||
"""The precedence the two handlers apply, isolated from aiohttp."""
|
||||
|
||||
@staticmethod
|
||||
def _resolve_responses(adapter, *, stored, key):
|
||||
# gateway/platforms/api_server.py::_handle_responses
|
||||
return stored or adapter._declared_conversation_session(key) or "minted-uuid"
|
||||
|
||||
@staticmethod
|
||||
def _resolve_runs(adapter, *, body_id, stored, key):
|
||||
# gateway/platforms/api_server.py::_handle_runs
|
||||
return (
|
||||
(body_id or stored)
|
||||
or adapter._declared_conversation_session(key)
|
||||
or "minted-run-id"
|
||||
)
|
||||
|
||||
def test_response_chain_still_outranks_the_declared_key(self, adapter):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert self._resolve_responses(a, stored="sess-chained", key=KEY) == "sess-chained"
|
||||
|
||||
def test_declared_key_outranks_a_minted_id(self, adapter):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert self._resolve_responses(a, stored=None, key=KEY) == "sess-live"
|
||||
|
||||
def test_undeclared_request_still_mints(self, adapter):
|
||||
a, _ = adapter
|
||||
assert self._resolve_responses(a, stored=None, key=None) == "minted-uuid"
|
||||
|
||||
def test_runs_body_session_id_still_wins(self, adapter):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert self._resolve_runs(a, body_id="explicit", stored=None, key=KEY) == "explicit"
|
||||
|
||||
def test_runs_declared_key_outranks_the_run_id(self, adapter):
|
||||
a, db = adapter
|
||||
_seed(db, "sess-live")
|
||||
assert self._resolve_runs(a, body_id=None, stored=None, key=KEY) == "sess-live"
|
||||
|
||||
def test_runs_undeclared_still_uses_the_run_id(self, adapter):
|
||||
a, _ = adapter
|
||||
assert self._resolve_runs(a, body_id=None, stored=None, key=None) == "minted-run-id"
|
||||
|
||||
|
||||
class TestRunAgentOptIn:
|
||||
"""Only the two routes that resolve a declared id record one."""
|
||||
|
||||
def test_bind_targets_the_rotated_session(self, adapter, monkeypatch):
|
||||
"""The finally block binds ``agent.session_id``, not the id it started on."""
|
||||
a, _ = adapter
|
||||
calls = []
|
||||
monkeypatch.setattr(
|
||||
a, "_bind_declared_conversation", lambda sid, key: calls.append((sid, key))
|
||||
)
|
||||
agent = types.SimpleNamespace(session_id="sess-rotated")
|
||||
a._bind_declared_conversation(
|
||||
getattr(agent, "session_id", None) or "sess-initial", KEY
|
||||
)
|
||||
assert calls == [("sess-rotated", KEY)]
|
||||
|
||||
def test_bind_falls_back_when_the_agent_never_started(self, adapter, monkeypatch):
|
||||
a, _ = adapter
|
||||
calls = []
|
||||
monkeypatch.setattr(
|
||||
a, "_bind_declared_conversation", lambda sid, key: calls.append((sid, key))
|
||||
)
|
||||
agent = types.SimpleNamespace(session_id=None)
|
||||
a._bind_declared_conversation(
|
||||
getattr(agent, "session_id", None) or "sess-initial", KEY
|
||||
)
|
||||
assert calls == [("sess-initial", KEY)]
|
||||
Reference in New Issue
Block a user