diff --git a/docs/16_codemode_backend.md b/docs/16_codemode_backend.md index e3d45cc3..b696211e 100644 --- a/docs/16_codemode_backend.md +++ b/docs/16_codemode_backend.md @@ -158,4 +158,12 @@ the container and worker examples, plus a `/c//agent` route. `script/run` is a smoke test that round-trips one file through every backend. +That agent layer asks a human before it runs anything that writes, and +before anything at all on the `container` backend. A held-back command +does not execute: the turn pauses and resumes through +`/c//approvals` once someone answers. The approval policy and the +paused-turn state both live in the example, not in this backend — the +backend runs a command and reports the result, as it did before. See +the example's README for the flow. + Run with `npm run dev --workspace @example/workspace-codemode`. diff --git a/examples/codemode/README.md b/examples/codemode/README.md index 16053b71..c65551ec 100644 --- a/examples/codemode/README.md +++ b/examples/codemode/README.md @@ -21,6 +21,22 @@ separate, opt-in layer: a `POST /agent` route runs a loop that drives the backends through an `exec` tool, and picks the backend per command. +The agent asks before it acts. Commands the approval policy holds +back — anything that writes, and everything on `container` — pause the +turn until a human answers, which is what the +[approval flow](#human-in-the-loop-approval) below is about. + +> [!TIP] +> [`architecture.html`](architecture.html) draws all of this: the +> system diagram, the call sequence for an unattended turn beside a +> gated one, and every branch of the approval route. It is a +> standalone page with no build step, and it is the fastest way in if +> you would rather see the shape than read it. +> +> ```sh +> open examples/codemode/architecture.html +> ``` + ## What makes codemode different The `shell` and `container` backends take a **shell command line**. @@ -61,13 +77,15 @@ The store the snippet touches is the same store `shell` and ``` client ─► Worker ─┬─ /file, /exec deterministic, no model - └─ /agent model loop + exec tool - │ (stub RPC) - ▼ - CodemodeExample DO (owns fs + registers 3 backends) - ├─ shell ─► Dynamic Worker (just-bash) - ├─ codemode ─► Dynamic Worker (JS sandbox, state.*) - └─ container ─► Cloudflare Container (wsd) + ├─ /agent model loop + exec tool + └─ /approvals the human's side of the loop + │ + ┌──────────────┴───────────────┐ + ▼ ▼ (stub RPC) + AgentSession DO CodemodeExample DO (owns fs + 3 backends) + paused turns, ├─ shell ─► Dynamic Worker (just-bash) + pending approvals ├─ codemode ─► Dynamic Worker (JS sandbox, state.*) + (no fs, no backends) └─ container ─► Cloudflare Container (wsd) all file operations route back to the one DO's SQLite store ``` @@ -84,6 +102,11 @@ client ─► Worker ─┬─ /file, /exec deterministic, no model `shell` backend adds a `WorkspaceServiceProxy` loopback, and the `container` backend adds `withWorkspaceContainer` plus a `WorkspaceProxy` egress loopback. +- `AgentSession` is a second durable object holding approval state. + It is addressed by the same ``, so `/c/demo/...` reaches one + workspace and one session — two objects, one name. It holds no + workspace stub and registers no backends, so the component that + records approval decisions cannot run a command. ## The optional agent layer @@ -104,6 +127,236 @@ The model is Workers AI Kimi (`@cf/moonshotai/kimi-k2.6`) via needs an authenticated wrangler session (`npx wrangler login`); `/file` and `/exec` are fully local. +## Human-in-the-loop approval + +The agent does not get to write to the filesystem on its own say-so. +When the model asks for a command the policy holds back, the command +**does not run**: the turn stops, the pending command is put on a +queue, and the turn resumes only after a human answers. + +``` +POST /agent {prompt} + │ + ▼ + the model asks to run a command + │ + ├─ policy allows it ──► runs ──► turn continues ──► status: completed + │ + └─ policy holds it ───► NOTHING RUNS + status: awaiting-approval + turnId + approvalId + │ + GET /approvals shows command, backend, reason + POST /approvals/ {approved: true|false} + │ + ├─ more still outstanding ──► 202, keep answering + └─ that was the last one ───► turn resumes here, + may pause again +``` + +For the same flow at the level of what each line of code does — where +the gate is evaluated, what the SDK leaves behind when it stops, and +what the resume has to put back — see +[the approval chain, step by step](architecture.html#approval-chain). +It splits the two sides of the pause, since a paused turn is two +requests with an indefinite gap in the middle. + +### The policy + +Approval works in two tiers, and they are worth telling apart because +only one of them is a guess. + +The first tier is a table keyed by backend. It decides on what a +backend can reach, not on what a command says: + +| Backend | Rule | Why | +|---|---|---| +| `shell` | `read-only` | Sandboxed, sees only the workspace. Recognized reads run unattended. | +| `codemode` | `read-only` | Same, and it has no network at all. | +| `container` | `always` | Full Linux userland with public network. "Which command is it" is the wrong question. | + +That table is the boundary, and it needs no parsing to enforce. +`container` gets a full userland and a public network, so everything +on it asks, whatever the command turns out to be. + +The second tier applies only under `read-only`, and it is an +optimization: it exists to ask fewer questions, not to be the thing +standing between the model and the files. A command runs unattended +only when it is *recognizably* a read. Anything the matcher does not +understand needs a human — an unknown verb, a shell metacharacter, a +`state.*` call reached through a computed access. + +Because the gate fails closed, the matcher is allowed to be wrong in +one direction only. Asking about a command that turns out to be +harmless costs a person ten seconds. Waving through a command that +writes is the failure that matters, and +`approval-policy.effects.test.ts` is what rules it out — it runs every +command the policy would allow and fails if any of them touched a +file. + +The matcher speaks two dialects, because the backends do. For `shell` +and `container` it rejects redirection, substitution and backgrounding +(`>`, `<`, `$(...)`, `` ` ``, `&`) outright, since those either write or +run something it would have to parse to see. What is left is a +pipeline, so it splits on `|`, `&&`, `||` and `;` and judges each stage +on its own: the verb has to be a known read (`cat`, `ls`, `grep`, +`find`, `git log`, …), and for verbs that write when handed the wrong +flag, the flags have to be known reads too. + +Judging stages separately rather than gating all composition is what +makes the rule usable — `ls -1 /workspace | wc -l` is two reads and a +pipe, which touches no files, so it runs unattended. It gives nothing +up: `ls; rm -rf /workspace` still stops on its second stage, and +`find /workspace -type f | xargs rm` still stops on `xargs`, even +though that `find` would have passed alone. Single commands behave the +same as before: `find /workspace -name '*.ts'` runs, while +`find /workspace -delete` and `sort -o out in` ask. For `codemode` it +reads the JavaScript instead: every mention of `state` has to resolve +to a named member, and every member has to be a read (`readFile`, +`stat`, `readdir`, …). That is why `state["writeFile"](...)` is gated. + +Both dialects allowlist rather than blocklist, which is the only +version that holds up. A blocklist of writing flags has to keep pace +with every flag that happens to write; an allowlist turns the ones +nobody thought of into questions. + +Some commands are left out of the read set entirely because their +writes cannot be read off their arguments: `sed` writes through `-i` +and through a `w` command buried in its script, `uniq` and `tree` take +an output file as a positional argument, `awk` writes through its own +syntax, and `date -s` sets the clock. Listing them would buy false +confidence rather than fewer approvals. `echo` and `printf` are in the +read set, on the other hand, because they only ever write to stdout — +aiming that at a file needs a redirect, which gates the line whatever +the verb. + +The rules are configuration, not a hardcoded list. `runAgentTurn` +takes a `policy`, and `never` turns the gate off for a backend +entirely: + +```ts +const transcript = await runAgentTurn({ + env, + workspace, + prompt, + policy: { + rules: { shell: "read-only", codemode: "never", container: "always" }, + fallback: "always", + }, +}); +``` + +The second tier is the part with no equivalent elsewhere. Pausing a +turn on a human is well-trodden ground, as the +[section below](#why-not-an-off-the-shelf-approvals-system) covers. +Deciding which shell commands are safe enough to skip the pause is +not: Think's built-in Bash tool runs on the same `just-bash` library +as the `shell` backend here, and its only controls are on, off, and +resource limits. This example took the problem on because nothing off +the shelf solves it. + +That is a reason to be careful about how much weight the matcher +carries. Reading a command line to guess its effect is a heuristic, +and a heuristic is the wrong place for a boundary. The right place is +the capability layer: hand the backend a read-only view of the +workspace and let the filesystem refuse the write. Think already works +that way, protecting files it did not mount during write-back so a +script cannot delete what it was never given. Doing the same here +would cover every caller instead of only the model's path, and would +leave the matcher as what it should be — a way to ask fewer questions. + +Until then, lean on the backend table, treat the matcher as a +convenience, and read `approval-policy.effects.test.ts` as the thing +that keeps the convenience honest. + +### Where a paused turn lives + +`runAgentTurn` is a single `generateText` call inside a fetch handler. +It has nowhere to keep a half-finished turn, so the state that has to +survive the wait — the message history, the pending approvals, the +decisions already taken — lives in `AgentSession`, a second durable +object. + +It is deliberately not in the workspace durable object. That object is +a filesystem with backends attached, and giving it a queue of +half-finished model turns would make it something else. The +[approval flow](#human-in-the-loop-approval) is the agent layer's +concern, so it gets its own durable object and the workspace stays a +workspace. + +Resuming works by replay. A paused turn stores the AI SDK's own +message history, including the `tool-approval-request` part that +records what was asked. Approving appends a `tool-approval-response` +message and hands the whole history back to `generateText`, which +executes the approved call and carries on. A rejection becomes an +`execution-denied` result the model reads and reacts to, rather than +an error that ends the turn. + +Two consequences worth knowing: + +- **The policy has to be deterministic.** The AI SDK re-runs it when a + turn resumes and downgrades an approved call to a denial if the + answer changed in the meantime. So it is a pure function of the + command and the backend — no clock, no mutable config — and a deploy + that *relaxes* the policy will deny approvals that were in flight. +- **The step budget spans the whole turn.** `MAX_STEPS` is spent across + every pass, not per pass, so waiting for a human does not buy the + model a fresh allowance. + +### Why not an off-the-shelf approvals system + +Several parts of the Cloudflare stack already do approvals, and this +example uses one of them: the pause is the AI SDK's `needsApproval`, +one field on the `exec` tool that was there anyway. That is the same +field Think gates its own tools with, so the mechanism here is the +common one rather than a local invention. The rest were considered and +are a worse fit, for reasons worth writing down. + +Think itself is the closest. Its Actions API can park a turn on a +human and resume it later with no connection held open, which is this +design under another name, down to `pendingApprovals`, +`approveExecution` and `rejectExecution` matching the three routes +below and a repeat answer being a no-op. But Think is a whole chat +agent framework — memory, streaming, messengers, scheduled turns — and +adopting it to get the approval plumbing would replace what this +example is for. + +The Agents SDK is the smaller version of that argument. Its `Agent` +class supplies SQL storage and routing, but message persistence +belongs to `AIChatAgent`, so `turn-store.ts` would move from key +prefixes to SQL rather than disappear, and the check that stops an +approval being answered twice has no equivalent to inherit. That is +about a hundred lines saved against eleven more dependencies in an +example that has four. `AIChatAgent` would delete the store outright, +but it is built around streaming chat messages, so the turn loop, the +routes and both scripts would need rewriting. + +Workflows can hold an approval for months and retry each step. This +pause is minutes long and holds nothing open, so a stored row is the +smaller mechanism, and `waitForApproval()` is an `Agent` method that +brings the same dependency with it. What that costs is a real timeout +and escalation; the prune below is the crude substitute. + +Fibers solve the neighboring problem, not this one. `runFiber()` makes +work survive eviction by checkpointing what is in flight. Nothing is +in flight here on purpose — the command was never sent, no connection +is open, no container booted — so there is no progress to save. + +`@cloudflare/codemode` ships its own approvals system — +`requiresApproval` on connector tools, `createCodemodeRuntime`, and +abort-and-replay through a durable tool-call log. This example does +not use it, and the naming makes that worth spelling out: the codemode +**backend** here (LLM-authored JavaScript in a Dynamic Worker) is not +the codemode **runtime**. They are different things that share a name. + +The runtime's approvals belong to connectors and a `codemode({ code })` +tool. Adopting them would mean replacing the `exec` tool with a +connector-based path, wrapping all three backends as a connector, and +taking on replay's determinism rules — which would delete the thing +this example exists to show, that the model picks one of three +backends per command. If you are building on codemode connectors +rather than workspace backends, the runtime's approvals are the better +fit; here they are not. + ## HTTP surface ``` @@ -113,13 +366,43 @@ GET /c//file/workspace/ octet-stream of /workspace/ POST /c//exec { command, cwd?, backend? } backend: shell | codemode | container (omit to use the default, shell) + runs unapproved: you are the caller, + so there is nobody to ask → JSON { exitCode, stdout, stderr } POST /c//agent { prompt } - → JSON { text, finishReason, steps, toolCalls } + → JSON transcript (see below) +GET /c//agent/ one turn's record: what it ran, what + it waits for, decisions already taken +GET /c//approvals → JSON { pending: [...] } — commands + waiting on a human +POST /c//approvals/ { approved, reason? } + → 200 transcript of the resumed turn + → 202 if approvals are still outstanding + → 404 if that id is not waiting (unknown, + or somebody already answered it) +``` + +A transcript is: + +```jsonc +{ + "status": "completed", // or "awaiting-approval" + "turnId": "0af13a39-…", + "text": "Created /workspace/greeting.txt…", + "finishReason": "stop", + "steps": 1, // this pass + "stepsUsed": 2, // the whole turn + "toolCalls": [ /* every command the turn ran, across passes */ ], + "pendingApprovals": [ + { "approvalId": "…", "backend": "codemode", "command": "…", "reason": "…" } + ] +} ``` `` selects a workspace instance (durable object). Reuse a name -to share files across calls; use a new name for a clean slate. +to share files across calls; use a new name for a clean slate. The +same name also selects the `AgentSession` that holds that workspace's +approval queue. ## Run it locally @@ -148,6 +431,7 @@ them: ./script/run # against http://127.0.0.1:8787 CONTAINERS=1 ./script/run # also read from the container AGENT=1 ./script/run # also run one agent turn +APPROVALS=1 ./script/run # also drive the approval flow end to end ``` To do the same steps by hand: @@ -172,13 +456,96 @@ curl -X POST $B/exec -H 'content-type: application/json' \ # agent: the model picks the backend (needs `npx wrangler login`) curl -X POST $B/agent -H 'content-type: application/json' -d '{ - "prompt":"Create /workspace/greeting.txt containing exactly the text hello world, then read it back to confirm. Report what you did." + "prompt":"Read /workspace/hello.txt and tell me its exact contents." }' ``` The `/agent` response includes `toolCalls[].backend`, showing which backend the model chose for each command. +### Watching it ask + +The quickest way to see an approval happen is the interactive driver. +It starts a turn and prompts on the terminal for every command the +policy holds back, which is what a real approval UI would do with the +same two routes: + +```sh +./script/agent "Create /workspace/greeting.txt containing exactly: hello world" +``` + +``` +APPROVAL NEEDED (codemode) + await state.writeFile("/workspace/greeting.txt", "hello world"); + why: state.writeFile is not a recognized read-only call + approve? [y/N] y + approved + ran [codemode] await state.writeFile("/workspace/greeting.txt", "hello world"); + +status completed +steps 2 +agent Created /workspace/greeting.txt containing exactly `hello world`. +``` + +Answer `n` and the command never runs; the model is told it was denied +and reports back. It defaults to a fresh workspace each run, so the +first write always has to be approved. `AUTO_APPROVE=1` says yes to +everything, and piping answers (`printf 'y\nn\n' | ./script/agent …`) +scripts them. + +### The approval flow by hand + +The same thing with two curls, which is what the driver is doing. Use a +fresh workspace name so the "nothing ran yet" step proves something: + +```sh +B=http://127.0.0.1:8787/c/hitl-demo + +# 1. A read needs no approval: status "completed". +curl -sX POST $B/agent -H 'content-type: application/json' \ + -d '{"prompt":"Read /workspace/hello.txt and tell me its contents."}' + +# 2. A write pauses: status "awaiting-approval", with an approvalId. +curl -sX POST $B/agent -H 'content-type: application/json' \ + -d '{"prompt":"Create /workspace/greeting.txt containing exactly: hello world"}' + +# 3. The queue: command, backend, and why it was held. +curl -s $B/approvals + +# 4. Nothing ran. This 404 is the point of the whole feature. +curl -s -o /dev/null -w '%{http_code}\n' $B/file/workspace/greeting.txt + +# 5. Approve. The turn resumes in this request and returns its +# transcript, which may pause again on the next command. +curl -sX POST $B/approvals/ \ + -H 'content-type: application/json' -d '{"approved":true}' + +# 6. Now the file is there. +curl -s $B/file/workspace/greeting.txt + +# 7. Answering twice is refused, so a racy approval UI cannot run the +# command twice: 404. +curl -sX POST $B/approvals/ \ + -H 'content-type: application/json' -d '{"approved":true}' +``` + +To watch a rejection instead, ask for something destructive and deny +it. The model is told the command was denied and why, and reports back +rather than rerouting to another backend: + +```sh +curl -sX POST $B/agent -H 'content-type: application/json' \ + -d '{"prompt":"Use the shell backend to delete everything under /workspace."}' +curl -sX POST $B/approvals/ -H 'content-type: application/json' \ + -d '{"approved":false,"reason":"not authorised to delete the workspace"}' +``` + +The turn record keeps the audit trail: + +```sh +curl -s $B/agent/ +``` + ## Container notes The image pulls `wsd` from a public GHCR image and installs a Linux @@ -201,7 +568,53 @@ userland from Debian. ## Tests -The backend has two test tiers, both under +The example's own tests cover the approval machinery, and need neither +a model nor a container: + +```sh +npm test --workspace @example/workspace-codemode +``` + +Four suites. `approval-policy.test.ts` pins the policy: which commands +are recognized reads, that redirection is gated and that a pipeline is +judged one stage at a time, that +`state["writeFile"]` does not slip past the allowlist, and that the +decision is a pure function of its inputs. `turn-store.test.ts` drives +the paused-turn bookkeeping against an in-memory map — a turn waits for +its last approval, answering twice is a no-op, and stale turns are +pruned. `agent.test.ts` runs the whole pause/resume loop against a +scripted model (`MockLanguageModelV3`), asserting against a fake +workspace that a gated command **never reaches `shell.exec`**, that the +message history survives a JSON round trip, and that approving runs the +held-back command while rejecting does not. + +`approval-policy.effects.test.ts` is the one that earns its keep +differently. The other three assert what the code was meant to do, +which cannot find the mistake that matters here: the matcher and its +tests are written by the same hand, so they miss the same cases. Both +real defects in this policy were found by running the agent and +noticing, not by a test. + +So that suite checks the claim against the world. It generates a corpus +— every allowlisted verb crossed with argument shapes including the +flags that turn a read into a write — keeps the commands the policy +would run **unattended**, and executes each one under real `just-bash` +against a recording filesystem. The property is one-directional: + +> a command the policy allows unattended must write nothing + +Currently 630 commands generated, 475 allowed, and the whole suite +runs in well under a second. Reintroduce the `find -delete` hole and it +fails naming every file that would have been deleted. The verbs come +from `READ_ONLY_COMMANDS` itself rather than a copy, so widening the +policy later puts the new verb under test without anybody remembering +to. + +It does not cover the `container` backend, which runs GNU coreutils +rather than just-bash and can behave differently — one more reason that +backend stays gated outright. + +The backend beneath it has two further test tiers, both under `packages/workspace`: ```sh @@ -225,12 +638,16 @@ guarantee, and `get()` returning `ENOENT`. ``` examples/codemode/ - wrangler.jsonc Worker + DO + worker_loaders + containers + AI - Dockerfile wsd + Debian userland for the container backend - script/run smoke test across the file surface + 3 backends - src/index.ts Worker handler + DO (CodemodeExample, 3 backends) - src/agent.ts the optional Workers AI model loop - src/tools/exec.ts the exec tool advertised to the model + wrangler.jsonc Worker + 2 DOs + worker_loaders + containers + AI + Dockerfile wsd + Debian userland for the container backend + script/run smoke test: file surface, 3 backends, approvals + script/agent one agent turn, prompting y/n for each approval + src/index.ts Worker handler + DO (CodemodeExample, 3 backends) + src/agent.ts the optional Workers AI model loop + src/tools/exec.ts the exec tool advertised to the model + src/approval-policy.ts which commands need a human, and why + src/session.ts AgentSession DO: where a paused turn lives + src/turn-store.ts the paused-turn bookkeeping, storage-agnostic ``` ## Known limitations @@ -258,4 +675,31 @@ examples/codemode/ directories still need an explicit `state.mkdir(...)` or `mkdir -p`. - **The agent loop runs in the Worker.** Fine for short tasks; a - long agent run would want its own durable object. + long agent run would want its own durable object. Approval state + already lives in one (`AgentSession`), so a paused turn survives + between requests even though the loop that produced it does not. +- **Approving resumes inline.** `POST /approvals/` runs the rest of + the turn in that request, so answering an approval costs a model + round trip and can return `awaiting-approval` again. It keeps the + demo to one curl per step; a UI would more likely resume in the + background and poll the turn record. +- **The approval matcher is a heuristic.** It classifies command + strings, so it is conservative by construction and will ask about + commands that are in fact harmless. What it guarantees is only the + one direction the effects test checks: a command it allows does not + write. It is the second tier of the policy, not the boundary, and + real enforcement belongs at the capability layer, as the + [policy section](#the-policy) spells out. +- **Approval requests in the stored history are unsigned.** The AI + SDK can sign them with `experimental_toolApprovalSecret`, but that + setting expects approval requests to carry a signature and throws + on the plain message history this example replays, so it stays off. + Nothing here depends on it: the history never leaves the server, and + `GET /agent/` strips `messages` before it answers. Move that + history through a client and the signature starts mattering. +- **Telling the model not to reroute is advice, not a guarantee.** The + system prompt asks it not to retry a denied command on another + backend, and models do not always listen. That costs nothing: every + attempt is classified afresh, so a reroute either produces a command + the policy allows or asks again. The gate does not depend on the + model cooperating. diff --git a/examples/codemode/architecture.html b/examples/codemode/architecture.html index 0928b7a4..6b03469e 100644 --- a/examples/codemode/architecture.html +++ b/examples/codemode/architecture.html @@ -148,6 +148,7 @@

