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/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/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/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 780e907559..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 @@ -569,23 +571,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..df25af1c11 100644 --- a/packages/agent-loop/src/agent.ts +++ b/packages/agent-loop/src/agent.ts @@ -26,6 +26,25 @@ 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 + /** + * 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). */ @@ -176,6 +195,34 @@ 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 + // 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 + // 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 +265,16 @@ export class ReactLoopAgent implements Agent { setAbort: controller => void (this.currentAbort = controller), 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 + // 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..b8f43a0a9b 100644 --- a/packages/agent-loop/src/loop.ts +++ b/packages/agent-loop/src/loop.ts @@ -107,6 +107,34 @@ 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 + /** + * 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 + /** + * 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 +176,49 @@ 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`. 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() + if (!agent.inbox.hasQueued) { + 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`. 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() + if (!agent.inbox.hasQueued) { + 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 +240,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 +403,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: handle.cancelReason() } + closeStep() + break + } + let stepOutcome: { hadToolCalls: boolean; finish: FinishReason } | { error: Error } try { stepOutcome = await runStep(ctx, agent, turn, step, abort.signal) @@ -382,6 +476,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: handle.cancelReason() } + 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..a4e9e13a1c --- /dev/null +++ b/packages/agent-loop/tests/cancel.spec.ts @@ -0,0 +1,337 @@ +/** + * 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, + * 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 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: 'from turn-start' }]) + }) + + 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 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: 'from continuation' }]) + }) + + 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('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. + // 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) + 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() }, } }