From c4bc6e0e38c3a9ed78c09e91b58009532831341c Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Sat, 20 Jun 2026 04:51:32 +0800 Subject: [PATCH 1/5] feat(agent): add queue-aware Agent.cancel() primitive MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit abort() only kills the in-flight step, so a queued-but-not-yet-started prompt ran to completion after a cancel and a prompt accepted right after could be batched into the cancelled turn (the loop merges queued messages into one turn). This closes TODO(rfc010-cancel-prestep) with a distinct cancel() verb. cancel() clears the queued + steering FIFOs, aborts the in-flight step, and drives a turn-scoped marker on the LoopHandle that the driver checks at EVERY point a turn could start or continue: - right after the idle wait (window 1): drop the about-to-run turn and settle whenIdle() waiters directly (no running→idle transition fires, and no agent/status is emitted, so an ACP listener can't see a spurious idle that resolves a freshly-queued prompt as cancelled); - after the synchronous setStatus('running') emit (window 2): a running listener can cancel in the gap before runTurn; - in the step-start window (before runStep, after setAbort): a synchronous turn-start/step-start listener can cancel before any AbortController exists; - at the continuation gate: a cancel during the continuation waterfall (the finished step's controller already cleared) ends the turn aborted. The marker is ARMED only when there is something to cancel (running, an in-flight step, or queued/steering work) — an idle no-op cancel cannot leave it set to drop a later prompt — and RESET unconditionally once per loop iteration, so it governs exactly one turn and never leaks onto the next prompt (even when a send() lands in the cancelled turn's flush window). ACP session/cancel now maps to agent.cancel() (keeping the synchronous settlePrompt). Teardown/disconnect still use abort('disposed') until PR D, so the ACP README narrows the remaining best-effort window to teardown only. Tests (agent-loop/cancel.spec.ts) cover every window unit-level (the F1 hang guard: a whenIdle() waiter registered before a pre-step cancel resolves; the F2 leak guard: idle cancel then a prompt runs; mid-step, continuation, both pre-step windows, turn-start-listener, steering-cleared, marker-reset). ACP turns.spec.ts adds the through-bridge tests with NO intervening whenIdle (idle cancel→prompt runs; mid-stream cancel→immediate next prompt runs) and updates the stale pre-step test to the queue-aware guarantee. The existing cancel snapshot golden is byte-identical (it drives the new cancel() path end-to-end through the real subprocess), so no new golden is needed. 100% coverage. --- docs/architecture.md | 1 + packages/acp/README.md | 4 +- packages/acp/src/index.ts | 30 ++- packages/acp/tests/turns.spec.ts | 63 ++++- packages/agent-loop/README.md | 2 + packages/agent-loop/src/agent.ts | 41 ++++ packages/agent-loop/src/inbox.ts | 10 + packages/agent-loop/src/loop.ts | 82 ++++++- packages/agent-loop/tests/cancel.spec.ts | 279 +++++++++++++++++++++++ packages/agent/README.md | 3 +- packages/agent/src/types.ts | 19 ++ packages/agent/tests/agent.spec.ts | 1 + 12 files changed, 504 insertions(+), 31 deletions(-) create mode 100644 packages/agent-loop/tests/cancel.spec.ts diff --git a/docs/architecture.md b/docs/architecture.md index 81a478b714..e28a384eb3 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -109,6 +109,7 @@ Tool schemas are deliberately **part of the assembly**: "what the model is told - `steer(content)` — mid-turn injection, drained **between steps**; behaves like `send` when idle - `inject(content)` — in-session context (`context/message` event); the next request sees it (Claude Code attachment / system-reminder analog). An inject made while the agent is *running* joins the open turn; an inject while *idle* is wrapped in a one-shot turn (`turn/start{trigger:injection}` → `context/message` → `turn/end`) so every event stays turn-enclosed (see [the turn-enclosure invariant](rfc/implemented/2026-06-15-turn-enclosure-invariant.md)). - `abort(reason)` — aborts the in-flight step via `AbortSignal` +- `cancel(reason)` — the broad cancel: clears queued + steering work, aborts the in-flight step, and drops a turn about to start (the pre-step window) so a queued-but-not-started prompt never runs and cannot be batched into the cancelled turn. `abort()` is the narrower step-only verb; `cancel()` is what a UI/ACP `session/cancel` maps to. - `whenIdle()` — resolves once the agent reaches quiescence after settling out of `running` (resolves immediately when already idle; awaits the loop exit when disposed). The teardown signal: `abort()` then `await whenIdle()` guarantees the in-flight turn has fully stopped. Observes the transition without disposing the agent. - `session`, `status`, `options` diff --git a/packages/acp/README.md b/packages/acp/README.md index 7db1ee8785..0bacaeef07 100644 --- a/packages/acp/README.md +++ b/packages/acp/README.md @@ -27,7 +27,7 @@ It is a **client-driver / UI plugin**, the structured analogue of the readline ` | `session/new` | `ctx.agents.create({ sessionId, meta:{cwd} })` | creates a new session/agent; N concurrent sessions are allowed, keyed by id; `cwd` must be absolute (it becomes the session's workspace — see Per-session cwd); non-empty `additionalDirectories` and `mcpServers` rejected | | `session/load` | `ctx.agents.resume(...)` | replays the persisted event log to the client as `session/update` — the USER side (`user/message` → `user_message_chunk`), assistant text/reasoning (`assistant/chunk`), and tool calls/results (`tool/call` + `tool/result`). Re-loading an already-live id is rejected; the id's load slot is reserved (`loadingIds`) BEFORE the async resume so a pipelined load of the SAME id can't leak a second agent (distinct ids load concurrently). The resumed session keeps its PERSISTED header `cwd`, so its bash tools run in the original workspace; the requested `cwd` must be absolute and match the persisted `cwd`. After the async resume a `closed` re-check refuses to install a record if the bridge tore down mid-load | | `session/prompt` | `agent.send()` | supports ACP `text` and `resource_link` blocks; rejects image/audio/embedded resource and empty prompts; one in-flight prompt PER session (independent); settles on the OWNING turn's end (a turn that ends in `error` rejects the RPC) | -| `session/cancel` | `agent.abort()` | aborts a running step + settles the prompt `cancelled` for ONLY that session — a cancel never touches another session's stream or prompt (see limitation below) | +| `session/cancel` | `agent.cancel()` | the queue-aware cancel: aborts a running step, clears queued + steering work, and drops a turn about to start, then settles the prompt `cancelled` — for ONLY that session (a cancel never touches another session's stream or prompt) | | `session/update` | `session/event` | `agent_message_chunk` (text-delta), `agent_thought_chunk` (reasoning-delta), `user_message_chunk` (load replay), `tool_call`/`tool_call_update` (title/kind/rawInput/content owned by the TOOL via `presentCall`/`presentResult` — see Tool-call presentation) | ## Multi-session @@ -66,7 +66,7 @@ Teardown reaches quiescence: for EVERY live session settle any pending prompt as ## Known limitations (tracked TODOs) - **`TODO(rfc010-permission-gate)`** — the `tools/execute` permission gate (`session/request_permission`) is NOT implemented; tools run with the executor's full authority. The `agent→sessionId` reverse map is in place so the gate can route a permission request (which receives only `exec.agent`) back to its originating session. [ACP support](../../docs/rfc/proposed/2026-06-14-acp-agent-client-protocol.md) and [ACP multi-session](../../docs/rfc/proposed/2026-06-14-acp-multi-session.md) stay `proposed` until the gate (and per-session permission ownership) land. -- **`TODO(rfc010-cancel-prestep)`** — `session/cancel` (and teardown/disconnect) is honest RPC/UI cancellation plus best-effort abort: a *running* step is aborted, but a turn that is queued-but-not-yet-started (the gap before `agent.abort()` has an `AbortController` to signal) may still run to completion. This same window means disposal/disconnect can return while one short queued turn per session still runs, and a prompt accepted right after a pre-step cancel can be batched into the cancelled turn (the loop merges queued messages into one turn). A loop-level queue-aware cancel will close this; the single-in-flight-per-session rule bounds the worst case to one extra prompt per session. +- **`TODO(rfc010-cancel-prestep)`** — `session/cancel` is now the queue-aware `agent.cancel()` (a running step is aborted, queued + steering work is cleared, and a turn about to start is dropped), so a queued-but-not-yet-started prompt no longer runs and a later prompt cannot be batched into the cancelled turn. **Teardown/disconnect still use the older `agent.abort('disposed')` + `whenIdle()`**, so the best-effort window remains there: disposal/disconnect can return while one short queued turn per session still runs. PR D's per-agent disposer switches teardown to the queue-aware path and closes this; the single-in-flight-per-session rule bounds the worst case to one extra prompt per session until then. - **`TODO(rfc010-agent-disposal)`** — the factory (`ctx.agents.create`/`resume`) returns no per-agent disposer, so teardown aborts+drains each agent but cannot individually unregister it; on a bare client disconnect (no host dispose) the idled agents linger in `ctx.agents` until the host context disposes. A reconnect spins up a fresh context, so this strands no work; a per-agent disposal seam is the follow-up. - **`additionalDirectories`** — rejected. A session operates in its single `cwd` (see Per-session cwd); widening the tool/filesystem scope to extra roots is a separate sandbox concern, not yet implemented. diff --git a/packages/acp/src/index.ts b/packages/acp/src/index.ts index 780e907559..040f4b10d1 100644 --- a/packages/acp/src/index.ts +++ b/packages/acp/src/index.ts @@ -569,23 +569,19 @@ export function apply(ctx: Context, config: AcpConfig): void { cancel(params: CancelNotification): Promise { const rec = sessions.get(params.sessionId) if (rec === undefined) return Promise.resolve() - // RFC 010: session/cancel maps to agent.abort(reason). This aborts a - // RUNNING step (the turn ends 'aborted' → 'cancelled' via turn-end). - // It aborts and settles ONLY this session's agent/prompt — a cancel in - // one session never touches another's stream or pending prompt (RFC 011 - // isolation). It also settles the in-flight prompt as cancelled directly, - // in case the abort lands in the pre-step window (queued-but-not-started) - // where abort() has no AbortController to signal — see the README - // TODO(rfc010-cancel-prestep): a not-yet-started queued turn may still - // run to completion until a loop-level cancel lands. Best-effort abort - // plus honest RPC/UI cancellation. A secondary consequence of that same - // gap: because the loop batches all queued messages into one turn, a - // prompt accepted right after a pre-step cancel can be merged into the - // same turn as the cancelled one — that turn then carries both prompts' - // text and the new prompt settles for it. Both are closed by the same - // queue-aware loop cancel; the single-in-flight rule bounds the blast - // radius to one extra prompt. - rec.agent.abort('session/cancel') + // session/cancel maps to the queue-aware agent.cancel(reason): it aborts + // a RUNNING step, clears the queued + steering FIFOs, and drops a + // turn that is about to start (the pre-step window) — so a queued-but- + // not-yet-started prompt never runs, and a prompt accepted right after + // cannot be batched into the cancelled turn. Scoped to THIS session's + // agent — a cancel in one session never touches another's stream or + // pending prompt (RFC 011 isolation). We ALSO settle the in-flight prompt + // as cancelled directly here: do NOT rely on the resulting turn/end to + // settle it, because cancel() may drop the turn before any turn/end is + // emitted, and removing this direct settle would move the RPC's + // resolution onto the settleFromLog/agent-status path, changing its + // timing. + rec.agent.cancel('session/cancel') settlePrompt(rec, 'cancelled') return Promise.resolve() }, diff --git a/packages/acp/tests/turns.spec.ts b/packages/acp/tests/turns.spec.ts index 19c3be2556..014dfecf91 100644 --- a/packages/acp/tests/turns.spec.ts +++ b/packages/acp/tests/turns.spec.ts @@ -327,21 +327,64 @@ describe('acp bridge — turn outcomes', () => { expect(res.stopReason).toBe('cancelled') }) - it('cancel in the pre-step window still settles the prompt cancelled exactly once', async () => { - // No script entry is consumed before cancel: cancel immediately after the - // prompt is sent, before the model step starts. The prompt must still - // settle cancelled (best-effort abort + settle), not hang. - harness = await makeBridgeHarness({ storageDir, script: [textResponse('late')] }) + it('cancel right after prompt settles cancelled and leaves the agent idle, no leaked turn', async () => { + // Over the async JSON-RPC transport the loop usually wakes before cancel + // arrives, so this is a running/mid-step cancel (the synchronous pre-step + // DROP is unit-tested in agent-loop/cancel.spec.ts). The ACP-level guarantee: + // the prompt settles cancelled, the agent reaches idle, and no second/leaked + // turn runs afterward. + harness = await makeBridgeHarness({ storageDir, script: [textResponse('answer'), textResponse('leaked')] }) const sessionId = await newSession(harness) const promptDone = harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] }) await harness.client.cancel({ sessionId }) const res = await promptDone expect(res.stopReason).toBe('cancelled') - // The queued turn may still start after the cancel cleared the in-flight - // slot (the documented TODO(rfc010-cancel-prestep) best-effort window): its - // turn-start then fires with no prompt to tag, and the bridge does nothing. - // Let it run to completion and assert nothing re-settles (no throw, no hang). - await harness.ctx.agents.get(sessionId)!.whenIdle() + const agent = harness.ctx.agents.get(sessionId)! + await agent.whenIdle() + // At most ONE turn ran (the cancelled one) — the cancel cleared the queue, so + // no second turn was batched or leaked. (A best-effort abort that left queued + // work could have started a second turn.) + const turnStarts = agent.session.events.filter(e => e.type === 'turn/start').length + expect(turnStarts).toBeLessThanOrEqual(1) + }) + + it('idle session/cancel then session/prompt runs the prompt (no intervening whenIdle)', async () => { + // The ACP bridge settles the cancel RPC synchronously and accepts the next + // prompt WITHOUT awaiting quiescence — so this drives cancel→prompt with NO + // whenIdle() between, the production race. An idle cancel must be a no-op that + // does NOT drop the following prompt. + harness = await makeBridgeHarness({ storageDir, script: [textResponse('real answer')] }) + const sessionId = await newSession(harness) + // Cancel while idle (no prompt in flight) — a no-op. + await harness.client.cancel({ sessionId }) + // Immediately prompt, no whenIdle() between. + const res = await harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] }) + expect(res.stopReason).toBe('end_turn') + const text = harness.updates + .filter(u => u.sessionUpdate === 'agent_message_chunk') + .map(u => (u.content.type === 'text' ? u.content.text : '')) + .join('') + expect(text).toContain('real answer') + }) + + it('mid-stream cancel then an IMMEDIATE next prompt runs (no intervening whenIdle)', async () => { + // Cancel a running turn, then send the next prompt WITHOUT awaiting quiescence + // (the synchronous-settle path). The new prompt must run — the cancel marker + // must not leak onto it. + harness = await makeBridgeHarness({ storageDir, script: ['hang', textResponse('next answer')] }) + const sessionId = await newSession(harness) + const a = harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'A' }] }) + await new Promise(r => setTimeout(r, 30)) + await harness.client.cancel({ sessionId }) + expect((await a).stopReason).toBe('cancelled') + // Immediately — no whenIdle() — send the next prompt. + const b = await harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'B' }] }) + expect(b.stopReason).toBe('end_turn') + const text = harness.updates + .filter(u => u.sessionUpdate === 'agent_message_chunk') + .map(u => (u.content.type === 'text' ? u.content.text : '')) + .join('') + expect(text).toContain('next answer') }) it('a cancelled turn\'s late turn/end does not settle the NEXT prompt', async () => { diff --git a/packages/agent-loop/README.md b/packages/agent-loop/README.md index 308e5f022b..e5e3a03d2d 100644 --- a/packages/agent-loop/README.md +++ b/packages/agent-loop/README.md @@ -66,6 +66,8 @@ forever: Error containment: a throwing plugin ends the **turn**, never the loop. Dispose mid-turn emits `agent/status('disposed')` and ends with reason `disposed`. A step that hits the model's output-token ceiling makes the turn end `max-tokens` (the rule: any `max-tokens` step in the turn surfaces as `max-tokens`; `disposed`/`aborted`/`error` still take precedence) — distinct from a clean `completed` stop. +Cancellation: `agent.abort()` aborts only the in-flight step; `agent.cancel()` is the broad verb — it clears the queued + steering FIFOs, aborts the in-flight step, and drives a turn-scoped marker the driver checks at every point a turn could start or continue (right after the idle wait, after the `running` flip, before each step, and at the continuation gate) so a turn about to start is dropped. A cancelled turn ends `aborted`; a queued-but-not-started prompt never runs and cannot be batched into the cancelled turn. The marker is reset once per loop iteration, so a cancel governs exactly one turn and never leaks onto a later prompt. + ### What is NOT here Everything that goes beyond "call the model, run the tools, repeat" belongs to plugins listening on the event taxonomy: diff --git a/packages/agent-loop/src/agent.ts b/packages/agent-loop/src/agent.ts index e964685b86..72bc1b0455 100644 --- a/packages/agent-loop/src/agent.ts +++ b/packages/agent-loop/src/agent.ts @@ -26,6 +26,14 @@ export class ReactLoopAgent implements Agent { private _status: AgentStatus = 'idle' private currentAbort: AbortController | undefined + /** + * Turn-scoped cancel marker, set by {@link cancel} and read/cleared by the + * driver loop (via the LoopHandle) at every point a turn could start or + * continue. Armed ONLY when there is something to cancel (a running turn, an + * in-flight step, or queued/steering work), so an idle no-op cancel cannot + * leave it set to wrongly drop a later prompt. + */ + private cancelRequested = false private disposed: Promise private resolveDisposed!: () => void /** Resolves when the driver loop has fully exited (tests/disposal). */ @@ -176,6 +184,30 @@ export class ReactLoopAgent implements Agent { this.currentAbort?.abort(reason ?? 'aborted') } + cancel(reason?: string): void { + // Arm-gate: only mark a cancellation when there is actually work to cancel — + // a running turn, an in-flight step, or queued/steering work. An idle cancel + // with nothing pending is a true no-op; arming the marker then would wrongly + // drop the NEXT legitimate prompt (the marker is consumed only at the loop's + // turn-decision points, which an idle parked loop does not reach until woken + // by a real send()). Note the gate canNOT be `status === 'running'` alone: + // the pre-step window (a send() queued but the loop not yet flipped to + // running) has status `idle` with `hasQueued` true, and the marker exists + // precisely to cover it. + if (this._status === 'running' || this.currentAbort !== undefined || this.inbox.hasQueued || this.inbox.hasSteering) { + this.cancelRequested = true + } + // Drop all pending queued + steering work (un-started prompts never run; the + // cancelled turn's steering is not re-enqueued). Cleared directly even when + // the loop is parked in waitForQueued — there is no turn to stop and nothing + // left for the parked loop to run, so no wake is needed. + this.inbox.clear() + // Interrupt an in-flight step immediately (the running turn observes the + // abort and ends `aborted`). The marker covers the windows where no step is + // running (pre-step, continuation). + this.currentAbort?.abort(reason ?? 'cancelled') + } + /** * Resolve once the agent has reached quiescence after settling out of * `running`. If it is already disposed, awaits {@link done} (the loop-exit @@ -218,6 +250,15 @@ export class ReactLoopAgent implements Agent { setAbort: controller => void (this.currentAbort = controller), disposed: this.disposed, isDisposed: () => this._status === 'disposed', + isCancelled: () => this.cancelRequested, + clearCancel: () => { this.cancelRequested = false }, + // Settle whenIdle() waiters WITHOUT a status transition — the pre-step + // cancel-skip path drops the about-to-run turn and re-parks without ever + // flipping running→idle, so a waiter registered in the pre-step window + // (status idle, hasQueued was true) would otherwise hang. This emits no + // agent/status, so an ACP agent/status listener never sees a spurious idle + // that would resolve a freshly-queued prompt as cancelled. + settleIdle: () => { this.settleIdleWaiters() }, }) // The disposer must be infallible: it runs inside the fiber's LIFO // disposal chain, where a throw would skip later disposers (e.g. the diff --git a/packages/agent-loop/src/inbox.ts b/packages/agent-loop/src/inbox.ts index 7aabad0166..a7b2e64e2c 100644 --- a/packages/agent-loop/src/inbox.ts +++ b/packages/agent-loop/src/inbox.ts @@ -52,6 +52,16 @@ export class Inbox { return this.steeringMessages.splice(0) } + /** + * Discard all pending messages (queued + steering) without delivering them — + * used by `cancel()`, which drops un-started work rather than draining it into + * a turn. Unlike `drainQueued`/`drainSteering`, the messages are thrown away. + */ + clear(): void { + this.queuedMessages.length = 0 + this.steeringMessages.length = 0 + } + /** Wait until a queued message arrives or `cancel` resolves. */ waitForQueued(cancel: Promise): Promise { if (this.hasQueued) return Promise.resolve() diff --git a/packages/agent-loop/src/loop.ts b/packages/agent-loop/src/loop.ts index dc1cbcf278..906467e6b4 100644 --- a/packages/agent-loop/src/loop.ts +++ b/packages/agent-loop/src/loop.ts @@ -107,6 +107,26 @@ export interface LoopHandle { /** Resolves when the agent is disposed — unblocks the idle wait. */ disposed: Promise isDisposed(): boolean + /** + * Whether a `cancel()` is pending for the current turn. The driver checks this + * at every decision point where a turn could start or continue (right after + * the idle wait, after the `running` flip, before each step, and at the + * continuation gate) and drops the about-to-run / continuing turn. Reset once + * per loop iteration via {@link clearCancel} after the turn returns, so the + * marker governs exactly one cancellation and never leaks to a later prompt. + */ + isCancelled(): boolean + /** Clear the cancel marker (called once per iteration after the turn returns). */ + clearCancel(): void + /** + * Settle pending `whenIdle()` waiters WITHOUT a status transition. Used by the + * pre-step cancel-skip path: it drops the about-to-run turn and re-parks at the + * idle wait, so no `running→idle` transition fires to settle a `whenIdle()` + * waiter that was registered in the pre-step window — this settles it directly + * (it emits no `agent/status`, so an ACP `agent/status` listener never sees a + * spurious idle that would resolve a freshly-queued prompt as cancelled). + */ + settleIdle(): void } /** @@ -148,7 +168,33 @@ export async function runLoop(ctx: Context, agent: ReactLoopAgent, handle: LoopH await agent.inbox.waitForQueued(handle.disposed) if (handle.isDisposed()) break + // Pre-step cancel (window 1): a `cancel()` landed after a `send()` woke the + // idle wait but before we flip to `running`. Drop the about-to-run turn: the + // queued/steering work is already cleared by `cancel()`, and we settle any + // `whenIdle()` waiter DIRECTLY (no status transition fires here, so the + // running→idle settle never runs) WITHOUT emitting `agent/status` (an ACP + // listener must not see a spurious idle that resolves a freshly-queued prompt + // as cancelled). Clear the marker and re-park. + if (handle.isCancelled()) { + handle.clearCancel() + handle.settleIdle() + continue + } + handle.setStatus('running') + + // Pre-step cancel (window 2): `setStatus('running')` emits `agent/status` + // SYNCHRONOUSLY, so a `running` listener can `cancel()` in the gap between the + // check above and `runTurn`. `cancel()` already cleared the queued FIFO, so + // drop the turn before it starts (runTurn would otherwise throw on an empty + // queue) and transition back to idle — `running` was already emitted, so a + // real `idle` transition (which also settles waiters) balances the status. + if (handle.isCancelled()) { + handle.clearCancel() + handle.setStatus('idle') + continue + } + // Re-derive the turn number from the log each iteration (do NOT keep a local // counter): an idle `agent.inject()` can append its own one-shot turn while // the loop waits above, so the next real turn must continue from whatever @@ -170,8 +216,18 @@ export async function runLoop(ctx: Context, agent: ReactLoopAgent, handle: LoopH } catch { /* contained: a throwing agent/error listener must not kill the driver */ } } + // Reset the cancel marker UNCONDITIONALLY here, after the turn returns and + // before the next iteration's idle wait. NOT gated on the idle transition + // below: a `send()` that lands during the cancelled turn's flush window makes + // `hasQueued` true at the `setStatus('idle')` guard, so an idle-gated reset + // would never fire and the stale marker would wrongly drop that next prompt's + // turn. Resetting per iteration scopes the marker to exactly the turn that was + // cancelled. + handle.clearCancel() + // Steering that arrived too late to join this turn (turn-end listeners, - // flush) becomes a queued message — it must never be stranded. + // flush) becomes a queued message — it must never be stranded. (A cancelled + // turn already cleared its steering, so there is nothing to re-enqueue.) for (const message of agent.inbox.drainSteering()) { agent.inbox.enqueue(message) } @@ -323,6 +379,20 @@ async function runTurn(ctx: Context, agent: ReactLoopAgent, handle: LoopHandle, const abort = new AbortController() handle.setAbort(abort) + // Cancel landing in the step-start window: a synchronous `agent/turn-start` + // or `agent/step-start` listener (both fire before this point) can have + // called `cancel()`, and `runStep` would otherwise run a full extra step + // with no AbortController having observed it. Check the marker AFTER + // setAbort (so the next-iteration drain sees a clean controller) and before + // `runStep`: drop the step, end the turn `aborted`. closeStep balances the + // already-appended step/start. + if (handle.isCancelled()) { + handle.setAbort(undefined) + reason = { kind: 'aborted', reason: 'cancelled' } + closeStep() + break + } + let stepOutcome: { hadToolCalls: boolean; finish: FinishReason } | { error: Error } try { stepOutcome = await runStep(ctx, agent, turn, step, abort.signal) @@ -382,6 +452,16 @@ async function runTurn(ctx: Context, agent: ReactLoopAgent, handle: LoopHandle, // next iteration's drain records it. if (!shouldContinue && agent.inbox.hasSteering) shouldContinue = true + // A cancel that landed during the continuation window — after the step's + // AbortController was cleared (setAbort(undefined)) but before the next + // step starts — has no controller to observe it, so the turn-scoped marker + // ends the turn here. cancel() also cleared the steering FIFO, so the + // override above did not re-arm continuation. + if (handle.isCancelled()) { + reason = { kind: 'aborted', reason: 'cancelled' } + break + } + if (!shouldContinue || handle.isDisposed()) { /* v8 ignore next -- disposal during continuation-decision window is a narrow race; error-path disposal is covered elsewhere */ if (handle.isDisposed()) reason = { kind: 'disposed' } diff --git a/packages/agent-loop/tests/cancel.spec.ts b/packages/agent-loop/tests/cancel.spec.ts new file mode 100644 index 0000000000..64b969811d --- /dev/null +++ b/packages/agent-loop/tests/cancel.spec.ts @@ -0,0 +1,279 @@ +/** + * Tests for the queue-aware `Agent.cancel()` primitive (PR C). `cancel()` is the + * broad verb — it clears queued + steering work, aborts an in-flight step, and + * drops a turn about to start — whereas `abort()` kills only the current step. + * These tests exercise every window where a cancel can land (idle, pre-step, + * mid-step, continuation) and the marker's arm/reset rules that keep a cancel + * from leaking to a later prompt or hanging `whenIdle()`. + * + * @module dsh-agent-loop/tests/cancel + */ + +import { describe, expect, it } from 'vitest' +import { Context } from 'cordis' +import LlmService from '@deepseek-ai/dsh-llm' +import SessionStore, { TurnEndReason } from '@deepseek-ai/dsh-session' +import SystemPrompt from '@deepseek-ai/dsh-system-prompt' +import ToolRegistry from '@deepseek-ai/dsh-tools' +import AgentRegistry from '@deepseek-ai/dsh-agent' +import AgentLoop, { ReactLoopAgent } from '@deepseek-ai/dsh-agent-loop' +import { MockAdapter, textResponse } from './mock-adapter.ts' + +async function harness(adapter: MockAdapter) { + const ctx = new Context() + await ctx.plugin(LlmService) + await ctx.plugin(SessionStore) + await ctx.plugin(SystemPrompt) + await ctx.plugin(ToolRegistry) + await ctx.plugin(AgentRegistry) + await ctx.plugin(AgentLoop, { agents: [] }) + ctx.llm.registerAdapter(['mock'], adapter) + return ctx +} + +function send(agent: ReactLoopAgent, text: string) { + agent.send([{ type: 'text', text }]) +} + +/** Resolve on the agent's next idle transition (event-based, not status poll). */ +function waitForIdle(ctx: Context, agent: ReactLoopAgent): Promise { + return new Promise((resolve) => { + const dispose = ctx.on('agent/status', (subject, status) => { + if (subject === agent && status === 'idle') { dispose(); resolve() } + }) + }) +} + +/** All user-message texts recorded in the log (to assert what actually ran). */ +function userTexts(agent: ReactLoopAgent): string[] { + return agent.session.events + .filter(e => e.type === 'user/message') + .flatMap(e => e.type === 'user/message' ? e.data.content : []) + .flatMap(b => b.type === 'text' ? [b.text] : []) +} + +describe('Agent.cancel()', () => { + it('cancel() on an idle agent with nothing queued is a no-op; the next prompt runs (F2 leak guard)', async () => { + const adapter = new MockAdapter([textResponse('reply')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + // The loop is parked at the idle wait with nothing queued. A cancel here must + // NOT arm the marker — otherwise the next legitimate prompt would be dropped. + agent.cancel('nothing to cancel') + + send(agent, 'real prompt') + await waitForIdle(ctx, agent) + + // The prompt ran: its user message is in the log and one turn completed. + expect(userTexts(agent)).toEqual(['real prompt']) + expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true) + }) + + it('pre-step cancel drops the about-to-start turn (no turn is opened)', async () => { + const adapter = new MockAdapter([textResponse('should not run')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + // send() queues synchronously (status still idle, loop microtask not yet + // resumed). Cancel in that pre-step window: the queued turn must not run. + send(agent, 'drop me') + agent.cancel('pre-step') + + // Give the loop a chance to wake and process the cancel. + await new Promise(r => setTimeout(r, 30)) + + // No turn was opened — the queued prompt was dropped, never recorded. + expect(userTexts(agent)).toEqual([]) + expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false) + expect(agent.status).toBe('idle') + }) + + it('a whenIdle() waiter registered BEFORE a pre-step cancel resolves (F1 hang guard)', async () => { + const adapter = new MockAdapter([textResponse('x')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + // Queue work, then register a whenIdle() waiter while in the pre-step window + // (status idle, hasQueued true) — it does NOT take the fast path. Then cancel. + // The skip path must settle this waiter directly (no running→idle transition + // ever fires), or it would hang forever. + send(agent, 'q') + const idle = agent.whenIdle() + agent.cancel('pre-step') + + // Must resolve (not hang). A timeout makes the failure a clear test failure. + await Promise.race([ + idle, + new Promise((_r, reject) => setTimeout(() => { reject(new Error('whenIdle hung after pre-step cancel')) }, 1000)), + ]) + expect(agent.status).toBe('idle') + }) + + it('cancel() mid-step aborts the in-flight model call; the turn ends aborted', async () => { + const adapter = new MockAdapter(['hang']) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + const reasons: TurnEndReason[] = [] + ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason)) + + send(agent, 'go') + await new Promise(r => setTimeout(r, 30)) + expect(agent.status).toBe('running') + agent.cancel('mid-step') + await waitForIdle(ctx, agent) + + expect(reasons).toEqual([{ kind: 'aborted', reason: 'mid-step' }]) + }) + + it('cancel() with no reason defaults to "cancelled" when aborting an in-flight step', async () => { + const adapter = new MockAdapter(['hang']) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + const reasons: TurnEndReason[] = [] + ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason)) + + send(agent, 'go') + await new Promise(r => setTimeout(r, 30)) + agent.cancel() // no reason → default 'cancelled' + await waitForIdle(ctx, agent) + + expect(reasons).toEqual([{ kind: 'aborted', reason: 'cancelled' }]) + }) + + it('a prompt sent AFTER a cancelled turn settles runs normally (marker reset)', async () => { + const adapter = new MockAdapter(['hang', textResponse('second reply')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + // First turn hangs; cancel it mid-step. + send(agent, 'first') + await new Promise(r => setTimeout(r, 30)) + agent.cancel('cancel first') + await waitForIdle(ctx, agent) + + // The marker must have been reset after the cancelled turn — a fresh prompt + // runs to completion rather than being dropped by a stale marker. + send(agent, 'second') + await waitForIdle(ctx, agent) + + expect(userTexts(agent)).toContain('second') + // The second turn completed (its reply was streamed). + const reasons = agent.session.events.filter(e => e.type === 'turn/end') + expect(reasons.length).toBe(2) + }) + + it('cancel from a synchronous agent/turn-start listener drops the step (step-start window)', async () => { + const adapter = new MockAdapter([textResponse('should not stream')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + // A turn-start listener fires BEFORE any AbortController is installed for the + // step. Cancelling there must still drop the step (the turn-scoped marker, + // not abort(), is what catches this) — no model step runs. + let streamed = false + ctx.on('agent/stream-chunk', () => { streamed = true }) + const dispose = ctx.on('agent/turn-start', (subject) => { + if (subject === agent) agent.cancel('from turn-start') + }) + + const reasons: TurnEndReason[] = [] + ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason)) + + send(agent, 'go') + await waitForIdle(ctx, agent) + dispose() + + // No step streamed (the model never ran), and the turn ended aborted. + expect(streamed).toBe(false) + expect(reasons).toEqual([{ kind: 'aborted', reason: 'cancelled' }]) + }) + + it('cancel during the continuation window ends the turn aborted and runs no further step', async () => { + // A continuation-waterfall listener cancels DURING the continuation decision + // (the finished step's AbortController is already cleared), and votes to + // continue — but the turn-scoped marker checked right after must end the turn + // `aborted` and run NO second step. + const adapter = new MockAdapter([textResponse('one'), textResponse('two')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + let steps = 0 + ctx.on('agent/step-start', () => { steps += 1 }) + const reasons: TurnEndReason[] = [] + ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason)) + + let continued = false + ctx.on('agent/turn-continuation', async (subject, _turn, _default, next) => { + if (subject === agent && !continued) { + continued = true + agent.cancel('from continuation') + return true // vote to continue — the post-waterfall marker check must override + } + return next() + }) + + send(agent, 'go') + await waitForIdle(ctx, agent) + + // Only ONE step ran (the second was cancelled in the continuation window), + // and the turn ended aborted. + expect(steps).toBe(1) + expect(reasons).toEqual([{ kind: 'aborted', reason: 'cancelled' }]) + }) + + it('cancel from a synchronous agent/status(running) listener drops the turn (window 2)', async () => { + const adapter = new MockAdapter([textResponse('should not run')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + // setStatus('running') emits agent/status SYNCHRONOUSLY, so a running + // listener can cancel in the gap between the loop's pre-step check and + // runTurn. The second check (after the running flip) must drop the turn — + // runTurn would otherwise throw on the now-empty queue. + let streamed = false + ctx.on('agent/stream-chunk', () => { streamed = true }) + const dispose = ctx.on('agent/status', (subject, status) => { + if (subject === agent && status === 'running') agent.cancel('from running listener') + }) + + send(agent, 'go') + await waitForIdle(ctx, agent) + dispose() + + // No turn opened, no step streamed, and a later prompt still runs (the marker + // was reset). + expect(streamed).toBe(false) + expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false) + }) + + it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => { + const adapter = new MockAdapter(['hang']) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + send(agent, 'go') + await new Promise(r => setTimeout(r, 30)) + expect(agent.status).toBe('running') + // Steer (joins the running turn's steering FIFO), then cancel: the steering + // must be dropped, NOT re-enqueued as a new queued turn. + agent.steer([{ type: 'text', text: 'steer text' }]) + agent.cancel('cancel with steering') + await waitForIdle(ctx, agent) + + // After the cancelled turn settles, the agent is idle with NO follow-up turn + // started from the dropped steering. + await new Promise(r => setTimeout(r, 30)) + expect(agent.status).toBe('idle') + const turnStarts = agent.session.events.filter(e => e.type === 'turn/start') + expect(turnStarts.length).toBe(1) // only the original (cancelled) turn + // The steering text was dropped — it never reached the log. + const flat = agent.session.events + .filter(e => e.type === 'steering/message') + .flatMap(e => e.type === 'steering/message' ? e.data.content : []) + .flatMap(b => b.type === 'text' ? [b.text] : []) + expect(flat).not.toContain('steer text') + }) +}) diff --git a/packages/agent/README.md b/packages/agent/README.md index 4380cd756d..d815761ae4 100644 --- a/packages/agent/README.md +++ b/packages/agent/README.md @@ -54,7 +54,8 @@ The handle every plugin programs against: - `agent.send(content, options?)` — queue a message; starts a turn when idle - `agent.steer(content, options?)` — steer a running turn (inject between steps); behaves like `send` when idle - `agent.inject(content, options?)` — inject in-session context (context/message event); the next request sees it. Does not run the model. While a turn is open it joins that turn; while idle it is wrapped in a one-shot `injection` turn so every event stays turn-enclosed ([the turn-enclosure invariant](../../docs/rfc/implemented/2026-06-15-turn-enclosure-invariant.md)) -- `agent.abort(reason?)` — abort the in-flight step +- `agent.abort(reason?)` — abort the in-flight step (the narrow, step-only verb) +- `agent.cancel(reason?)` — cancel ALL pending work: clears the queued + steering FIFOs, aborts the in-flight step, and drops a turn about to start (the pre-step window) so a queued-but-not-started prompt never runs. A UI/ACP `session/cancel` maps to this. Idle with nothing pending → a safe no-op. - `agent.whenIdle()` — resolve once the agent reaches quiescence after settling out of `running` (idle → immediately; disposed → awaits the loop exit), the signal a teardown awaits (`abort()` then `await whenIdle()`). Observes the transition without disposing the agent. - `agent.session`, `agent.status`, `agent.options`, `agent.id` diff --git a/packages/agent/src/types.ts b/packages/agent/src/types.ts index 762d325068..0ebf582bda 100644 --- a/packages/agent/src/types.ts +++ b/packages/agent/src/types.ts @@ -81,6 +81,25 @@ export interface Agent { /** Abort the in-flight step (if any); the turn ends with reason 'aborted'. */ abort(reason?: string): void + /** + * Cancel ALL pending work for the agent — the narrower {@link abort} kills + * only the in-flight step. `cancel()`: + * + * - clears the queued FIFO (un-started prompts never run) and the steering + * FIFO (steering for the cancelled turn is dropped, not re-enqueued); + * - aborts the in-flight step if one is running (the turn ends `aborted`); + * - drops a turn that is about to start (a `cancel()` landing in the + * pre-step window — after a `send()` queued but before the loop flips to + * `running`, or after `running` is emitted but before the first step) so + * that queued prompt does not run and cannot be batched into the cancelled + * turn. + * + * After `cancel()`, `whenIdle()` resolves on the post-cancel quiescent state. + * `cancel()` on an idle agent with nothing queued or running is a safe no-op + * — it does NOT arm anything that would drop a later legitimate prompt. + */ + cancel(reason?: string): void + /** * Resolve once the agent has reached quiescence after settling out of * `running`, or immediately if it is already idle with no queued work. The diff --git a/packages/agent/tests/agent.spec.ts b/packages/agent/tests/agent.spec.ts index 1994c785a6..28072154f6 100644 --- a/packages/agent/tests/agent.spec.ts +++ b/packages/agent/tests/agent.spec.ts @@ -14,6 +14,7 @@ function stubAgent(rawId: string): Agent { steer() {}, inject() {}, abort() {}, + cancel() {}, whenIdle() { return Promise.resolve() }, } } From 9ee22bc6f69c9912c327c4cafbdd55e3d04c0823 Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Sat, 20 Jun 2026 05:10:16 +0800 Subject: [PATCH 2/5] fix(agent): don't resolve whenIdle() early on pre-step cancel + requeue (Codex review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Codex's converge pass found a quiescence-contract violation: a whenIdle() waiter registered for prompt A, then cancel() clears A, then prompt B is queued BEFORE the loop resumes from the idle wait. The window-1 cancel branch called settleIdle() UNCONDITIONALLY, resolving the waiter while B was still queued-and-unrun — whenIdle() resolved with zero events, then B ran afterward. Fix: in window 1, only settleIdle() + re-park when NO new work is queued. If a send() raced in after the cancel, the marker was for the cancelled work only — clear it and fall through to run the new prompt's turn, letting THAT turn's running→idle settle the waiter (so whenIdle() waits for B to actually run). Adds a regression test reproducing the exact interleaving (send A → whenIdle → cancel → send B): whenIdle() now resolves only after B's turn ran (B's user message + a turn/end in the log), and A was dropped. --- packages/agent-loop/src/loop.ts | 24 +++++++++++++++-------- packages/agent-loop/tests/cancel.spec.ts | 25 ++++++++++++++++++++++++ 2 files changed, 41 insertions(+), 8 deletions(-) diff --git a/packages/agent-loop/src/loop.ts b/packages/agent-loop/src/loop.ts index 906467e6b4..1c78f49beb 100644 --- a/packages/agent-loop/src/loop.ts +++ b/packages/agent-loop/src/loop.ts @@ -169,16 +169,24 @@ export async function runLoop(ctx: Context, agent: ReactLoopAgent, handle: LoopH if (handle.isDisposed()) break // Pre-step cancel (window 1): a `cancel()` landed after a `send()` woke the - // idle wait but before we flip to `running`. Drop the about-to-run turn: the - // queued/steering work is already cleared by `cancel()`, and we settle any - // `whenIdle()` waiter DIRECTLY (no status transition fires here, so the - // running→idle settle never runs) WITHOUT emitting `agent/status` (an ACP - // listener must not see a spurious idle that resolves a freshly-queued prompt - // as cancelled). Clear the marker and re-park. + // idle wait but before we flip to `running`. The cancelled queued/steering + // work is already cleared by `cancel()`. Clear the marker, then: + // - if NOTHING new is queued, drop the about-to-run turn and re-park, + // settling any `whenIdle()` waiter DIRECTLY (no running→idle transition + // fires here to settle it) and WITHOUT emitting `agent/status` (an ACP + // listener must not see a spurious idle that resolves a freshly-queued + // prompt as cancelled); + // - if a NEW prompt was queued AFTER the cancel (a send() that raced in + // before the loop resumed), the marker was for the cancelled work only — + // fall through and run the new prompt's turn. Do NOT settle waiters here: + // a whenIdle() waiter must wait for that new turn's running→idle, not + // resolve before it runs (the quiescence contract). if (handle.isCancelled()) { handle.clearCancel() - handle.settleIdle() - continue + if (!agent.inbox.hasQueued) { + handle.settleIdle() + continue + } } handle.setStatus('running') diff --git a/packages/agent-loop/tests/cancel.spec.ts b/packages/agent-loop/tests/cancel.spec.ts index 64b969811d..b8c7a4f78e 100644 --- a/packages/agent-loop/tests/cancel.spec.ts +++ b/packages/agent-loop/tests/cancel.spec.ts @@ -249,6 +249,31 @@ describe('Agent.cancel()', () => { expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false) }) + it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => { + // The subtle race: a whenIdle() waiter is registered for prompt A; cancel() + // clears A; prompt B is queued BEFORE the loop resumes from the idle wait. + // The window-1 cancel branch must NOT settle the waiter while B is still + // queued-and-unrun — whenIdle() must wait for B's turn to actually run and + // settle (the quiescence contract), not resolve before B's first event. + const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + send(agent, 'A') // queues A (status still idle, loop microtask pending) + const idle = agent.whenIdle() // registers a waiter (idle + hasQueued → no fast path) + agent.cancel('drop A') // arms marker, clears A + send(agent, 'B') // B races in before the loop resumes + + // whenIdle() must resolve only AFTER B's turn fully ran — by which point B's + // user message and a turn/end are in the log. (Before the fix it resolved + // immediately, with zero events, then B ran afterward.) + await idle + expect(userTexts(agent)).toContain('B') + expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true) + // A was dropped (never ran); only B's turn is recorded. + expect(userTexts(agent)).not.toContain('A') + }) + it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => { const adapter = new MockAdapter(['hang']) const ctx = await harness(adapter) From 6a1e381e38c6683d59dff18271efd5951cd78846 Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Sat, 20 Jun 2026 11:51:32 +0800 Subject: [PATCH 3/5] fix(agent): carry cancel(reason) through the marker-only windows (review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A reviewer found that `cancel(reason)` only preserved the caller's reason when an active AbortController observed it (the mid-step path, via `abort.signal.reason`). The marker-only windows (step-start at loop.ts and the continuation gate) hardcoded `reason: 'cancelled'`, so the logged `turn/end` reason was race-dependent on WHERE the cancel landed and the public `cancel(reason?)` parameter was half-effective. Capture the resolved reason (`reason ?? 'cancelled'`) on the agent when the marker is armed, expose it on the LoopHandle as `cancelReason()`, and use it in both marker branches so a turn dropped without a live controller records the SAME `{kind:'aborted', reason}` the mid-step path produces. The two existing window tests asserted `reason: 'cancelled'` while passing `'from turn-start'` / `'from continuation'` — they documented the bug. Updated both to assert the caller's reason (behavior + test changed together, per AGENTS.md "tests document behavior, not golden truth"). Also fixes two stale docs the PR's contract change left behind: the module-level ACP mapping comment and `codec.ts` both still said `session/cancel -> agent.abort()`. --- packages/acp/src/codec.ts | 2 +- packages/acp/src/index.ts | 4 +++- packages/agent-loop/src/agent.ts | 16 ++++++++++++++++ packages/agent-loop/src/loop.ts | 12 ++++++++++-- packages/agent-loop/tests/cancel.spec.ts | 11 +++++++---- 5 files changed, 37 insertions(+), 8 deletions(-) diff --git a/packages/acp/src/codec.ts b/packages/acp/src/codec.ts index dd28bc5ae6..5f5a53f529 100644 --- a/packages/acp/src/codec.ts +++ b/packages/acp/src/codec.ts @@ -26,7 +26,7 @@ import type { ContentBlock as AcpContentBlock, StopReason } from '@agentclientpr * * - `completed` → `end_turn` (the model chose to stop) * - `max-tokens` → `max_tokens` (cut off at the output-token ceiling) - * - `aborted` → `cancelled` (an `agent.abort()`, e.g. from `session/cancel`) + * - `aborted` → `cancelled` (a step abort or a queue-aware `agent.cancel()`, e.g. from `session/cancel`) * - `error` → `end_turn` (defensive fallback only: the bridge REJECTS the * `session/prompt` RPC on an error turn BEFORE calling this, so * a client sees a JSON-RPC error, not a stop reason — see diff --git a/packages/acp/src/index.ts b/packages/acp/src/index.ts index 040f4b10d1..91553669ca 100644 --- a/packages/acp/src/index.ts +++ b/packages/acp/src/index.ts @@ -13,7 +13,9 @@ * - `session/load` → `ctx.agents.resume(...)` then replay the event log * - `session/prompt` → `agent.send()`, settle on the owning turn's end (a turn * that ends in `error` rejects the RPC) - * - `session/cancel` → `agent.abort()` + settle the in-flight prompt + * - `session/cancel` → `agent.cancel()` (the queue-aware cancel: aborts a + * running step, clears queued + steering work, and drops a + * turn about to start) + settle the in-flight prompt * * Multi-session (RFC 011): N concurrent sessions per connection, each mapped to * its own `ReactLoopAgent`. Sessions are keyed by id in `sessions` (forward) with an diff --git a/packages/agent-loop/src/agent.ts b/packages/agent-loop/src/agent.ts index 72bc1b0455..df25af1c11 100644 --- a/packages/agent-loop/src/agent.ts +++ b/packages/agent-loop/src/agent.ts @@ -34,6 +34,17 @@ export class ReactLoopAgent implements Agent { * leave it set to wrongly drop a later prompt. */ private cancelRequested = false + /** + * The resolved reason for the pending {@link cancel} (`reason ?? 'cancelled'`), + * read by the driver loop's marker branches so a turn dropped in a + * marker-only window (pre-step / continuation, where no `AbortController` + * carries the reason) ends with the SAME `{kind:'aborted', reason}` the + * mid-step abort path produces from `abort.signal.reason`. Without this the + * caller's `cancel(reason)` would be silently replaced by the literal + * 'cancelled' whenever the cancel landed outside a running step — making the + * logged reason race-dependent and the public `reason?` param half-effective. + */ + private cancelReason = 'cancelled' private disposed: Promise private resolveDisposed!: () => void /** Resolves when the driver loop has fully exited (tests/disposal). */ @@ -196,6 +207,10 @@ export class ReactLoopAgent implements Agent { // precisely to cover it. if (this._status === 'running' || this.currentAbort !== undefined || this.inbox.hasQueued || this.inbox.hasSteering) { this.cancelRequested = true + // Capture the resolved reason for the marker-only windows (pre-step / + // continuation). The mid-step path reads it from abort.signal.reason + // below; the marker path reads it via the LoopHandle's cancelReason(). + this.cancelReason = reason ?? 'cancelled' } // Drop all pending queued + steering work (un-started prompts never run; the // cancelled turn's steering is not re-enqueued). Cleared directly even when @@ -251,6 +266,7 @@ export class ReactLoopAgent implements Agent { disposed: this.disposed, isDisposed: () => this._status === 'disposed', isCancelled: () => this.cancelRequested, + cancelReason: () => this.cancelReason, clearCancel: () => { this.cancelRequested = false }, // Settle whenIdle() waiters WITHOUT a status transition — the pre-step // cancel-skip path drops the about-to-run turn and re-parks without ever diff --git a/packages/agent-loop/src/loop.ts b/packages/agent-loop/src/loop.ts index 1c78f49beb..3a6ac76a74 100644 --- a/packages/agent-loop/src/loop.ts +++ b/packages/agent-loop/src/loop.ts @@ -116,6 +116,14 @@ export interface LoopHandle { * marker governs exactly one cancellation and never leaks to a later prompt. */ isCancelled(): boolean + /** + * The resolved reason for the pending cancel (`reason ?? 'cancelled'`), read + * by the marker branches (pre-step / continuation) so a turn dropped where no + * `AbortController` carries the reason still records the caller's + * `cancel(reason)` value — matching the mid-step abort path. Only meaningful + * when {@link isCancelled} is true. + */ + cancelReason(): string /** Clear the cancel marker (called once per iteration after the turn returns). */ clearCancel(): void /** @@ -396,7 +404,7 @@ async function runTurn(ctx: Context, agent: ReactLoopAgent, handle: LoopHandle, // already-appended step/start. if (handle.isCancelled()) { handle.setAbort(undefined) - reason = { kind: 'aborted', reason: 'cancelled' } + reason = { kind: 'aborted', reason: handle.cancelReason() } closeStep() break } @@ -466,7 +474,7 @@ async function runTurn(ctx: Context, agent: ReactLoopAgent, handle: LoopHandle, // ends the turn here. cancel() also cleared the steering FIFO, so the // override above did not re-arm continuation. if (handle.isCancelled()) { - reason = { kind: 'aborted', reason: 'cancelled' } + reason = { kind: 'aborted', reason: handle.cancelReason() } break } diff --git a/packages/agent-loop/tests/cancel.spec.ts b/packages/agent-loop/tests/cancel.spec.ts index b8c7a4f78e..ab05c4e25b 100644 --- a/packages/agent-loop/tests/cancel.spec.ts +++ b/packages/agent-loop/tests/cancel.spec.ts @@ -186,9 +186,11 @@ describe('Agent.cancel()', () => { await waitForIdle(ctx, agent) dispose() - // No step streamed (the model never ran), and the turn ended aborted. + // No step streamed (the model never ran), and the turn ended aborted with + // the CALLER's reason — the marker carries `cancel(reason)` through even + // though no AbortController observed it in this window. expect(streamed).toBe(false) - expect(reasons).toEqual([{ kind: 'aborted', reason: 'cancelled' }]) + expect(reasons).toEqual([{ kind: 'aborted', reason: 'from turn-start' }]) }) it('cancel during the continuation window ends the turn aborted and runs no further step', async () => { @@ -219,9 +221,10 @@ describe('Agent.cancel()', () => { await waitForIdle(ctx, agent) // Only ONE step ran (the second was cancelled in the continuation window), - // and the turn ended aborted. + // and the turn ended aborted with the CALLER's reason (carried by the + // marker, since the finished step's AbortController was already cleared). expect(steps).toBe(1) - expect(reasons).toEqual([{ kind: 'aborted', reason: 'cancelled' }]) + expect(reasons).toEqual([{ kind: 'aborted', reason: 'from continuation' }]) }) it('cancel from a synchronous agent/status(running) listener drops the turn (window 2)', async () => { From f58b031465c7def8a9b0809209cc62db4f0168ed Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Sat, 20 Jun 2026 12:57:32 +0800 Subject: [PATCH 4/5] fix(agent): close the window-2 early-whenIdle race + sync cancellation RFC docs (review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A reviewer found that window 2 (a cancel from a synchronous agent/status('running') listener) had the same early-whenIdle() race that window 1 already guards: it unconditionally `setStatus('idle')` + continue, which settles `whenIdle()` waiters — so if the running listener cancels AND queues replacement work, the waiter resolves while the replacement is still queued-and-unrun (the next iteration runs it later, but the caller already observed quiescence). Mirror window 1: after clearing the marker, only `setStatus('idle')` when nothing new is queued; otherwise fall through to run the queued replacement (status is already `running`), so `whenIdle()` resolves on that turn's running→idle. Regression test reproduces the reviewer's interleaving (running listener cancels A, sends B; whenIdle() resolves only after B ran). Also syncs the cancellation contract in the two ACP RFCs that describe the live behavior: `session/cancel` is the queue-aware `agent.cancel()` (drops an about-to-start turn), not the old best-effort `agent.abort()` pre-step limitation. --- .../2026-06-14-acp-agent-client-protocol.md | 4 +-- .../proposed/2026-06-14-acp-multi-session.md | 2 +- packages/agent-loop/src/loop.ts | 20 +++++++++---- packages/agent-loop/tests/cancel.spec.ts | 30 +++++++++++++++++++ 4 files changed, 47 insertions(+), 9 deletions(-) diff --git a/docs/rfc/proposed/2026-06-14-acp-agent-client-protocol.md b/docs/rfc/proposed/2026-06-14-acp-agent-client-protocol.md index 3e68024fa9..3dc779ffa4 100644 --- a/docs/rfc/proposed/2026-06-14-acp-agent-client-protocol.md +++ b/docs/rfc/proposed/2026-06-14-acp-agent-client-protocol.md @@ -3,7 +3,7 @@ Status: proposed -> **Implementation status (MVP landed):** steps 1, 2, 3, 4, 6, 7, 8 are implemented in `packages/acp` + `examples/acp-agent`. **Step 5 (the `session/request_permission` permission gate) is deferred** — the bridge ships a pass-through (tools run with the executor's full authority) marked `TODO(rfc010-permission-gate)`, and lays down only the `WeakMap` ownership seam the gate will build on. Status stays `proposed` until the gate lands. One further best-effort limitation is tracked as `TODO(rfc010-cancel-prestep)`: `session/cancel` aborts a running step and settles the RPC as `cancelled`, but a turn still queued (not yet started) when the cancel arrives may execute before the abort takes effect, pending a loop-level pre-step cancel. **Per-session `cwd` is now honored** (lifting the original "launch the server in the workspace root" restriction — see § Deferred): `session/new` accepts any absolute `cwd`, and `session/load` requires the request `cwd` to match the persisted session `cwd` so the editor and bash executor agree on the workspace. +> **Implementation status (MVP landed):** steps 1, 2, 3, 4, 6, 7, 8 are implemented in `packages/acp` + `examples/acp-agent`. **Step 5 (the `session/request_permission` permission gate) is deferred** — the bridge ships a pass-through (tools run with the executor's full authority) marked `TODO(rfc010-permission-gate)`, and lays down only the `WeakMap` ownership seam the gate will build on. Status stays `proposed` until the gate lands. `session/cancel` is the queue-aware `agent.cancel()`: it aborts a running step, clears queued + steering work, and drops a turn that is about to start, so a queued-but-not-yet-started prompt never runs and a later prompt cannot be batched into the cancelled turn. **Per-session `cwd` is now honored** (lifting the original "launch the server in the workspace root" restriction — see § Deferred): `session/new` accepts any absolute `cwd`, and `session/load` requires the request `cwd` to match the persisted session `cwd` so the editor and bash executor agree on the workspace. ## Problem @@ -33,7 +33,7 @@ The mapping between ACP and existing harness seams — each row names the seam a | `session/update: tool_call` (pending→in_progress) | `session/event` `tool/call` | demux via a Session→sessionId map; `kind` inferred from the tool name | | `session/update: tool_call_update` (completed/failed) | `session/event` `tool/result` | a throwing `tools/execute` yields NO `tool/result` → fail the pending tool UI from `agent/error`/turn-end | | `session/request_permission {sessionId, toolCall, options}` | prepended `tools/execute` listener | no-op unless `exec.agent` is ACP-owned; await the outcome; `selected/allow_*` → `next()`; `reject_*`/`cancelled` → veto `ToolExecutionResult{isError}` | -| `session/cancel` (notification) | `agent.abort(reason)` | settle the in-flight prompt as `cancelled`; resolve any pending permission as `cancelled` exactly once | +| `session/cancel` (notification) | `agent.cancel(reason)` | the queue-aware cancel (abort running step, clear queued + steering, drop an about-to-start turn); settle the in-flight prompt as `cancelled`; resolve any pending permission as `cancelled` exactly once | The permission gate is the first real consumer of the `tools/execute` veto seam (the documented "single veto/sandbox/permission seam" plus the deferred "Permission system" TODO in [docs/architecture.md](../../architecture.md)). It is a single global listener registered with `prepend: true` so it runs before any other tool wrapper. `ToolExecution.agent` is optional and the `Agent` interface carries no origin marker, so the bridge tracks ownership itself: it records each agent it creates in a `WeakMap` and the gate no-ops (calls `next()` immediately) for any `exec.agent` it does not own — non-ACP agents and the no-agent case pass straight through. For an owned agent it resolves the session, issues `session/request_permission`, and stores the pending resolver on that session's record so the outcome — or a `session/cancel`/connection-close — settles it exactly once. diff --git a/docs/rfc/proposed/2026-06-14-acp-multi-session.md b/docs/rfc/proposed/2026-06-14-acp-multi-session.md index 8cc2237eff..25fe24c972 100644 --- a/docs/rfc/proposed/2026-06-14-acp-multi-session.md +++ b/docs/rfc/proposed/2026-06-14-acp-multi-session.md @@ -18,7 +18,7 @@ The harness core already supports many agents (`AgentRegistry.list()` and `Agent - Lift the single-session guard in `session/new`; allow N live sessions, each mapped to its own `ReactLoopAgent`. - The bridge's `sessionId→agent` and `Session→sessionId` maps (introduced single-entry by [the ACP support RFC](2026-06-14-acp-agent-client-protocol.md)) become true multi-entry, plus a third `agent→sessionId` reverse map: the `tools/execute` permission gate receives only `exec.agent` (no sessionId), so it needs an O(1) reverse lookup to find the owning session. Every `agent/*` event and every `session/event` is demuxed strictly by id, so two sessions streaming at once never interleave their `session/update` notifications. - Per-session prompt queues: [the ACP support RFC](2026-06-14-acp-agent-client-protocol.md)'s single-entry in-flight-prompt state becomes multi-entry — one in-flight prompt *per session*, tracked per `sessionId`. -- Per-session cancel routing: `session/cancel` aborts only its own session's agent and settles only that session's in-flight prompt. `agent.abort()` drives a per-agent `AbortController`, so the per-session `exec.signal` is the natural isolation fence. +- Per-session cancel routing: `session/cancel` cancels only its own session's agent (via the queue-aware `agent.cancel()`) and settles only that session's in-flight prompt. The cancel is scoped to that one agent — a per-agent `AbortController` for the running step plus the agent's own queued/steering FIFOs — so it never touches another session's stream or pending prompt. - Per-session permission ownership: a `session/request_permission` and its outcome are bound to the originating session via the reverse map, so a permission prompt or a cancel in one session can never resolve another session's pending permission. ## Plan diff --git a/packages/agent-loop/src/loop.ts b/packages/agent-loop/src/loop.ts index 3a6ac76a74..b8f43a0a9b 100644 --- a/packages/agent-loop/src/loop.ts +++ b/packages/agent-loop/src/loop.ts @@ -201,14 +201,22 @@ export async function runLoop(ctx: Context, agent: ReactLoopAgent, handle: LoopH // Pre-step cancel (window 2): `setStatus('running')` emits `agent/status` // SYNCHRONOUSLY, so a `running` listener can `cancel()` in the gap between the - // check above and `runTurn`. `cancel()` already cleared the queued FIFO, so - // drop the turn before it starts (runTurn would otherwise throw on an empty - // queue) and transition back to idle — `running` was already emitted, so a - // real `idle` transition (which also settles waiters) balances the status. + // check above and `runTurn`. Mirror window 1: clear the marker, then + // - if NOTHING new is queued, drop the about-to-run turn and transition + // back to `idle` (`running` was already emitted, so a real idle + // transition balances the status AND settles `whenIdle()` waiters); + // - if a NEW prompt was queued AFTER the cancel (a `running` listener that + // cancels then sends), the marker was for the cancelled work only — fall + // through and run the new prompt's turn (status is already `running`), so + // a `whenIdle()` waiter resolves on THAT turn's running→idle, not before + // it runs. Settling here would resolve quiescence while the replacement + // is still queued and unrun (the same early-resolve race window 1 fixes). if (handle.isCancelled()) { handle.clearCancel() - handle.setStatus('idle') - continue + if (!agent.inbox.hasQueued) { + handle.setStatus('idle') + continue + } } // Re-derive the turn number from the log each iteration (do NOT keep a local diff --git a/packages/agent-loop/tests/cancel.spec.ts b/packages/agent-loop/tests/cancel.spec.ts index ab05c4e25b..9392b4b3d0 100644 --- a/packages/agent-loop/tests/cancel.spec.ts +++ b/packages/agent-loop/tests/cancel.spec.ts @@ -252,6 +252,36 @@ describe('Agent.cancel()', () => { expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false) }) + it('window 2: whenIdle() does NOT resolve early when a running listener cancels then queues replacement work', async () => { + // The window-1 early-resolve race has a window-2 twin: a synchronous + // agent/status('running') listener cancels the about-to-run turn AND queues a + // replacement. window 2 must NOT settle waiters (via setStatus('idle')) while + // the replacement is still queued-and-unrun — it must fall through and run it, + // so whenIdle() resolves on the replacement turn's running→idle, not before. + const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create('a1', { model: 'mock' }) + + let replaced = false + const dispose = ctx.on('agent/status', (subject, status) => { + if (subject !== agent || status !== 'running' || replaced) return + replaced = true + agent.cancel('drop A') + send(agent, 'B') + }) + + send(agent, 'A') + const idle = agent.whenIdle() + await idle + dispose() + + // whenIdle() resolved only AFTER B's turn ran: B's user message + a turn/end + // are in the log, and A was dropped. + expect(userTexts(agent)).toContain('B') + expect(userTexts(agent)).not.toContain('A') + expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true) + }) + it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => { // The subtle race: a whenIdle() waiter is registered for prompt A; cancel() // clears A; prompt B is queued BEFORE the loop resumes from the idle wait. From 16304872e134a4e1f712973c7d5d3c0624c1ff4e Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Sat, 20 Jun 2026 13:51:30 +0800 Subject: [PATCH 5/5] docs(agent-loop): drop PR-letter ref from cancel test header (review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The cancel.spec.ts module doc named "PR C", narrating the change's origin — process/history a reader of the current test does not need. Per the repo doc-current-state convention, describe only what the suite tests. --- packages/agent-loop/tests/cancel.spec.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/agent-loop/tests/cancel.spec.ts b/packages/agent-loop/tests/cancel.spec.ts index 9392b4b3d0..a4e9e13a1c 100644 --- a/packages/agent-loop/tests/cancel.spec.ts +++ b/packages/agent-loop/tests/cancel.spec.ts @@ -1,5 +1,5 @@ /** - * Tests for the queue-aware `Agent.cancel()` primitive (PR C). `cancel()` is the + * Tests for the queue-aware `Agent.cancel()` primitive. `cancel()` is the * broad verb — it clears queued + steering work, aborts an in-flight step, and * drops a turn about to start — whereas `abort()` kills only the current step. * These tests exercise every window where a cancel can land (idle, pre-step,