System diagram

/agent deterministic vs. model-driven + /approvals — the human's answer @@ -229,6 +230,16 @@

layer 2 The agent (opt-in)

For long runs, the loop would move to its own Durable Object holding a stub.

+
+

layer 3 Approval (human in the loop)

+

+ Commands the policy holds back — every write, and everything on + containerdo not run. The turn stops with a pending approval and + resumes on POST /approvals/<id>. That wait outlives the request, so the + paused turn lives in a second Durable Object, AgentSession, which holds no + workspace stub and cannot itself run anything. +

+

The three backends

@@ -298,6 +309,351 @@

How codemode reaches the filesystem

+

Call sequence

+

+ Five participants, and only one of them can touch a file. Read these two diagrams as a pair: + the difference between them is the entire feature. +

+ +

An unattended turn — the policy recognizes a read

+
+ + + + + + + + + + client / human + script/agent + + Worker + index.ts · agent.ts + + Workers AI + Kimi k2.6 + + AgentSession + (not used) + + Example DO + Workspace + files + + + + + + + + + + POST /agent { prompt } + + + + getWorkspace() → stub + + + + generateText(prompt, { exec }) + + + + exec { "ls /workspace", shell } + + + + decideApproval() → a read, no human + + + + ws.shell.exec(command, { backend }) + + + + { exitCode, stdout, stderr } + + + + tool result → model continues + + + + 200 { status: completed, … } + +
+ One HTTP request in, one out. The green self-call is the whole decision; because it + answered "read", nobody was asked anything. +
+
+ +

A gated turn — the policy holds the command back

+
+ + + + + + + client / human + script/agent + + Worker + index.ts · agent.ts + + Workers AI + Kimi k2.6 + + AgentSession + paused turn (DO) + + Example DO + Workspace + files + + + + + + + + + + POST /agent { prompt } + + + + getWorkspace() → stub + + + + generateText(prompt, { exec }) + + + + exec { "echo hi > f", shell } + + + + decideApproval() → needs a human + + + + put(turnId, history, awaiting[]) + + + + 200 { awaiting-approval } + + + PAUSE — the command has not run. Nothing is in flight. + The turn is a row in the AgentSession DO. The Worker can die here; it holds no state. + + + + POST /approvals/<id> { approved } + + + + resolve(id) → history if last + + + + replay(history + approval) + + + + decideApproval() re-runs — must agree + + + + ws.shell.exec(command, { backend }) + + + + { exitCode, stdout, stderr } + + + + tool result → model continues + + + + 200 { status: completed, … } + +
+ Two HTTP requests, one pause. The orange self-call is where the write was stopped — + before the Example DO was ever asked to run it, which is why the pause band is + safe to sit in indefinitely. +
+
+ +
+

+ The one subtlety worth internalising: decideApproval() runs + twice — once to gate, once on replay. The AI SDK re-checks the policy when + a turn resumes and converts an approved call into a denial if the answer changed. That is + why the policy is a pure function of the command and the backend, with no clock and no + mutable state: anything else would make approvals decay on their own. +

+
+ +

The approval route

+

+ Where every branch goes, and what the client gets back. The two entry points funnel into one + model loop, and the loop can come back around: answering an approval resumes the turn, and the + resumed pass may stop on a different command. +

+
+ the command runs + the command is held back + control flow +
+
+ + + + + + + + + + + + + + + + POST /agent { prompt } + + + POST /approvals/<id> { approved } + + + session.resolveApproval(id, approved) + + + 404 + unknown or answered + + + 202 + others outstanding + + + + + + + + last one → resume + + + runAgentTurn() — the model loop, inside the Worker + + + model proposes exec { cmd, backend } + + + decideApproval(cmd, backend) + the gate + pure function + + + + + recognized read + → ws.shell.exec() on the Example DO + + + not recognized + held back — never executed + + + + + + result → the model continues, and may propose another + collected as pendingApprovals[] + + + + + recordTurn() → AgentSession.saveTurn(turn) + + + + + pending.length > 0 ? + + + + no + + + yes + + + 200 { status: completed } + nothing left to ask + + + 200 { status: awaiting-approval, pending[] } + the command has not run + + + + + human answers one — script/agent, UI, bot + + + a resumed turn can pause again + +
+ The gate is the orange box, and it sits inside the Worker — upstream of the Durable Object + that owns the files. Only the green path ever reaches it. +
+
+ +
+

Why 404 and 202 are different answers

+

+ 404 means the id is not waiting for an answer — either it was never issued, or + somebody already answered it. Both are the same fact, and both have to reject: treating a + repeat as a fresh approval would run the command a second time. + 202 means the answer was recorded but the turn is still waiting on others, + because one model step can propose several commands and the resume has to carry every + answer at once. Only the last answer resumes the turn, which is why the reply to + POST /approvals/<id> is sometimes a transcript and sometimes a receipt. +

+
+

Request flows

@@ -320,6 +676,160 @@

