fix(gateway): fail closed when a pinned route moves under a delegation completion - #113716
wangtaotaotao95 wants to merge 2 commits into
Conversation
…n completion `_resolve_async_delegation_session` awaits the spawning session's DB row and then repoints the routing key unconditionally on the non-compression branch. That await is the race window: /new or /stop can revoke the run generation and a concurrent /resume can repoint the key while the lookup is suspended, after which the stale resolution overwrites the newer route with the pin it snapshotted before the boundary. Re-check the run token after the await — threaded from the routed turn, since `_hmwa_resolve_session` never received it — and pass the session the caller resolved against to `switch_session` as a compare-and-swap condition, so the route only moves while it still points there. Both boundaries now drop the injection instead of publishing it. Refs NousResearch#113690
|
Independent verification on the PR head (2ee68ad): the pinned-route race suite plus the delegation binding suite pass locally, 9/9 on Linux (canonical runner). Cleaner than the sibling approach: the CAS lives inside switch_session as an optional expected_session_id rather than a parallel method, so all callers share one conditional path and stale resolutions fail closed with a logged refusal. Note for maintainers: #113692 covers adjacent ground with switch_session_if_current; worth converging on one shape. |
andrexibiza
left a comment
There was a problem hiding this comment.
I traced this exact head through the real inbound run-generation claim, _hmwa_resolve_session, the async delegation resolver, AsyncSessionStore, the synchronous SessionStore mutation, the new event-synchronized regression, #113690, the competing #113692 carrier, and the established dispatch-time ownership lineage from merged #60871. The expected-session CAS is a real improvement and closes the concurrent-route-replacement half of the reporter's reproducer without weakening the captured owner. There is still one blocking authority race in the current shape.
P1 — the run-generation proof is still separated from the effect by another scheduling boundary.
The new check correctly catches a /stop//new that lands while session_db.get_session() is suspended. But the proof is consumed only once, immediately after that first await. The non-compression mutation then goes through AsyncSessionStore, whose generic wrapper is await asyncio.to_thread(attr, *args, **kwargs). That creates a second interference window after the generation check and before the synchronous SessionStore.switch_session() mutation.
A concrete execution is:
- Turn owns generation
G; the pinned-row lookup returns. _is_session_run_current(key, G)succeeds.- Before the offloaded
switch_session()acquires the store lock,/stopinvalidates the run toG+1without changing the route's current session id — the pure-revoke case from #113690. switch_session(... expected_session_id=prior_session_id)still sees the expected route id and commits the stale repin.
The expected-session CAS cannot detect that boundary because route identity and run generation are different authority dimensions. The same root exists on the compression side: after the one generation check, _resolve_compression_lineage_target() can perform several more awaits, and the target_session_id == session_entry.session_id fast path can return a stale entry without any final generation validation at all. So the PR narrows the original race window, but it does not yet make the advertised invariant — a revoked completion cannot publish after the boundary — true for the whole mutation/delivery path.
The closure needs the generation proof to be effect-bound, not another preflight observation: validate the expected route and expected run generation under one coordination/commit boundary that cannot interleave with generation invalidation, then perform/authorize the route transition (or fail closed). A second check immediately before another await is still TOCTOU; post-hoc rollback is also unsafe because a newer route decision may have landed by then. Please add a deterministic witness that suspends the route commit after the new generation check, invalidates the generation while the commit is held, then releases it and proves both that the routing key stayed put and that the completion was not returned for injection. Mirror that witness across the compression/identity fast path as the other side of the same invariant.
The surrounding topology matters here:
- #113690 / tobific owns the current race report and its
none / revoke / replaceinvariant. - #113692 / KoNit-K is a competing partial carrier. Its route CAS closes
replace, but the reporter's purerevokecase was independently shown still failing because the route can remain unchanged. This PR correctly adds generation to the proof; it should be the place where that proof is carried all the way to the effect rather than stopping one await early. - #60871 / teknium1, salvaging #57535 / nankingjing with authorship preserved, established
parent_session_id/gateway_session_idas the durable dispatch-time owner. This PR preserves that stronger identity invariant, which is important. - #107812 / salch-cred is not an equivalent substitute: its current-session retargeting changes ownership and has an existing P1 for recreating the #57498 cross-session contamination class.
- Closed-unmerged #64530 / richkapp is useful historical compression-lineage/CAS design evidence, but it is not landed provenance by itself.
On verification: the author supplied useful RED→GREEN local evidence, and another contributor reports 9/9 focused Linux tests on this exact head. Hosted acceptance is nevertheless absent. This branch has one surviving commit (2ee68ade52ccf753dc8ba78e6991d9f465e9eb23), and its Nix run 35178430405, Docker run 35178430468, and CI run 35178430621 all concluded action_required; the CI run created zero jobs. With one surviving commit, exact-head and every-commit acceptance are the same object here, and that object is 0/1 hosted-green.
Current repository state at review: main / actual merge base is 6005aa1fd9aac8b1024ace50fec8cd1c85a04bae; this carrier is 1 ahead / 0 behind and GitHub reports it mergeable. The branch is nicely scoped and the expected-session CAS is worth keeping. The remaining fix is to make the generation proof survive to the actual effect boundary rather than merely checking it earlier in the async path.
| pinned_row = await session_db.get_session(pinned_session_id) | ||
| except Exception: | ||
| logger.debug("Async-delegation parent lookup failed for %s", pinned_session_id, exc_info=True) | ||
| if not self._is_session_run_current(generation_key, expected_generation): |
There was a problem hiding this comment.
P1 — this generation check is still TOCTOU with respect to the route mutation. It proves G here, but the non-compression path later does await self.async_session_store.switch_session(...); AsyncSessionStore offloads the synchronous store call through asyncio.to_thread, so /stop can invalidate G -> G+1 after this check and before the store mutation. In the pure-revoke case the route id can remain unchanged, so expected_session_id still matches and the stale completion repins anyway. The compression path has the same root because several lineage awaits (and the target == current early return) occur after this lone generation check. Please bind expected generation + expected route to the effect/commit boundary itself, and add an event-synchronized regression that invalidates after this check but before route commit; another pre-await check is not sufficient.
|
Thanks — noted on #113692. I'm happy to converge on whichever shape maintainers prefer; folding the condition into One coverage difference worth flagging while both are open, so the choice is made on coverage rather than style: A pure session-id CAS catches the case where the route itself moved. #113690's second case is different:
So this PR carries the run-generation guard in addition to the CAS ( Not blocking either way — just making the difference explicit. |
The run token was sampled once after the spawning-session lookup, but the route mutation then went through AsyncSessionStore — another await, and an offloaded call. A /stop or /new landing in that second window left the routing key unchanged, so the expected-session CAS still matched and the stale pin committed anyway. The compression branch had the same root with several more awaits, and the identity fast path returned an entry for injection without sampling the token at all. Sample the token at the effect instead: switch_session_if_current and advance_compression_session take an `authorize` predicate evaluated under a new store authority lock, and every generation bump takes that same lock. A routing transition cannot move onto the event loop — a repoint is structural, so it rewrites the whole routing index rather than taking the single-entry fast path — which is why a shared authority, not a synchronous commit, is what makes the proof and the mutation one boundary. The identity fast path samples the token in the same block that returns, with no await in between. The authority is an RLock: invalidate nests begin. Tests: the parametrized boundary suite grows an identity/compression witness that injects the boundary at the instant the transition samples the token. Red on base (1 passed, 4 failed), green after. Refs NousResearch#113690
|
You're right, and thanks for the line-level trace — I reproduced it before changing anything. The check was consumed once after Fixed by making the proof effect-bound rather than preflight, in
One thing worth recording, because it changed the shape of the fix: a synchronous commit on the loop is not available here. A repoint is structural, so I also hit a self-deadlock implementing it: Witness, as asked — Green after ( On convergence with #113692: I kept the CAS inside the store rather than a parallel body and named the method
|
…new won the race The non-compression branch of _resolve_async_delegation_session awaited the spawning-session row lookup and then unconditionally called switch_session(), so a run invalidated (/stop) or a route replaced (/new, /resume) while the lookup was pending was overwritten by the stale completion's pin. - Snapshot the routing key's run generation before the await and re-check it before mutating the route; invalidated -> drop the injection, route untouched. - switch_session(expected_session_id=...) turns the /resume primitive into a CAS for this caller, mirroring advance_compression_session: a route that moved past the snapshot wins over the late completion. - _current_session_run_generation extracted from _is_session_run_current. Slimmer redo of #113692 (@KoNit-K) and #113716 (@wangtaotaotao95): same direction, without the second routing-authority lock and authorize callbacks. Fixes #113690 Co-authored-by: KoNit-K <konit.block@protonmail.com> Co-authored-by: wangtaotaotao95 <wangtaotaotao95@users.noreply.github.com>
|
Thanks @wangtaotaotao95. Your change was salvaged into #114765 with your authorship preserved (co-authored); #114765 — fix(gateway): async completion no longer re-pins a route after /stop or /new won the race (#113690, salvage #113692, #113716) — is now merged on |
What does this PR do?
Fixes the race in #113690: a pinned async-delegation completion could overwrite a routing key that had already moved, because the resolution awaits the spawning session's DB row and then repoints the key unconditionally.
gateway/run_notifications.py::_resolve_async_delegation_sessiondoes:That
awaitis the window./newor/stopcan revoke the in-flight run generation and a concurrent/resumecan repoint the key while the lookup is suspended; the stale resolution then publishes the pin it snapshotted before the boundary. A revoked completion must not be able to write routing at all — the resolved entry is what every later delivery in that turn is addressed to.The compression branch was already safe: it goes through
advance_compression_session, which is a compare-and-swap on the expected current session. The non-compression branch had no such condition, andswitch_sessionaccepted none.Related Issue
Fixes #113690
Type of Change
Changes Made
gateway/run_notifications.py—_resolve_async_delegation_sessionnow takes the caller's run token (optional), snapshots it before the row lookup, and re-checks it after the await; a revoked run drops the injection instead of moving the route. The non-compression repoint passes the session it resolved against as a condition.gateway/session.py—switch_sessiontakes an optionalexpected_session_idand refuses the repoint when the route no longer points there (returnsNone, i.e. fail closed). Existing two-argument callers (/resume, handoff, startup) are unchanged.gateway/run_turn.py—_hmwa_resolve_session(and its call from_handle_message_with_agent) forwards_quick_key/run_generation, which the resolver never received. The parameter is optional, so the eval harness caller is unaffected.gateway/run_agent_cache.py— adds_current_session_run_generation, so a caller that does not pass its own token still snapshots a comparable one.tests/gateway/test_pinned_route_race.py— new invariant test over a realSessionStoreand the real generation primitives.tests/gateway/test_async_delegation_session_binding.py— the "live spawning session rebinds" assertion now expects the conditional repoint. The behaviour it asserts is the same; the contract got stricter.Both boundaries fail closed: the route keeps whatever the newer decision set, and the completion is dropped.
How to Test
The test suspends the session-row lookup on an
asyncio.Event, applies a boundary, then releases it — no wall-clock sleeps:none— nothing moves; the pin repoints the route and is returned.revoke— the run generation is invalidated while the lookup is pending; the route must stay put and the result must beNone.replace— the run is revoked and the key is repointed to a replacement session; the replacement must survive and the result must beNone.Proof it is red on base, at
6005aa1fd9:With the fix:
3 passed.No regressions. Full gateway suite (
tests/gateway/, 901 files):The 6 failures are identical on unpatched
6005aa1fd9(397 passed, 6 failedon the same six files) — they are host-environment failures (test_gateway_trust_env,test_slack,test_feishu,test_wecom,test_telegram_thread_fallback,test_multiplex_residue_paritytry to reach the network / observe this host's system proxy). The 31 test files that reference any changed symbol pass:356 passed.ruff checkandruff format --checkare clean on the new test file.Checklist
Code
fix(scope):,feat(scope):, etc.)Documentation & Housekeeping
cli-config.yaml.exampleif I added/changed config keys — N/ACONTRIBUTING.mdorAGENTS.mdif I changed architecture or workflows — N/AScreenshots / Logs
Base (
6005aa1fd9) vs. fix, same command: