From fc1aaecd74be39dabc8598b62f79973b210c752b Mon Sep 17 00:00:00 2001 From: linhdmn Date: Mon, 21 Sep 2026 15:33:38 +0700 Subject: [PATCH] =?UTF-8?q?feat(xdev):=20the=20loop=20can=20write=20and=20?= =?UTF-8?q?verify=20=E2=80=94=20run=5Ftests/write=5Ffile=20as=20real=20xde?= =?UTF-8?q?v=20turns?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Asked to "fix the off-by-one in parseConfig and run the tests", the loop rotated four tools and changed nothing: `write_file` returned `written: false` with a message saying `stub`, and `run_tests` returned `passed: 0, output: ""`. Two of the four v1 tools were fiction, and the approval gate in front of them was theatre. ## internal/xdev — the executor client Speaks xdev's `rpc` JSONL protocol: `{"type":"ready"}` first (validated against protocol version 1, so a future v2 fails loudly instead of mis-parsing), then id-tagged prompts answered by a stream of events and one terminal response. Four properties are load-bearing, and each has a test against a real child process rather than a mock: - **the ready frame is a hard gate** — a child that dies without one fails Start instead of blocking a run forever; - **events are read until the response** — a client that waits for the response without draining events deadlocks behind a full pipe, which xdev's own Serve comment names; - **one turn at a time** — xdev refuses a second concurrent prompt, so serialising here turns a server-side error into a client-side wait; - **the step's context kills the child** — a turn that outlives its budget must not leave a half-read stream and a sandbox still mutating a workspace. ## The tools `run_tests` and `write_file` each run as one xdev turn. That is §4.1's contract in code: agentloop owns the loop and the bound, xdev owns the turn — the file mutation, the test run, the model call inside it. Three outcomes are made distinguishable, because "no sandbox" must never read as "the tests passed": | Situation | Result | |---|---| | no executor configured | Success, `ran`/`written` false, message says so | | the turn failed | Success=false with the child's reason — an observation the loop keeps | | the turn ran | Success, `ran`/`written` true, the turn's text in `output` | `write_file` with no path fails closed and says why. ## The seam that was missing on the loop side `nextToolDefault` passed only `{step, goal}`, so every write was rejected for a missing path before the sandbox saw it. It now gives each tool the arguments that make it a real call — the readers get the goal as query text, `run_tests` gets the workspace, `write_file` gets a target. A real planner should name the file; this exists so the write path is exercisable end to end until it does, and the comment says so. ## Wiring AGENTLOOP_XDEV_BIN default "xdev" AGENTLOOP_XDEV_DIR workspace turns run in (default: a fresh temp dir, because a sandbox pointed at the server's cwd could edit agentloop itself) AGENTLOOP_XDEV_OFF disable One child is shared across runs on purpose (a per-run child would pay startup per step and lose the session between them) and it is still one turn at a time. A child that will not start is not fatal: the registry is built without a sandbox and the two tools report no executor — the same shape as LeanKG being down. The M6 eval factory keeps **no** sandbox and no knowledge client: the deploy gate stays deterministic and offline, and still passes 4/4. ## Verified live Against a sandbox speaking the protocol: step 2 run_tests -> ran: true, output: "ran `go test ./...`: 12 passed, 0 failed" step 3 write_file -> written: true, path: agentloop--step-3.txt With `AGENTLOOP_XDEV_OFF=1` the same run reports `ran: false` / `written: false` and names the missing executor. Against the real xdev binary the client reaches the ready frame and Start then fails on that machine's malformed models.yml — which is the correct shape: Start fails, and the server logs a warning and runs without a sandbox rather than dying. ## Checks `make check` green: gofmt, vet, 14/14 packages (7 new tests in internal/xdev, 3 new in internal/tools), lint 0 issues, PRD OK, selftest 12/12. Docs: PRD §4 marks three tools built and one stub; §13 records this; USAGE's config table gains the three variables and §9 splits what reads/writes today from what still does not. --- cmd/agentloop/main.go | 62 ++++- docs/PRD.md | 3 +- docs/USAGE.md | 22 +- internal/loop/runner.go | 27 +- internal/tools/registry_impl.go | 189 ++++++++++++-- internal/tools/registry_impl_test.go | 202 ++++++++++++++- internal/xdev/xdev.go | 373 +++++++++++++++++++++++++++ internal/xdev/xdev_test.go | 239 +++++++++++++++++ 8 files changed, 1081 insertions(+), 36 deletions(-) create mode 100644 internal/xdev/xdev.go create mode 100644 internal/xdev/xdev_test.go diff --git a/cmd/agentloop/main.go b/cmd/agentloop/main.go index 40e0fc8..b2c6d37 100644 --- a/cmd/agentloop/main.go +++ b/cmd/agentloop/main.go @@ -23,6 +23,7 @@ import ( "github.com/FreePeak/agentloop/internal/onegw" "github.com/FreePeak/agentloop/internal/planner" "github.com/FreePeak/agentloop/internal/tools" + "github.com/FreePeak/agentloop/internal/xdev" ) // Server holds the in-memory run store, the runner/gate registry, @@ -36,6 +37,9 @@ type Server struct { tools tools.ToolRegistry evalRunner *eval.Runner model *onegw.Client + // sandboxDir is the workspace sandboxed turns run in; empty when no + // executor is configured. + sandboxDir string } // NewServer creates a Server with the v1 tool set and empty stores. @@ -56,12 +60,14 @@ type Server struct { // and no model: the M6 deploy gate must stay deterministic and offline, so it // must not depend on LeanKG or onegw being up (PRD §11.4). func NewServer() *Server { - reg := tools.NewRegistryWithKnowledge(leankgFromEnv()) + sandboxDir := os.Getenv("AGENTLOOP_XDEV_DIR") + reg := tools.NewRegistryWithKnowledge(leankgFromEnv()).WithSandbox(xdevFromEnv(context.Background(), sandboxDir), sandboxDir) return &Server{ - runs: make(map[string]loop.RunResult), - runners: make(map[string]*loop.LoopRunner), - gates: make(map[string]*loop.ApprovalGate), - tools: reg, + runs: make(map[string]loop.RunResult), + runners: make(map[string]*loop.LoopRunner), + gates: make(map[string]*loop.ApprovalGate), + tools: reg, + sandboxDir: sandboxDir, model: onegw.New( envOr("AGENTLOOP_ONEGW_URL", "http://127.0.0.1:8080"), os.Getenv("AGENTLOOP_ONEGW_KEY"), @@ -73,6 +79,9 @@ func NewServer() *Server { // bare runner here made the adversarial case unscoreable — no // gate means no pause, and a pause is its whole premise. g := budget.New(cfg.CostBudget, float64(loop.DailyCeilingMult)*cfg.CostBudget) + // No sandbox and no knowledge client: the M6 deploy gate must + // stay deterministic and offline (PRD §11.4), so it must not + // depend on xdev or LeanKG being up. reg := tools.NewRegistry() gate := loop.NewApprovalGate() cfg.Gate = gate @@ -103,6 +112,48 @@ func leankgFromEnv() *leankg.Client { return leankg.New(envOr("AGENTLOOP_LEANKG_URL", "http://127.0.0.1:8090")) } +// xdevFromEnv starts ONE sandbox child and returns it, or nil when no +// executor is wanted. +// +// AGENTLOOP_XDEV_BIN default "xdev" +// AGENTLOOP_XDEV_DIR the workspace turns run in; default is a fresh +// temp dir, because a sandbox pointed at the +// server's own cwd can edit agentloop itself +// AGENTLOOP_XDEV_OFF any value disables it +// +// One child is shared by every run on purpose: `xdev rpc` keeps a session +// and a model conversation, so a per-run child would pay startup on every +// step and lose the thread between them. It is still one turn at a time — +// the client serialises prompts. +// +// A child that will not start is not fatal: the registry is built without +// a sandbox and `run_tests`/`write_file` report that no executor is +// configured. That is the same shape as LeanKG being down, and it keeps a +// deploy without xdev able to run the read-only paths. +func xdevFromEnv(ctx context.Context, dir string) *xdev.Client { + if os.Getenv("AGENTLOOP_XDEV_OFF") != "" { + return nil + } + if dir == "" { + d, err := os.MkdirTemp("", "agentloop-sandbox-") + if err != nil { + log.Printf("sandbox: no workspace: %v", err) + return nil + } + dir = d + } + c, err := xdev.Start(ctx, xdev.Config{ + Command: envOr("AGENTLOOP_XDEV_BIN", "xdev"), + Dir: dir, + }) + if err != nil { + log.Printf("sandbox: xdev unavailable, run_tests and write_file will report no executor: %v", err) + return nil + } + log.Printf("sandbox: xdev rpc started in %s", dir) + return c +} + type runRequest struct { Goal string `json:"goal"` Context string `json:"context"` @@ -139,6 +190,7 @@ func (s *Server) submitRun(w http.ResponseWriter, r *http.Request) { Goal: body.Goal, Context: body.Context, Gate: gate, + SandboxDir: s.sandboxDir, // M8: the loop can call the portfolio's gateway for synthesis. // The eval runner deliberately gets no model — the deploy gate // must stay deterministic and offline. diff --git a/docs/PRD.md b/docs/PRD.md index 7e38a1b..1616bb7 100644 --- a/docs/PRD.md +++ b/docs/PRD.md @@ -136,7 +136,7 @@ Rules (Ch.6, `design.md` §5): typed envelope `ToolResult(success, data, message | 3 | `run_tests` | `xdev rpc` (JSONL over stdio) running in a **restricted** `--add-dir` workspace | **yes** (sandboxed) | the verification half of the write-test-fix loop (Ch.14, ≤3 attempts); never a raw shell tool in v1 | | 4 | `write_file` | `xdev rpc` file tools | **yes** | read twin = `query`; approval gate by policy (§7.3); idempotency key on every call (§4.2) | -**Implementation status of that table** (so a reader can tell built from planned): `query` is **real** — one `POST /api/v1/query` through `internal/leankg`, wired from `AGENTLOOP_LEANKG_URL`; the tool reports the retrieval rung that answered in `Metadata`, and a LeanKG outage is recorded as a failed step, not a crash. `web_search`, `run_tests` and `write_file` are still stubs and say `stub` in their result message. The M6 eval runner deliberately gets a registry with **no** knowledge client, so the deploy gate stays deterministic and offline. +**Implementation status of that table** (so a reader can tell built from planned): `query` is **real** — one `POST /api/v1/query` through `internal/leankg`, wired from `AGENTLOOP_LEANKG_URL`; the tool reports the retrieval rung that answered in `Metadata`, and a LeanKG outage is recorded as a failed step, not a crash. `run_tests` and `write_file` are **real** too: each runs as one **xdev turn** through `internal/xdev` (the `rpc` JSONL protocol), wired from `AGENTLOOP_XDEV_{BIN,DIR,OFF}` — which is §4.1's contract in code, *agentloop drives xdev as a tool executor for one already-planned step* — and with no sandbox reachable both report *no executor configured* with `written`/`ran` false rather than a success nobody earned. `web_search` is the last stub. The M6 eval runner deliberately gets a registry with **no** knowledge client and **no** sandbox, so the deploy gate stays deterministic and offline. **Divergence from `design.md` §18, stated on purpose:** `design.md`'s starting set is *"3 reads + 1 search + 1 ticket/incident writer"*. This PRD ships `write_file` + `run_tests` instead of the ticket writer, because the loop's own verification primitive (`run_tests`) is what makes the evaluate phase real, and because a ticket writer is a template concern (Appendix C row 6) rather than loop infrastructure. That is the one place the PRD knowingly overrides the architecture of record; everything else in §4 is narrower than `design.md`, not different from it. @@ -730,6 +730,7 @@ Written the way an unfriendly reviewer would write it, then answered. Every find **Read next.** §13.1 (scope → milestones), §17 (defaults), §18 (where to discount the source), §22 (this document's own weaknesses). +* Last updated: 2026-09-21 (The loop can finally **write and verify**. `internal/xdev` speaks xdev's `rpc` JSONL protocol (ready-frame version gate, event-before-response interleaving, one turn at a time, child killed when the step's budget expires) and `run_tests`/`write_file` run as one xdev turn each, in a sandboxed workspace from `AGENTLOOP_XDEV_DIR`. With no sandbox the two report *no executor configured* and `written`/`ran` stay false — "no sandbox" can never read as "the tests passed". `nextToolDefault` now gives each tool the arguments it needs to be a real call, so a write has a target instead of failing closed on a missing path. `web_search` is the last stub.) * Last updated: 2026-09-21 (The M6 eval gate was reporting **0.5 — deploy blocked — for days** while `go test ./...` was green. Two causes, both real bugs: the eval factory built a *gateless, plannerless* runner, so the adversarial case's premise ("the gate holds") was unreachable by construction; and the score functions asserted step counts calibrated against that gateless 9-step rotation, so the honest gated behaviour — three read steps then a hold on the writer — scored 0.3. The factory now builds the service's runner (minus the model, as §11.4 requires) and the scores describe the category's expected behaviour rather than a step count. Two tests close the hole that let it rot: `TestEval_DefaultSuiteIsGreen` asserts the *default* suite passes (every prior test used its own factory or its own score fn — nothing pinned the real one) and `TestEval_DefaultSuiteCanFail` requires the adversarial case to fail against a gateless runner. Verified by breaking `Categorize` and watching the endpoint block at 0.75.) * Last updated: 2026-09-21 (CI exists: `.github/workflows/ci.yml` runs gofmt, `go vet`, `go test`, `golangci-lint` (config pinned in `.golangci.yml`) and `docs/check-prd.py` **plus its `--selftest`** on every PR — the "deploys blocked on the suite" half of M6 is no longer aspirational (#30 closes #27); `make check` runs the same five steps locally. **Caveat recorded here rather than discovered later:** `FreePeak/agentloop` is private, and GitHub-hosted runners are billed — until the org's spending limit is raised, both jobs fail at dispatch with a billing error that says nothing about the code (seen on PR #31). `make check` is the fallback that keeps the gate honest in the meantime. Getting the lint job to a clean baseline exposed real code, not just style: an unused `currentTier` field, an unused `maxLandmarkTokens` budget that nothing enforced (recorded as §9.1's fourth accepted ceiling instead, since landmarks are never evicted), and a `Categorize` switch staticcheck flagged.) * Last updated: 2026-09-21 (Docs synced to `b492cc2`. §13's M5 row recorded the `Categorize` fix (#25) and its stale test count; the M6 row now says out loud that the "deploys blocked on the full-suite gate" half is **not** enforced — there is no CI in this repo (#27); the §13 task-record line stopped claiming the repo has no commits/remote and now records the tracker state: priority bands P0–P3 defined and every open issue labelled (#28, half done), with issue closure still not PR-linked. `docs/USAGE.md` lost a duplicated line and its "does nothing" summary was split into what reads today vs what still does not. `docs/check-prd.py` stops false-failing on §-refs into other docs — the cause of the long-standing dual-§-ref FAIL, whose refs pointed at `docs/JEV-INTEGRATION.md`, not this PRD.) diff --git a/docs/USAGE.md b/docs/USAGE.md index 84a5a2d..61f9f98 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -37,7 +37,7 @@ Read this before you plan work around it. As of this writing: | HITL approval gate wired into the runner | **implemented and working**: the gate holds a write, and approving it **resumes** the run. The read-only three tools never interrupt; `write_file` always holds — §3, §9 | | HTTP API, admin console, eval harness | **implemented** | | Model calls to onegw | **partial** — one outbound call, at synthesis (§2, §9). The loop's *planning* still makes none | -| The four built-in tools | **one real, three stubs** — `query` reaches LeanKG and returns its envelope; `web_search`, `run_tests`, `write_file` return canned results whose message says `stub` (`internal/tools/registry_impl.go`) | +| The four built-in tools | **three real, one stub** — `query` reaches LeanKG; `run_tests` and `write_file` each run as one **xdev turn** in a sandboxed workspace (so the loop can actually write and verify); `web_search` still returns a canned result whose message says `stub` | | The planner | **deterministic**, no model calls; model-driven planning is the documented production path | | M7 multi-agent (`internal/supervisor`) | **gated shut** by design — refused unless one of [PRD §10](PRD.md#10-multi-agent-stance)'s four conditions is met | @@ -102,6 +102,9 @@ Deployment facts, not compiled defaults: | `AGENTLOOP_ONEGW_COMBO` | `dev` | combo name, used as the wire `model` | | `AGENTLOOP_LEANKG_URL` | `http://127.0.0.1:8090` | LeanKG REST root behind the `query` tool | | `AGENTLOOP_LEANKG_OFF` | *(unset)* | any value disables the knowledge client; `query` then says no service is configured | +| `AGENTLOOP_XDEV_BIN` | `xdev` | the sandbox binary; agentloop speaks its `rpc` JSONL protocol | +| `AGENTLOOP_XDEV_DIR` | a fresh temp dir | the workspace `write_file`/`run_tests` turns run in | +| `AGENTLOOP_XDEV_OFF` | *(unset)* | any value disables the sandbox; those two tools then report no executor | Take the onegw key from `onegw.toml`'s `[auth] [[auth.keys]]`; the combo must exist there too, since the client sends whatever name you give it and onegw @@ -293,11 +296,18 @@ Stated plainly, so nobody discovers it the hard way: stored state *and* signals the runner's kill channel, which `Run()` checks at every step boundary. The stop is bounded by one step, not instant — a step already in flight finishes first. -- **Three of the four tools are stubs.** `query` is **real**: it reaches LeanKG - `POST /api/v1/query` and reports the retrieval rung that answered. - `web_search` has no client, and `run_tests`/`write_file` should go through xdev - rpc in a restricted workspace. Each stub says so in its result message, so a - step that did nothing cannot be read as one that worked. +- **`web_search` is the last stub.** It has no client, and says so in its result + message, so a step that did nothing cannot be read as one that worked. +- **`write_file` and `run_tests` are real, and they need a sandbox.** Each runs + as **one xdev turn** (`internal/xdev` speaks xdev's `rpc` JSONL protocol). + agentloop owns the loop and the bound; xdev owns the turn — the file mutation, + the test run, the model call inside it (PRD §4.3). With no xdev binary + reachable, both tools report *no executor configured* and `Data[written]` / + `Data[ran]` stay false, so "no sandbox" can never be read as "the write + succeeded" or "the tests passed". +- **The turn is bounded by the step, not by xdev.** If a step's budget expires + mid-turn the client kills the child, rather than leaving a half-read stream + and a sandbox still mutating a workspace. - **Retrieval needs a LeanKG to point at.** `AGENTLOOP_LEANKG_URL` (default `http://127.0.0.1:8090`); `AGENTLOOP_LEANKG_OFF=1` disables the client, and the tool then says no knowledge service is configured. A LeanKG that is *down* is diff --git a/internal/loop/runner.go b/internal/loop/runner.go index 43352b6..c58cf61 100644 --- a/internal/loop/runner.go +++ b/internal/loop/runner.go @@ -77,6 +77,7 @@ type RunnerConfig struct { ConfidenceFloor float64 EscalationThreshold float64 Gate *ApprovalGate // M5: fail-closed HITL gate (nil = bypass, tests) + SandboxDir string // workspace a sandboxed turn runs in (empty = tool's own default) Model ModelClient // M8: outbound model transport (nil = deterministic synthesis only) } @@ -638,10 +639,34 @@ func synthesizePartial(steps []StepRecord) string { // nextToolDefault returns the tool and args for a step. M1 uses a // deterministic rotation over the 4 v1 tools for testing; // production replaces this with the planner's output. +// nextToolDefault is the deterministic stand-in for model-driven tool +// selection (PRD §4.3: "agentloop owns the loop, xdev owns the turn"). +// It rotates the v1 surface and gives each tool the arguments it needs to +// be a real call rather than a shape: +// +// - the readers get the goal as their query text; +// - run_tests gets the workspace to verify; +// - write_file gets a target path, because a write with no target now +// fails closed — the previous rotation passed only step/goal, so every +// write was rejected before the model inside the sandbox ever saw it. +// +// ponytail: the path is derived from cfg.RunID, not from what the run +// learned. A real planner names the file; this exists so the write path is +// exercisable end to end until that lands. func nextToolDefault(step int, cfg RunnerConfig) (string, map[string]any) { tools := []string{"query", "web_search", "run_tests", "write_file"} name := tools[step%len(tools)] - return name, map[string]any{"step": step, "goal": cfg.Goal} + args := map[string]any{"step": step, "goal": cfg.Goal} + switch name { + case "run_tests": + if cfg.SandboxDir != "" { + args["cwd"] = cfg.SandboxDir + } + case "write_file": + args["path"] = fmt.Sprintf("agentloop-%s-step-%d.txt", cfg.RunID, step) + args["content"] = fmt.Sprintf("step %d of run %s: %s\n", step, cfg.RunID, cfg.Goal) + } + return name, args } // dedupKey is the agentloop idempotency fingerprint: diff --git a/internal/tools/registry_impl.go b/internal/tools/registry_impl.go index 2726404..24cc1a7 100644 --- a/internal/tools/registry_impl.go +++ b/internal/tools/registry_impl.go @@ -3,17 +3,21 @@ // registry's Execute routes through AgentBase.execute(tool, args) // (PRD §4.3 — agentloop owns the loop, xdev owns the turn). // -// One tool is real today: `query` reaches LeanKG when a knowledge -// client is injected (cmd/agentloop wires it from the environment). -// The other three are still stubs, and they say so in their message -// rather than returning a success the loop cannot tell apart from work. +// Two of the four are real: `query` reaches LeanKG when a knowledge +// client is injected, and `run_tests`/`write_file` run as one xdev turn +// when a sandbox client is injected (cmd/agentloop wires both from the +// environment). `web_search` has no client yet, and it says so in its +// message rather than returning a success the loop cannot tell apart +// from work. package tools import ( "context" "fmt" + "strings" "github.com/FreePeak/agentloop/internal/leankg" + "github.com/FreePeak/agentloop/internal/xdev" ) // Registry is an in-memory ToolRegistry pre-loaded with the 4 v1 tools. @@ -24,6 +28,14 @@ type Registry struct { // instead of answering with an empty hit list, which would read as // "the graph has nothing" (P100 honest degradation). Knowledge *leankg.Client + // Sandbox is the xdev client behind `run_tests` and `write_file`. Nil + // means no executor is configured, and those two tools say so rather + // than reporting success for work nobody did. This is the seam PRD + // §4.3 draws: agentloop owns the loop, xdev owns the turn. + Sandbox *xdev.Client + // SandboxDir is the workspace a turn runs in. Every write and every + // test run is confined to it. + SandboxDir string } // NewRegistry returns a Registry with the default v1 tool set and no @@ -38,6 +50,16 @@ func NewRegistryWithKnowledge(c *leankg.Client) *Registry { return &Registry{Tools: append([]Tool(nil), DefaultTools...), Knowledge: c} } +// WithSandbox returns a copy of the Registry whose `run_tests` and +// `write_file` tools run as xdev turns in dir. Passing a nil client leaves +// them reporting "no executor configured", which is the honest answer. +func (r *Registry) WithSandbox(c *xdev.Client, dir string) *Registry { + out := *r + out.Tools = append([]Tool(nil), r.Tools...) + out.Sandbox, out.SandboxDir = c, dir + return &out +} + // List returns the registered tools. func (r *Registry) List() []Tool { return r.Tools } @@ -76,18 +98,9 @@ func (r *Registry) Execute(ctx context.Context, name string, args map[string]any Message: "stub: web_search results would come from onegw provider kind=searxng", }, nil case "run_tests": - return ToolResult{ - Success: true, - Data: map[string]any{"passed": 0, "failed": 0, "output": ""}, - Message: "stub: run_tests executes via xdev rpc in a restricted --add-dir workspace", - }, nil + return r.runTests(ctx, args) case "write_file": - return ToolResult{ - Success: true, - Data: map[string]any{"path": "", "written": false}, - Message: "stub: write_file executes via xdev rpc in a restricted --add-dir workspace; nothing was written", - Metadata: map[string]string{"idempotency_key": "pending"}, - }, nil + return r.writeFile(ctx, args) } return ToolResult{}, fmt.Errorf("unknown tool: %s", name) } @@ -145,3 +158,149 @@ func (r *Registry) queryKnowledge(ctx context.Context, args map[string]any) (Too }, }, nil } + +// runTests runs the workspace's tests as one xdev turn. +// +// It is a prompt to the sandbox rather than a fork/exec of `go test`, +// because the whole reason xdev is in this architecture is that it owns +// the turn: the model inside it decides how to run the suite, reads the +// failures and can fix them inside the same step (PRD §4.1 — "agentloop +// drives xdev as a tool executor for ONE already-planned step"). agentloop +// still owns the bound: the turn's context is the step's. +// +// The result is deliberately honest about which of three things happened: +// no executor, a failed turn, or a turn that ran. A caller can never +// mistake "no sandbox" for "tests passed". +func (r *Registry) runTests(ctx context.Context, args map[string]any) (ToolResult, error) { + if r.Sandbox == nil { + return ToolResult{ + Success: true, + Data: map[string]any{"ran": false, "output": ""}, + Message: "no executor configured (AGENTLOOP_XDEV_OFF); the step ran and no tests were executed", + }, nil + } + cwd, _ := args["cwd"].(string) + extra, _ := args["args"].([]any) + prompt := testPrompt(cwd, extra) + + turn, err := r.Sandbox.Prompt(ctx, prompt) + if err != nil { + // A sandbox that failed is an observation, not a crash: the run + // keeps its steps and the reason is recorded (same rule as LeanKG). + return ToolResult{ + Success: false, + Data: map[string]any{"ran": false, "error": err.Error()}, + Message: fmt.Sprintf("xdev test turn failed: %v", err), + Metadata: map[string]string{"sandbox": "xdev"}, + }, nil + } + return ToolResult{ + Success: true, + Data: map[string]any{ + "ran": true, + "output": turn.Text, + "stop_reason": turn.StopReason, + "model": turn.Model, + "events": len(turn.Events), + "duration_ms": turn.Duration.Milliseconds(), + }, + Message: summarise(turn.Text), + Metadata: map[string]string{"sandbox": "xdev", "stop_reason": turn.StopReason}, + }, nil +} + +// writeFile asks the sandbox to make a file change, as one xdev turn. +// +// The approval gate is upstream (Categorize puts write_file in +// CatApprove), so by the time this runs a human has already said yes — +// which is why this function does not re-prompt. The idempotency key rides +// back in Metadata because the caller (the runner) persists it. +func (r *Registry) writeFile(ctx context.Context, args map[string]any) (ToolResult, error) { + path, _ := args["path"].(string) + content, _ := args["content"].(string) + if path == "" { + // Fail closed and loudly: a write with no target is how a run + // scribbles on a workspace it cannot name. + return ToolResult{ + Success: false, + Data: map[string]any{"path": "", "written": false}, + Message: "write_file requires a path", + }, nil + } + if r.Sandbox == nil { + return ToolResult{ + Success: true, + Data: map[string]any{"path": path, "written": false}, + Message: "no executor configured (AGENTLOOP_XDEV_OFF); the step ran and nothing was written", + Metadata: map[string]string{"idempotency_key": "pending"}, + }, nil + } + + turn, err := r.Sandbox.Prompt(ctx, writePrompt(path, content)) + if err != nil { + return ToolResult{ + Success: false, + Data: map[string]any{"path": path, "written": false, "error": err.Error()}, + Message: fmt.Sprintf("xdev write turn failed: %v", err), + Metadata: map[string]string{"sandbox": "xdev", "idempotency_key": "pending"}, + }, nil + } + return ToolResult{ + Success: true, + Data: map[string]any{ + "path": path, + "written": true, + "output": turn.Text, + "stop_reason": turn.StopReason, + "duration_ms": turn.Duration.Milliseconds(), + }, + Message: summarise(turn.Text), + Metadata: map[string]string{"sandbox": "xdev", "stop_reason": turn.StopReason, "idempotency_key": "pending"}, + }, nil +} + +// testPrompt frames the verification half of the write-test-fix loop. +// One instruction, one job: run the suite, report what failed. +func testPrompt(cwd string, extra []any) string { + var b strings.Builder + b.WriteString("Run this workspace's test suite and report the result.") + if cwd != "" { + fmt.Fprintf(&b, " Work in %s.", cwd) + } + if len(extra) > 0 { + parts := make([]string, 0, len(extra)) + for _, a := range extra { + if s, ok := a.(string); ok { + parts = append(parts, s) + } + } + if len(parts) > 0 { + fmt.Fprintf(&b, " Pass these arguments: %s.", strings.Join(parts, " ")) + } + } + b.WriteString(" Report: the command you ran, the number of tests that passed and failed, " + + "and the exact failure output for anything that failed. Do not change any file.") + return b.String() +} + +// writePrompt frames one file change. It names the path once and gives the +// content verbatim: an agent asked to "improve" a file will, and this step +// already decided what the change is. +func writePrompt(path, content string) string { + return fmt.Sprintf("Write exactly this content to the file %q, creating it if it does not exist. "+ + "Make no other change, and report only whether the write succeeded.\n\n%s", path, content) +} + +// summarise trims a turn's report to the first non-empty line, so a step +// record stays readable while the full text stays in Data. +func summarise(text string) string { + for _, line := range strings.Split(text, "\n") { + if s := strings.TrimSpace(line); s != "" { + if len(s) > 200 { + return s[:200] + "…" + } + return s + } + } + return "turn produced no output" +} diff --git a/internal/tools/registry_impl_test.go b/internal/tools/registry_impl_test.go index 44b356d..0b0d6e4 100644 --- a/internal/tools/registry_impl_test.go +++ b/internal/tools/registry_impl_test.go @@ -5,9 +5,14 @@ import ( "encoding/json" "net/http" "net/http/httptest" + "os" + "os/exec" + "path/filepath" "testing" + "time" "github.com/FreePeak/agentloop/internal/leankg" + "github.com/FreePeak/agentloop/internal/xdev" ) // The registry must carry LeanKG's own envelope back to the loop, and @@ -109,18 +114,61 @@ func TestQuery_FallsBackToTheGoalForQueryText(t *testing.T) { } } -// The other three tools are still stubs. They must say so rather than -// return a success the loop cannot tell apart from real work. -func TestStubsAdvertiseThemselves(t *testing.T) { +// `web_search` has no client. It must say so rather than return a success +// the loop cannot tell apart from real work. +func TestWebSearchAdvertisesItselfAsAStub(t *testing.T) { reg := NewRegistry() - for _, name := range []string{"web_search", "run_tests", "write_file"} { - got, err := reg.Execute(context.Background(), name, map[string]any{}) + got, err := reg.Execute(context.Background(), "web_search", map[string]any{}) + if err != nil { + t.Fatalf("Execute() error: %v", err) + } + if !contains(got.Message, "stub") { + t.Errorf("Message = %q, want it to say stub", got.Message) + } +} + +// `run_tests` and `write_file` with no executor must report that, not a +// success. "No sandbox" and "the tests passed" must never look alike. +func TestSandboxToolsReportNoExecutor(t *testing.T) { + reg := NewRegistry() // no sandbox client + cases := []struct { + name string + args map[string]any + want string + }{ + {"run_tests", map[string]any{}, "no executor configured"}, + {"write_file", map[string]any{"path": "a.txt", "content": "x"}, "no executor configured"}, + } + for _, c := range cases { + got, err := reg.Execute(context.Background(), c.name, c.args) if err != nil { - t.Fatalf("%s: Execute() error: %v", name, err) + t.Fatalf("%s: Execute() error: %v", c.name, err) + } + if !contains(got.Message, c.want) { + t.Errorf("%s: Message = %q, want it to mention %q", c.name, got.Message, c.want) } - if !contains(got.Message, "stub") { - t.Errorf("%s: Message = %q, want it to say stub", name, got.Message) + if written, ok := got.Data["written"].(bool); ok && written { + t.Errorf("%s: Data[written] = true with no executor", c.name) } + if ran, ok := got.Data["ran"].(bool); ok && ran { + t.Errorf("%s: Data[ran] = true with no executor", c.name) + } + } +} + +// A write with no target fails closed: a run must not scribble on a +// workspace it cannot name. +func TestWriteFileRequiresAPath(t *testing.T) { + reg := NewRegistry() + got, err := reg.Execute(context.Background(), "write_file", map[string]any{"content": "x"}) + if err != nil { + t.Fatalf("Execute() error: %v", err) + } + if got.Success { + t.Error("write_file with no path reported success") + } + if !contains(got.Message, "path") { + t.Errorf("Message = %q, want it to name the missing path", got.Message) } } @@ -132,3 +180,141 @@ func contains(s, sub string) bool { } return false } + +// With a sandbox client, `write_file` and `run_tests` become real turns. +// The fixture is a process speaking xdev's wire protocol, so this asserts +// the whole path: registry -> client -> child -> result. +func TestSandboxToolsRunRealTurns(t *testing.T) { + bin := xdevFixture(t) + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + defer cancel() + c, err := xdev.Start(ctx, xdev.Config{Command: bin}) + if err != nil { + t.Fatalf("xdev.Start() error: %v", err) + } + defer c.Close() + + reg := NewRegistry().WithSandbox(c, t.TempDir()) + + // write_file: the prompt must carry the path and the content verbatim. + got, err := reg.Execute(ctx, "write_file", map[string]any{"path": "b.txt", "content": "hello"}) + if err != nil { + t.Fatalf("write_file: %v", err) + } + if !got.Success { + t.Fatalf("write_file failed: %s", got.Message) + } + if written, _ := got.Data["written"].(bool); !written { + t.Errorf("write_file: written = %v, want true with a live sandbox", got.Data["written"]) + } + if got.Metadata["sandbox"] != "xdev" { + t.Errorf("write_file metadata = %v, want sandbox=xdev", got.Metadata) + } + + // run_tests: same path, different prompt. + got, err = reg.Execute(ctx, "run_tests", map[string]any{}) + if err != nil { + t.Fatalf("run_tests: %v", err) + } + if ran, _ := got.Data["ran"].(bool); !ran { + t.Errorf("run_tests: ran = %v, want true with a live sandbox", got.Data["ran"]) + } + if !contains(got.Message, "did:") { + t.Errorf("run_tests Message = %q, want the turn's report", got.Message) + } +} + +// A sandbox turn that fails is an observation with the reason recorded — +// not a hard error that would abort the run. +func TestSandboxTurnFailureIsAnObservation(t *testing.T) { + bin := xdevFixture(t) + t.Setenv("FAKE_MODE", "fail-turn") + + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + defer cancel() + c, err := xdev.Start(ctx, xdev.Config{Command: bin}) + if err != nil { + t.Fatalf("xdev.Start() error: %v", err) + } + defer c.Close() + + reg := NewRegistry().WithSandbox(c, t.TempDir()) + got, err := reg.Execute(ctx, "run_tests", map[string]any{}) + if err != nil { + t.Fatalf("Execute() returned a hard error: %v (it must be an observation)", err) + } + if got.Success { + t.Error("Success = true on a failed turn") + } + if !contains(got.Message, "model exploded") { + t.Errorf("Message = %q, want the child's reason", got.Message) + } +} + +// xdevFixture builds the same wire-protocol child internal/xdev's tests +// use, so this package's tests exercise the real client against a real +// process rather than a stub. +func xdevFixture(t *testing.T) string { + t.Helper() + dir := t.TempDir() + bin := filepath.Join(dir, "fake-xdev") + src := filepath.Join(dir, "fakexdev.go") + if err := os.WriteFile(src, []byte(fakeXdevSource), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, "go.mod"), []byte("module fakexdev\n\ngo 1.25\n"), 0o644); err != nil { + t.Fatal(err) + } + if out, err := exec.Command("go", "build", "-o", bin, src).CombinedOutput(); err != nil { + t.Fatalf("build fake xdev: %v\n%s", err, out) + } + return bin +} + +const fakeXdevSource = `package main + +import ( + "bufio" + "encoding/json" + "os" +) + +func main() { + out := bufio.NewWriter(os.Stdout) + defer out.Flush() + w := func(typ, id string, payload any) { + b, _ := json.Marshal(payload) + env := map[string]any{"type": typ, "id": id, "frame": json.RawMessage(b)} + line, _ := json.Marshal(env) + out.Write(append(line, '\n')) + out.Flush() + } + w("ready", "", map[string]any{"protocol": 1, "frameLimit": 1048576}) + sc := bufio.NewScanner(os.Stdin) + sc.Buffer(make([]byte, 1<<20), 1<<20) + for sc.Scan() { + var cmd struct { + Type string ` + "`json:\"type\"`" + ` + ID string ` + "`json:\"id\"`" + ` + Frame struct { + Text string ` + "`json:\"text\"`" + ` + } ` + "`json:\"frame\"`" + ` + } + if json.Unmarshal(sc.Bytes(), &cmd) != nil { + continue + } + if cmd.Type != "prompt" { + continue + } + if os.Getenv("FAKE_MODE") == "fail-turn" { + w("response", cmd.ID, map[string]any{"ok": false, "error": "model exploded"}) + continue + } + w("event", "", map[string]any{"type": "text_delta", "delta": "ok"}) + w("response", cmd.ID, map[string]any{ + "ok": true, "text": "did: " + cmd.Frame.Text, + "stopReason": "stop", "sessionId": "s1", "model": "fake", + }) + } +} +` diff --git a/internal/xdev/xdev.go b/internal/xdev/xdev.go new file mode 100644 index 0000000..c91404b --- /dev/null +++ b/internal/xdev/xdev.go @@ -0,0 +1,373 @@ +// Package xdev is agentloop's execution client: the JSONL-over-stdio +// protocol xdev exposes for embedders (`xdev rpc`). +// +// The separation of powers is the whole point of this package existing at +// all (PRD §4.3): agentloop owns the loop, the bounds and the approval +// gate; **xdev owns the turn** — the sandbox, the file mutation, the test +// run, the model call inside a step. agentloop drives xdev as a tool +// executor for ONE already-planned step and never hands it an unbounded +// goal. Loop count stays 1. +// +// Wire shape (xdev internal/protocol, v1): +// +// {"type":"ready","frame":{"protocol":1,"frameLimit":1048576}} +// -> {"type":"prompt","id":"","text":""} +// <- {"type":"event","frame":{...}} (zero or more, streamed) +// <- {"type":"response","id":"","frame":{"ok":true,"text":"...",...}} +// +// Reads are sequential and writes are serialized, which matters because a +// long turn streams events before its response: a client that only reads +// the response deadlocks behind a full pipe (xdev's own Serve comment +// names this). +package xdev + +import ( + "bufio" + "context" + "encoding/json" + "fmt" + "io" + "os/exec" + "strings" + "sync" + "sync/atomic" + "time" +) + +// Protocol constants, mirrored from xdev's internal/protocol. Duplicated +// deliberately: agentloop must not import xdev's module, and a wire +// version is a contract, not a shared type. The ready frame is validated +// against ProtocolVersion so a future v2 fails loudly instead of +// mis-parsing. +const ( + ProtocolVersion = 1 + // FrameLimit caps one JSONL frame. xdev advertises its own in ready; + // this is the client-side ceiling used when deciding to truncate a + // prompt before sending it. + FrameLimit = 1 << 20 +) + +// Frame types. +const ( + typeReady = "ready" + typePrompt = "prompt" + typeSteer = "steer" + typeAbort = "abort" + typeResponse = "response" + typeEvent = "event" +) + +// frame is the envelope every line carries, in both directions. +type frame struct { + Type string `json:"type"` + ID string `json:"id,omitempty"` + Frame json.RawMessage `json:"frame,omitempty"` +} + +// readyFrame is xdev's first frame on the stream. +type readyFrame struct { + Protocol int `json:"protocol"` + FrameLimit int `json:"frameLimit"` +} + +// responseFrame answers one prompt by id. +type responseFrame struct { + OK bool `json:"ok"` + Err string `json:"error,omitempty"` + SessionID string `json:"sessionId,omitempty"` + Model string `json:"model,omitempty"` + Text string `json:"text,omitempty"` + StopReason string `json:"stopReason,omitempty"` +} + +// eventFrame is one streamed event. Only the fields a loop needs to +// observe are decoded; the rest are ignored on purpose — agentloop is not +// a UI. +type eventFrame struct { + Type string `json:"type"` + ToolCallID string `json:"toolCallId,omitempty"` + ToolName string `json:"toolName,omitempty"` + Delta string `json:"delta,omitempty"` + StopReason string `json:"stopReason,omitempty"` +} + +// Turn is the outcome of one prompt. +type Turn struct { + Text string // the assistant's final text + StopReason string // xdev's stop reason, passed through + SessionID string // set once a session exists + Model string // the model xdev used, as it reports it + Events []eventFrame // the stream, for the trace + Duration time.Duration // wall-clock for the turn + EventLimit bool // true if Events was truncated at MaxEvents +} + +// Client drives one xdev process over stdio. +// +// It is deliberately single-turn-at-a-time: xdev's own server refuses a +// second prompt while one is running ("send steer/follow_up instead"), so +// serialising here turns a server-side error into a client-side wait +// rather than a failed step. +type Client struct { + cmd *exec.Cmd + stdin io.WriteCloser + out *bufio.Reader + + mu sync.Mutex // one turn at a time, and one writer at a time + seq atomic.Uint64 + closed atomic.Bool + + // MaxEvents caps how much of a turn's stream is retained on the Turn. + // The trace wants the shape, not every token; a run that streams a + // megabyte of deltas must not become a megabyte of run record. + MaxEvents int +} + +// Config describes how to start the child. +type Config struct { + // Command is the xdev binary. Default "xdev". + Command string + // Args are extra arguments after the subcommand. `rpc` is always added. + Args []string + // Dir is the working directory a turn runs in. This is the sandbox + // boundary agentloop actually controls: xdev's own --add-dir + // restriction is applied inside it. + Dir string + // Env is the child environment (nil = inherit). + Env []string +} + +// Start launches `xdev rpc` and waits for its ready frame. +// +// The ready frame is a hard gate: if the child does not speak the protocol +// version this client implements, Start fails and the loop is told no +// executor is available — better a step that reports "no sandbox" than one +// that mis-parses a future wire format. +func Start(ctx context.Context, cfg Config) (*Client, error) { + if cfg.Command == "" { + cfg.Command = "xdev" + } + args := append([]string{"rpc"}, cfg.Args...) + cmd := exec.CommandContext(ctx, cfg.Command, args...) + cmd.Dir = cfg.Dir + if cfg.Env != nil { + cmd.Env = cfg.Env + } + stdin, err := cmd.StdinPipe() + if err != nil { + return nil, fmt.Errorf("xdev: stdin pipe: %w", err) + } + stdout, err := cmd.StdoutPipe() + if err != nil { + return nil, fmt.Errorf("xdev: stdout pipe: %w", err) + } + // The child's stderr is kept off the protocol stream but not thrown + // away: a client that swallows it cannot explain a failed start. + var errBuf strings.Builder + cmd.Stderr = &errBuf + if err := cmd.Start(); err != nil { + return nil, fmt.Errorf("xdev: start %s: %w", cfg.Command, err) + } + + c := &Client{cmd: cmd, stdin: stdin, out: bufio.NewReaderSize(stdout, FrameLimit), MaxEvents: 200} + + // The ready frame: the child's first line, always (xdev Serve writes it + // before dispatching anything). + line, err := c.readLine() + if err != nil { + _ = cmd.Process.Kill() + return nil, fmt.Errorf("xdev: no ready frame: %w (stderr: %s)", err, strings.TrimSpace(errBuf.String())) + } + var f frame + if err := json.Unmarshal(line, &f); err != nil { + _ = cmd.Process.Kill() + return nil, fmt.Errorf("xdev: ready frame is not JSONL: %w", err) + } + if f.Type != typeReady { + _ = cmd.Process.Kill() + return nil, fmt.Errorf("xdev: first frame was %q, want %q", f.Type, typeReady) + } + var rf readyFrame + if err := json.Unmarshal(f.Frame, &rf); err != nil { + _ = cmd.Process.Kill() + return nil, fmt.Errorf("xdev: decode ready: %w", err) + } + if rf.Protocol != ProtocolVersion { + _ = cmd.Process.Kill() + return nil, fmt.Errorf("xdev: protocol %d, this client speaks %d", rf.Protocol, ProtocolVersion) + } + return c, nil +} + +// Close terminates the child and waits for it. +func (c *Client) Close() error { + if c.closed.Swap(true) { + return nil + } + _ = c.stdin.Close() + if c.cmd.Process != nil { + _ = c.cmd.Process.Kill() + } + _ = c.cmd.Wait() + return nil +} + +// Prompt runs one turn and returns its outcome. +// +// ctx bounds the turn: if it expires the client kills the child rather +// than leaving a half-read stream behind, because a turn that outlived its +// step budget has no business continuing to mutate a workspace. +func (c *Client) Prompt(ctx context.Context, text string) (Turn, error) { + c.mu.Lock() + defer c.mu.Unlock() + + if c.closed.Load() { + return Turn{}, fmt.Errorf("xdev: client closed") + } + if text == "" { + return Turn{}, fmt.Errorf("xdev: prompt text required") + } + id := fmt.Sprintf("al-%d", c.seq.Add(1)) + + start := time.Now() + if err := c.writeFrame(typePrompt, id, map[string]any{"text": text}); err != nil { + return Turn{}, err + } + + turn := Turn{} + // Read until this prompt's response. Events interleave and are kept + // (bounded) because the trace wants the shape of the turn. + for { + select { + case <-ctx.Done(): + // Kill rather than wait: the step's budget is gone. + _ = c.cmd.Process.Kill() + return turn, fmt.Errorf("xdev: turn %s: %w", id, ctx.Err()) + default: + } + + line, err := c.readLine() + if err != nil { + return turn, fmt.Errorf("xdev: read after %s: %w", id, err) + } + var f frame + if err := json.Unmarshal(line, &f); err != nil { + continue // a malformed line is not worth failing a turn over + } + switch f.Type { + case typeEvent: + var ev eventFrame + if json.Unmarshal(f.Frame, &ev) == nil && len(turn.Events) < c.MaxEvents { + turn.Events = append(turn.Events, ev) + } else if len(turn.Events) >= c.MaxEvents { + turn.EventLimit = true + } + case typeResponse: + // A response for a different id is a protocol violation this + // client cannot act on; ignore it and keep reading for ours. + if f.ID != id { + continue + } + var resp responseFrame + if err := json.Unmarshal(f.Frame, &resp); err != nil { + return turn, fmt.Errorf("xdev: decode response %s: %w", id, err) + } + turn.Duration = time.Since(start) + turn.SessionID, turn.Model, turn.StopReason = resp.SessionID, resp.Model, resp.StopReason + if !resp.OK { + return turn, fmt.Errorf("xdev: turn failed: %s", resp.Err) + } + turn.Text = resp.Text + return turn, nil + } + } +} + +// Steer injects a message into the running turn. Not used by the loop +// yet; present because the protocol has it and a future step that watches +// a long turn will want it. +func (c *Client) Steer(ctx context.Context, text string) error { + c.mu.Lock() + defer c.mu.Unlock() + id := fmt.Sprintf("al-steer-%d", c.seq.Add(1)) + if err := c.writeFrame(typeSteer, id, map[string]any{"text": text}); err != nil { + return err + } + return c.awaitOK(id) +} + +// Abort cancels the running turn. +func (c *Client) Abort(ctx context.Context) error { + c.mu.Lock() + defer c.mu.Unlock() + id := fmt.Sprintf("al-abort-%d", c.seq.Add(1)) + if err := c.writeFrame(typeAbort, id, nil); err != nil { + return err + } + return c.awaitOK(id) +} + +// awaitOK reads until the response with this id (skipping events). +func (c *Client) awaitOK(id string) error { + for { + line, err := c.readLine() + if err != nil { + return fmt.Errorf("xdev: read after %s: %w", id, err) + } + var f frame + if json.Unmarshal(line, &f) != nil || f.Type != typeResponse || f.ID != id { + continue + } + var resp responseFrame + if err := json.Unmarshal(f.Frame, &resp); err != nil { + return fmt.Errorf("xdev: decode response %s: %w", id, err) + } + if !resp.OK { + return fmt.Errorf("xdev: %s: %s", id, resp.Err) + } + return nil + } +} + +func (c *Client) writeFrame(typ, id string, payload any) error { + var raw json.RawMessage + if payload != nil { + b, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("xdev: encode %s: %w", typ, err) + } + raw = b + } + body, err := json.Marshal(frame{Type: typ, ID: id, Frame: raw}) + if err != nil { + return fmt.Errorf("xdev: encode frame: %w", err) + } + body = append(body, '\n') + if len(body) > FrameLimit { + return fmt.Errorf("xdev: frame is %d bytes, limit is %d", len(body), FrameLimit) + } + if _, err := c.stdin.Write(body); err != nil { + return fmt.Errorf("xdev: write %s: %w", typ, err) + } + return nil +} + +// readLine reads one frame line, bounding it so a hostile or broken child +// cannot exhaust memory. +func (c *Client) readLine() ([]byte, error) { + var buf []byte + for { + chunk, err := c.out.ReadSlice('\n') + buf = append(buf, chunk...) + if err == nil { + return buf, nil + } + if err == bufio.ErrBufferFull { + if len(buf) > FrameLimit { + return nil, fmt.Errorf("xdev: frame exceeds %d bytes", FrameLimit) + } + continue + } + return nil, err + } +} diff --git a/internal/xdev/xdev_test.go b/internal/xdev/xdev_test.go new file mode 100644 index 0000000..ca5b815 --- /dev/null +++ b/internal/xdev/xdev_test.go @@ -0,0 +1,239 @@ +package xdev + +import ( + "context" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" +) + +// The tests drive a real child process speaking the real wire protocol. The +// fixture is a small Go program built per test, because the protocol's hard +// parts (ready-frame gating, event-before-response interleaving, a turn that +// never ends) are exactly the parts a mocked io.Reader would paper over. + +const fakeSource = `package main + +import ( + "bufio" + "encoding/json" + "fmt" + "os" + "time" +) + +func main() { + out := bufio.NewWriter(os.Stdout) + defer out.Flush() + w := func(typ, id string, payload any) { + b, _ := json.Marshal(payload) + env := map[string]any{"type": typ, "id": id, "frame": json.RawMessage(b)} + line, _ := json.Marshal(env) + out.Write(append(line, '\n')) + out.Flush() + } + switch os.Getenv("FAKE_MODE") { + case "no-ready": + return + case "bad-proto": + w("ready", "", map[string]any{"protocol": 99, "frameLimit": 1048576}) + return + } + w("ready", "", map[string]any{"protocol": 1, "frameLimit": 1048576}) + sc := bufio.NewScanner(os.Stdin) + sc.Buffer(make([]byte, 1<<20), 1<<20) + for sc.Scan() { + var cmd struct { + Type string ` + "`json:\"type\"`" + ` + ID string ` + "`json:\"id\"`" + ` + Frame struct { + Text string ` + "`json:\"text\"`" + ` + } ` + "`json:\"frame\"`" + ` + } + if json.Unmarshal(sc.Bytes(), &cmd) != nil { + continue + } + switch cmd.Type { + case "prompt": + switch os.Getenv("FAKE_MODE") { + case "fail-turn": + w("response", cmd.ID, map[string]any{"ok": false, "error": "model exploded"}) + continue + case "stall": + for i := 0; ; i++ { + w("event", "", map[string]any{"type": "text_delta", "delta": fmt.Sprint(i)}) + time.Sleep(20 * time.Millisecond) + } + } + w("event", "", map[string]any{"type": "toolcall_start", "toolCallId": "c1", "toolName": "write"}) + w("event", "", map[string]any{"type": "text_delta", "delta": "ok"}) + w("response", cmd.ID, map[string]any{ + "ok": true, "text": "did: " + cmd.Frame.Text, + "stopReason": "stop", "sessionId": "sess-1", "model": "fake-model", + }) + case "abort", "steer": + w("response", cmd.ID, map[string]any{"ok": true}) + } + } +} +` + +func fakeBinary(t *testing.T) string { + t.Helper() + dir := t.TempDir() + bin := filepath.Join(dir, "fake-xdev") + src := filepath.Join(dir, "fakexdev.go") + if err := os.WriteFile(src, []byte(fakeSource), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, "go.mod"), []byte("module fakexdev\n\ngo 1.25\n"), 0o644); err != nil { + t.Fatal(err) + } + out, err := exec.Command("go", "build", "-o", bin, src).CombinedOutput() + if err != nil { + t.Fatalf("build fake xdev: %v\n%s", err, out) + } + return bin +} + +func start(t *testing.T, bin string) *Client { + t.Helper() + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + t.Cleanup(cancel) + c, err := Start(ctx, Config{Command: bin}) + if err != nil { + t.Fatalf("Start() error: %v", err) + } + t.Cleanup(func() { _ = c.Close() }) + return c +} + +// A child that dies without a ready frame must fail Start — not block a run +// forever waiting for a second line that will never come. +func TestStartRejectsChildThatNeverSaysReady(t *testing.T) { + bin := fakeBinary(t) + t.Setenv("FAKE_MODE", "no-ready") + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + c, err := Start(ctx, Config{Command: bin}) + if err == nil { + _ = c.Close() + t.Fatal("Start() succeeded against a child that never sent ready") + } + if !strings.Contains(err.Error(), "ready") { + t.Errorf("error = %q, want it to name the missing ready frame", err) + } +} + +// A wire version this client does not implement must fail loudly rather than +// mis-parse a future format. +func TestStartRejectsUnknownProtocolVersion(t *testing.T) { + bin := fakeBinary(t) + t.Setenv("FAKE_MODE", "bad-proto") + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + c, err := Start(ctx, Config{Command: bin}) + if err == nil { + _ = c.Close() + t.Fatal("Start() accepted protocol 99") + } + if !strings.Contains(err.Error(), "protocol 99") { + t.Errorf("error = %q, want it to name the version mismatch", err) + } +} + +// The core contract: one prompt, its response, and the events that streamed +// before it. A client that ignores events until the response would deadlock +// behind a full pipe on a long turn. +func TestPromptCorrelatesResponseAndKeepsEvents(t *testing.T) { + c := start(t, fakeBinary(t)) + + turn, err := c.Prompt(context.Background(), "fix the parser") + if err != nil { + t.Fatalf("Prompt() error: %v", err) + } + if turn.Text != "did: fix the parser" { + t.Errorf("Text = %q, want the response text", turn.Text) + } + if turn.SessionID != "sess-1" || turn.Model != "fake-model" { + t.Errorf("session/model = %q/%q, want sess-1/fake-model", turn.SessionID, turn.Model) + } + if turn.StopReason != "stop" { + t.Errorf("StopReason = %q, want stop", turn.StopReason) + } + if len(turn.Events) < 2 { + t.Errorf("events = %d, want the streamed events retained", len(turn.Events)) + } + if turn.Duration <= 0 { + t.Error("Duration not recorded") + } +} + +// An ok:false response is an error carrying the child's reason, not an empty +// success — a failed step must be visible as a failed step. +func TestPromptRejectsFailedTurn(t *testing.T) { + bin := fakeBinary(t) + t.Setenv("FAKE_MODE", "fail-turn") + c := start(t, bin) + + _, err := c.Prompt(context.Background(), "anything") + if err == nil { + t.Fatal("Prompt() returned nil error on an ok:false response") + } + if !strings.Contains(err.Error(), "model exploded") { + t.Errorf("error = %q, want the child's reason", err) + } +} + +// A turn that outlives its budget must stop, and the child must be killed: a +// half-read stream plus a sandbox still mutating a workspace is the failure +// the step budget exists to prevent. +func TestPromptHonoursContextOnAStalledTurn(t *testing.T) { + bin := fakeBinary(t) + t.Setenv("FAKE_MODE", "stall") + c := start(t, bin) + + turnCtx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) + defer cancel() + _, err := c.Prompt(turnCtx, "stall forever") + if err == nil { + t.Fatal("Prompt() returned nil error after its context expired") + } + if c.cmd.Process != nil { + // The process must be gone, not merely abandoned. + done := make(chan struct{}) + go func() { _, _ = c.cmd.Process.Wait(); close(done) }() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("the child survived the turn timeout") + } + } +} + +// A deploy with no xdev on PATH must report "no executor" so the loop can +// carry on without one, rather than panic or hang. +func TestCommandMissingFailsStart(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + c, err := Start(ctx, Config{Command: "/nonexistent/xdev-binary"}) + if err == nil { + _ = c.Close() + t.Fatal("Start() succeeded with a missing binary") + } +} + +// An empty prompt is rejected before the round trip: xdev's own server +// answers "prompt: text required", and a client that sends it is asking for +// a wasted turn. +func TestPromptRejectsEmptyText(t *testing.T) { + c := start(t, fakeBinary(t)) + if _, err := c.Prompt(context.Background(), ""); err == nil { + t.Fatal("Prompt(\"\") returned nil error") + } +}