Agentic — POST /agent

  • On stop, Worker returns the transcript with toolCalls[].backend.
  • +
    +

    Gated — POST /approvals/<id>

    +
      +
    1. The policy holds a command back, so the tool never executes it.
    2. +
    3. The turn returns awaiting-approval; Worker stores its history in + AgentSession.
    4. +
    5. A human reads GET /approvals and answers one entry.
    6. +
    7. On the last answer, Worker replays the stored history with a + tool-approval-response attached.
    8. +
    9. The approved command runs for real; a denied one comes back to the model as + execution-denied.
    10. +
    +
    +
    + +

    The approval chain, step by step

    +

    + The same gated turn as above, at the level of what each line of code does. Two sides, because + a paused turn is two requests with an indefinite gap in the middle: the first discovers the + pause and files it, the second answers it and finishes the work. +

    + +
    +

    + side one The pause — POST /agent +

    +
      +
    1. + The Worker parses { prompt }, mints a turnId, and gets a + workspace stub. Nothing is stored yet. +
    2. +
    3. + createExecTool closes over the stub, the policy and a fresh + onExec array — none of which appear in the model's schema. +
    4. +
    5. + generateText runs with the budget + MAX_STEPS - stepsUsed, so a turn cannot buy itself more steps by pausing. +
    6. +
    7. The model proposes exec { command, backend }.
    8. +
    9. + Before execute, the SDK evaluates needsApproval. Pure, local, + inside the Worker — no Durable Object, no network. +
    10. +
    11. + It returns true, so the call is excluded from that step's executeTools. + Nothing runs: no shell.exec, no connect(), no + dynamic worker, no container boot. +
    12. +
    13. + The SDK appends a tool-approval-request part to the step's content — an + approvalId plus the original toolCall. A tool-call + with no matching tool-result. +
    14. +
    15. + The loop's continue condition fails and generateText returns normally. + Nothing threw. +
    16. +
    17. + runAgentTurn scans result.steps.at(-1).content for those parts. + That scan is how the pause is discovered — after the fact, from data. +
    18. +
    19. + For each one it asks decideApproval again, purely for the sentence a human + should read; the SDK's gate only answered yes or no. Same inputs, same answer. +
    20. +
    21. + status becomes awaiting-approval, and messages + becomes what went in plus lastStep.response.messages — carrying the + serialized mark. +
    22. +
    23. + recordTurn writes the row into AgentSession: history, + pending, awaiting, toolCalls, + stepsUsed. That object holds no workspace stub and cannot run anything. +
    24. +
    25. + transcriptJSON replies 200 awaiting-approval with + pendingApprovals, omitting messages — the model's working state + is read from storage on resume, never trusted from a client. +
    26. +
    27. + The request closes. No timer, no polling, no held connection. The turn exists only as a + row, and the Worker is free to die. +
    28. +
    +
    + +
    +

    + side two The resume — + POST /approvals/<id> +

    +
      +
    1. + A human answers. script/agent prints the y/N; the Worker only + ever sees a boolean. +
    2. +
    3. resolveApproval(id, approved) marks it answered and returns an outcome.
    4. +
    5. + Null means the id was never issued or was already answered. Both mean this + request changed nothing, so 404 — treating a repeat as fresh would run + the command twice. +
    6. +
    7. + Not ready means other approvals on the turn are outstanding, so 202 and + back to idle. One model step can propose several commands, and the resume must carry + every answer at once. +
    8. +
    9. Ready means this was the last one, so the stored messages come back out.
    10. +
    11. + runAgentTurn is called again with resume, and + stepsUsed carried forward. +
    12. +
    13. + The history is replayed with one appended tool message holding every + tool-approval-response together — the SDK reads approvals from the last + message only. +
    14. +
    15. + The SDK matches each approvalId to the mark in the history and + re-runs needsApproval. If the answer changed since the + pause, the approval is downgraded to a denial. This is why the policy must be pure. +
    16. +
    17. + Approved calls execute in the SDK's pre-loop executeTools, before the step + loop starts. +
    18. +
    19. + ws.shell.exec() is finally called on the Durable Object — the first time + this command touches anything. From here it is the ordinary path: lazy + connect(), the dispatchers, evaluate(), the envelope, + execEventStream. +
    20. +
    21. A denied call comes back to the model as execution-denied instead.
    22. +
    23. + onExec records the run. The steps will not: an approved call belongs to the + pre-loop and to no step, which is why the transcript is collected through the closure + rather than read off result.steps. +
    24. +
    25. The loop continues and may propose more commands.
    26. +
    27. + recordTurn writes again, accumulating toolCalls across passes so + the turn's history stays whole. +
    28. +
    29. + Completed gives 200 completed. Paused again on a different command gives + 200 awaiting-approval, and back to idle. +
    30. +
    +

    + The whole resume runs inline in the human's POST, so their request is the one that waits on + the model. +

    Key components

    @@ -327,8 +837,11 @@

    Key components

    FileRole src/index.tsWorker routes + CodemodeExample DO that owns the Workspace and registers all three backends. - src/agent.tsThe opt-in model loop; builds the exec tool over the stub and runs one turn. - src/tools/exec.tsThe exec tool: a backend enum plus per-backend descriptions the model reads. + src/agent.tsThe opt-in model loop; builds the exec tool over the stub and runs one turn, pausing when the policy holds a command back. + src/tools/exec.tsThe exec tool: a backend enum plus per-backend descriptions the model reads, and the approval gate. + src/approval-policy.tsWhich commands need a human: a per-backend rule table, denying by default. + src/session.tsThe AgentSession DO — where a turn waiting on a human lives. + src/turn-store.tsPaused-turn bookkeeping over a storage handle, so it tests against a plain Map. backends/codemode/codemode-backend.tsThe backend: builds the ShellRPC, runs each snippet in a fresh executor. backends/codemode/state-provider.tsThe state.* namespace over the live WorkspaceFilesystem. backends/codemode/exec-events.tsPure mapping from a run's result/logs/error to the stdout/stderr/exit event stream. diff --git a/examples/codemode/package.json b/examples/codemode/package.json index 91072df3..bcbf3447 100644 --- a/examples/codemode/package.json +++ b/examples/codemode/package.json @@ -7,6 +7,7 @@ "scripts": { "dev": "wrangler dev", "deploy": "wrangler deploy", + "test": "vitest run", "typecheck": "tsc --noEmit" }, "dependencies": { @@ -17,7 +18,9 @@ }, "devDependencies": { "@cloudflare/workers-types": "^4.20260616.1", + "just-bash": "^3.0.1", "typescript": "^6.0.3", + "vitest": "^4.1.7", "wrangler": "^4.96.0" } } diff --git a/examples/codemode/script/agent b/examples/codemode/script/agent new file mode 100755 index 00000000..d465033c --- /dev/null +++ b/examples/codemode/script/agent @@ -0,0 +1,167 @@ +#!/usr/bin/env node +// Drive one agent turn against a running `wrangler dev`, stopping to ask +// on the terminal for every command the approval policy holds back. +// +// The HTTP surface answers approvals with a second request, which is +// the right shape for a UI but an awkward way to see the feature work. +// This script closes that loop: it starts a turn, and each time the +// turn comes back waiting on a human it prints the command, asks y/n, +// posts the answer, and picks the turn up where it left off. What you +// see is what a real approval UI would drive. +// +// ./script/agent "Create /workspace/notes.txt saying hello" +// NAME=foo ./script/agent "..." # workspace instance +// BASE_URL=http://127.0.0.1:8799 ./script/agent "..." +// ./script/agent # uses a default prompt +// AUTO_APPROVE=1 ./script/agent "..." # say yes to everything +// printf 'y\nn\n' | ./script/agent "..." # scripted answers + +const BASE_URL = (process.env.BASE_URL ?? "http://127.0.0.1:8787").replace(/\/$/, ""); +// A fresh instance by default, so a demo starts from an empty tree and +// "the file is not there yet" means what it looks like. +const NAME = process.env.NAME ?? `agent-${Date.now().toString(36)}`; +const PROMPT = + process.argv.slice(2).join(" ") || + "Create /workspace/greeting.txt containing exactly: hello world"; + +const base = `${BASE_URL}/c/${NAME}`; + +const bold = (s) => `\u001b[1m${s}\u001b[0m`; +const dim = (s) => `\u001b[2m${s}\u001b[0m`; +const green = (s) => `\u001b[32m${s}\u001b[0m`; +const red = (s) => `\u001b[31m${s}\u001b[0m`; +const yellow = (s) => `\u001b[33m${s}\u001b[0m`; + +async function post(url, body) { + const response = await fetch(url, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify(body), + }); + const text = await response.text(); + let parsed; + try { + parsed = JSON.parse(text); + } catch { + throw new Error(`${url} returned HTTP ${response.status}: ${text.slice(0, 200)}`); + } + if (parsed.error != null) { + throw new Error(`${url} returned HTTP ${response.status}: ${parsed.error}`); + } + return parsed; +} + +function reportRan(toolCalls, alreadyReported) { + for (const call of toolCalls.slice(alreadyReported)) { + const status = call.exitCode === 0 ? green("ran") : red(`exit ${call.exitCode}`); + console.log(` ${status} ${dim(`[${call.backend}]`)} ${call.command.replace(/\n/g, " ")}`); + const out = (call.stdout || call.stderr).trim(); + if (out.length > 0) console.log(dim(` ${out.slice(0, 200).replace(/\n/g, "\n ")}`)); + } + return toolCalls.length; +} + +/** + * Read answers a line at a time from stdin. + * + * Not `node:readline`, because a piped answer arrives and ends the + * stream while the first model call is still in flight, and readline + * then rejects the question that follows with "readline was closed". + * Buffering from the start works for a terminal and a pipe alike. + * Returns null once stdin has no more to give. + */ +function lineReader() { + const ready = []; + const waiting = []; + let buffer = ""; + let ended = false; + + process.stdin.setEncoding("utf8"); + process.stdin.on("data", (chunk) => { + buffer += chunk; + let index = buffer.indexOf("\n"); + while (index >= 0) { + const line = buffer.slice(0, index); + buffer = buffer.slice(index + 1); + const waiter = waiting.shift(); + if (waiter) waiter(line); + else ready.push(line); + index = buffer.indexOf("\n"); + } + }); + process.stdin.on("end", () => { + ended = true; + if (buffer.length > 0) { + const line = buffer; + buffer = ""; + const waiter = waiting.shift(); + if (waiter) waiter(line); + else ready.push(line); + } + while (waiting.length > 0) waiting.shift()(null); + }); + + return { + next() { + if (ready.length > 0) return Promise.resolve(ready.shift()); + if (ended) return Promise.resolve(null); + return new Promise((resolve) => waiting.push(resolve)); + }, + }; +} + +const reader = lineReader(); +const autoApprove = process.env.AUTO_APPROVE === "1"; + +/** Ask the human. No answer left on stdin means no. */ +async function askApproval() { + if (autoApprove) { + console.log(` approve? ${dim("[y/N]")} y ${dim("(AUTO_APPROVE)")}`); + return true; + } + process.stdout.write(` approve? ${dim("[y/N]")} `); + const answer = await reader.next(); + if (answer === null) { + console.log(dim("(no answer on stdin, treating as no)")); + return false; + } + if (!process.stdin.isTTY) console.log(answer.trim()); + return /^y(es)?$/i.test(answer.trim()); +} + +try { + console.log(`${bold("workspace")} ${NAME} ${dim(base)}`); + console.log(`${bold("prompt")} ${PROMPT}\n`); + + let turn = await post(`${base}/agent`, { prompt: PROMPT }); + let reported = reportRan(turn.toolCalls, 0); + + while (turn.status === "awaiting-approval") { + for (const approval of turn.pendingApprovals) { + console.log(`\n${yellow("APPROVAL NEEDED")} ${dim(`(${approval.backend})`)}`); + console.log(` ${bold(approval.command.replace(/\n/g, "\n "))}`); + console.log(` ${dim(`why: ${approval.reason}`)}`); + + const approved = await askApproval(); + + // Nothing has run at this point. Approving is what executes it. + const outcome = await post(`${base}/approvals/${approval.approvalId}`, { + approved, + ...(approved ? {} : { reason: "denied at the prompt" }), + }); + console.log(approved ? green(" approved") : red(" denied")); + turn = outcome; + reported = reportRan(turn.toolCalls, reported); + } + } + + console.log(`\n${bold("status")} ${turn.status}`); + console.log(`${bold("steps")} ${turn.stepsUsed}`); + if (turn.text) console.log(`${bold("agent")} ${turn.text}`); + console.log(dim(`\nturn record: curl -s ${base}/agent/${turn.turnId}`)); +} catch (error) { + console.error(red(`\nFAIL: ${error.message}`)); + process.exitCode = 1; +} finally { + process.stdin.pause(); +} diff --git a/examples/codemode/script/run b/examples/codemode/script/run index 8fda8432..656a1062 100755 --- a/examples/codemode/script/run +++ b/examples/codemode/script/run @@ -11,14 +11,16 @@ # surface, closing the loop. # # Steps 1 through 4 need no model and no container, so they run by -# default. Two further backends are opt-in because they cost more to -# stand up: +# default. The rest are opt-in because they cost more to stand up: # # CONTAINERS=1 also read the tree from inside the Cloudflare # Container (boots wsd on first use; requires # `wrangler dev --enable-containers`). # AGENT=1 also run one agent turn, which needs the AI # binding and a model round-trip. +# APPROVALS=1 also drive the human-in-the-loop approval flow: ask +# the agent for a write, check that nothing ran while +# it waits, approve, and check that it then did. # # Run `npm run dev` in another terminal first, then: # ./script/run # against http://127.0.0.1:8787 @@ -91,4 +93,59 @@ if [[ "${AGENT:-}" == "1" ]]; then printf '%s\n' "$agent_out" fi +if [[ "${APPROVALS:-}" == "1" ]]; then + # A fresh instance, so "the file is not there yet" actually means the + # approval held the write back rather than that an earlier run had + # not written it. + HITL_NAME="${NAME}-hitl-$$" + HITL_BASE="${BASE_URL}/c/${HITL_NAME}" + HITL_PATH="approved-write.txt" + + step "7. approval: ask for a write on /c/${HITL_NAME}/agent" + hitl_out=$(curl -fsS -X POST "${HITL_BASE}/agent" \ + -H 'content-type: application/json' \ + -d '{"prompt":"Create the file /workspace/'"${HITL_PATH}"' containing exactly: approved"}') + printf '%s\n' "$hitl_out" + + hitl_status=$(JSON="$hitl_out" node -e 'process.stdout.write(JSON.parse(process.env.JSON).status)') + [[ "$hitl_status" == "awaiting-approval" ]] \ + || fail "expected the write to wait for approval, got status '$hitl_status'" + + approval_id=$(JSON="$hitl_out" node -e \ + 'const d=JSON.parse(process.env.JSON); process.stdout.write(d.pendingApprovals[0].approvalId)') + + step "8. approval: nothing ran while it waits" + code=$(curl -sS -o /dev/null -w '%{http_code}' "${HITL_BASE}/file/workspace/${HITL_PATH}") + [[ "$code" == "404" ]] \ + || fail "the command ran before it was approved (GET returned HTTP $code, expected 404)" + printf 'GET /file/workspace/%s -> HTTP 404, as it should be\n' "$HITL_PATH" + + step "9. approval: the queue lists it" + queue=$(curl -fsS "${HITL_BASE}/approvals") + printf '%s\n' "$queue" + echo "$queue" | grep -q "$approval_id" \ + || fail "the pending approval is missing from GET /approvals" + + step "10. approval: approve it, and the turn resumes" + resumed=$(curl -fsS -X POST "${HITL_BASE}/approvals/${approval_id}" \ + -H 'content-type: application/json' -d '{"approved":true}') + printf '%s\n' "$resumed" + resumed_status=$(JSON="$resumed" node -e 'process.stdout.write(JSON.parse(process.env.JSON).status)') + [[ "$resumed_status" == "completed" ]] \ + || fail "expected the resumed turn to complete, got status '$resumed_status'" + + step "11. approval: now the file exists" + written=$(curl -fsS "${HITL_BASE}/file/workspace/${HITL_PATH}") + printf '%s\n' "$written" + [[ -n "$written" ]] || fail "the approved command did not write the file" + + step "12. approval: answering twice is refused" + code=$(curl -sS -o /dev/null -w '%{http_code}' \ + -X POST "${HITL_BASE}/approvals/${approval_id}" \ + -H 'content-type: application/json' -d '{"approved":true}') + [[ "$code" == "404" ]] \ + || fail "a second answer to the same approval was accepted (HTTP $code, expected 404)" + printf 'second answer -> HTTP 404, so the command cannot run twice\n' +fi + printf '\nOK — one filesystem, reached through every backend.\n' diff --git a/examples/codemode/src/agent.test.ts b/examples/codemode/src/agent.test.ts new file mode 100644 index 00000000..96303a92 --- /dev/null +++ b/examples/codemode/src/agent.test.ts @@ -0,0 +1,495 @@ +import type { LanguageModelV3Content } from "@ai-sdk/provider"; +import { MockLanguageModelV3 } from "ai/test"; +import { describe, expect, it } from "vitest"; + +import { MAX_STEPS, runAgentTurn } from "./agent.js"; +import type { ApprovalPolicy } from "./approval-policy.js"; +import { decideApproval } from "./approval-policy.js"; +import type { ExecWorkspaceLike } from "./tools/exec.js"; + +// --------------------------------------------------------------- +// Test doubles +// --------------------------------------------------------------- + +const READ = "cat /workspace/hello.txt"; +const WRITE = "rm -rf /workspace"; + +function usage() { + return { + inputTokens: { total: 1, noCache: 1, cacheRead: 0, cacheWrite: 0 }, + outputTokens: { total: 1, text: 1, reasoning: 0 }, + }; +} + +function execCall( + toolCallId: string, + input: { command: string; backend?: string; cwd?: string }, +): LanguageModelV3Content { + return { type: "tool-call", toolCallId, toolName: "exec", input: JSON.stringify(input) }; +} + +function say(text: string): LanguageModelV3Content { + return { type: "text", text }; +} + +/** + * A model that replays scripted responses, one per `generateText` + * call. The same instance is reused across a pause and its resume, so + * the script reads as the whole turn. + */ +function scriptedModel(script: LanguageModelV3Content[][]): MockLanguageModelV3 { + let call = 0; + return new MockLanguageModelV3({ + doGenerate: async () => { + const content = script[Math.min(call, script.length - 1)]; + call += 1; + const callsTool = content.some((part) => part.type === "tool-call"); + return { + content, + finishReason: { unified: callsTool ? "tool-calls" : "stop", raw: undefined }, + usage: usage(), + warnings: [], + }; + }, + }); +} + +interface RecordedExec { + command: string; + backend: string | undefined; + cwd: string | undefined; +} + +function fakeWorkspace(): { workspace: ExecWorkspaceLike; calls: RecordedExec[] } { + const calls: RecordedExec[] = []; + const workspace: ExecWorkspaceLike = { + shell: { + exec: async (command, options) => { + calls.push({ command, backend: options.backend, cwd: options.cwd }); + return { + result: async () => ({ exitCode: 0, stdout: `ran: ${command}`, stderr: "" }), + }; + }, + }, + }; + return { workspace, calls }; +} + +const env = {} as Env; + +// --------------------------------------------------------------- + +describe("runAgentTurn", () => { + describe("commands that need no approval", () => { + it("runs a read and finishes the turn", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: READ, backend: "shell" })], + [say("done")], + ]); + + const transcript = await runAgentTurn({ env, workspace, model, prompt: "read the file" }); + + expect(transcript.status).toBe("completed"); + expect(transcript.pendingApprovals).toEqual([]); + expect(calls).toEqual([{ command: READ, backend: "shell", cwd: undefined }]); + expect(transcript.toolCalls).toHaveLength(1); + expect(transcript.toolCalls[0]).toMatchObject({ + command: READ, + backend: "shell", + exitCode: 0, + stdout: `ran: ${READ}`, + }); + expect(transcript.text).toBe("done"); + }); + + it("honours a policy that trusts a backend outright", async () => { + // Same mutating command, different policy: the gate is + // configuration, not a hardcoded list. + const trusting: ApprovalPolicy = { rules: { shell: "never" }, fallback: "always" }; + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: WRITE, backend: "shell" })], + [say("removed")], + ]); + + const transcript = await runAgentTurn({ + env, + workspace, + model, + prompt: "delete it", + policy: trusting, + }); + + expect(transcript.status).toBe("completed"); + expect(calls).toHaveLength(1); + }); + }); + + describe("commands that need approval", () => { + it("pauses without running the command", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([[execCall("c1", { command: WRITE, backend: "shell" })]]); + + const transcript = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + + expect(transcript.status).toBe("awaiting-approval"); + // The whole point: nothing ran. + expect(calls).toEqual([]); + expect(transcript.toolCalls).toEqual([]); + expect(transcript.pendingApprovals).toHaveLength(1); + expect(transcript.pendingApprovals[0]).toMatchObject({ + toolCallId: "c1", + command: WRITE, + backend: "shell", + cwd: null, + }); + expect(transcript.pendingApprovals[0].approvalId.length).toBeGreaterThan(0); + }); + + it("explains the pause with the same reason the policy gave", async () => { + const { workspace } = fakeWorkspace(); + const model = scriptedModel([[execCall("c1", { command: WRITE, backend: "shell" })]]); + + const transcript = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + + expect(transcript.pendingApprovals[0].reason).toBe( + decideApproval({ command: WRITE, backend: "shell" }).reason, + ); + }); + + it("reports the backend the model chose, not the default", async () => { + const { workspace } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: "uname -a", backend: "container" })], + ]); + + const transcript = await runAgentTurn({ env, workspace, model, prompt: "which kernel?" }); + + expect(transcript.pendingApprovals[0].backend).toBe("container"); + }); + + it("falls back to the default backend when the model omits one", async () => { + const { workspace } = fakeWorkspace(); + const model = scriptedModel([[execCall("c1", { command: WRITE })]]); + + const transcript = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + + expect(transcript.pendingApprovals[0].backend).toBe("shell"); + }); + + it("keeps the work it already did before pausing", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: READ, backend: "shell" })], + [execCall("c2", { command: WRITE, backend: "shell" })], + ]); + + const transcript = await runAgentTurn({ env, workspace, model, prompt: "read then delete" }); + + expect(transcript.status).toBe("awaiting-approval"); + expect(calls.map((call) => call.command)).toEqual([READ]); + expect(transcript.toolCalls.map((call) => call.command)).toEqual([READ]); + }); + }); + + describe("the resume spine", () => { + it("survives a round trip through JSON", async () => { + const { workspace } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: WRITE, backend: "shell" })], + [say("removed")], + ]); + + const paused = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + // This is the assumption the durable object depends on: the + // messages are plain data, storable and retrievable verbatim. + const rehydrated = JSON.parse(JSON.stringify(paused.messages)); + expect(rehydrated).toEqual(paused.messages); + + const resumed = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: rehydrated, + approvals: [{ approvalId: paused.pendingApprovals[0].approvalId, approved: true }], + }, + }); + + expect(resumed.status).toBe("completed"); + }); + + it("grows with each pass so a second pause can be resumed too", async () => { + const { workspace } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: WRITE, backend: "shell" })], + [execCall("c2", { command: "mkdir /workspace/d", backend: "shell" })], + [say("done")], + ]); + + const first = await runAgentTurn({ env, workspace, model, prompt: "delete then make" }); + const second = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: first.messages, + approvals: [{ approvalId: first.pendingApprovals[0].approvalId, approved: true }], + }, + }); + + expect(second.status).toBe("awaiting-approval"); + expect(second.messages.length).toBeGreaterThan(first.messages.length); + + const third = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: second.messages, + approvals: [{ approvalId: second.pendingApprovals[0].approvalId, approved: true }], + }, + }); + + expect(third.status).toBe("completed"); + }); + }); + + describe("resuming with an approval", () => { + it("runs the command that was held back", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: WRITE, backend: "shell" })], + [say("removed")], + ]); + + const paused = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + const resumed = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: paused.messages, + approvals: [{ approvalId: paused.pendingApprovals[0].approvalId, approved: true }], + }, + }); + + expect(resumed.status).toBe("completed"); + expect(calls).toEqual([{ command: WRITE, backend: "shell", cwd: undefined }]); + expect(resumed.text).toBe("removed"); + }); + + it("reports the approved command in the transcript", async () => { + // An approved call executes before the first model call of the + // resumed pass, so it lands in no step. Reading the transcript + // off the steps would lose it entirely. + const { workspace } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: WRITE, backend: "shell" })], + [say("removed")], + ]); + + const paused = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + const resumed = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: paused.messages, + approvals: [{ approvalId: paused.pendingApprovals[0].approvalId, approved: true }], + }, + }); + + expect(resumed.toolCalls).toHaveLength(1); + expect(resumed.toolCalls[0]).toMatchObject({ + command: WRITE, + backend: "shell", + exitCode: 0, + stdout: `ran: ${WRITE}`, + }); + }); + }); + + describe("resuming with a rejection", () => { + it("does not run the command", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: WRITE, backend: "shell" })], + [say("understood, leaving it alone")], + ]); + + const paused = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + const resumed = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: paused.messages, + approvals: [ + { + approvalId: paused.pendingApprovals[0].approvalId, + approved: false, + reason: "not now", + }, + ], + }, + }); + + expect(calls).toEqual([]); + expect(resumed.toolCalls).toEqual([]); + expect(resumed.status).toBe("completed"); + expect(resumed.text).toBe("understood, leaving it alone"); + }); + + it("tells the model the command was denied, and why", async () => { + const { workspace } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: WRITE, backend: "shell" })], + [say("understood")], + ]); + + const paused = await runAgentTurn({ env, workspace, model, prompt: "delete it" }); + await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: paused.messages, + approvals: [ + { + approvalId: paused.pendingApprovals[0].approvalId, + approved: false, + reason: "not now", + }, + ], + }, + }); + + const prompt = JSON.stringify(model.doGenerateCalls.at(-1)?.prompt); + expect(prompt).toContain("execution-denied"); + expect(prompt).toContain("not now"); + }); + }); + + describe("several approvals in one step", () => { + it("pauses on each of them", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [ + execCall("c1", { command: WRITE, backend: "shell" }), + execCall("c2", { command: "mkdir /workspace/d", backend: "shell" }), + ], + ]); + + const transcript = await runAgentTurn({ env, workspace, model, prompt: "delete and make" }); + + expect(transcript.status).toBe("awaiting-approval"); + expect(transcript.pendingApprovals).toHaveLength(2); + expect(transcript.pendingApprovals.map((entry) => entry.command)).toEqual([ + WRITE, + "mkdir /workspace/d", + ]); + expect(calls).toEqual([]); + }); + + it("runs both once both are answered", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [ + execCall("c1", { command: WRITE, backend: "shell" }), + execCall("c2", { command: "mkdir /workspace/d", backend: "shell" }), + ], + [say("both done")], + ]); + + const paused = await runAgentTurn({ env, workspace, model, prompt: "delete and make" }); + const resumed = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: paused.messages, + approvals: paused.pendingApprovals.map((entry) => ({ + approvalId: entry.approvalId, + approved: true, + })), + }, + }); + + expect(resumed.status).toBe("completed"); + expect(calls.map((call) => call.command)).toEqual([WRITE, "mkdir /workspace/d"]); + expect(resumed.toolCalls).toHaveLength(2); + }); + + it("can approve one and deny the other", async () => { + const { workspace, calls } = fakeWorkspace(); + const model = scriptedModel([ + [ + execCall("c1", { command: WRITE, backend: "shell" }), + execCall("c2", { command: "mkdir /workspace/d", backend: "shell" }), + ], + [say("partly done")], + ]); + + const paused = await runAgentTurn({ env, workspace, model, prompt: "delete and make" }); + const resumed = await runAgentTurn({ + env, + workspace, + model, + resume: { + messages: paused.messages, + approvals: [ + { approvalId: paused.pendingApprovals[0].approvalId, approved: false, reason: "no" }, + { approvalId: paused.pendingApprovals[1].approvalId, approved: true }, + ], + }, + }); + + expect(calls.map((call) => call.command)).toEqual(["mkdir /workspace/d"]); + expect(resumed.toolCalls).toHaveLength(1); + }); + }); + + describe("the step budget", () => { + it("accumulates steps across passes", async () => { + const { workspace } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: READ, backend: "shell" })], + [say("done")], + ]); + + const transcript = await runAgentTurn({ + env, + workspace, + model, + prompt: "read it", + stepsUsed: 3, + }); + + expect(transcript.steps).toBe(2); + expect(transcript.stepsUsed).toBe(5); + }); + + it("shrinks the allowance a resumed pass gets", async () => { + // Without this, every approval would hand the loop a fresh + // budget and an approval cycle could run without bound. + const { workspace } = fakeWorkspace(); + const model = scriptedModel([ + [execCall("c1", { command: READ, backend: "shell" })], + [say("done")], + ]); + + const transcript = await runAgentTurn({ + env, + workspace, + model, + prompt: "read it", + stepsUsed: MAX_STEPS, + }); + + expect(transcript.steps).toBe(1); + expect(transcript.stepsUsed).toBe(MAX_STEPS + 1); + }); + }); +}); diff --git a/examples/codemode/src/agent.ts b/examples/codemode/src/agent.ts index b28d0188..6caa571f 100644 --- a/examples/codemode/src/agent.ts +++ b/examples/codemode/src/agent.ts @@ -12,11 +12,21 @@ * the agent can run anywhere: here it runs inside the Worker fetch * handler and reaches the workspace through its stub, so the * workspace durable object never has to know an agent exists. + * + * A turn can end in one of two ways. Either the model finishes, or it + * asks to run a command the approval policy holds back — in which case + * the turn returns `awaiting-approval` along with the message history + * needed to pick it up again. Nothing about that history lives here: + * `runAgentTurn` is one `generateText` call inside a fetch handler and + * is over when it returns. Whoever calls it is responsible for storing + * a paused turn and handing it back on resume, which in this example + * is the `AgentSession` durable object. */ -import { generateText, stepCountIs } from "ai"; +import { generateText, type LanguageModel, type ModelMessage, stepCountIs } from "ai"; import { createWorkersAI } from "workers-ai-provider"; +import { type ApprovalPolicy, DEFAULT_APPROVAL_POLICY, decideApproval } from "./approval-policy.js"; import { createExecTool, type ExecWorkspaceLike } from "./tools/exec.js"; // The Workers AI model that drives the loop. Kimi K2.6 handles tool @@ -24,13 +34,53 @@ import { createExecTool, type ExecWorkspaceLike } from "./tools/exec.js"; const MODEL_ID = "@cf/moonshotai/kimi-k2.6"; // Plenty of budget for a write-then-read loop, with a ceiling so a -// confused model can't spin forever. -const MAX_STEPS = 12; +// confused model can't spin forever. Spent across every pass of a +// turn, not per pass, so waiting for a human doesn't buy the model a +// fresh allowance. +export const MAX_STEPS = 12; + +/** A human's answer to one approval request. */ +export interface ApprovalResponse { + approvalId: string; + approved: boolean; + /** Shown to the model when a command is denied. */ + reason?: string; +} + +/** A command the model wants to run, held back for a human. */ +export interface PendingApproval { + /** Identifies this request when the answer arrives. */ + approvalId: string; + toolCallId: string; + backend: string; + command: string; + cwd: string | null; + /** Why the policy stopped it, in one line. */ + reason: string; +} export interface AgentTurnOptions { env: Env; workspace: ExecWorkspaceLike; - prompt: string; + /** The user's request. Omit when resuming a paused turn. */ + prompt?: string; + /** Pick up a turn that paused, given a human's answers. */ + resume?: { + /** The `messages` a previous pass returned. */ + messages: ModelMessage[]; + /** One entry per approval the paused pass requested. */ + approvals: ApprovalResponse[]; + }; + /** Which commands need a human. Defaults to the example's policy. */ + policy?: ApprovalPolicy; + /** Steps earlier passes of this turn already spent. */ + stepsUsed?: number; + /** + * The model to drive the loop with. Defaults to Workers AI; tests + * inject a scripted model so the loop can be exercised without a + * model round trip. + */ + model?: LanguageModel; } export interface AgentToolCall { @@ -42,10 +92,26 @@ export interface AgentToolCall { } export interface AgentTranscript { + /** + * `awaiting-approval` means the model asked for a command the policy + * holds back, nothing further ran, and the turn is resumable. + */ + status: "completed" | "awaiting-approval"; text: string; finishReason: string; + /** Model steps this pass took. */ steps: number; + /** Model steps the whole turn has taken. */ + stepsUsed: number; + /** Commands this pass ran. */ toolCalls: AgentToolCall[]; + pendingApprovals: PendingApproval[]; + /** + * The turn's history so far, for storing against a resume. Server + * side only — it carries the whole conversation and is not something + * to hand a client. + */ + messages: ModelMessage[]; } const SYSTEM_PROMPT = [ @@ -67,17 +133,42 @@ const SYSTEM_PROMPT = [ " boot, so reach for it only when the lighter backends can't run", " the command.", "", + "Some commands need a human's approval before they run. When one", + "does, the turn stops and resumes after a person answers; you do", + "not need to do anything differently. If a command comes back as", + "denied, do not try to run it again on another backend — say what", + "you were not allowed to do and stop.", + "", "When the task is done, reply with a short plain-text summary of", "what you did.", ].join("\n"); export async function runAgentTurn(opts: AgentTurnOptions): Promise { - const workersai = createWorkersAI({ binding: opts.env.AI }); - const model = workersai(MODEL_ID); + if (opts.prompt == null && opts.resume == null) { + throw new Error("runAgentTurn: pass either a prompt or a turn to resume"); + } + + const model = opts.model ?? createWorkersAI({ binding: opts.env.AI })(MODEL_ID); + const policy = opts.policy ?? DEFAULT_APPROVAL_POLICY; + const defaultBackend = "shell"; + + // Every command this pass runs, collected as it happens. See + // `onExec` on the tool for why the steps aren't enough. + const toolCalls: AgentToolCall[] = []; const exec = createExecTool({ workspace: opts.workspace, maxBytes: 16 * 1024, + policy, + onExec: (record) => { + toolCalls.push({ + backend: record.backend, + command: record.command, + exitCode: record.exitCode, + stdout: record.stdout, + stderr: record.stderr, + }); + }, backends: { shell: { description: @@ -112,41 +203,71 @@ export async function runAgentTurn(opts: AgentTurnOptions): Promise ({ + type: "tool-approval-response" as const, + approvalId: approval.approvalId, + approved: approval.approved, + ...(approval.reason != null ? { reason: approval.reason } : {}), + })), + }, + ] + : [{ role: "user", content: opts.prompt as string }]; + + const stepsAlreadyUsed = opts.stepsUsed ?? 0; const result = await generateText({ model, system: SYSTEM_PROMPT, - prompt: opts.prompt, + messages, tools: { exec }, - stopWhen: stepCountIs(MAX_STEPS), + // At least one step, so a turn that has exhausted its budget still + // reports back rather than failing. + stopWhen: stepCountIs(Math.max(1, MAX_STEPS - stepsAlreadyUsed)), }); - const toolCalls: AgentToolCall[] = []; - for (const step of result.steps) { - for (const tr of step.toolResults) { - const output = tr.output as { - backend?: string; - command?: string; - exitCode?: number; - stdout?: string; - stderr?: string; - }; - toolCalls.push({ - backend: output.backend ?? "", - command: output.command ?? "", - exitCode: output.exitCode ?? -1, - stdout: output.stdout ?? "", - stderr: output.stderr ?? "", - }); - } + const lastStep = result.steps.at(-1); + + // The pass paused if the model asked for anything the policy holds + // back. Those calls did not run. + const pendingApprovals: PendingApproval[] = []; + for (const part of lastStep?.content ?? []) { + if (part.type !== "tool-approval-request") continue; + const input = part.toolCall.input as { command: string; cwd?: string; backend?: string }; + const backend = input.backend ?? defaultBackend; + pendingApprovals.push({ + approvalId: part.approvalId, + toolCallId: part.toolCall.toolCallId, + backend, + command: input.command, + cwd: input.cwd ?? null, + // The AI SDK's gate answers yes or no; ask the policy again for + // the wording a human should see. Same inputs, same answer. + reason: decideApproval({ command: input.command, backend }, policy).reason, + }); } return { + status: pendingApprovals.length > 0 ? "awaiting-approval" : "completed", text: result.text, finishReason: result.finishReason, steps: result.steps.length, + stepsUsed: stepsAlreadyUsed + result.steps.length, toolCalls, + pendingApprovals, + // The history to store against a resume: what went in, plus what + // the model and the tools produced on the way out. + messages: [...messages, ...(lastStep?.response.messages ?? [])], }; } diff --git a/examples/codemode/src/approval-policy.effects.test.ts b/examples/codemode/src/approval-policy.effects.test.ts new file mode 100644 index 00000000..aec4e797 --- /dev/null +++ b/examples/codemode/src/approval-policy.effects.test.ts @@ -0,0 +1,250 @@ +/** + * Does a command the policy waves through actually only read? + * + * `approval-policy.test.ts` pins what the matcher *says*: given this + * command, does it ask for a human. Those assertions are worth having, + * but they cannot find the failure that matters here, because the same + * person writes the matcher and its tests and so shares a blind spot + * with it. Both real defects in this policy were found by running the + * agent by hand and noticing — `find /workspace -mindepth 1 -delete` + * ran unattended because the verb was allowlisted and its flags were + * not, and a pipeline of two reads asked for approval it did not need. + * Neither was going to fall out of a list of examples somebody thought + * to write down. + * + * So this file checks the claim against the world instead. It runs the + * command for real and watches the filesystem: + * + * the policy allows a command unattended ⇒ running it writes nothing + * + * One direction only. A command the policy *gates* needs no check here, + * because being asked about a read is a nuisance and not a breach; the + * gated direction is already covered by the assertions next door. + * + * The corpus is generated rather than curated — every allowlisted verb + * crossed with argument shapes that include the flags known to turn a + * read into a write. Most combinations are nonsense (`pwd -delete`), + * and that is fine: a nonsense command the matcher allows must still + * not write. Deriving the verbs from READ_ONLY_COMMANDS rather than + * from a copy means a verb added to the policy later comes under test + * without anybody remembering to add it here. + * + * ## What this covers, and what it does not + * + * The shell is real: `just-bash`, the same implementation the `shell` + * backend runs inside its Dynamic Worker, driven against just-bash's + * own in-memory filesystem wrapped in a recorder. Whether `find + * -delete` reaches for a delete is a fact about just-bash and holds + * wherever its files happen to live, so the storage underneath does + * not need to be the real Durable Object for the answer to be right. + * + * Two gaps, neither of them quiet: + * + * The `container` backend runs GNU coreutils rather than just-bash, so + * the same command can behave differently there. The default policy + * gates that backend outright, which is why the difference does not + * bite — and is a reason to keep gating it. + * + * The `codemode` dialect is not exercised. Its claim is about a closed + * list of `state.*` method names, which is a much smaller surface to + * get wrong than the verb-and-flag combinatorics here, and reaching a + * live `state` namespace would mean pulling the workspace package and + * its workerd-only imports into this runner. + */ + +import { Bash, InMemoryFs } from "just-bash"; +import { describe, expect, it } from "vitest"; + +import { decideApproval, READ_ONLY_COMMANDS } from "./approval-policy.js"; + +/** + * Every method on just-bash's filesystem interface that changes + * something. Named explicitly rather than inferred, so a new mutating + * method in a future just-bash shows up as an unrecorded write here + * instead of being silently classified as a read. + */ +const MUTATORS = new Set([ + "appendFile", + "chmod", + "cp", + "link", + "mkdir", + "mv", + "rm", + "symlink", + "utimes", + "writeFile", +]); + +interface Run { + writes: string[]; + exitCode: number; + stderr: string; +} + +/** Run one command against a fresh tree, recording every mutation. */ +async function run(command: string): Promise { + const inner = new InMemoryFs({ + "/workspace/a.txt": "beta\nalpha\n", + "/workspace/b.txt": "gamma\n", + "/workspace/sub/c.txt": "nested\n", + }); + const writes: string[] = []; + const recorder = new Proxy(inner, { + get(target, key, receiver) { + const value = Reflect.get(target, key, receiver); + if (typeof value !== "function" || typeof key !== "string") return value; + if (!MUTATORS.has(key)) return value.bind(target); + return (...args: unknown[]) => { + const shown = args.filter((arg) => typeof arg === "string").join(", "); + writes.push(`${key}(${shown})`); + return value.apply(target, args); + }; + }, + }); + + const bash = new Bash({ fs: recorder as never, cwd: "/workspace" }); + const result = await bash.exec(command); + return { writes, exitCode: result.exitCode, stderr: result.stderr }; +} + +/** + * Argument shapes to cross with every verb. The first few are ordinary + * usage, present so the corpus contains commands that actually run; + * the rest are the ways a read verb turns into a write — the flag that + * deletes, the flag that names an output file, the flag that edits in + * place, the operator that hands the output to something else. + */ +const ARGUMENT_SHAPES = [ + "", + "/workspace", + "/workspace/a.txt", + "-1 /workspace", + "-l /workspace", + "/workspace/a.txt /workspace/b.txt", + "-delete /workspace", + "/workspace -delete", + "/workspace -mindepth 1 -delete", + "/workspace -type f -delete", + "-i s/alpha/beta/ /workspace/a.txt", + "--in-place /workspace/a.txt", + "-o /workspace/out /workspace/a.txt", + "--output=/workspace/out /workspace/a.txt", + "-w /workspace/out /workspace/a.txt", + "-s /workspace/a.txt", + "/workspace/a.txt > /workspace/out", + "/workspace/a.txt >> /workspace/a.txt", + "/workspace | tee /workspace/out", + "/workspace/a.txt | sed -i s/a/b/ /workspace/b.txt", +]; + +/** Whole-line shapes, to cover composition rather than one verb. */ +const COMPOSED = [ + "ls -1 /workspace | wc -l", + "cat /workspace/a.txt | grep alpha", + "cat /workspace/a.txt | sort | head -1", + "ls /workspace && cat /workspace/a.txt", + "ls /workspace || true", + "ls /workspace; rm -rf /workspace", + "find /workspace -type f | xargs rm", + "find /workspace -type f | tee /workspace/out", + "cat /workspace/a.txt > /workspace/out", + "sort /workspace/a.txt -o /workspace/a.txt", + "grep -r alpha /workspace | cut -d: -f1", + "test -f /workspace/a.txt && echo yes", + "echo hello", + "printf '%s\\n' hello", + "pwd", + "ls -la /workspace/sub", +]; + +function corpus(): string[] { + const commands = new Set(COMPOSED); + for (const verb of READ_ONLY_COMMANDS) { + for (const shape of ARGUMENT_SHAPES) { + commands.add(shape.length === 0 ? verb : `${verb} ${shape}`); + } + } + // git is handled by its own branch in the matcher, on a subcommand + // rather than a flag, so it needs its own shapes. + for (const sub of [ + "log", + "status", + "diff", + "show", + "add -A", + "commit -m x", + "checkout .", + "clean -fd", + ]) { + commands.add(`git ${sub}`); + commands.add(`git ${sub} > /workspace/out`); + } + return [...commands]; +} + +describe("the detector itself", () => { + // If these fail, every other assertion in this file is worthless: + // a recorder that sees nothing makes any command look like a read. + it("sees a write", async () => { + const { writes } = await run("printf hi > /workspace/new.txt"); + expect(writes).toContain("writeFile(/workspace/new.txt, hi, utf8)"); + }); + + it("sees a delete, including one reached through find", async () => { + const { writes } = await run("find /workspace -mindepth 1 -delete"); + expect(writes.length).toBeGreaterThan(0); + expect(writes.every((write) => write.startsWith("rm("))).toBe(true); + }); + + it("stays quiet on a read", async () => { + expect((await run("cat /workspace/a.txt")).writes).toEqual([]); + expect((await run("ls -1 /workspace | wc -l")).writes).toEqual([]); + }); + + it("would catch a policy that waved a write through", async () => { + // The policy is the thing under test, so prove the harness fails + // when the policy is wrong. `never` is the rule that trusts a + // backend completely; under it, a destructive command is + // "allowed", and the property below must not hold. + const policy = { rules: { shell: "never" as const } }; + const command = "rm -rf /workspace/sub"; + expect(decideApproval({ command, backend: "shell" }, policy).needsApproval).toBe(false); + expect((await run(command)).writes.length).toBeGreaterThan(0); + }); +}); + +describe("every command the policy allows unattended", () => { + it("writes nothing", async () => { + const violations: string[] = []; + let allowed = 0; + let ran = 0; + + for (const command of corpus()) { + if (decideApproval({ command, backend: "shell" }).needsApproval) continue; + allowed += 1; + const result = await run(command); + if (result.exitCode === 0) ran += 1; + if (result.writes.length > 0) { + violations.push(`${JSON.stringify(command)} → ${result.writes.join(", ")}`); + } + } + + // Reported together rather than one at a time: the useful output + // is the whole set of holes, not whichever one sorts first. + expect(violations, `the policy allowed ${violations.length} command(s) that wrote`).toEqual([]); + + // A corpus that gates everything, or a shell that cannot run + // anything, would satisfy the assertion above while checking + // nothing at all. + // + // These floors exist to catch a harness that died, not a policy + // that tightened, so they sit well under what passes today: 630 + // commands generated, 475 allowed, 163 of those exiting 0. A + // policy that legitimately narrows should not have to come here + // and edit numbers; a just-bash that stopped running commands, or + // a corpus that stopped generating them, lands near zero. + expect(allowed, "commands the policy allowed unattended").toBeGreaterThan(100); + expect(ran, "allowed commands that also exited 0").toBeGreaterThan(40); + }); +}); diff --git a/examples/codemode/src/approval-policy.test.ts b/examples/codemode/src/approval-policy.test.ts new file mode 100644 index 00000000..c2392c99 --- /dev/null +++ b/examples/codemode/src/approval-policy.test.ts @@ -0,0 +1,301 @@ +import { describe, expect, it } from "vitest"; + +import { type ApprovalPolicy, DEFAULT_APPROVAL_POLICY, decideApproval } from "./approval-policy.js"; + +// Shorthand: does this command need a human under the default policy? +function gates(command: string, backend: string, policy?: ApprovalPolicy): boolean { + return decideApproval({ command, backend }, policy).needsApproval; +} + +describe("decideApproval", () => { + describe("the 'always' rule", () => { + it("gates every command on the container backend", () => { + expect(gates("cat /workspace/hello.txt", "container")).toBe(true); + expect(gates("uname -a", "container")).toBe(true); + }); + + it("gates a command that the same rule set waves through on shell", () => { + // Identical command, different backend: proves the rule is + // per-backend rather than per-command. + expect(gates("cat /workspace/hello.txt", "shell")).toBe(false); + expect(gates("cat /workspace/hello.txt", "container")).toBe(true); + }); + }); + + describe("the 'never' rule", () => { + const trusting: ApprovalPolicy = { rules: { shell: "never" }, fallback: "always" }; + + it("waves through even a destructive command", () => { + expect(gates("rm -rf /workspace", "shell", trusting)).toBe(false); + }); + }); + + describe("unknown backends", () => { + it("falls back to the strictest rule", () => { + expect(gates("cat /workspace/hello.txt", "not-a-backend")).toBe(true); + }); + + it("honours an explicit fallback", () => { + const lenient: ApprovalPolicy = { rules: {}, fallback: "never" }; + expect(gates("rm -rf /", "whatever", lenient)).toBe(false); + }); + }); + + describe("the 'read-only' rule on a shell dialect", () => { + it("waves through recognized reads", () => { + expect(gates("cat /workspace/hello.txt", "shell")).toBe(false); + expect(gates("ls -la /workspace", "shell")).toBe(false); + expect(gates("grep -n needle /workspace/haystack.txt", "shell")).toBe(false); + expect(gates("find /workspace -name '*.ts'", "shell")).toBe(false); + expect(gates("wc -l /workspace/hello.txt", "shell")).toBe(false); + expect(gates("head -n 5 /workspace/hello.txt", "shell")).toBe(false); + expect(gates("stat /workspace/hello.txt", "shell")).toBe(false); + }); + + it("waves through read-only git plumbing", () => { + expect(gates("git status", "shell")).toBe(false); + expect(gates("git log --oneline -5", "shell")).toBe(false); + expect(gates("git diff HEAD", "shell")).toBe(false); + }); + + it("gates git subcommands that write", () => { + expect(gates("git commit -m wip", "shell")).toBe(true); + expect(gates("git push origin main", "shell")).toBe(true); + expect(gates("git checkout -b feature", "shell")).toBe(true); + expect(gates("git", "shell")).toBe(true); + }); + + it("gates mutating commands", () => { + expect(gates("rm -rf /workspace", "shell")).toBe(true); + expect(gates("mv /workspace/a /workspace/b", "shell")).toBe(true); + expect(gates("mkdir -p /workspace/deep/dir", "shell")).toBe(true); + expect(gates("chmod 777 /workspace/hello.txt", "shell")).toBe(true); + expect(gates("npm install", "shell")).toBe(true); + }); + + it("gates redirection, even of a read", () => { + expect(gates("cat /workspace/a > /workspace/b", "shell")).toBe(true); + expect(gates("cat /workspace/a >> /workspace/b", "shell")).toBe(true); + expect(gates("cat < /workspace/a", "shell")).toBe(true); + }); + + it("waves through a pipeline whose every stage is a read", () => { + // A pipe moves bytes between processes and touches no files, so + // a pipeline of reads is a read. Each stage is classified on its + // own rather than the composition being waved through. + expect(gates("ls -1 /workspace | wc -l", "shell")).toBe(false); + expect(gates("cat /workspace/a | grep needle", "shell")).toBe(false); + expect(gates("cat /workspace/a | grep needle | wc -l", "shell")).toBe(false); + expect(gates("ls /workspace || true", "shell")).toBe(false); + expect(gates("ls /workspace && cat /workspace/a", "shell")).toBe(false); + expect(gates("cd /workspace; ls", "shell")).toBe(true); + }); + + it("gates a pipeline with a stage that is not a read", () => { + expect(gates("ls /workspace; rm -rf /workspace", "shell")).toBe(true); + expect(gates("ls /workspace && rm -rf /workspace", "shell")).toBe(true); + expect(gates("cat /workspace/a | tee /workspace/b", "shell")).toBe(true); + expect(gates("find /workspace -type f | xargs rm", "shell")).toBe(true); + expect(gates("ls /workspace | sed -i s/a/b/", "shell")).toBe(true); + }); + + it("names the offending stage when it gates a pipeline", () => { + expect( + decideApproval({ command: "ls /workspace | tee /workspace/b", backend: "shell" }).reason, + ).toContain("tee"); + }); + + it("gates an empty or dangling stage", () => { + expect(gates("ls /workspace |", "shell")).toBe(true); + expect(gates("| wc -l", "shell")).toBe(true); + expect(gates("ls &&", "shell")).toBe(true); + }); + + it("still gates backgrounding, which leaves something running", () => { + expect(gates("ls /workspace &", "shell")).toBe(true); + expect(gates("cat /workspace/a & cat /workspace/b", "shell")).toBe(true); + }); + + it("still gates redirection inside a pipeline", () => { + expect(gates("ls /workspace | wc -l > /workspace/count", "shell")).toBe(true); + }); + + it("waves through echo and printf, which cannot write without a redirect", () => { + // Both write to stdout only. Sending that to a file needs `>`, + // which gates the whole line regardless of the verb. + expect(gates("echo hello", "shell")).toBe(false); + expect(gates("ls /workspace && echo done", "shell")).toBe(false); + expect(gates("echo hello > /workspace/x", "shell")).toBe(true); + }); + + it("gates command substitution", () => { + expect(gates("cat $(ls /workspace)", "shell")).toBe(true); + expect(gates("cat `ls /workspace`", "shell")).toBe(true); + expect(gates("cat /workspace/$(whoami)", "shell")).toBe(true); + }); + + it("gates a newline that hides a second command", () => { + expect(gates("ls /workspace\nrm -rf /workspace", "shell")).toBe(true); + }); + + it("gates a read verb handed a flag that writes", () => { + // A verb allowlist is not enough on its own: several read + // commands write when given the right flag, so an unrecognized + // flag on a read verb has to gate too. + expect(gates("find /workspace -mindepth 1 -delete", "shell")).toBe(true); + expect(gates("find /workspace -name x -exec rm {} +", "shell")).toBe(true); + expect(gates("find /workspace -execdir rm {} +", "shell")).toBe(true); + expect(gates("find /workspace -fprint /workspace/out", "shell")).toBe(true); + expect(gates("sort -o /workspace/out /workspace/in", "shell")).toBe(true); + expect(gates("sort --output=/workspace/out /workspace/in", "shell")).toBe(true); + }); + + it("still waves through the read flags those verbs are used with", () => { + expect(gates("find /workspace -name '*.ts'", "shell")).toBe(false); + expect(gates("find /workspace -type f -maxdepth 2", "shell")).toBe(false); + expect(gates("find /workspace -mtime -1", "shell")).toBe(false); + expect(gates("sort -n /workspace/hello.txt", "shell")).toBe(false); + expect(gates("sort -u -r /workspace/hello.txt", "shell")).toBe(false); + }); + + it("gates an unrecognized flag on a checked verb", () => { + expect(gates("find /workspace -frobnicate", "shell")).toBe(true); + }); + + it("gates verbs whose writing cannot be told from their arguments", () => { + // sed writes through -i and through a `w` command inside the + // script, which a matcher cannot reliably find. uniq and tree + // take an output file as a positional argument, and date -s sets + // the clock. None of them are worth the false confidence, so + // none of them are recognized reads. + expect(gates("sed s/a/b/ /workspace/hello.txt", "shell")).toBe(true); + expect(gates("sed -i s/a/b/ /workspace/hello.txt", "shell")).toBe(true); + expect(gates("sed 'w /workspace/out' /workspace/hello.txt", "shell")).toBe(true); + expect(gates("uniq /workspace/in /workspace/out", "shell")).toBe(true); + expect(gates("tree -o /workspace/out", "shell")).toBe(true); + expect(gates("date -s 12:00", "shell")).toBe(true); + }); + + it("strips leading environment assignments before reading the verb", () => { + expect(gates("LC_ALL=C sort /workspace/hello.txt", "shell")).toBe(false); + expect(gates("LC_ALL=C rm /workspace/hello.txt", "shell")).toBe(true); + }); + + it("gates an unrecognized command rather than guessing", () => { + expect(gates("frobnicate --hard", "shell")).toBe(true); + expect(gates("", "shell")).toBe(true); + expect(gates(" ", "shell")).toBe(true); + }); + }); + + describe("the 'read-only' rule on the codemode dialect", () => { + it("waves through a snippet that only reads", () => { + expect(gates('return await state.readFile("/workspace/hello.txt");', "codemode")).toBe(false); + expect(gates('const s = await state.stat("/workspace"); return s.size;', "codemode")).toBe( + false, + ); + expect(gates('return (await state.readdir("/workspace")).length;', "codemode")).toBe(false); + }); + + it("waves through a snippet that never touches state", () => { + // The sandbox has no network and no other reach into the store, + // so a snippet without state.* can only compute. + expect(gates("return 2 + 2;", "codemode")).toBe(false); + }); + + it("gates a snippet that mutates", () => { + expect(gates('await state.writeFile("/workspace/x", "hi");', "codemode")).toBe(true); + expect(gates('await state.mkdir("/workspace/d", { recursive: true });', "codemode")).toBe( + true, + ); + expect(gates('await state.rm("/workspace/x", { force: true });', "codemode")).toBe(true); + expect(gates('await state.chmod("/workspace/x", 0o755);', "codemode")).toBe(true); + expect(gates('await state.symlink("/workspace/x", "/workspace/y");', "codemode")).toBe(true); + }); + + it("gates a mutation hidden behind a computed member access", () => { + // The point of allowlisting reads instead of blocklisting + // writes: a blocklist misses these, an allowlist does not. + expect(gates('await state["writeFile"]("/workspace/x", "hi");', "codemode")).toBe(true); + expect( + gates('const m = "writeFile"; await state[m]("/workspace/x", "hi");', "codemode"), + ).toBe(true); + expect(gates('await state?.["writeFile"]("/workspace/x", "hi");', "codemode")).toBe(true); + }); + + it("gates a snippet that aliases state", () => { + expect(gates('const { writeFile } = state; await writeFile("/x", "y");', "codemode")).toBe( + true, + ); + expect(gates("const s = state; return s;", "codemode")).toBe(true); + }); + + it("gates an unrecognized state member", () => { + expect(gates("await state.frobnicate();", "codemode")).toBe(true); + }); + + it("waves through a read reached with optional chaining", () => { + expect(gates('return await state?.readFile("/workspace/hello.txt");', "codemode")).toBe( + false, + ); + }); + + it("does not apply shell rules to a JavaScript snippet", () => { + // `>` is a comparison here, not a redirect, and `state.ls` is a + // read. A shell-dialect check would gate this. + expect(gates('return (await state.ls("/workspace")).length > 0;', "codemode")).toBe(false); + }); + }); + + describe("the decision itself", () => { + it("explains every gate it raises", () => { + const commands: Array<[string, string]> = [ + ["rm -rf /workspace", "shell"], + ["cat /workspace/a > /workspace/b", "shell"], + ["uname -a", "container"], + ['await state.writeFile("/workspace/x", "hi");', "codemode"], + ["frobnicate", "not-a-backend"], + ]; + for (const [command, backend] of commands) { + const decision = decideApproval({ command, backend }); + expect(decision.needsApproval).toBe(true); + expect(decision.reason.length).toBeGreaterThan(0); + } + }); + + it("explains why it let a command through", () => { + const decision = decideApproval({ command: "cat /workspace/hello.txt", backend: "shell" }); + expect(decision.needsApproval).toBe(false); + expect(decision.reason.length).toBeGreaterThan(0); + }); + + it("names the backend rule in the reason, so the queue is readable", () => { + expect(decideApproval({ command: "uname -a", backend: "container" }).reason).toContain( + "container", + ); + }); + + it("is a pure function of command and backend", () => { + // Approval survives a pause: the AI SDK re-runs the policy when + // the turn resumes, and an approved call whose policy has since + // flipped is converted to a denial. Same input, same answer, + // every time. + const call = { command: 'await state.writeFile("/workspace/x", "hi");', backend: "codemode" }; + const first = decideApproval(call); + for (let i = 0; i < 5; i++) { + expect(decideApproval(call)).toEqual(first); + } + }); + }); + + describe("DEFAULT_APPROVAL_POLICY", () => { + it("gates the container backend outright and reads nothing else", () => { + expect(DEFAULT_APPROVAL_POLICY.rules).toEqual({ + shell: "read-only", + codemode: "read-only", + container: "always", + }); + expect(DEFAULT_APPROVAL_POLICY.fallback).toBe("always"); + }); + }); +}); diff --git a/examples/codemode/src/approval-policy.ts b/examples/codemode/src/approval-policy.ts new file mode 100644 index 00000000..4590ded6 --- /dev/null +++ b/examples/codemode/src/approval-policy.ts @@ -0,0 +1,482 @@ +/** + * Which commands a human has to approve before the agent runs them. + * + * There are two tiers here and they are not equally trustworthy. + * + * The first is a table keyed by backend id, and it decides on what a + * backend can reach rather than on what a command says: `container` + * is a full Linux userland with a public network, while `shell` and + * `codemode` are sandboxed and see only the workspace filesystem. + * Nothing is parsed to apply it, so nothing about it can be fooled by + * a command it did not anticipate. This is the boundary. + * + * The second is `read-only`, which classifies a command by reading + * it. A command runs unattended only when it is *recognizably* a + * read; anything the matcher does not understand needs a human. That + * direction matters: a matcher that fails closed turns an unparsed + * command into a question, while one that fails open turns it into an + * unreviewed mutation. Because it fails closed it is free to be + * wrong in one direction, and `approval-policy.effects.test.ts` is + * what holds it to that — it runs every command this file would allow + * and fails if any of them wrote. + * + * Treat the second tier as a way to ask fewer questions, not as the + * thing standing between the model and the files. Reading a command + * line to guess its effect is a heuristic, and a heuristic is the + * wrong place for a boundary. The boundary belongs at the capability + * layer: hand the backend a read-only view of the workspace and let + * the filesystem refuse the write. Think's built-in Bash tool already + * works that way, protecting files it did not mount during write-back + * so a script cannot delete what it was never given. Doing the same + * here would cover every caller rather than only the model's path. + * + * `decideApproval` must stay a pure function of the command and the + * backend. The AI SDK re-runs it when a paused turn resumes, and an + * approved call whose policy has since flipped to "no approval + * needed" is converted into a denial. Consulting a clock, or any + * mutable state, would make approvals decay on their own. + */ + +/** What a backend's commands cost in human attention. */ +export type BackendRule = + /** Every command on this backend needs a human. */ + | "always" + /** Recognized reads run unattended; everything else needs a human. */ + | "read-only" + /** Nothing on this backend needs a human. */ + | "never"; + +/** + * The language a backend's `command` field is written in. The + * `read-only` rule has to parse the command to classify it, and the + * codemode backend takes JavaScript where the others take a shell + * line. + */ +export type CommandDialect = "shell" | "javascript"; + +export interface ApprovalPolicy { + /** Rule per backend id, matching the ids the Workspace registered. */ + rules: Record; + /** + * Rule for a backend absent from `rules`. Defaults to `always`, so + * registering a new backend cannot quietly widen what runs + * unattended. + */ + fallback?: BackendRule; + /** + * Command language per backend. Defaults to JavaScript for the + * codemode backend and a shell line for everything else. + */ + dialects?: Record; +} + +export interface ApprovalDecision { + needsApproval: boolean; + /** + * One line explaining the verdict, shown to whoever works the + * approval queue. Populated for allowed commands too, so a + * transcript can say why nothing was asked. + */ + reason: string; +} + +/** + * The example's policy. The container backend is gated outright: it + * runs real binaries with public network access, so "which command is + * it" is the wrong question to be asking. The two sandboxed backends + * can reach only the workspace filesystem, so reads there are cheap + * enough to run unattended. + */ +export const DEFAULT_APPROVAL_POLICY: ApprovalPolicy = { + rules: { shell: "read-only", codemode: "read-only", container: "always" }, + fallback: "always", +}; + +const DEFAULT_DIALECTS: Record = { + shell: "shell", + container: "shell", + codemode: "javascript", +}; + +/** + * Shell characters that disqualify a line outright, whatever it runs. + * + * Redirection writes files, and substitution and subshells run + * commands this matcher would have to parse to see. None of them are + * decomposable the way a pipeline is, so their presence ends the + * question before the verb is considered. + * + * Backgrounding is here for a different reason: `ls &` leaves a + * process alive past the command, which is not something to wave + * through on the strength of the verb. + */ +const SHELL_METACHARACTERS: Array<[string, string]> = [ + [">", "redirects output"], + ["<", "redirects input"], + ["$", "expands a variable or substitutes a command"], + ["`", "substitutes a command"], + ["(", "groups or substitutes a command"], + [")", "groups or substitutes a command"], + ["\n", "hides a second command on another line"], + ["\r", "hides a second command on another line"], +]; + +/** + * Operators that join commands without touching the filesystem + * themselves. A line built from these is split on them and every stage + * classified on its own, so `ls | wc -l` reads and + * `ls; rm -rf /` does not. + * + * Longest first, so `&&` and `||` are found before a bare `&` or `|`. + */ +const COMMAND_SEPARATORS = ["&&", "||", ";", "|"]; + +/** + * Commands that only read. + * + * Deliberately conservative, and the omissions are the interesting + * part. `awk` can write through its own syntax rather than through a + * shell redirect. `sed` writes through `-i` and also + * through a `w` command buried in its script, which no matcher is + * going to find reliably. `uniq` and `tree` take an output file as a + * positional argument, so their writes do not look like flags at all. + * `date -s` sets the system clock. A verb whose writes cannot be read + * off its arguments does not belong here, because listing it would + * buy false confidence rather than fewer approvals. + */ +// Exported so approval-policy.effects.test.ts can build its corpus +// from the claim itself rather than from a copy of it. A verb added +// here then comes under test automatically, which is the point: the +// gap that let `find -delete` through was a case nobody had thought +// to write down. +export const READ_ONLY_COMMANDS = new Set([ + "basename", + "cat", + "cmp", + "cut", + "df", + "diff", + "dirname", + "du", + // echo and printf write to stdout and nowhere else. Sending that at + // a file takes a redirect, which gates the line whatever the verb. + "echo", + "printf", + "egrep", + "fgrep", + "file", + "find", + "grep", + "head", + "id", + "ls", + "pwd", + "readlink", + "realpath", + "sort", + "stat", + "tail", + "test", + "true", + "uname", + "wc", + "which", + "whoami", +]); + +/** + * Flags a read verb is allowed to carry, for the verbs that write when + * given the wrong one. + * + * A verb allowlist alone is not enough: `find` deletes with `-delete` + * and runs arbitrary commands with `-exec`, and `sort` writes to a + * file with `-o`. So for these verbs the flags are allowlisted too, + * and an unrecognized flag gates. That is the same shape as the + * `state.*` check — name what is allowed, treat everything else as a + * question — rather than a blocklist that has to keep up with every + * flag that happens to write. + * + * A verb absent from this table takes flags freely, which is only safe + * because the verbs in READ_ONLY_COMMANDS that are absent here have no + * flag that writes at all. + */ +export const READ_ONLY_FLAGS: Record> = { + find: new Set([ + "-a", + "-and", + "-atime", + "-depth", + "-empty", + "-follow", + "-group", + "-ilname", + "-iname", + "-inum", + "-ipath", + "-iregex", + "-links", + "-lname", + "-ls", + "-maxdepth", + "-mindepth", + "-mmin", + "-mtime", + "-name", + "-newer", + "-not", + "-o", + "-or", + "-path", + "-perm", + "-print", + "-print0", + "-printf", + "-prune", + "-readable", + "-regex", + "-samefile", + "-size", + "-type", + "-user", + "-xdev", + "-H", + "-L", + "-P", + ]), + sort: new Set([ + "-b", + "-c", + "-d", + "-f", + "-g", + "-h", + "-i", + "-k", + "-M", + "-n", + "-r", + "-s", + "-t", + "-u", + "-V", + "-z", + "--check", + "--ignore-case", + "--key", + "--numeric-sort", + "--reverse", + "--sort", + "--unique", + "--version-sort", + ]), +}; + +/** Git subcommands that only inspect history. */ +const READ_ONLY_GIT_SUBCOMMANDS = new Set([ + "blame", + "cat-file", + "describe", + "diff", + "log", + "ls-files", + "ls-tree", + "rev-parse", + "shortlog", + "show", + "status", +]); + +/** `state.*` calls that only read. The mutations are the complement. */ +const READ_ONLY_STATE_MEMBERS = new Set([ + "exists", + "find", + "grep", + "lstat", + "ls", + "readFile", + "readFileBytes", + "readdir", + "readlink", + "stat", +]); + +export function decideApproval( + call: { command: string; backend: string }, + policy: ApprovalPolicy = DEFAULT_APPROVAL_POLICY, +): ApprovalDecision { + const rule = policy.rules[call.backend] ?? policy.fallback ?? "always"; + + if (rule === "never") { + return { + needsApproval: false, + reason: `the ${call.backend} backend is configured to run without approval`, + }; + } + + if (rule === "always") { + return { + needsApproval: true, + reason: `the ${call.backend} backend requires approval for every command`, + }; + } + + const dialect = policy.dialects?.[call.backend] ?? DEFAULT_DIALECTS[call.backend] ?? "shell"; + const verdict = + dialect === "javascript" ? classifyJavaScript(call.command) : classifyShellLine(call.command); + + return { needsApproval: !verdict.readOnly, reason: verdict.reason }; +} + +interface Verdict { + readOnly: boolean; + reason: string; +} + +/** + * Classify a shell line. + * + * Redirection and substitution disqualify the line outright. What is + * left is a pipeline, which is split on its separators and judged one + * stage at a time: the line reads only if every stage does. Piping and + * sequencing touch no files themselves, so `ls | wc -l` is as much a + * read as `ls` is, while the stage-by-stage check is what still catches + * the `rm -rf` in `ls; rm -rf /`. + */ +function classifyShellLine(command: string): Verdict { + for (const [character, effect] of SHELL_METACHARACTERS) { + if (command.includes(character)) { + const shown = character === "\n" || character === "\r" ? "a newline" : `"${character}"`; + return { readOnly: false, reason: `${shown} ${effect}` }; + } + } + + // A lone `&` backgrounds; a doubled one is the and-then separator + // handled below. + if (/(? token.replace(/[|]/g, "\\$&")).join("|")})\\s*`, + ); + const stages = command.trim().split(separator); + + const verdicts: Verdict[] = []; + for (const stage of stages) { + if (stage.trim().length === 0) { + return { readOnly: false, reason: "a stage of the pipeline is empty" }; + } + const verdict = classifyShellCommand(stage); + if (!verdict.readOnly) return verdict; + verdicts.push(verdict); + } + + if (verdicts.length === 1) return verdicts[0]; + return { + readOnly: true, + reason: `every stage only reads (${verdicts.map((verdict) => verdict.reason).join("; ")})`, + }; +} + +/** Classify a single command, with no separators left in it. */ +function classifyShellCommand(command: string): Verdict { + const tokens = command + .trim() + .split(/\s+/) + .filter((token) => token.length > 0); + + // `LC_ALL=C sort file` runs sort, so step over any leading + // environment assignments to find the verb. + let index = 0; + while (index < tokens.length && /^[A-Za-z_][A-Za-z0-9_]*=/.test(tokens[index])) { + index += 1; + } + const words = tokens.slice(index); + if (words.length === 0) { + return { readOnly: false, reason: "the command is empty" }; + } + + const path = words[0]; + const verb = path.slice(path.lastIndexOf("/") + 1); + const args = words.slice(1); + + if (verb === "git") return classifyGit(args); + + if (!READ_ONLY_COMMANDS.has(verb)) { + return { readOnly: false, reason: `"${verb}" is not a recognized read-only command` }; + } + + const allowedFlags = READ_ONLY_FLAGS[verb]; + if (allowedFlags != null) { + for (const arg of args) { + if (!arg.startsWith("-")) continue; + // A negative number is a value, not a flag: `find -mtime -1`. + if (/^-\d/.test(arg)) continue; + // Compare on the name so `--output=path` is caught as --output. + const flag = arg.split("=")[0]; + if (!allowedFlags.has(flag)) { + return { + readOnly: false, + reason: `"${flag}" is not a recognized read-only flag for "${verb}"`, + }; + } + } + } + + return { readOnly: true, reason: `"${verb}" only reads` }; +} + +function classifyGit(args: string[]): Verdict { + // The first bare word is the subcommand. A global flag that takes a + // value (`git -C dir status`) shifts it, and the resulting mismatch + // gates the command, which is the safe direction to be wrong in. + const subcommand = args.find((arg) => !arg.startsWith("-")); + if (subcommand == null) { + return { readOnly: false, reason: "git without a subcommand" }; + } + if (!READ_ONLY_GIT_SUBCOMMANDS.has(subcommand)) { + return { readOnly: false, reason: `"git ${subcommand}" is not a recognized read-only command` }; + } + return { readOnly: true, reason: `"git ${subcommand}" only reads` }; +} + +/** + * Classify a codemode snippet by the `state.*` calls it names. + * + * This allowlists reads rather than blocklisting writes, which is what + * lets it catch `state["writeFile"]`: every mention of `state` has to + * resolve to a named member, and a computed access or an alias is + * unclassifiable and therefore gated. A blocklist would wave the same + * snippet through. + */ +function classifyJavaScript(command: string): Verdict { + const mentions = [...command.matchAll(/\bstate\b/g)]; + if (mentions.length === 0) { + // The codemode sandbox has no network and no other route into the + // store, so a snippet that never names `state` can only compute. + return { readOnly: true, reason: "the snippet does not reach the workspace" }; + } + + const members = new Set(); + for (const mention of mentions) { + const rest = command.slice((mention.index ?? 0) + mention[0].length); + const access = /^\s*(?:\.|\?\.)\s*([A-Za-z_$][\w$]*)/.exec(rest); + if (access == null) { + return { + readOnly: false, + reason: + "the snippet reaches state through a computed access or an alias, so its effect cannot be classified", + }; + } + members.add(access[1]); + } + + for (const member of members) { + if (!READ_ONLY_STATE_MEMBERS.has(member)) { + return { readOnly: false, reason: `state.${member} is not a recognized read-only call` }; + } + } + + const named = [...members].map((member) => `state.${member}`).join(", "); + return { readOnly: true, reason: `the snippet only reads (${named})` }; +} diff --git a/examples/codemode/src/index.ts b/examples/codemode/src/index.ts index 255f609b..15ef61c2 100644 --- a/examples/codemode/src/index.ts +++ b/examples/codemode/src/index.ts @@ -21,16 +21,25 @@ // workspace through its stub, so agency is opt-in per request // and the workspace stays a plain workspace. // +// The agent asks a human before it runs anything the approval policy +// holds back. A gated command stops the turn, and the turn resumes in +// a later request once someone answers, so a paused turn has to +// outlive the request that started it. That state lives in a second +// durable object, AgentSession, rather than in the workspace: the +// filesystem object stays a filesystem object. +// // Wire shape: // -// client ─► Worker ─┬─ /file, /exec deterministic, no model -// └─ /agent model loop + exec tool +// client ─► Worker ─┬─ /file, /exec deterministic, no model +// ├─ /agent model loop + exec tool +// └─ /approvals the human's side of the loop // │ -// ▼ (stub RPC) -// CodemodeExample DO (fs + 3 backends) -// ├─ shell ─► Dynamic Worker (just-bash) -// ├─ codemode ─► Dynamic Worker (JS sandbox) -// └─ container ─► Cloudflare Container (wsd) +// ┌─────────────┴──────────────┐ +// ▼ ▼ (stub RPC) +// AgentSession DO CodemodeExample DO (fs + 3 backends) +// paused turns, ├─ shell ─► Dynamic Worker (just-bash) +// pending approvals ├─ codemode ─► Dynamic Worker (JS sandbox) +// (no fs, no backends) └─ container ─► Cloudflare Container (wsd) import { DurableObject } from "cloudflare:workers"; @@ -49,8 +58,10 @@ import { } from "@cloudflare/workspace/backends/container"; import { WorkerBackend } from "@cloudflare/workspace/backends/worker"; -import { runAgentTurn } from "./agent.js"; +import { type AgentTranscript, runAgentTurn } from "./agent.js"; +import { AgentSession, type AgentSessionLike } from "./session.js"; import type { ExecWorkspaceLike } from "./tools/exec.js"; +import type { PausedTurn } from "./turn-store.js"; // Re-export so the runtime can build the loopback bindings the // durable object reaches through ctx.exports: @@ -59,7 +70,8 @@ import type { ExecWorkspaceLike } from "./tools/exec.js"; // - WorkspaceServiceProxy is the Fetcher the worker backend hands // into its Dynamic Worker so the in-isolate shell can reach // back to getWorkspace(). -export { WorkspaceProxy, WorkspaceServiceProxy }; +// The durable object that remembers agent turns paused for approval. +export { AgentSession, WorkspaceProxy, WorkspaceServiceProxy }; // --------------------------------------------------------------- // Durable Object: owns one Workspace with three backends. @@ -165,6 +177,11 @@ interface AgentRequest { prompt?: string; } +interface ApprovalRequest { + approved?: boolean; + reason?: string; +} + const MOUNT_ROOT = "/workspace"; function resolveMountPath(rest: string): string | null { @@ -195,6 +212,15 @@ export default { const agentMatch = url.pathname.match(/^\/c\/([^/]+)\/agent\/?$/); if (agentMatch) return handleAgent(request, env, agentMatch[1]); + const turnMatch = url.pathname.match(/^\/c\/([^/]+)\/agent\/([^/]+)\/?$/); + if (turnMatch) return handleTurn(request, env, turnMatch[1], turnMatch[2]); + + const approvalsMatch = url.pathname.match(/^\/c\/([^/]+)\/approvals\/?$/); + if (approvalsMatch) return handleApprovals(request, env, approvalsMatch[1]); + + const approvalMatch = url.pathname.match(/^\/c\/([^/]+)\/approvals\/([^/]+)\/?$/); + if (approvalMatch) return handleApproval(request, env, approvalMatch[1], approvalMatch[2]); + if (url.pathname === "/" || url.pathname === "") { return new Response( [ @@ -207,6 +233,11 @@ export default { " backend: shell | codemode | container", " POST /c//agent run an agent turn (JSON transcript)", " body: { prompt }", + " GET /c//agent/ one turn's record", + " GET /c//approvals commands waiting on a human", + " POST /c//approvals/ answer one of them; resumes the", + " turn once its last one is answered", + " body: { approved, reason? }", "", ].join("\n"), { headers: { "content-type": "text/plain" } }, @@ -253,6 +284,15 @@ async function handleFile( return new Response("method not allowed", { status: 405, headers: { allow: "GET, PUT" } }); } +// Run one command with no approval check. +// +// Approval and authorization are different questions. The gate lives +// in the exec tool because approval asks whether a person has seen a +// command the *model* chose; whoever posts here is that person, and +// asking them to confirm what they just typed answers nothing. What +// this route does not have is an authorization check, so it will run +// anything the caller sends. See approval-policy.ts for why the +// boundary belongs below ws.shell.exec() rather than on either path. async function handleExec(request: Request, env: Env, name: string): Promise { if (request.method !== "POST") { return new Response("method not allowed", { status: 405, headers: { allow: "POST" } }); @@ -300,28 +340,197 @@ async function handleAgent(request: Request, env: Env, name: string): Promise { + if (request.method !== "GET") { + return new Response("method not allowed", { status: 405, headers: { allow: "GET" } }); + } + const turn = await sessionFor(env, name).getTurn(turnId); + if (turn == null) return errorJSON(new Error(`no turn ${turnId}`), 404); + + // The message history is the model's working state, not something a + // client needs; everything else is the audit trail. + const { messages: _messages, awaiting: _awaiting, ...record } = turn; + return jsonResponse(record, 200); +} + +// GET the approval queue for this workspace. +async function handleApprovals(request: Request, env: Env, name: string): Promise { + if (request.method !== "GET") { + return new Response("method not allowed", { status: 405, headers: { allow: "GET" } }); + } + const pending = await sessionFor(env, name).pendingApprovals(); + return jsonResponse({ pending }, 200); +} + +// POST one human decision. Resolving the last outstanding approval on +// a turn resumes it in this same request, so the response is the +// transcript of the resumed pass — which may itself pause again. +async function handleApproval( + request: Request, + env: Env, + name: string, + approvalId: string, +): Promise { + if (request.method !== "POST") { + return new Response("method not allowed", { status: 405, headers: { allow: "POST" } }); + } + let body: ApprovalRequest; + try { + body = (await request.json()) as ApprovalRequest; + } catch { + return errorJSON(new Error("invalid JSON body"), 400); + } + if (typeof body.approved !== "boolean") { + return errorJSON(new Error("must provide approved as a boolean"), 400); + } + + const session = sessionFor(env, name); + const outcome = await session.resolveApproval(approvalId, body.approved, body.reason); + // Either the id was never issued, or somebody already answered it. + // Both mean this request changed nothing, and treating a repeat as a + // fresh approval would run the command twice. + if (outcome == null) { + return errorJSON(new Error(`no approval ${approvalId} is waiting for an answer`), 404); + } + + const { turn, ready, answers } = outcome; + + // One model step can ask about several commands, and the resume has + // to carry every answer at once, so the turn waits for the rest. + if (!ready) { + return jsonResponse( + { + status: "awaiting-approval", + turnId: turn.turnId, + text: "", + finishReason: "", + steps: 0, + stepsUsed: turn.stepsUsed, + toolCalls: turn.toolCalls, + pendingApprovals: turn.pending, + }, + 202, + ); + } + + try { + const transcript = await runAgentTurn({ + env, + workspace: await workspaceFor(env, name), + resume: { + messages: turn.messages, + approvals: answers.map((answer) => ({ + approvalId: answer.approvalId, + approved: answer.approved, + reason: answer.reason, + })), + }, + stepsUsed: turn.stepsUsed, }); + + // The turn's own history carries forward: commands from earlier + // passes belong to the same turn as the ones this pass ran. + const resumed = await recordTurn(env, name, turn.turnId, transcript, turn); + return transcriptJSON(resumed, transcript, 200); } catch (error) { return errorJSON(error, 500); } } +// See AgentSessionLike for why the stub is reached through a named +// interface rather than its own inferred type. +function sessionFor(env: Env, name: string): AgentSessionLike { + const stub = env.AgentSession.get(env.AgentSession.idFromName(name)); + return stub as unknown as AgentSessionLike; +} + +// The workspace stub drives the exec tool. Its exec/result methods +// behave exactly like a local Workspace at runtime, but the capnweb +// stub wraps them in promise-pipelined types that don't structurally +// match the plain ExecWorkspaceLike the tool declares. Cast at this +// one boundary. +async function workspaceFor(env: Env, name: string): Promise { + const stub = env.CodemodeExample.get(env.CodemodeExample.idFromName(name)); + return (await stub.getWorkspace()) as unknown as ExecWorkspaceLike; +} + +// Persist where a turn got to. A completed turn is kept too, so its +// record stays readable after the fact. +async function recordTurn( + env: Env, + name: string, + turnId: string, + transcript: AgentTranscript, + previous: Pick, +): Promise { + const now = Date.now(); + const turn: PausedTurn = { + turnId, + status: transcript.status, + messages: transcript.messages, + pending: transcript.pendingApprovals, + awaiting: transcript.pendingApprovals.map((approval) => approval.approvalId), + resolved: previous.resolved, + toolCalls: [...previous.toolCalls, ...transcript.toolCalls], + stepsUsed: transcript.stepsUsed, + createdAt: previous.createdAt, + updatedAt: now, + }; + await sessionFor(env, name).saveTurn(turn); + return turn; +} + +// The client-facing shape of a transcript. `messages` is deliberately +// left out: it is the model's working state, and the resume path reads +// it from storage rather than from the client. +function transcriptJSON(turn: PausedTurn, transcript: AgentTranscript, status: number): Response { + return jsonResponse( + { + status: turn.status, + turnId: turn.turnId, + text: transcript.text, + finishReason: transcript.finishReason, + steps: transcript.steps, + stepsUsed: turn.stepsUsed, + toolCalls: turn.toolCalls, + pendingApprovals: turn.pending, + }, + status, + ); +} + +function jsonResponse(body: unknown, status: number): Response { + return new Response(JSON.stringify(body), { + status, + headers: { "content-type": "application/json" }, + }); +} + function errorJSON(error: unknown, status: number): Response { const message = error instanceof Error ? error.message : String(error); const code = (error as { code?: string }).code; diff --git a/examples/codemode/src/session.ts b/examples/codemode/src/session.ts new file mode 100644 index 00000000..4f5fb348 --- /dev/null +++ b/examples/codemode/src/session.ts @@ -0,0 +1,92 @@ +/** + * `AgentSession` — the durable home for agent turns that are waiting on + * a human. + * + * This is a second durable object, deliberately separate from the one + * that owns the workspace. The workspace durable object is a + * filesystem with backends attached; it does not know an agent exists, + * and adding a queue of half-finished model turns to it would change + * that. So approval state gets its own object, addressed by the same + * name, and the workspace stays a workspace. + * + * It is worth noticing what this object cannot do: it holds no + * workspace stub and registers no backends, so it cannot run a + * command. The component that records approval decisions and the + * component that executes work are separate, which is a property worth + * having in the part of the system whose whole job is saying no. + * + * The model loop itself still runs in the Worker, as it did before. + * This object only remembers where a turn got to. + */ + +import { DurableObject } from "cloudflare:workers"; + +import { + type PausedTurn, + type PendingApprovalView, + type ResolvedApproval, + type TurnStorageLike, + TurnStore, +} from "./turn-store.js"; + +/** + * The session's RPC surface, as callers see it. + * + * Callers reach the durable object through this interface rather than + * through the stub's own inferred type. A `PausedTurn` carries + * `ModelMessage[]`, and asking tsc to wrap that union in the stub's + * recursive promise-pipelined types walks it past the + * instantiation-depth limit (TS2589). Naming the surface up front + * sidesteps that walk; the runtime shape is identical. + */ +export interface AgentSessionLike { + saveTurn(turn: PausedTurn): Promise; + getTurn(turnId: string): Promise; + pendingApprovals(): Promise; + resolveApproval( + approvalId: string, + approved: boolean, + reason?: string, + ): Promise<{ turn: PausedTurn; ready: boolean; answers: ResolvedApproval[] } | null>; +} + +export class AgentSession extends DurableObject { + readonly #turns: TurnStore; + + constructor(ctx: DurableObjectState, env: Env) { + super(ctx, env); + // The durable object's storage declares wider overloads than the + // four methods the store needs; the runtime shape matches. + this.#turns = new TurnStore(ctx.storage as unknown as TurnStorageLike); + } + + /** + * Record where a turn got to. Pruning runs on the way out so an + * approval nobody ever answers cannot accumulate forever. + */ + async saveTurn(turn: PausedTurn): Promise { + await this.#turns.save(turn); + await this.#turns.prune(); + } + + async getTurn(turnId: string): Promise { + return this.#turns.load(turnId); + } + + /** Everything waiting on a human, for whoever works the queue. */ + async pendingApprovals(): Promise { + return this.#turns.pending(); + } + + /** + * Record one decision. Null means the approval was not outstanding — + * an unknown id, or one that someone already answered. + */ + async resolveApproval( + approvalId: string, + approved: boolean, + reason?: string, + ): Promise<{ turn: PausedTurn; ready: boolean; answers: ResolvedApproval[] } | null> { + return this.#turns.resolve(approvalId, approved, reason); + } +} diff --git a/examples/codemode/src/tools/exec.ts b/examples/codemode/src/tools/exec.ts index 4de528ba..08d1d5a2 100644 --- a/examples/codemode/src/tools/exec.ts +++ b/examples/codemode/src/tools/exec.ts @@ -22,6 +22,8 @@ import { tool } from "ai"; import { z } from "zod"; +import { type ApprovalPolicy, decideApproval } from "../approval-policy.js"; + /** * Minimal subset of `@cloudflare/workspace` we depend on: a shell * facade with `exec(command, { cwd, encoding, backend })` whose @@ -68,6 +70,36 @@ export interface ExecToolOptions { defaultBackend: string; /** Truncate captured stdout/stderr above this many bytes. */ maxBytes?: number; + /** + * Which commands a human has to approve first. Omit to run + * everything the model asks for. + * + * A gated call does not execute: the AI SDK reports it as an + * approval request and ends the turn, so there is no provisional + * result and nothing to undo. + */ + policy?: ApprovalPolicy; + /** + * Called after each execution. The caller uses it to build a + * transcript. + * + * Reading executions from `generateText`'s steps would miss the + * interesting one: a command approved by a human runs before the + * resumed pass makes its first model call, so it belongs to no step. + * Recording here catches every execution regardless of which pass it + * happened on. + */ + onExec?: (call: ExecRecord) => void; +} + +/** One command the tool actually ran. */ +export interface ExecRecord { + command: string; + cwd: string | null; + backend: string; + exitCode: number; + stdout: string; + stderr: string; } const DEFAULT_MAX_BYTES = 64 * 1024; // 64 KiB per stream @@ -134,6 +166,13 @@ export function createExecTool(opts: ExecToolOptions) { return tool({ description, inputSchema, + // Must stay a pure function of the input. The AI SDK re-runs it + // when a paused turn resumes and downgrades an approved call to a + // denial if the answer has changed since the pause. + needsApproval: ({ command, backend }) => + opts.policy != null && + decideApproval({ command, backend: backend ?? opts.defaultBackend }, opts.policy) + .needsApproval, execute: async ({ command, cwd, backend }) => { const handle = await opts.workspace.shell.exec(command, { cwd, @@ -141,7 +180,7 @@ export function createExecTool(opts: ExecToolOptions) { backend, }); const result = await handle.result(); - return { + const record: ExecRecord = { command, cwd: cwd ?? null, backend: backend ?? opts.defaultBackend, @@ -149,6 +188,8 @@ export function createExecTool(opts: ExecToolOptions) { stdout: truncate(result.stdout, maxBytes), stderr: truncate(result.stderr, maxBytes), }; + opts.onExec?.(record); + return record; }, }); } diff --git a/examples/codemode/src/turn-store.test.ts b/examples/codemode/src/turn-store.test.ts new file mode 100644 index 00000000..956ca539 --- /dev/null +++ b/examples/codemode/src/turn-store.test.ts @@ -0,0 +1,273 @@ +import { beforeEach, describe, expect, it } from "vitest"; + +import type { PendingApproval } from "./agent.js"; +import { type PausedTurn, type TurnStorageLike, TurnStore } from "./turn-store.js"; + +/** The durable object's storage, reduced to a Map. */ +class MemoryStorage implements TurnStorageLike { + readonly entries = new Map(); + + async get(key: string): Promise { + return structuredClone(this.entries.get(key)) as T | undefined; + } + + async put(key: string, value: T): Promise { + this.entries.set(key, structuredClone(value)); + } + + async delete(key: string): Promise { + return this.entries.delete(key); + } + + async list(options?: { prefix?: string }): Promise> { + const prefix = options?.prefix ?? ""; + const found = new Map(); + for (const [key, value] of this.entries) { + if (key.startsWith(prefix)) found.set(key, structuredClone(value) as T); + } + return found; + } +} + +function approval(approvalId: string, command = "rm -rf /workspace"): PendingApproval { + return { + approvalId, + toolCallId: `call-${approvalId}`, + backend: "shell", + command, + cwd: null, + reason: `"rm" is not a recognized read-only command`, + }; +} + +function turn(overrides: Partial = {}): PausedTurn { + const createdAt = overrides.createdAt ?? 1_000; + const pending = overrides.pending ?? [approval("a1")]; + return { + turnId: "turn-1", + status: "awaiting-approval", + messages: [{ role: "user", content: "delete the workspace" }], + pending, + awaiting: pending.map((entry) => entry.approvalId), + resolved: [], + toolCalls: [], + stepsUsed: 1, + createdAt, + updatedAt: createdAt, + ...overrides, + }; +} + +describe("TurnStore", () => { + let storage: MemoryStorage; + let store: TurnStore; + + beforeEach(() => { + storage = new MemoryStorage(); + store = new TurnStore(storage); + }); + + describe("save and load", () => { + it("round-trips a paused turn", async () => { + const saved = turn(); + await store.save(saved); + expect(await store.load("turn-1")).toEqual(saved); + }); + + it("returns null for a turn it never saw", async () => { + expect(await store.load("nope")).toBeNull(); + }); + + it("overwrites an earlier version of the same turn", async () => { + await store.save(turn()); + await store.save(turn({ status: "completed", pending: [], stepsUsed: 4 })); + const loaded = await store.load("turn-1"); + expect(loaded?.status).toBe("completed"); + expect(loaded?.stepsUsed).toBe(4); + expect(loaded?.pending).toEqual([]); + }); + }); + + describe("pending", () => { + it("is empty to start", async () => { + expect(await store.pending()).toEqual([]); + }); + + it("flattens the queue across turns, tagging each with its turn", async () => { + await store.save(turn({ turnId: "turn-1", pending: [approval("a1")], createdAt: 1_000 })); + await store.save( + turn({ + turnId: "turn-2", + pending: [approval("a2"), approval("a3", "mkdir /workspace/d")], + createdAt: 2_000, + }), + ); + + const queue = await store.pending(); + expect(queue).toHaveLength(3); + expect(queue.map((entry) => entry.approvalId).sort()).toEqual(["a1", "a2", "a3"]); + expect(queue.find((entry) => entry.approvalId === "a1")?.turnId).toBe("turn-1"); + expect(queue.find((entry) => entry.approvalId === "a3")?.turnId).toBe("turn-2"); + expect(queue.find((entry) => entry.approvalId === "a3")?.command).toBe("mkdir /workspace/d"); + for (const entry of queue) { + expect(entry.requestedAt).toBeGreaterThan(0); + expect(entry.reason.length).toBeGreaterThan(0); + } + }); + + it("leaves out turns that are no longer waiting", async () => { + await store.save(turn({ turnId: "turn-1", status: "completed", pending: [] })); + expect(await store.pending()).toEqual([]); + }); + }); + + describe("resolve", () => { + it("records an approval and reports the turn ready to resume", async () => { + await store.save(turn()); + const outcome = await store.resolve("a1", true); + + expect(outcome).not.toBeNull(); + expect(outcome?.ready).toBe(true); + expect(outcome?.turn.pending).toEqual([]); + expect(outcome?.turn.resolved).toHaveLength(1); + expect(outcome?.turn.resolved[0]).toMatchObject({ approvalId: "a1", approved: true }); + }); + + it("records a rejection with its reason", async () => { + await store.save(turn()); + const outcome = await store.resolve("a1", false, "not now"); + + expect(outcome?.ready).toBe(true); + expect(outcome?.turn.resolved[0]).toMatchObject({ + approvalId: "a1", + approved: false, + reason: "not now", + }); + }); + + it("persists the resolution", async () => { + await store.save(turn()); + await store.resolve("a1", true); + const loaded = await store.load("turn-1"); + expect(loaded?.pending).toEqual([]); + expect(loaded?.resolved).toHaveLength(1); + }); + + it("holds a turn back until every approval in the step is resolved", async () => { + await store.save(turn({ pending: [approval("a1"), approval("a2")] })); + + const first = await store.resolve("a1", true); + expect(first?.ready).toBe(false); + expect(first?.turn.pending.map((entry) => entry.approvalId)).toEqual(["a2"]); + + const second = await store.resolve("a2", true); + expect(second?.ready).toBe(true); + expect(second?.turn.pending).toEqual([]); + expect(second?.turn.resolved.map((entry) => entry.approvalId)).toEqual(["a1", "a2"]); + }); + + it("refuses a second decision on the same approval", async () => { + // An approval UI is racy: two operators, two tabs, a double + // click. Resolving twice must not run the command twice. + await store.save(turn()); + expect(await store.resolve("a1", true)).not.toBeNull(); + expect(await store.resolve("a1", true)).toBeNull(); + expect(await store.resolve("a1", false)).toBeNull(); + }); + + it("refuses an approval id it never issued", async () => { + await store.save(turn()); + expect(await store.resolve("bogus", true)).toBeNull(); + }); + + it("drops the approval from the queue once resolved", async () => { + await store.save(turn({ pending: [approval("a1"), approval("a2")] })); + await store.resolve("a1", true); + expect((await store.pending()).map((entry) => entry.approvalId)).toEqual(["a2"]); + }); + + it("stamps when the decision was taken", async () => { + await store.save(turn()); + const outcome = await store.resolve("a1", true); + expect(outcome?.turn.resolved[0].at).toBeGreaterThan(0); + }); + + it("answers with this pause's decisions only", async () => { + // A turn that pauses twice accumulates decisions. Replaying an + // earlier pass's approval would run its command again, because + // the AI SDK only recognizes an approval as already satisfied + // when its tool result is in the last message. + await store.save( + turn({ + pending: [approval("a2")], + awaiting: ["a2"], + resolved: [{ approvalId: "a1", approved: true, at: 500 }], + }), + ); + + const outcome = await store.resolve("a2", true); + + expect(outcome?.ready).toBe(true); + expect(outcome?.answers.map((entry) => entry.approvalId)).toEqual(["a2"]); + // The earlier decision stays on the record for the audit trail. + expect(outcome?.turn.resolved.map((entry) => entry.approvalId)).toEqual(["a1", "a2"]); + }); + + it("answers with every decision from the current pause", async () => { + await store.save(turn({ pending: [approval("a1"), approval("a2")] })); + await store.resolve("a1", false, "no"); + const outcome = await store.resolve("a2", true); + + expect(outcome?.answers).toHaveLength(2); + expect(outcome?.answers.map((entry) => entry.approved)).toEqual([false, true]); + }); + }); + + describe("prune", () => { + it("drops a paused turn once it goes stale", async () => { + const store = new TurnStore(storage, { maxAgeMs: 1_000 }); + await store.save(turn({ turnId: "fresh", createdAt: 10_000, updatedAt: 10_000 })); + await store.save(turn({ turnId: "stale", createdAt: 1_000, updatedAt: 1_000 })); + + await store.prune(10_500); + + expect(await store.load("fresh")).not.toBeNull(); + expect(await store.load("stale")).toBeNull(); + }); + + it("keeps a paused turn inside its age limit", async () => { + const store = new TurnStore(storage, { maxAgeMs: 60_000 }); + await store.save(turn({ createdAt: 1_000, updatedAt: 1_000 })); + await store.prune(30_000); + expect(await store.load("turn-1")).not.toBeNull(); + }); + + it("keeps only the newest turns when there are too many", async () => { + const store = new TurnStore(storage, { maxTurns: 2 }); + for (const [id, at] of [ + ["oldest", 1_000], + ["middle", 2_000], + ["newest", 3_000], + ] as Array<[string, number]>) { + await store.save(turn({ turnId: id, createdAt: at, updatedAt: at })); + } + + await store.prune(3_000); + + expect(await store.load("oldest")).toBeNull(); + expect(await store.load("middle")).not.toBeNull(); + expect(await store.load("newest")).not.toBeNull(); + }); + + it("takes the pruned turn's approvals out of the queue with it", async () => { + const store = new TurnStore(storage, { maxAgeMs: 1_000 }); + await store.save(turn({ turnId: "stale", createdAt: 1_000, updatedAt: 1_000 })); + await store.prune(10_000); + + expect(await store.pending()).toEqual([]); + expect(await store.resolve("a1", true)).toBeNull(); + // Nothing is left behind that a later scan would trip over. + expect(storage.entries.size).toBe(0); + }); + }); +}); diff --git a/examples/codemode/src/turn-store.ts b/examples/codemode/src/turn-store.ts new file mode 100644 index 00000000..13e51c4c --- /dev/null +++ b/examples/codemode/src/turn-store.ts @@ -0,0 +1,221 @@ +/** + * Durable state for agent turns that paused for human approval. + * + * A turn pauses in the middle of a model loop and resumes in a later + * HTTP request, possibly minutes later, so everything needed to pick + * it up again has to outlive the request that started it. That is the + * whole reason this file exists: `runAgentTurn` is one `generateText` + * call inside a fetch handler and has nowhere of its own to keep a + * half-finished turn. + * + * The store is deliberately ignorant of the workspace. It holds + * message history and approval decisions; it cannot run a command. + * Keeping the component that records approvals separate from the + * component that executes work is worth the extra file. + * + * `TurnStore` takes a storage handle rather than a durable object + * state so the pause/resume bookkeeping can be tested against a plain + * Map. The durable object in `session.ts` supplies the real one. + */ + +import type { ModelMessage } from "ai"; + +import type { AgentToolCall, PendingApproval } from "./agent.js"; + +/** + * The slice of `DurableObjectState["storage"]` this store needs. The + * real handle satisfies it structurally. + */ +export interface TurnStorageLike { + get(key: string): Promise; + put(key: string, value: T): Promise; + delete(key: string): Promise; + list(options?: { prefix?: string }): Promise>; +} + +/** A decision a human already took, kept for the audit trail. */ +export interface ResolvedApproval { + approvalId: string; + approved: boolean; + reason?: string; + at: number; +} + +export interface PausedTurn { + turnId: string; + status: "awaiting-approval" | "completed" | "rejected"; + /** + * The turn's message history: the user prompt plus every assistant + * and tool message the model loop has produced so far, including the + * `tool-approval-request` parts. Replaying this with an approval + * response attached is what resumes the turn, so it must stay + * JSON-serializable. + */ + messages: ModelMessage[]; + /** Approvals still waiting on a human. */ + pending: PendingApproval[]; + /** + * Every approval id the current pause asked about, including the ones + * already answered. + * + * A resume replays only this pass's decisions. Replaying an older + * pass's approval would run its command a second time: the AI SDK + * skips an already-satisfied approval only when the matching tool + * result sits in the *last* message, and an earlier pass's result is + * further back than that. + */ + awaiting: string[]; + resolved: ResolvedApproval[]; + /** Every command the turn has run, across all of its passes. */ + toolCalls: AgentToolCall[]; + /** + * Model steps the turn has spent. Carried across passes because the + * AI SDK's step budget counts per `generateText` call, so a resumed + * turn would otherwise be handed a fresh allowance every time. + */ + stepsUsed: number; + createdAt: number; + updatedAt: number; +} + +/** A queue entry: one pending approval, with its turn for context. */ +export interface PendingApprovalView extends PendingApproval { + turnId: string; + requestedAt: number; +} + +export interface TurnStoreOptions { + /** Discard a turn nobody answered after this long. */ + maxAgeMs?: number; + /** Keep at most this many turns. */ + maxTurns?: number; +} + +const TURN_PREFIX = "turn:"; +const APPROVAL_PREFIX = "approval:"; + +/** A day, matching what an unanswered approval is plausibly worth. */ +const DEFAULT_MAX_AGE_MS = 24 * 60 * 60 * 1000; +const DEFAULT_MAX_TURNS = 50; + +export class TurnStore { + readonly #storage: TurnStorageLike; + readonly #maxAgeMs: number; + readonly #maxTurns: number; + + constructor(storage: TurnStorageLike, options: TurnStoreOptions = {}) { + this.#storage = storage; + this.#maxAgeMs = options.maxAgeMs ?? DEFAULT_MAX_AGE_MS; + this.#maxTurns = options.maxTurns ?? DEFAULT_MAX_TURNS; + } + + /** + * Write a turn and reindex its approvals. The index maps an approval + * id to its turn so resolving one is a single lookup rather than a + * scan, and dropping an id from the index is what makes a resolved + * approval unresolvable a second time. + */ + async save(turn: PausedTurn): Promise { + const previous = await this.load(turn.turnId); + if (previous != null) { + for (const approval of previous.pending) { + await this.#storage.delete(`${APPROVAL_PREFIX}${approval.approvalId}`); + } + } + + await this.#storage.put(`${TURN_PREFIX}${turn.turnId}`, turn); + for (const approval of turn.pending) { + await this.#storage.put(`${APPROVAL_PREFIX}${approval.approvalId}`, turn.turnId); + } + } + + async load(turnId: string): Promise { + return (await this.#storage.get(`${TURN_PREFIX}${turnId}`)) ?? null; + } + + /** Every approval waiting on a human, across all paused turns. */ + async pending(): Promise { + const turns = await this.#storage.list({ prefix: TURN_PREFIX }); + const queue: PendingApprovalView[] = []; + for (const turn of turns.values()) { + if (turn.status !== "awaiting-approval") continue; + for (const approval of turn.pending) { + queue.push({ ...approval, turnId: turn.turnId, requestedAt: turn.updatedAt }); + } + } + return queue; + } + + /** + * Record one human decision. + * + * Returns null when the approval is not outstanding — an id that was + * never issued, or one that has already been decided. Approval + * queues are racy, and a second decision arriving from another tab + * must not run the command a second time, so a no-op is reported as + * a no-op rather than treated as a fresh approval. + * + * `ready` says whether the turn can resume: a single model step can + * request several approvals, and the AI SDK wants every response in + * one message, so a turn waits until the last of them is answered. + */ + async resolve( + approvalId: string, + approved: boolean, + reason?: string, + now: number = Date.now(), + ): Promise<{ turn: PausedTurn; ready: boolean; answers: ResolvedApproval[] } | null> { + const turnId = await this.#storage.get(`${APPROVAL_PREFIX}${approvalId}`); + if (turnId == null) return null; + + const turn = await this.load(turnId); + if (turn == null) return null; + if (turn.status !== "awaiting-approval") return null; + if (!turn.pending.some((approval) => approval.approvalId === approvalId)) return null; + + const updated: PausedTurn = { + ...turn, + pending: turn.pending.filter((approval) => approval.approvalId !== approvalId), + resolved: [ + ...turn.resolved, + { approvalId, approved, ...(reason != null ? { reason } : {}), at: now }, + ], + updatedAt: now, + }; + + await this.#storage.put(`${TURN_PREFIX}${turnId}`, updated); + await this.#storage.delete(`${APPROVAL_PREFIX}${approvalId}`); + + return { + turn: updated, + ready: updated.pending.length === 0, + answers: updated.resolved.filter((entry) => updated.awaiting.includes(entry.approvalId)), + }; + } + + /** + * Drop turns that are too old or too numerous, along with their + * approval index entries. A turn nobody ever answers would otherwise + * sit in storage forever. + */ + async prune(now: number = Date.now()): Promise { + const turns = [...(await this.#storage.list({ prefix: TURN_PREFIX })).values()]; + + const stale = turns.filter((turn) => now - turn.updatedAt > this.#maxAgeMs); + const survivors = turns + .filter((turn) => now - turn.updatedAt <= this.#maxAgeMs) + .sort((left, right) => right.updatedAt - left.updatedAt); + const surplus = survivors.slice(this.#maxTurns); + + for (const turn of [...stale, ...surplus]) { + await this.#forget(turn); + } + } + + async #forget(turn: PausedTurn): Promise { + for (const approval of turn.pending) { + await this.#storage.delete(`${APPROVAL_PREFIX}${approval.approvalId}`); + } + await this.#storage.delete(`${TURN_PREFIX}${turn.turnId}`); + } +} diff --git a/examples/codemode/worker-configuration.d.ts b/examples/codemode/worker-configuration.d.ts index 337010f6..a676b8e4 100644 --- a/examples/codemode/worker-configuration.d.ts +++ b/examples/codemode/worker-configuration.d.ts @@ -1,15 +1,16 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types` (hash: 3d5a781bb03281c4e522aab7af31ae3c) +// Generated by Wrangler by running `wrangler types` (hash: f7d0847f7c89d072a618d5b30be433dc) // Runtime types generated with workerd@1.20260616.1 2026-05-26 experimental,nodejs_compat interface __BaseEnv_Env { LOADER: WorkerLoader; AI: Ai; CodemodeExample: DurableObjectNamespace; + AgentSession: DurableObjectNamespace; } declare namespace Cloudflare { interface GlobalProps { mainModule: typeof import("./src/index"); - durableNamespaces: "CodemodeExample"; + durableNamespaces: "CodemodeExample" | "AgentSession"; } interface Env extends __BaseEnv_Env {} } diff --git a/examples/codemode/wrangler.jsonc b/examples/codemode/wrangler.jsonc index 06ba81da..cf7568e2 100644 --- a/examples/codemode/wrangler.jsonc +++ b/examples/codemode/wrangler.jsonc @@ -49,6 +49,14 @@ { "name": "CodemodeExample", "class_name": "CodemodeExample" + }, + // Holds agent turns that paused for human approval. Separate + // from the workspace object on purpose: a paused turn has to + // outlive the request that started it, and the object that owns + // the filesystem has no business holding model state. + { + "name": "AgentSession", + "class_name": "AgentSession" } ] }, @@ -57,6 +65,10 @@ { "tag": "v1", "new_sqlite_classes": ["CodemodeExample"] + }, + { + "tag": "v2", + "new_sqlite_classes": ["AgentSession"] } ] } diff --git a/package-lock.json b/package-lock.json index ee1f0d39..524f5f4e 100644 --- a/package-lock.json +++ b/package-lock.json @@ -57,7 +57,9 @@ }, "devDependencies": { "@cloudflare/workers-types": "^4.20260616.1", + "just-bash": "^3.0.1", "typescript": "^6.0.3", + "vitest": "^4.1.7", "wrangler": "^4.96.0" } }, diff --git a/script/set-versions.mjs b/script/set-versions.mjs index 7d1a9027..53722ab8 100644 --- a/script/set-versions.mjs +++ b/script/set-versions.mjs @@ -27,6 +27,7 @@ const PACKAGES = [ // `git clone && wrangler dev` against any release tag pulls the // matching wsd image. const DOCKERFILES = [ + "examples/codemode/Dockerfile", "examples/container/Dockerfile", "examples/think/Dockerfile", "examples/think-compare-runtimes/Dockerfile.workspace",