diff --git a/README.md b/README.md index 8d1df76..88eaa89 100644 --- a/README.md +++ b/README.md @@ -259,6 +259,23 @@ Rules worth knowing before you write one: code is the risk to design for: `factory.yaml` scores edits to CI or lint configuration, and test files that lose more lines than they gain, as not low, so such a "fix" is held for a person. +- **A read-only review panel can run in parallel.** `parallel: + [review-correctness, review-security, review-maintainability]` runs those + phases at once, as one step of the chain. It is opt-in and narrow: one + group per definition, two or more consecutive agent phases, each with its + own role, every role `writes: []`, and no member's `if:` guard reading a + field a sibling reports. Each member runs in its own ephemeral HOME and + session and is handed the envelope from BEFORE the group, never a + sibling's; results merge in declared order, and the next phase gets the + last member's envelope. Any worktree change while the group runs — + including one a crashed member left — rolls back and aborts the attempt. + Rejections resolve after every member finishes: the first member in + declared order that rejected with budget left dispatches its repair target + and charges only its own budget, and then the whole group runs again, so + earlier approvals are re-judged. The group runs at most 1 + the sum of its + members' budgets times. `examples/definitions/factory-parallel.yaml` is the + stock factory with its panel grouped; `factory.yaml` itself stays + sequential until the parallel panel has been watched on real work. - **Validation happens before anything runs.** `jig def validate ` is the same check the store applies at save time, offline. diff --git a/cmd/jig/factory_test.go b/cmd/jig/factory_test.go index be0ddc6..12b61e2 100644 --- a/cmd/jig/factory_test.go +++ b/cmd/jig/factory_test.go @@ -9,6 +9,7 @@ import ( "os" "os/exec" "path/filepath" + "reflect" "strings" "testing" @@ -155,3 +156,35 @@ func TestTheFactoryRepairsRedCIThroughItsWholeReviewPanel(t *testing.T) { } } } + +func parseStock(t *testing.T, name string) *protocol.DefinitionSpec { + t.Helper() + source, err := os.ReadFile(stockDefinition(name)) + if err != nil { + t.Fatal(err) + } + spec, err := protocol.ParseDefinition(source) + if err != nil { + t.Fatal(err) + } + return spec +} + +// The stock factory keeps its sequential panel (plan KTD7); the parallel +// example groups exactly its three reviewers and is otherwise the same +// definition, so the two cannot drift apart. +func TestTheParallelFactoryIsTheStockFactoryWithItsPanelGrouped(t *testing.T) { + stock := parseStock(t, "factory.yaml") + parallel := parseStock(t, "factory-parallel.yaml") + if stock.Parallel != nil { + t.Fatalf("factory.yaml declares parallel: %v; the stock factory stays sequential", stock.Parallel) + } + want := protocol.ParallelGroup{"review-correctness", "review-security", "review-maintainability"} + if !reflect.DeepEqual(parallel.Parallel, want) { + t.Fatalf("factory-parallel.yaml groups %v, want %v", parallel.Parallel, want) + } + parallel.Name, parallel.Parallel = stock.Name, nil + if !reflect.DeepEqual(parallel, stock) { + t.Fatal("factory-parallel.yaml differs from factory.yaml beyond its name and its parallel group") + } +} diff --git a/examples/definitions/factory-parallel.yaml b/examples/definitions/factory-parallel.yaml new file mode 100644 index 0000000..2df74b4 --- /dev/null +++ b/examples/definitions/factory-parallel.yaml @@ -0,0 +1,404 @@ +# factory-parallel.yaml — factory.yaml with its reviewer panel run in +# parallel. Everything but the name and the `parallel:` line after `phases` +# is the stock factory (a test holds the two in step), and the notes after +# THE PARALLEL PANEL apply unchanged. The stock factory stays sequential +# until this panel's semantics have been watched on real work (plan KTD7). +# +# jig run --def examples/definitions/factory-parallel.yaml "add rate limiting to the API" +# +# THE PARALLEL PANEL. `parallel: [review-correctness, review-security, +# review-maintainability]` runs the three reviewers at once, each in its own +# ephemeral HOME and session. A group may hold only consecutive agent phases +# whose roles are read-only (`writes: []`) and distinct; any change to the +# worktree while the panel runs aborts the attempt. What changes against the +# sequential panel: +# +# - Prompts. Every reviewer is handed the envelope from before the panel +# (the tests), never a sibling's review. The phase after the panel gets +# the last reviewer's envelope, as it would sequentially. +# - Rejections. Edges resolve after all three finish, in declared order: +# the FIRST reviewer that rejected with budget left sends its blocking +# list to the builder, and then the WHOLE panel runs again — so an +# earlier approval is re-judged against the revision, and a later +# reviewer's rejection waits for the next panel run (its budget is not +# charged until it is the one dispatching). A reviewer whose budget is +# spent applies its `exhausted` policy: fail-job ends the attempt, +# proceed moves on to the next reviewer. +# - Bounds. The panel runs at most 1 + the sum of its reviewers' budgets +# times (here 1 + 3 + 3 + 3 = 10). Every run spends one send per reviewer +# at least, from the same attempt-wide send budget. +# +# THE RISK GATE. `classify-risk` is a code phase with `reports_fields: true`: +# its last output line is a JSON object, and jig merges those fields into the +# envelope view the guards and the publish hold read. The script is +# deterministic — paths touched and lines changed, no model — so the same +# diff always scores the same. Anything but `risk: low` runs one more review +# and trips `publish.hold_when`, which parks the work: the attempt ends +# accepted_unpublished with the hold on record, the worktree and branch are +# retained, and the operator's publish retry is the human sign-off. Low-risk +# work publishes on its own. Both predicates say `risk != low` rather than +# `risk == high` so they fail closed: work the classifier never scored — a +# missing or unexpected `risk` — is held, not shipped. +# +# THE CI GATE. `publish.ci.wait` makes green CI part of `accepted`: after the +# pull request exists, the worker polls its head's checks for up to 30 +# minutes. Timed-out CI ends the job accepted_unpublished with the pending +# checks named; push a fix to the branch and the publish retry re-judges CI +# on the new head. +# +# THE CI REPAIR. `on_fail` repairs red CI inside the attempt, up to two +# rounds. The failing checks and the tail of each failed Actions job's log +# go to `build`; then EVERY phase after it runs again — commit, test, the +# three reviewers, the risk classifier, the extra review — and only a fix +# that passes all of them and is not held is pushed to the same pull +# request. The classifier scores a change that edits CI or lint +# configuration, or removes more test lines than it adds, as not low, so a +# "fix" that weakens a check instead of the code is held for a person. A +# round that fails, is held, or changes nothing ends the job +# accepted_unpublished, as red CI did before; the publish retry judges CI +# but never repairs. +# +# THE ROSTER. Builders and reviewers are `opus`; the planner and the extra +# high-risk reviewer are `sonnet`. `effort` is medium for building to a plan +# and high for reviewers hunting what the build missed. `budget_usd` caps what +# one send may spend; the cap is enforced by the Claude Code CLI, and the +# Codex runtime refuses a capped role rather than run it uncapped. +# +# EDIT BEFORE USE: the test command (it detects common entry points and +# fails loudly when it finds none), the builder's `writes` allowlist (`**` is +# the weakest boundary jig can enforce — narrow it, R10), and the risk +# classifier's path pattern and line threshold, which should name what is +# dangerous in YOUR repository. +name: factory-parallel +roster: + planner: + model: sonnet + budget_usd: 2 + system_prompt: | + You are a read-only planner working in a git repository. You never + create, modify, or delete any file; you produce a plan the builder + will execute, and you report exclusively through your final JSON + envelope. + user_prompt: | + Task: {{prompt}} + + Read enough of this repository to plan the change: which files to + touch, in what order, how the result will be verified, and what + "done" looks like — the concrete condition the reviewers will check. + + Your final JSON must also include: + - "notes_for_next_agent": the plan itself, concrete enough to + execute without re-deriving it + env: [PATH] + writes: [] + builder: + model: opus + effort: medium + budget_usd: 8 + system_prompt: | + You are a builder working in a git repository. You execute the plan + you are handed, you write only what the task requires, and you report + honestly: claim only files you actually created or modified. + + You may be run several times in one job. The first time you are + handed a plan. Every later time you are handed a REVIEW that rejected + your work — its "blocking" list is what you must fix, and nothing + else. Never weaken or delete a test to satisfy a reviewer. + + Before you finish, write your commit message — one subject line, then + a blank line, then why — to the file $HOME/commit-message. It is your + own ephemeral home, outside the repository; it is never committed as a + file, and jig uses it as the message of the commit that carries your + work. + user_prompt: | + Task: {{prompt}} + + The previous phase's report is above. If it is a review with + "approved": false, fix every blocking item it names. Otherwise, + execute the plan. + + Write $HOME/commit-message, then report. + + Your final JSON must also include: + - "changed_files": the repo-relative paths of every file you + created or modified + - "artifacts": the same paths, for existence verification + - "revised": true ONLY if you were fixing a review's blocking + items on this run; false when you were executing the plan + env: [PATH] + writes: ["**"] + correctness-reviewer: + model: opus + effort: high + budget_usd: 4 + system_prompt: | + You are a read-only correctness reviewer. You judge whether the change + does what the task asked, completely and without regressions, against + the repository's own conventions. You never create, modify, or delete + any file. + + A finding is a requirement you checked, with "met" true or false. Your + verdict must agree with your findings: never approve while a finding + is unmet or a blocking item stands. Never report "revised" — that + field belongs to the builder. List only problems you would block the + merge for; for each, give the file and line, why it is wrong, and how + to show it fails. + + A change that deletes or weakens tests, lint rules, or CI configuration + is blocking unless the task asked for exactly that: a fix must make the + checks pass, not make them check less. This matters most when the + builder is repairing red CI. + + Your "status" is "success" whenever you completed the review. The + verdict travels in "approved", never in "status": a rejection you + reported correctly is a phase that worked. + user_prompt: | + Task: {{prompt}} + + Review the change in this worktree — `git diff HEAD~1` shows the + committed work, and any uncommitted changes are a revision in + progress. + + Your final JSON must also include: + - "approved": true when the change is complete and correct + - "findings": a list of {"requirement": "", + "met": true|false, "evidence": ""} + - "blocking": a list of strings naming exactly what must change + before approval (an empty list when you approve) + env: [PATH] + writes: [] + security-reviewer: + model: opus + effort: high + budget_usd: 4 + system_prompt: | + You are a read-only security reviewer. You look only for security + defects the change introduces or leaves in place: injection, broken + authentication or authorization, secrets in code, unsafe + deserialization, path traversal, unvalidated input reaching a + dangerous sink, and dependencies pulled in without need. You never + create, modify, or delete any file, and you do not review style or + correctness — other reviewers own those. + + A finding is a requirement you checked, with "met" true or false. Your + verdict must agree with your findings: never approve while a finding + is unmet or a blocking item stands. Never report "revised". + + Your "status" is "success" whenever you completed the review. The + verdict travels in "approved", never in "status". + user_prompt: | + Task: {{prompt}} + + Review the change in this worktree for security defects — `git diff + HEAD~1` shows the committed work, and any uncommitted changes are a + revision in progress. Read the code paths the change touches far + enough to know where its inputs come from. + + Your final JSON must also include: + - "approved": true when you found nothing you would block for + - "findings": a list of {"requirement": "", + "met": true|false, "evidence": ""} + - "blocking": a list of strings naming exactly what must change + before approval (an empty list when you approve) + env: [PATH] + writes: [] + maintainability-reviewer: + model: opus + effort: high + budget_usd: 4 + system_prompt: | + You are a read-only maintainability reviewer. You judge whether the + change reads like the surrounding code — naming, structure, comment + density, idiom — and whether its tests would catch the bug it fixes + or the regression it risks. You never create, modify, or delete any + file, and you do not re-review correctness or security. + + Block only for what a maintainer would refuse to merge; a preference + is a note in "findings" with "met": true, not a blocking item. Your + verdict must agree with your findings. Never report "revised". + + Your "status" is "success" whenever you completed the review. The + verdict travels in "approved", never in "status". + user_prompt: | + Task: {{prompt}} + + Review the change in this worktree — `git diff HEAD~1` shows the + committed work, and any uncommitted changes are a revision in + progress. + + Your final JSON must also include: + - "approved": true when a maintainer would merge this as-is + - "findings": a list of {"requirement": "", + "met": true|false, "evidence": ""} + - "blocking": a list of strings naming exactly what must change + before approval (an empty list when you approve) + env: [PATH] + writes: [] + risk-reviewer: + model: sonnet + effort: high + budget_usd: 3 + system_prompt: | + You are a read-only reviewer brought in because a deterministic + classifier scored this change HIGH RISK — it touches sensitive paths + or is large. Your job is to write the briefing the human approver will + read: what the change does to the risky paths, what could go wrong in + production, and what to check before releasing it. You never create, + modify, or delete any file, and you do not block: the person does. + + Your "status" is "success" whenever you completed the briefing. + user_prompt: | + Task: {{prompt}} + + The classifier's report is above; its "risky_paths" are where to + look first. `git diff HEAD~1` shows the committed work. + + Your final JSON must also include: + - "approved": true (you inform; the release decision is human) + - "briefing": a short plain-language summary for the approver + - "release_checks": a list of strings — what to verify before + releasing this change + env: [PATH] + writes: [] +phases: + - name: plan + kind: agent + owner: planner + - name: build + kind: agent + owner: builder + gates: + - {name: artifacts_exist} + - {name: diff_matches_claims} + # The commit phases stage with `git add -A` deliberately: by the time a + # code phase runs, the write boundary has already verified every path the + # agent changed against its allowlist and rolled back anything else (R10), + # so what is left in the worktree is exactly the authorized work. + - name: commit-build + kind: code + owner: builder + command: | + set -e + message="$HOME/commit-message" + git add -A + if git diff --cached --quiet; then + echo "commit-build: nothing to commit" + exit 0 + fi + if [ -s "$message" ]; then + git -c user.name=jig -c user.email=jig@localhost \ + -c core.hooksPath=/dev/null commit -q -F "$message" + rm -f "$message" + else + echo 'commit-build: the build phase wrote no $HOME/commit-message; using the fallback' + git -c user.name=jig -c user.email=jig@localhost \ + -c core.hooksPath=/dev/null commit -q -m "jig: build phase" + fi + - name: test + kind: code + owner: builder + command: > + if [ -f Justfile ] || [ -f justfile ]; then just test; + elif [ -f Makefile ]; then make test; + elif [ -f go.mod ]; then go test ./...; + elif [ -f package.json ]; then npm test --silent; + else echo "factory: no test command detected — edit the test phase"; exit 1; fi + on_fail: {run: build, then: rerun-self, budget: 3, exhausted: fail-job} + - name: review-correctness + kind: agent + owner: correctness-reviewer + gates: + - {name: verdict_consistent} + on_fail: {when: "approved == false", run: build, then: rerun-self, budget: 3, exhausted: fail-job} + - name: review-security + kind: agent + owner: security-reviewer + gates: + - {name: verdict_consistent} + on_fail: {when: "approved == false", run: build, then: rerun-self, budget: 3, exhausted: fail-job} + - name: review-maintainability + kind: agent + owner: maintainability-reviewer + gates: + - {name: verdict_consistent} + on_fail: {when: "approved == false", run: build, then: rerun-self, budget: 3, exhausted: fail-job} + - name: commit-revision + kind: code + owner: builder + if: revised + command: | + set -e + message="$HOME/commit-message" + git add -A + if git diff --cached --quiet; then + echo "commit-revision: nothing to commit" + exit 0 + fi + if [ -s "$message" ]; then + git -c user.name=jig -c user.email=jig@localhost \ + -c core.hooksPath=/dev/null commit -q -F "$message" + rm -f "$message" + else + echo 'commit-revision: the builder wrote no $HOME/commit-message; using the fallback' + git -c user.name=jig -c user.email=jig@localhost \ + -c core.hooksPath=/dev/null commit -q -m "jig: revision after review" + fi + - name: retest + kind: code + owner: builder + if: revised + command: > + if [ -f Justfile ] || [ -f justfile ]; then just test; + elif [ -f Makefile ]; then make test; + elif [ -f go.mod ]; then go test ./...; + elif [ -f package.json ]; then npm test --silent; + else echo "factory: no test command detected — edit the retest phase"; exit 1; fi + # Deterministic risk classification. JIG_BASE_SHA is the commit the run was + # pinned to, so the diff from there is the whole change, however many + # commits it took. The last line printed is the report; everything before + # it is ordinary output. Fields a reporting phase sets are protected: a + # later agent envelope cannot overwrite `risk`. + - name: classify-risk + kind: code + reports_fields: true + command: | + set -eu + base="${JIG_BASE_SHA:?JIG_BASE_SHA is required}" + pattern='(^|/)(auth|authn|authz|login|session|secret|token|password|crypto|migrations?|schema|deploy|infra|terraform|helm|k8s)([./_-]|$)|(^|/)\.github/|(^|/)Dockerfile|(^|/)go\.mod$|(^|/)package(-lock)?\.json$' + # Checks a change could weaken to get CI green: lint and test runner + # configuration is always sensitive, and a test file that loses more + # lines than it gains is too. Adding tests is ordinary work. + checks='(^|/)\.(golangci|eslintrc|prettierrc|stylelintrc|flake8|pylintrc|rubocop)[^/]*$|(^|/)(jest|vitest|karma|playwright)\.config\.[a-z]+$|(^|/)(pytest\.ini|setup\.cfg|tox\.ini|\.pre-commit-config\.yaml)$' + tests='(_test\.go|_test\.py|\.(test|spec)\.[cm]?[jt]sx?|_spec\.rb)$|(^|/)(tests?|__tests__|spec)/' + threshold=400 + files=$(git diff --name-only "$base" HEAD) + lines=$(git diff --numstat "$base" HEAD | awk '{ total += $1 + $2 } END { print total + 0 }') + # The pattern reaches awk through the environment: `awk -v` would + # process its backslashes and turn `\.` into "any character". + shrunk=$(git diff --numstat "$base" HEAD | tests="$tests" awk '$3 ~ ENVIRON["tests"] && $2 > $1 { print $3 }') + risky=$( { printf '%s\n' "$files" | grep -E "$pattern|$checks" || true; printf '%s\n' "$shrunk"; } | sed '/^$/d' | sort -u) + risky_json=$(printf '%s\n' "$risky" | sed '/^$/d' | sed 's/\\/\\\\/g; s/"/\\"/g; s/.*/"&"/' | paste -sd, -) + risk=low + reason="no sensitive paths, $lines lines changed (threshold $threshold)" + if [ -n "$shrunk" ]; then + risk=high; reason="removes more test lines than it adds" + elif [ -n "$risky" ]; then + risk=high; reason="touches sensitive paths" + elif [ "$lines" -gt "$threshold" ]; then + risk=high; reason="$lines lines changed exceeds threshold $threshold" + fi + echo "classify-risk: $risk — $reason" + printf '{"risk":"%s","risk_reason":"%s","changed_lines":%s,"risky_paths":[%s]}\n' \ + "$risk" "$reason" "$lines" "$risky_json" + - name: review-risk + kind: agent + owner: risk-reviewer + if: "risk != low" +parallel: [review-correctness, review-security, review-maintainability] +acceptance: [all_phases_passed, verdict_consistent, diff_matches_claims] +publish: + hold_when: "risk != low" + ci: + wait: true + on_fail: {run: build, budget: 2} + timeout: 30m diff --git a/internal/engine/enginetest/runtime.go b/internal/engine/enginetest/runtime.go index 0c1c45b..7bb5355 100644 --- a/internal/engine/enginetest/runtime.go +++ b/internal/engine/enginetest/runtime.go @@ -6,6 +6,11 @@ // session identity, options) so tests can assert that corrections re-enter // the SAME live session, that fresh sessions appear after deaths, and that // a can-resume=false engine replays transcript digests into new sessions. +// +// The fake is safe for concurrent StartOrContinue calls, which a parallel +// group makes. Concurrent calls arrive in no fixed order, so steps can be +// routed: Route gives calls matching a predicate (a role's system prompt, +// say) their own queue, and Started/WaitFor let two calls rendezvous. package enginetest import ( @@ -44,6 +49,21 @@ type Step struct { ExitCode int // Usage is the send's reported accounting. Usage runtime.Usage + // Started, when set, is closed as the call starts, so another step can + // wait for this one to be running. + Started chan struct{} + // WaitFor, when set, holds the send (no events, no result) until it is + // closed or the send is killed. + WaitFor <-chan struct{} + // Delay holds the result this long after WaitFor released, or until the + // send is killed. + Delay time.Duration + // ProcessGroup is the process group the handle reports. + ProcessGroup int64 + // Do, when set, runs first as the call starts, with the call as + // recorded — what an agent does outside the worktree (its handoff notes) + // or a runtime that panics. It may block. + Do func(Call) } // Call records one StartOrContinue invocation for assertions. @@ -61,10 +81,25 @@ type Runtime struct { // the transcript-digest degraded path. CanResume bool - mutex sync.Mutex + mutex sync.Mutex + steps []Step + next int + routes []*route + calls []Call + kills int +} + +// route is a queue of steps reserved for the calls its predicate matches. +type route struct { + match func(Call) bool steps []Step next int - calls []Call +} + +// ForSystemPrompt matches the calls made under one system prompt — in +// practice, one role's calls. +func ForSystemPrompt(prompt string) func(Call) bool { + return func(call Call) bool { return call.Options.SystemPrompt == prompt } } // New builds a resume-capable scripted runtime. @@ -79,6 +114,22 @@ func (r *Runtime) Append(steps ...Step) { r.steps = append(r.steps, steps...) } +// Route reserves steps for the calls match accepts. A matching call takes +// the next unconsumed step of the first matching route that has one, and +// never falls through to the unrouted queue. +func (r *Runtime) Route(match func(Call) bool, steps ...Step) { + r.mutex.Lock() + defer r.mutex.Unlock() + r.routes = append(r.routes, &route{match: match, steps: steps}) +} + +// Kills reports how many sends were killed. +func (r *Runtime) Kills() int { + r.mutex.Lock() + defer r.mutex.Unlock() + return r.kills +} + // Calls returns a copy of every recorded call. func (r *Runtime) Calls() []Call { r.mutex.Lock() @@ -90,7 +141,11 @@ func (r *Runtime) Calls() []Call { func (r *Runtime) Remaining() int { r.mutex.Lock() defer r.mutex.Unlock() - return len(r.steps) - r.next + remaining := len(r.steps) - r.next + for _, route := range r.routes { + remaining += len(route.steps) - route.next + } + return remaining } // Probe reports the scripted capability record. @@ -106,12 +161,11 @@ func (r *Runtime) StartOrContinue( _ context.Context, session *runtime.Session, prompt string, opts runtime.Options, ) (runtime.Handle, error) { r.mutex.Lock() - if r.next >= len(r.steps) { + step, found := r.nextStep(Call{Prompt: prompt, SessionKey: session.Key, Options: opts}) + if !found { r.mutex.Unlock() - return nil, fmt.Errorf("scripted runtime: no step for call %d (prompt %.80q)", r.next+1, prompt) + return nil, fmt.Errorf("scripted runtime: no step for call %d (prompt %.80q)", len(r.calls)+1, prompt) } - step := r.steps[r.next] - r.next++ if session.NativeID == "" { session.NativeID = "fake-" + session.Key } @@ -123,8 +177,12 @@ func (r *Runtime) StartOrContinue( Send: session.Sends, Options: opts, }) + call := r.calls[len(r.calls)-1] r.mutex.Unlock() + if step.Do != nil { + step.Do(call) + } for path, content := range step.Files { target, err := resolvePath(path, opts) if err != nil { @@ -138,15 +196,40 @@ func (r *Runtime) StartOrContinue( } } + if step.Started != nil { + close(step.Started) + } handle := &handle{ - step: step, - events: make(chan runtime.Event, len(step.Events)+1), - killed: make(chan struct{}), + runtime: r, + step: step, + events: make(chan runtime.Event, len(step.Events)+1), + killed: make(chan struct{}), } go handle.run() return handle, nil } +// nextStep takes the next step for a call: its route's, when one matches, +// else the unrouted queue's. Callers hold the mutex. +func (r *Runtime) nextStep(call Call) (Step, bool) { + routed := false + for _, route := range r.routes { + if !route.match(call) { + continue + } + routed = true + if route.next < len(route.steps) { + route.next++ + return route.steps[route.next-1], true + } + } + if routed || r.next >= len(r.steps) { + return Step{}, false + } + r.next++ + return r.steps[r.next-1], true +} + func resolvePath(path string, opts runtime.Options) (string, error) { if home, found := strings.CutPrefix(path, "~/"); found { for _, entry := range opts.Env { @@ -163,6 +246,7 @@ func resolvePath(path string, opts runtime.Options) (string, error) { } type handle struct { + runtime *Runtime step Step events chan runtime.Event killed chan struct{} @@ -170,6 +254,22 @@ type handle struct { } func (h *handle) run() { + if h.step.WaitFor != nil { + select { + case <-h.step.WaitFor: + case <-h.killed: + close(h.events) + return + } + } + if h.step.Delay > 0 { + select { + case <-time.After(h.step.Delay): + case <-h.killed: + close(h.events) + return + } + } if h.step.Hang { // Emit nothing; the channel closes only when the group is "killed". <-h.killed @@ -191,10 +291,15 @@ func (h *handle) run() { } func (h *handle) Events() <-chan runtime.Event { return h.events } -func (h *handle) ProcessGroupID() int64 { return 0 } +func (h *handle) ProcessGroupID() int64 { return h.step.ProcessGroup } func (h *handle) Kill() error { - h.killOnce.Do(func() { close(h.killed) }) + h.killOnce.Do(func() { + h.runtime.mutex.Lock() + h.runtime.kills++ + h.runtime.mutex.Unlock() + close(h.killed) + }) return nil } diff --git a/internal/engine/export_test.go b/internal/engine/export_test.go new file mode 100644 index 0000000..583894d --- /dev/null +++ b/internal/engine/export_test.go @@ -0,0 +1,10 @@ +package engine + +import "github.com/StructuPath/jig/internal/worker" + +// FieldViewOf exposes a kept chain's merged field view to the external test +// package: the view guards and the publish hold read is part of what a +// parallel group must merge exactly as sequential execution would. +func FieldViewOf(continuation worker.Continuation) map[string]any { + return continuation.(*ciContinuation).e.fieldView +} diff --git a/internal/engine/homedir.go b/internal/engine/homedir.go index 296570e..57214e0 100644 --- a/internal/engine/homedir.go +++ b/internal/engine/homedir.go @@ -32,15 +32,25 @@ type attemptScratch struct { root string home string handoff string + // members holds one ephemeral HOME per parallel group member, beside + // the chain's HOME rather than inside it, created when a member first + // runs (parallel.go). + members string + // memberHandoff holds one private handoff directory per parallel group + // member, so concurrent members never read or overwrite each other's + // notes; the notes join handoff only when the group's run merges. + memberHandoff string } // createScratch materializes the attempt's scratch family under root, 0700. func createScratch(root, attemptID string) (*attemptScratch, error) { base := filepath.Join(root, attemptID) scratch := &attemptScratch{ - root: base, - home: filepath.Join(base, "home"), - handoff: filepath.Join(base, "handoff"), + root: base, + home: filepath.Join(base, "home"), + handoff: filepath.Join(base, "handoff"), + members: filepath.Join(base, "members"), + memberHandoff: filepath.Join(base, "member-handoff"), } for _, dir := range []string{scratch.home, scratch.handoff} { if err := os.MkdirAll(dir, 0o700); err != nil { diff --git a/internal/engine/parallel.go b/internal/engine/parallel.go new file mode 100644 index 0000000..e183a28 --- /dev/null +++ b/internal/engine/parallel.go @@ -0,0 +1,542 @@ +// parallel.go — the one opt-in parallel construct (plan U5, R8–R11): a +// definition's `parallel:` group of read-only agent phases runs as ONE step +// of the chain, its members concurrently. +// +// Each member runs in a private view of the execution (KTD6): its own +// session, transcript, results, deferred field merges, gate reports, and an +// ephemeral HOME of its own, seeded serially before any member starts, so +// concurrent agent CLIs never share credentials or `~/.claude.json`. What +// the members share with the chain is safe under concurrency by +// construction: the locked emitter, the read-only deadline, and the +// attempt-wide counters (sends, session keys, phase entries). +// +// The worktree is the group's, not any member's. One snapshot is taken +// before the members start and the boundary is enforced once after all of +// them stop, with an empty allowlist: every member is read-only, so ANY +// change — a member's write, a dead member's leftover — is a breach that +// rolls the tree back to the group snapshot and aborts the attempt. +// +// Members see the envelope that preceded the group, never a sibling's. +// Their results merge into the chain in declared order after the join, and +// the phase after the group receives the last member's envelope. Repair +// edges resolve after the join too, walking members in declared order: the +// first member whose edge fired with budget left dispatches its repair +// target with its own envelope, charges only its own budget, and the whole +// group runs again. Total group runs are at most 1 + the sum of member +// budgets. +package engine + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "fmt" + "io/fs" + "os" + "path/filepath" + "strconv" + "strings" + "sync" + + "github.com/StructuPath/jig/internal/protocol" + "github.com/StructuPath/jig/internal/runtime" +) + +// memberRun is one member's part in one group run. +type memberRun struct { + skipped bool + run phaseRun + view *execution +} + +// memberResult is what one member goroutine hands back to the runner. +type memberResult struct { + run phaseRun + died bool + panicked bool + panicValue any +} + +// memberTerminal is a member result that ends the attempt. +type memberTerminal struct { + index int + run phaseRun +} + +// runGroup runs the group as one chain step with R10's edge walk, and +// returns the envelope the next phase receives. +func (e *execution) runGroup( + ctx context.Context, members []protocol.PhaseSpec, previous *parsedEnvelope, edgeUses map[string]int, +) (*chainEnd, *parsedEnvelope) { + // Guards are judged once, on the group's first run. A repair between + // runs can change the field a guard reads, and a re-judged guard would + // let the repair skip the very reviewer that rejected it — recorded as + // skipped, which acceptance counts as passed. Sequential rerun-self never + // re-checks a guard either: a member that ran runs again, one that was + // skipped stays skipped. + var skipped []bool + for groupRun := 1; ; groupRun++ { + runs, end := e.runGroupOnce(ctx, members, previous, groupRun, skipped) + if end != nil { + return end, nil + } + if skipped == nil { + skipped = make([]bool, len(members)) + for i := range runs { + skipped[i] = runs[i].skipped + } + } + last := previous + rerun := false + for i, member := range members { + if runs[i].skipped { + continue + } + run := runs[i].run + if run.hasEnvelope { + last = run.envelopeRef() + } + if !edgeTriggered(member, run) { + // Not triggered: a member that passed is done; anything else + // ends the attempt here, as a failure would have run + // sequentially. Success is earned, never defaulted into. + if run.outcome != phasePassed { + return &chainEnd{endFailed, run.failure}, nil + } + continue + } + end, repaired, dispatched := e.followEdge(ctx, member, run, edgeUses) + if end != nil { + return end, nil + } + if dispatched { + previous = repaired + rerun = true + break + } + // Exhausted under proceed: on to the next member. + } + if !rerun { + return nil, last + } + } +} + +// runGroupOnce runs every member that is not skipped, concurrently, then +// enforces the boundary and merges. On the first run (firstSkipped nil) +// each member's guard decides; on every later run firstSkipped does. A +// non-nil chainEnd means nothing was merged. +func (e *execution) runGroupOnce( + ctx context.Context, members []protocol.PhaseSpec, previous *parsedEnvelope, groupRun int, + firstSkipped []bool, +) ([]memberRun, *chainEnd) { + if e.cancelled() { + return nil, &chainEnd{endCancelled, "cancelled before the parallel group"} + } + if e.ceilingExceeded() { + return nil, &chainEnd{endCeiling, ""} + } + names := make([]string, len(members)) + for i, member := range members { + names[i] = member.Name + } + e.emit.emit(protocol.EventLog, "", "parallel_group_start", + map[string]any{"members": names, "run": groupRun}) + + // Every guard is judged against the view that preceded the group's + // first run — validation guarantees no guard reads a field a sibling + // reports, and runGroup keeps a repair from re-judging one. + runs := make([]memberRun, len(members)) + var active []int + for i, member := range members { + skip := member.If != "" && !guardHolds(member.If, e.fieldView) + if firstSkipped != nil { + skip = firstSkipped[i] + } + if skip { + runs[i].skipped = true + e.emit.emit(protocol.EventLog, member.Name, "phase_skipped", + map[string]string{"guard": member.If}) + continue + } + active = append(active, i) + } + + if len(active) > 0 { + before, err := snapshotTree(ctx, e.attempt.WorktreePath) + if err != nil { + return nil, e.groupInfra(names, err.Error()) + } + // Seed serially, before any member starts. + for _, i := range active { + view, err := e.memberView(i, members[i]) + if err != nil { + return nil, e.groupInfra(names, err.Error()) + } + if err := view.beginGroupRun(e.scratch.handoff); err != nil { + return nil, e.groupInfra(names, err.Error()) + } + // Whatever a member left in its private handoff directory is + // published at the merge below or discarded here, never both. + defer os.RemoveAll(view.scratch.handoff) + owner := members[i].Owner + if err := view.seedRole(owner, e.spec.Roster[owner]); err != nil { + return nil, e.groupInfra(names, fmt.Sprintf( + "seed ephemeral HOME for member %q: %s", members[i].Name, err)) + } + runs[i].view = view + } + if end := e.runMembers(ctx, members, active, runs, previous, names, before); end != nil { + return nil, end + } + } + + for i, member := range members { + if runs[i].skipped { + e.recordResult(protocol.PhaseResult{ + Phase: member.Name, Kind: member.Kind, Status: phaseStatusSkipped, + }) + continue + } + view := runs[i].view + e.results = append(e.results, view.results...) + for _, merge := range view.pendingFields { + e.mergeAgentFields(merge.phase, merge.role, merge.fields) + } + for name, report := range view.gateReports { + e.gateReports[name] = report + } + // A member's handoff notes join the chain's handoff directory only + // now, in declared order, each under its own folder: no sibling + // could read them mid-run, and none can overwrite another's. + if err := publishMemberHandoff( + view.scratch.handoff, e.scratch.handoff, i, member.Name, view.handoffSeed); err != nil { + e.emit.emit(protocol.EventError, member.Name, "handoff_merge_failed", + map[string]string{"error": err.Error()}) + } + } + e.emit.emit(protocol.EventLog, "", "parallel_group_end", + map[string]any{"members": names, "run": groupRun}) + return runs, nil +} + +// runMembers fans the active members out in waves: every member runs +// concurrently; a member that died re-enters alone in the next wave, after +// its siblings finished, until its death budget is spent. The first result +// that ends the attempt stops every sibling — their sends are killed — and +// the runner waits for all of them before it enforces the boundary once. +func (e *execution) runMembers( + ctx context.Context, members []protocol.PhaseSpec, active []int, runs []memberRun, + previous *parsedEnvelope, names []string, before treeSnapshot, +) *chainEnd { + stop := make(chan struct{}) + groupCtx, cancelGroup := context.WithCancel(ctx) + defer cancelGroup() + var stopOnce sync.Once + stopMembers := func() { + stopOnce.Do(func() { + close(stop) + cancelGroup() + }) + } + + deaths := make([]int, len(members)) + var terminals []memberTerminal + var panicked bool + var panicValue any + for pending := active; len(pending) > 0; { + results := make([]memberResult, len(members)) + var wait sync.WaitGroup + for _, i := range pending { + view := runs[i].view + view.stop = stop + wait.Add(1) + go func(i int, view *execution) { + defer wait.Done() + defer func() { + if recovered := recover(); recovered != nil { + results[i] = memberResult{panicked: true, panicValue: recovered} + stopMembers() + } + }() + run, died := view.runAgentPhaseAttempt(groupCtx, members[i], previous) + results[i] = memberResult{run: run, died: died} + if endsAttempt(run.outcome) { + stopMembers() + } + }(i, view) + } + wait.Wait() + + var next []int + for _, i := range pending { + result := results[i] + switch { + case result.panicked: + if !panicked { + panicked, panicValue = true, result.panicValue + } + case endsAttempt(result.run.outcome): + terminals = append(terminals, memberTerminal{index: i, run: result.run}) + case result.died: + deaths[i]++ + runs[i].view.dropSession(members[i].Owner) + if deaths[i] > e.phaseCorrectionBudget(members[i]) { + runs[i].run = phaseRun{outcome: phaseFailed, failure: fmt.Sprintf( + "phase %q: agent died %d time(s): %s", members[i].Name, deaths[i], result.run.failure)} + continue + } + next = append(next, i) + default: + runs[i].run = result.run + } + } + if panicked || len(terminals) > 0 { + break + } + pending = next + } + + if panicked { + // A panic in a member is a panic in the chain, exactly as it would + // be sequentially — re-raised on the chain's own goroutine, where + // Execute's deferred cleanup still runs. Siblings kept running until + // the stop, so the worktree is enforced first: a sibling's write + // must not outlive the attempt just because another member crashed. + _, breaches, err := enforceBoundary( + context.WithoutCancel(ctx), e.attempt.WorktreePath, before, []string{}) + e.emit.emit(protocol.EventError, "", "parallel_group_panic", map[string]any{ + "parallel_group": names, "panic": fmt.Sprint(panicValue), + "breaches": breaches, "enforcement_error": errorText(err), + }) + panic(panicValue) + } + + // The one enforcement, detached from the caller's context: cancellation + // is one of the exits it guards, and on a dead context every git command + // would fail. + _, breaches, err := enforceBoundary(context.WithoutCancel(ctx), e.attempt.WorktreePath, before, []string{}) + if err != nil { + return e.groupInfra(names, "write-boundary enforcement: "+err.Error()) + } + if len(breaches) > 0 { + return e.groupBreach(members, active, runs, names, breaches) + } + if len(terminals) > 0 { + return groupTerminalEnd(terminals) + } + return nil +} + +func errorText(err error) string { + if err == nil { + return "" + } + return err.Error() +} + +// endsAttempt reports whether a member's outcome ends the attempt. +func endsAttempt(outcome phaseOutcome) bool { + switch outcome { + case phaseAborted, phaseCancelled, phaseCeiling, phaseSendBudget: + return true + } + return false +} + +// terminalRank orders the causes a group can end on: cancellation, then +// the ceiling, then the send budget. A breach outranks them all and is +// decided before this is consulted. +var terminalRank = map[phaseOutcome]int{ + phaseAborted: 0, + phaseCancelled: 1, + phaseCeiling: 2, + phaseSendBudget: 3, +} + +// groupTerminalEnd picks the attempt's end among the members' terminal +// results: by cause precedence, then by declared member order. +func groupTerminalEnd(terminals []memberTerminal) *chainEnd { + best := terminals[0] + for _, candidate := range terminals[1:] { + rank, bestRank := terminalRank[candidate.run.outcome], terminalRank[best.run.outcome] + if rank < bestRank || (rank == bestRank && candidate.index < best.index) { + best = candidate + } + } + return best.run.attemptEnd() +} + +// groupBreach records the abort. enforceBoundary has already rolled every +// change back to the group snapshot. No member's partial view is merged: +// each member that ran records one failed result, since the breach cannot +// be pinned on one of them. +func (e *execution) groupBreach( + members []protocol.PhaseSpec, active []int, runs []memberRun, names []string, breaches []breach, +) *chainEnd { + e.emit.emit(protocol.EventError, "", "write_boundary_breach", map[string]any{ + "parallel_group": names, "writes": []string{}, "breaches": breaches, + }) + for _, i := range active { + entry := 0 + if results := runs[i].view.results; len(results) > 0 { + entry = results[len(results)-1].PhaseAttempt + } + e.recordResult(protocol.PhaseResult{ + Phase: members[i].Name, Kind: members[i].Kind, Status: protocol.EnvelopeFail, + PhaseAttempt: entry, Error: "write boundary breach during the parallel group", + }) + } + paths := make([]string, 0, len(breaches)) + for _, item := range breaches { + paths = append(paths, item.Path+" — "+item.Outcome) + } + return &chainEnd{endAborted, fmt.Sprintf( + "parallel group [%s]: members are read-only, but %d path(s) changed during the group: %s", + strings.Join(names, ", "), len(breaches), strings.Join(paths, "; "))} +} + +func (e *execution) groupInfra(names []string, detail string) *chainEnd { + e.emit.emit(protocol.EventError, "", "engine_error", map[string]any{ + "parallel_group": names, "error": detail, + }) + return &chainEnd{endFailed, fmt.Sprintf("parallel group [%s]: %s", strings.Join(names, ", "), detail)} +} + +// memberView returns the member's private view, creating it — and its own +// ephemeral HOME under the attempt's scratch family — on first use. A view +// lives until the HOMEs are wiped, so a member keeps its session across +// group runs just as a sequential phase keeps its role's session across +// rerun-self. +func (e *execution) memberView(index int, member protocol.PhaseSpec) (*execution, error) { + if view := e.members[member.Name]; view != nil { + return view, nil + } + home := filepath.Join(e.scratch.members, strconv.Itoa(index)) + if err := os.MkdirAll(home, 0o700); err != nil { + return nil, fmt.Errorf("create ephemeral HOME for member %q: %w", member.Name, err) + } + scratch := *e.scratch + scratch.home = home + scratch.handoff = filepath.Join(e.scratch.memberHandoff, strconv.Itoa(index)) + view := &execution{ + runner: e.runner, + attempt: e.attempt, + spec: e.spec, + capability: e.capability, + emit: e.emit, + scratch: &scratch, + timeouts: e.timeouts, + deadline: e.deadline, + sessions: make(map[string]*runtime.Session), + transcripts: make(map[string][]exchange), + seededRoles: make(map[string]bool), + sessionKey: e.sessionKey, + counters: e.counters, + grouped: true, + } + if e.members == nil { + e.members = make(map[string]*execution) + } + e.members[member.Name] = view + return view, nil +} + +// beginGroupRun clears what one group run's merge reads, keeping what the +// member carries across runs: its session, transcript, and seeded HOME. +// +// Its private handoff directory is rebuilt every run from the chain's +// handoff directory — the notes every phase before the group left, which a +// sequential reviewer could read too — minus handoff/parallel, where members' +// own merged notes live: no member reads a sibling's notes. What was seeded +// is fingerprinted, so the join publishes only what this member wrote. +func (e *execution) beginGroupRun(shared string) error { + e.results = nil + e.pendingFields = nil + e.gateReports = make(map[string]protocol.GateReport) + e.touchedPaths = make(map[string]bool) + if err := os.RemoveAll(e.scratch.handoff); err != nil { + return fmt.Errorf("reset member handoff directory: %w", err) + } + if err := os.MkdirAll(e.scratch.handoff, 0o700); err != nil { + return fmt.Errorf("create member handoff directory: %w", err) + } + e.handoffSeed = make(map[string]string) + err := copyNotes(shared, e.scratch.handoff, func(relative string, digest string) bool { + if relative == memberNotesRoot || strings.HasPrefix(relative, memberNotesRoot+string(filepath.Separator)) { + return false + } + e.handoffSeed[relative] = digest + return true + }) + if err != nil { + return fmt.Errorf("seed member handoff directory: %w", err) + } + return nil +} + +// memberNotesRoot is where members' merged notes live inside the chain's +// handoff directory. +const memberNotesRoot = "parallel" + +// memberHandoffFolder is where a member's notes land inside the chain's +// handoff directory: handoff/parallel/, or the member's index when +// its phase name is not a plain path element. +func memberHandoffFolder(shared string, index int, phase string) string { + name := phase + if !filepath.IsLocal(name) || strings.ContainsAny(name, `/\`) { + name = "member-" + strconv.Itoa(index) + } + return filepath.Join(shared, memberNotesRoot, name) +} + +// publishMemberHandoff copies the notes a member wrote or changed — not the +// pre-group notes it was seeded with — into the chain's handoff directory, +// replacing what an earlier run of the same member put there. +func publishMemberHandoff(private, shared string, index int, phase string, seed map[string]string) error { + target := memberHandoffFolder(shared, index, phase) + if err := os.RemoveAll(target); err != nil { + return err + } + return copyNotes(private, target, func(relative string, digest string) bool { + seeded, found := seed[relative] + return !found || seeded != digest + }) +} + +// copyNotes copies the regular files under source into destination, +// creating directories as files need them, for every file keep accepts +// (given its path relative to source and a digest of its content). Nothing +// but regular files is copied: a symlink is how an agent would point the +// next phase at something outside the scratch family. A missing source +// copies nothing. +func copyNotes(source, destination string, keep func(relative, digest string) bool) error { + if _, err := os.Lstat(source); os.IsNotExist(err) { + return nil + } + return filepath.WalkDir(source, func(path string, entry fs.DirEntry, walkErr error) error { + if walkErr != nil { + return walkErr + } + if !entry.Type().IsRegular() { + return nil + } + relative, err := filepath.Rel(source, path) + if err != nil { + return err + } + body, err := os.ReadFile(path) + if err != nil { + return err + } + sum := sha256.Sum256(body) + if !keep(relative, hex.EncodeToString(sum[:])) { + return nil + } + target := filepath.Join(destination, relative) + if err := os.MkdirAll(filepath.Dir(target), 0o700); err != nil { + return err + } + return os.WriteFile(target, body, 0o600) + }) +} diff --git a/internal/engine/parallel_internal_test.go b/internal/engine/parallel_internal_test.go new file mode 100644 index 0000000..2ab8bb6 --- /dev/null +++ b/internal/engine/parallel_internal_test.go @@ -0,0 +1,87 @@ +package engine + +import "testing" + +// A group that ends on several members' terminal results reports ONE cause: +// cancellation over the ceiling over the send budget, and between equal +// causes the member declared first. (A breach outranks all three; the +// runner decides it before consulting this.) +func TestGroupTerminalCausePrecedence(t *testing.T) { + terminal := func(index int, outcome phaseOutcome, failure string) memberTerminal { + return memberTerminal{index: index, run: phaseRun{outcome: outcome, failure: failure}} + } + cases := []struct { + name string + terminals []memberTerminal + want endKind + diagnosis string + }{ + {"cancellation over the ceiling", []memberTerminal{ + terminal(0, phaseCeiling, "c"), terminal(2, phaseCancelled, "x")}, endCancelled, "x"}, + {"cancellation over the send budget", []memberTerminal{ + terminal(0, phaseSendBudget, "b"), terminal(1, phaseCancelled, "x")}, endCancelled, "x"}, + {"the ceiling over the send budget", []memberTerminal{ + terminal(0, phaseSendBudget, "b"), terminal(1, phaseCeiling, "c")}, endCeiling, "c"}, + {"equal causes go to the first declared member", []memberTerminal{ + terminal(2, phaseSendBudget, "late"), terminal(0, phaseSendBudget, "first")}, endSendBudget, "first"}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + end := groupTerminalEnd(tc.terminals) + if end.kind != tc.want || end.diagnostic != tc.diagnosis { + t.Fatalf("end = %+v, want kind %d with %q", *end, tc.want, tc.diagnosis) + } + }) + } +} + +// The send budget is one atomic check-and-take: however many members race +// for the last sends, exactly the budget is granted. +func TestConcurrentSendAcquisitionNeverOvershoots(t *testing.T) { + counters := &attemptCounters{phaseEntries: map[string]int{}} + const limit, racers = 50, 16 + granted := make(chan int, racers) + for r := 0; r < racers; r++ { + go func() { + count := 0 + for counters.acquireSend(limit) { + count++ + } + granted <- count + }() + } + total := 0 + for r := 0; r < racers; r++ { + total += <-granted + } + if total != limit || counters.sends.Load() != limit { + t.Fatalf("granted %d sends (counter %d), want exactly %d", total, counters.sends.Load(), limit) + } +} + +// The tightest race: many members released at once for the one send left. +// Repeated, because a check-then-add race loses only occasionally. +func TestExactlyOneMemberWinsTheLastSend(t *testing.T) { + for round := 0; round < 500; round++ { + counters := &attemptCounters{phaseEntries: map[string]int{}} + counters.sends.Store(9) + start := make(chan struct{}) + wins := make(chan bool, 32) + for r := 0; r < 32; r++ { + go func() { + <-start + wins <- counters.acquireSend(10) + }() + } + close(start) + won := 0 + for r := 0; r < 32; r++ { + if <-wins { + won++ + } + } + if won != 1 { + t.Fatalf("round %d: %d members took the last send, want exactly 1", round, won) + } + } +} diff --git a/internal/engine/parallel_test.go b/internal/engine/parallel_test.go new file mode 100644 index 0000000..60a80fe --- /dev/null +++ b/internal/engine/parallel_test.go @@ -0,0 +1,978 @@ +// parallel_test.go — the parallel read-only reviewer group (plan U5, R8–R11) +// on the scripted runtime. Members run concurrently in private views and +// merge in declared order; repair edges resolve after the join; any +// worktree change during the group is a breach. +package engine_test + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "reflect" + "regexp" + "strconv" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/StructuPath/jig/internal/engine" + "github.com/StructuPath/jig/internal/engine/enginetest" + "github.com/StructuPath/jig/internal/protocol" + "github.com/StructuPath/jig/internal/worker" +) + +// panelRoster: a builder, three read-only reviewers with distinct system +// prompts (the fake routes on them), and a read-only phase after the panel. +const panelRoster = ` +roster: + builder: + model: test-model + system_prompt: Build. + user_prompt: "Build the app." + writes: ["src/"] + correctness: + model: test-model + system_prompt: Review correctness. + user_prompt: "Review the work." + writes: [] + security: + model: test-model + system_prompt: Review security. + user_prompt: "Review the work." + writes: [] + maintainability: + model: test-model + system_prompt: Review maintainability. + user_prompt: "Review the work." + writes: [] + closer: + model: test-model + system_prompt: Close. + user_prompt: "Wrap up." + writes: [] +` + +// panelSnapshot is the panel chain with the given edges and group line. +// Every reviewer carries a verdict gate and a rejection edge back to build. +func panelSnapshot(edges map[string]string, group string) string { + edge := func(name string) string { + if custom, found := edges[name]; found { + return custom + } + return `{when: "approved == false", run: build, then: rerun-self, budget: 1, exhausted: fail-job}` + } + return "name: panel\n" + panelRoster + ` +phases: + - {name: build, kind: agent, owner: builder} + - name: review-correctness + kind: agent + owner: correctness + gates: [{name: verdict_consistent}] + on_fail: ` + edge("review-correctness") + ` + - name: review-security + kind: agent + owner: security + gates: [{name: verdict_consistent}] + on_fail: ` + edge("review-security") + ` + - name: review-maintainability + kind: agent + owner: maintainability + gates: [{name: verdict_consistent}] + on_fail: ` + edge("review-maintainability") + ` + - {name: wrap-up, kind: agent, owner: closer} +` + group + ` +acceptance: [all_phases_passed, verdict_consistent] +publish: + ci: + wait: true + on_fail: {run: build, budget: 1} +` +} + +const panelGroup = "parallel: [review-correctness, review-security, review-maintainability]" + +var ( + buildRole = enginetest.ForSystemPrompt("Build.") + correctnessRole = enginetest.ForSystemPrompt("Review correctness.") + securityRole = enginetest.ForSystemPrompt("Review security.") + maintainabilityRole = enginetest.ForSystemPrompt("Review maintainability.") + closerRole = enginetest.ForSystemPrompt("Close.") +) + +func approve(summary string, extra map[string]any) enginetest.Step { + fields := map[string]any{"approved": true, "blocking": []any{}} + for key, value := range extra { + fields[key] = value + } + return success(summary, fields) +} + +func reject(summary string) enginetest.Step { + return success(summary, map[string]any{"approved": false, "blocking": []any{summary}}) +} + +// normalizedSummary is the outcome's result with the wall-clock stamps +// removed, so a grouped and an ungrouped run can be compared exactly. +func normalizedSummary(t *testing.T, result string) map[string]any { + t.Helper() + var summary map[string]any + if err := json.Unmarshal([]byte(result), &summary); err != nil { + t.Fatalf("result does not parse: %v\n%s", err, result) + } + phases, _ := summary["phases"].([]any) + for _, phase := range phases { + entry := phase.(map[string]any) + delete(entry, "started_at") + delete(entry, "ended_at") + } + return summary +} + +func callsFor(fake *enginetest.Runtime, match func(enginetest.Call) bool) []enginetest.Call { + var matched []enginetest.Call + for _, call := range fake.Calls() { + if match(call) { + matched = append(matched, call) + } + } + return matched +} + +// scriptPanel scripts the panel identically for a grouped and an ungrouped +// run: prompt-independent reviewers, so the only thing that differs is how +// the engine runs them. Security's first emission omits its verdict, so its +// own gate correction goes to its own session inside the group. +func scriptPanel(fake *enginetest.Runtime) { + fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + fake.Route(correctnessRole, approve("correct", map[string]any{"score": 1.0, "correctness": "ok"})) + fake.Route(securityRole, + success("forgot the verdict", map[string]any{"score": 2.0}), + approve("secure", map[string]any{"score": 2.0, "security": "ok"})) + fake.Route(maintainabilityRole, approve("maintainable", map[string]any{"score": 3.0})) + fake.Route(closerRole, success("wrapped up", nil)) +} + +// R9: grouped and ungrouped runs of prompt-independent reviewers merge +// identically — results, gate reports, acceptance, and the field view — and +// every member's prompt carries the pre-group envelope, while the phase +// after the group gets the last member's envelope. +func TestAGroupMergesExactlyAsTheSameReviewersWouldSequentially(t *testing.T) { + run := func(group string) (*repairFixture, map[string]any, map[string]any) { + f := newRepairFixture(t, panelSnapshot(nil, group), nil) + scriptPanel(f.fake) + outcome := f.execute(t) + if outcome.State != protocol.AttemptAcceptedUnpublished || outcome.Continuation == nil { + t.Fatalf("%q: outcome = %s (%s), want accepted with a continuation", group, outcome.State, outcome.Error) + } + t.Cleanup(outcome.Continuation.Release) + if f.fake.Remaining() != 0 { + t.Fatalf("%q: unconsumed scripted steps: %d", group, f.fake.Remaining()) + } + return f, normalizedSummary(t, outcome.Result), engine.FieldViewOf(outcome.Continuation) + } + _, sequential, sequentialView := run("") + grouped, parallel, parallelView := run(panelGroup) + + if !reflect.DeepEqual(parallel, sequential) { + a, _ := json.MarshalIndent(parallel, "", " ") + b, _ := json.MarshalIndent(sequential, "", " ") + t.Fatalf("grouped summary differs from sequential\ngrouped:\n%s\nsequential:\n%s", a, b) + } + if !reflect.DeepEqual(parallelView, sequentialView) { + t.Fatalf("grouped field view %v differs from sequential %v", parallelView, sequentialView) + } + if parallelView["score"] != 3.0 { + t.Fatalf("score = %v, want the last member's 3 (declared-order merge)", parallelView["score"]) + } + + // Every member saw the builder's envelope, and no sibling's. + for _, member := range []func(enginetest.Call) bool{correctnessRole, securityRole, maintainabilityRole} { + first := callsFor(grouped.fake, member)[0] + if !strings.Contains(first.Prompt, "built the app") { + t.Fatalf("member prompt lacks the pre-group envelope: %.400q", first.Prompt) + } + for _, sibling := range []string{`"correct"`, `"secure"`, `"maintainable"`} { + if strings.Contains(first.Prompt, sibling) { + t.Fatalf("member prompt carries a sibling's envelope %s: %.400q", sibling, first.Prompt) + } + } + } + closer := callsFor(grouped.fake, closerRole)[0] + if !strings.Contains(closer.Prompt, `"maintainable"`) || strings.Contains(closer.Prompt, `"secure"`) { + t.Fatalf("the phase after the group must get the last member's envelope: %.400q", closer.Prompt) + } + // Security's gate correction stayed in its own session. + security := callsFor(grouped.fake, securityRole) + if len(security) != 2 || security[0].SessionKey != security[1].SessionKey { + t.Fatalf("security calls = %+v, want a correction in the same session", security) + } +} + +// executeWithin runs the fixture's attempt and fails the test if it does not +// return in time — the shape of a deadlock. +func executeWithin(t *testing.T, f *repairFixture, limit time.Duration) worker.Outcome { + t.Helper() + done := make(chan worker.Outcome, 1) + go func() { done <- f.runner.Execute(context.Background(), f.attempt) }() + select { + case outcome := <-done: + if outcome.Continuation != nil { + t.Cleanup(outcome.Continuation.Release) + } + return outcome + case <-time.After(limit): + t.Fatalf("the attempt did not finish within %s", limit) + return worker.Outcome{} + } +} + +// Members run concurrently: two reviewers that each hold their result until +// the other has started both finish, where one-at-a-time would deadlock. +// Session keys stay unique across members and the chain, and the trace's +// seq stays strictly increasing and gapless while members interleave. +func TestGroupMembersRunConcurrentlyUnderUniqueSessionsAndOneSequence(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), nil) + correctnessStarted, securityStarted := make(chan struct{}), make(chan struct{}) + correctness := approve("correct", nil) + correctness.Started, correctness.WaitFor = correctnessStarted, securityStarted + security := approve("secure", nil) + security.Started, security.WaitFor = securityStarted, correctnessStarted + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, correctness) + f.fake.Route(securityRole, security) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + + keys := map[string]string{} + for _, call := range f.fake.Calls() { + if owner, taken := keys[call.SessionKey]; taken && owner != call.Options.SystemPrompt { + t.Fatalf("session key %q is shared by %q and %q", call.SessionKey, owner, call.Options.SystemPrompt) + } + keys[call.SessionKey] = call.Options.SystemPrompt + } + if len(keys) != 5 { + t.Fatalf("session keys = %v, want one per role (5)", keys) + } + for i, event := range f.sink.all() { + if event.Seq != int64(i+1) { + t.Fatalf("event %d has seq %d: seq must be strictly increasing without gaps", i, event.Seq) + } + } +} + +func repairEdges(sink *recordingSink) []string { + var phases []string + for _, event := range sink.all() { + if event.Type == protocol.EventLog && event.Name == "repair_edge" { + phases = append(phases, event.Phase) + } + } + return phases +} + +// R10: when two members reject, only the first in declared order dispatches +// the builder, with its own envelope, and the whole group runs again. Only +// that member's budget is charged: the second member still has its one use +// when it rejects on the next run. +func TestOnlyTheFirstRejectingMemberDispatchesAndOnlyItsBudgetIsCharged(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), nil) + f.fake.Route(buildRole, + writes(map[string]string{"src/app.txt": "v1"}, "built the app"), + writes(map[string]string{"src/app.txt": "v2"}, "rebuilt once"), + writes(map[string]string{"src/app.txt": "v3"}, "rebuilt twice")) + f.fake.Route(correctnessRole, reject("correctness-blocker"), approve("correct", nil), approve("correct", nil)) + f.fake.Route(securityRole, reject("security-blocker-1"), reject("security-blocker-2"), approve("secure", nil)) + f.fake.Route(maintainabilityRole, + approve("maintainable", nil), approve("maintainable", nil), approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + if f.fake.Remaining() != 0 { + t.Fatalf("unconsumed scripted steps: %d", f.fake.Remaining()) + } + if got := repairEdges(f.sink); !reflect.DeepEqual(got, []string{"review-correctness", "review-security"}) { + t.Fatalf("repair edges = %v, want correctness's then (next run) security's", got) + } + if got := f.sink.count(protocol.EventLog, "parallel_group_start"); got != 3 { + t.Fatalf("group runs = %d, want 3", got) + } + builds := callsFor(f.fake, buildRole) + if !strings.Contains(builds[1].Prompt, "correctness-blocker") || strings.Contains(builds[1].Prompt, "security-blocker") { + t.Fatalf("the first dispatch must carry only the first rejecting member's envelope: %.500q", builds[1].Prompt) + } + if !strings.Contains(builds[2].Prompt, "security-blocker-2") { + t.Fatalf("the second dispatch must carry security's second rejection: %.500q", builds[2].Prompt) + } + // The rerun group sees the repaired envelope as its previous. + if second := callsFor(f.fake, correctnessRole)[1]; !strings.Contains(second.Prompt, "rebuilt once") { + t.Fatalf("the rerun group's members must see the repair's envelope: %.500q", second.Prompt) + } +} + +const ( + proceedEdge = `{when: "approved == false", run: build, then: rerun-self, budget: 1, exhausted: proceed}` +) + +// R10: the first member triggers with its budget spent under proceed, so the +// walk moves on and the second member dispatches. +func TestAnExhaustedProceedMemberLetsTheNextMemberDispatch(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(map[string]string{"review-correctness": proceedEdge}, panelGroup), nil) + f.fake.Route(buildRole, + writes(map[string]string{"src/app.txt": "v1"}, "built the app"), + writes(map[string]string{"src/app.txt": "v2"}, "rebuilt once"), + writes(map[string]string{"src/app.txt": "v3"}, "rebuilt twice")) + f.fake.Route(correctnessRole, reject("c-1"), reject("c-2"), reject("c-3")) + f.fake.Route(securityRole, approve("secure", nil), reject("security-blocker"), approve("secure", nil)) + f.fake.Route(maintainabilityRole, + approve("maintainable", nil), approve("maintainable", nil), approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + if got := repairEdges(f.sink); !reflect.DeepEqual(got, []string{"review-correctness", "review-security"}) { + t.Fatalf("repair edges = %v, want correctness's, then security's past correctness's exhaustion", got) + } + if got := f.sink.count(protocol.EventLog, "repair_exhausted"); got != 2 { + t.Fatalf("repair_exhausted events = %d, want 2 (correctness on runs 2 and 3)", got) + } + if builds := callsFor(f.fake, buildRole); !strings.Contains(builds[2].Prompt, "security-blocker") { + t.Fatalf("the second dispatch must carry security's rejection: %.500q", builds[2].Prompt) + } +} + +// R10: the first member triggers with its budget spent under fail-job: the +// attempt ends there, and no later member dispatches. +func TestAnExhaustedFailJobMemberEndsTheAttempt(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), nil) + f.fake.Route(buildRole, + writes(map[string]string{"src/app.txt": "v1"}, "built the app"), + writes(map[string]string{"src/app.txt": "v2"}, "rebuilt once")) + f.fake.Route(correctnessRole, reject("c-1"), reject("c-2")) + f.fake.Route(securityRole, approve("secure", nil), reject("security-blocker")) + f.fake.Route(maintainabilityRole, approve("maintainable", nil), approve("maintainable", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptFailed || !strings.Contains(outcome.Error, `"review-correctness": repair budget (1) exhausted`) { + t.Fatalf("outcome = %s (%s), want failed on correctness's exhausted budget", outcome.State, outcome.Error) + } + if got := len(callsFor(f.fake, buildRole)); got != 2 { + t.Fatalf("build ran %d time(s), want 2: nothing may dispatch past a fail-job exhaustion", got) + } +} + +// R10's bound: total group runs are at most one plus the sum of member +// budgets. Reviewers that never approve, under proceed, reach it exactly. +func TestGroupRunsAreBoundedByOnePlusTheSumOfMemberBudgets(t *testing.T) { + proceed := func(budget int) string { + return `{when: "approved == false", run: build, then: rerun-self, budget: ` + + strconv.Itoa(budget) + `, exhausted: proceed}` + } + f := newRepairFixture(t, panelSnapshot(map[string]string{ + "review-correctness": proceed(1), "review-security": proceed(2), "review-maintainability": proceed(1), + }, panelGroup), nil) + const runs = 1 + 1 + 2 + 1 + for i := 0; i < runs; i++ { + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": strconv.Itoa(i)}, "built")) + f.fake.Route(correctnessRole, reject("c")) + f.fake.Route(securityRole, reject("s")) + f.fake.Route(maintainabilityRole, reject("m")) + } + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 30*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted (every exhaustion proceeds)", outcome.State, outcome.Error) + } + if got := f.sink.count(protocol.EventLog, "parallel_group_start"); got != runs { + t.Fatalf("group runs = %d, want exactly %d", got, runs) + } + if f.fake.Remaining() != 0 { + t.Fatalf("unconsumed scripted steps: %d", f.fake.Remaining()) + } +} + +func phaseStatuses(t *testing.T, result string) map[string][]string { + t.Helper() + var summary struct { + Phases []protocol.PhaseResult `json:"phases"` + } + if err := json.Unmarshal([]byte(result), &summary); err != nil { + t.Fatalf("result does not parse: %v", err) + } + statuses := map[string][]string{} + for _, phase := range summary.Phases { + statuses[phase.Phase] = append(statuses[phase.Phase], phase.Status) + } + return statuses +} + +// R11: a member that writes breaches: the attempt aborts, the write is +// rolled back to the group snapshot, and no member's view is merged. +func TestAMemberWriteIsABreachThatAbortsWithNoPartialMerge(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), nil) + sneaky := approve("correct", nil) + sneaky.Files = map[string]string{"notes.txt": "a reviewer's scratch"} + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, sneaky) + f.fake.Route(securityRole, approve("secure", nil)) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptFailed || !strings.Contains(outcome.Error, "notes.txt") || + !strings.Contains(outcome.Error, "parallel group") { + t.Fatalf("outcome = %s (%s), want failed on a group breach naming notes.txt", outcome.State, outcome.Error) + } + if _, err := os.Stat(filepath.Join(f.repo, "notes.txt")); !os.IsNotExist(err) { + t.Fatalf("the breaching write survived: %v", err) + } + for _, member := range []string{"review-correctness", "review-security", "review-maintainability"} { + if got := phaseStatuses(t, outcome.Result)[member]; !reflect.DeepEqual(got, []string{protocol.EnvelopeFail}) { + t.Fatalf("%s results = %v, want one failed breach result and nothing merged", member, got) + } + } + if !f.sink.has(protocol.EventError, "write_boundary_breach") { + t.Fatal("no write_boundary_breach event") + } +} + +// R11: a member that dies leaving a write is a breach too — the member's +// death does not roll the tree back; the group's one enforcement finds it. +func TestADeadMembersLeftoverWriteIsABreach(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), nil) + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, + enginetest.Step{Crash: true, Files: map[string]string{"left-behind.txt": "half a thought"}}, + approve("correct", nil)) + f.fake.Route(securityRole, approve("secure", nil)) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptFailed || !strings.Contains(outcome.Error, "left-behind.txt") { + t.Fatalf("outcome = %s (%s), want failed on the dead member's leftover write", outcome.State, outcome.Error) + } + if _, err := os.Stat(filepath.Join(f.repo, "left-behind.txt")); !os.IsNotExist(err) { + t.Fatalf("the dead member's write survived: %v", err) + } +} + +func eventSeq(t *testing.T, sink *recordingSink, eventType, phase string, match func(protocol.Event) bool) int64 { + t.Helper() + for _, event := range sink.all() { + if event.Type == eventType && event.Phase == phase && (match == nil || match(event)) { + return event.Seq + } + } + t.Fatalf("no %s event for %s", eventType, phase) + return 0 +} + +// A member death re-enters only that member, after its siblings finish, and +// the siblings' results are kept rather than re-run. +func TestAMemberDeathReentersOnlyThatMemberAfterSiblingsFinish(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), nil) + correctnessStarted := make(chan struct{}) + slowSecurity := approve("secure", nil) + slowSecurity.WaitFor, slowSecurity.Delay = correctnessStarted, 300*time.Millisecond + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, + enginetest.Step{Crash: true, Started: correctnessStarted}, + approve("correct", nil)) + f.fake.Route(securityRole, slowSecurity) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + reentry := eventSeq(t, f.sink, protocol.EventPhaseStart, "review-correctness", func(event protocol.Event) bool { + return strings.Contains(string(event.Payload), `"phase_attempt":2`) + }) + for _, sibling := range []string{"review-security", "review-maintainability"} { + if end := eventSeq(t, f.sink, protocol.EventPhaseEnd, sibling, nil); end > reentry { + t.Fatalf("correctness re-entered (seq %d) before %s finished (seq %d)", reentry, sibling, end) + } + } + statuses := phaseStatuses(t, outcome.Result) + if got := statuses["review-correctness"]; !reflect.DeepEqual(got, []string{protocol.EnvelopeFail, protocol.EnvelopeSuccess}) { + t.Fatalf("correctness results = %v, want its death then its re-entry", got) + } + if got := statuses["review-security"]; !reflect.DeepEqual(got, []string{protocol.EnvelopeSuccess}) { + t.Fatalf("security results = %v, want its one kept result", got) + } + if calls := callsFor(f.fake, correctnessRole); calls[0].SessionKey == calls[1].SessionKey { + t.Fatalf("the re-entry must use a fresh session, got %q twice", calls[0].SessionKey) + } +} + +const twoMemberGroup = "parallel: [review-correctness, review-security]" + +// A repair that flips the field a member's guard reads must not skip the +// member that rejected it: guards are judged on the group's first run only. +// Sequentially, rerun-self never re-checks a guard, and a skipped phase +// counts as passed — so a re-judged guard would let a repair silence its own +// reviewer. A member skipped on the first run stays skipped. +func TestARepairCannotFlipAGuardToSkipTheReviewerThatRejectedIt(t *testing.T) { + snapshot := strings.Replace(panelSnapshot(nil, panelGroup), + " owner: security\n", " owner: security\n if: \"touches_auth == true\"\n", 1) + snapshot = strings.Replace(snapshot, + " owner: maintainability\n", " owner: maintainability\n if: \"touches_db == true\"\n", 1) + f := newRepairFixture(t, snapshot, nil) + built := writes(map[string]string{"src/app.txt": "v1"}, "built the app") + built.Text = envelope(map[string]any{"status": "success", "summary": "built the app", + "touches_auth": true, "touches_db": false}) + rebuilt := writes(map[string]string{"src/app.txt": "v2"}, "rebuilt") + rebuilt.Text = envelope(map[string]any{"status": "success", "summary": "moved the auth check away", + "touches_auth": false, "touches_db": true}) + f.fake.Route(buildRole, built, rebuilt) + f.fake.Route(correctnessRole, approve("correct", nil), approve("correct", nil)) + f.fake.Route(securityRole, reject("auth bypass"), approve("secure now", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + if got := len(callsFor(f.fake, securityRole)); got != 2 { + t.Fatalf("security ran %d time(s), want 2: the repair flipped its guard, and it must re-judge anyway", got) + } + if got := len(callsFor(f.fake, maintainabilityRole)); got != 0 { + t.Fatalf("maintainability ran %d time(s): skipped on the first run, it must stay skipped", got) + } + statuses := phaseStatuses(t, outcome.Result)["review-security"] + if statuses[len(statuses)-1] != protocol.EnvelopeSuccess { + t.Fatalf("security's last result = %v, want its approval, never a skip", statuses) + } +} + +var handoffLine = regexp.MustCompile(`Share working notes for later phases in: (\S+)`) + +func handoffDir(t *testing.T, prompt string) string { + t.Helper() + match := handoffLine.FindStringSubmatch(prompt) + if match == nil { + t.Errorf("prompt names no handoff directory: %.300q", prompt) + return "" + } + return match[1] +} + +// Members never share a handoff directory: a sibling cannot read another's +// notes mid-run, and two notes with the same name both survive, merged into +// the chain's handoff directory in declared order under per-member folders. +func TestMembersKeepPrivateHandoffNotesThatMergeAtTheJoin(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, twoMemberGroup), nil) + built := writes(map[string]string{"src/app.txt": "app"}, "built the app") + built.Do = func(call enginetest.Call) { + // Notes a phase before the group left, and a stand-in for a + // sibling's notes merged on an earlier group run. + dir := handoffDir(t, call.Prompt) + for name, body := range map[string]string{ + "plan.md": "the plan", filepath.Join("parallel", "review-security", "notes.md"): "sibling's", + } { + if err := os.MkdirAll(filepath.Dir(filepath.Join(dir, name)), 0o700); err != nil { + t.Errorf("mkdir: %v", err) + } + if err := os.WriteFile(filepath.Join(dir, name), []byte(body), 0o600); err != nil { + t.Errorf("write %s: %v", name, err) + } + } + } + seesPreGroupNotesOnly := func(dir string) { + if body, err := os.ReadFile(filepath.Join(dir, "plan.md")); err != nil || string(body) != "the plan" { + t.Errorf("a member cannot read the notes earlier phases left: %q, %v", body, err) + } + if _, err := os.Stat(filepath.Join(dir, "parallel")); !os.IsNotExist(err) { + t.Errorf("a member was seeded with members' merged notes: %v", err) + } + } + wrote := make(chan struct{}) + correctness := approve("correct", nil) + correctness.Do = func(call enginetest.Call) { + defer close(wrote) + dir := handoffDir(t, call.Prompt) + seesPreGroupNotesOnly(dir) + if err := os.WriteFile(filepath.Join(dir, "notes.md"), []byte("from correctness"), 0o600); err != nil { + t.Errorf("write correctness notes: %v", err) + } + } + security := approve("secure", nil) + security.Do = func(call enginetest.Call) { + <-wrote + dir := handoffDir(t, call.Prompt) + seesPreGroupNotesOnly(dir) + if _, err := os.Stat(filepath.Join(dir, "notes.md")); !os.IsNotExist(err) { + t.Errorf("security can see a sibling's notes mid-run in %s: %v", dir, err) + } + if err := os.WriteFile(filepath.Join(dir, "notes.md"), []byte("from security"), 0o600); err != nil { + t.Errorf("write security notes: %v", err) + } + // An edit to a seeded note is the member's own work, and is kept. + if err := os.WriteFile(filepath.Join(dir, "plan.md"), []byte("the plan, annotated"), 0o600); err != nil { + t.Errorf("annotate the plan: %v", err) + } + } + f.fake.Route(buildRole, built) + f.fake.Route(correctnessRole, correctness) + f.fake.Route(securityRole, security) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + shared := filepath.Join(f.scratch, "attempt-1", "handoff") + for member, want := range map[string]string{ + "review-correctness": "from correctness", "review-security": "from security", + } { + body, err := os.ReadFile(filepath.Join(shared, "parallel", member, "notes.md")) + if err != nil || string(body) != want { + t.Fatalf("%s's merged notes = %q, %v; want %q", member, body, err, want) + } + } + // Only what a member wrote or changed is published, not what it was + // seeded with: correctness left the plan alone, security annotated it. + if _, err := os.Stat(filepath.Join(shared, "parallel", "review-correctness", "plan.md")); !os.IsNotExist(err) { + t.Fatalf("correctness republished the pre-group notes it was seeded with: %v", err) + } + annotated, err := os.ReadFile(filepath.Join(shared, "parallel", "review-security", "plan.md")) + if err != nil || string(annotated) != "the plan, annotated" { + t.Fatalf("security's edit to a seeded note = %q, %v; want it published", annotated, err) + } + if body, err := os.ReadFile(filepath.Join(shared, "plan.md")); err != nil || string(body) != "the plan" { + t.Fatalf("the pre-group notes changed: %q, %v", body, err) + } + if got := handoffDir(t, callsFor(f.fake, closerRole)[0].Prompt); got != shared { + t.Fatalf("the phase after the group was pointed at %s, want the chain's handoff %s", got, shared) + } +} + +// A member that panics is a panic in the chain, but only after the group's +// one enforcement: a sibling's worktree write is rolled back first. +func TestAPanickingMemberStillHasItsSiblingsWritesRolledBack(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, twoMemberGroup), nil) + correctnessStarted := make(chan struct{}) + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, enginetest.Step{Hang: true, Started: correctnessStarted, + Files: map[string]string{"planted.txt": "written before the crash"}}) + f.fake.Route(securityRole, enginetest.Step{Do: func(enginetest.Call) { + <-correctnessStarted + panic("runtime blew up") + }}) + + recovered := func() (value any) { + defer func() { value = recover() }() + f.runner.Execute(context.Background(), f.attempt) + return nil + }() + if recovered != "runtime blew up" { + t.Fatalf("Execute recovered %v, want the member's panic re-raised", recovered) + } + if _, err := os.Stat(filepath.Join(f.repo, "planted.txt")); !os.IsNotExist(err) { + t.Fatalf("a sibling's write survived the panic: %v", err) + } + if !f.sink.has(protocol.EventError, "parallel_group_panic") { + t.Fatal("no parallel_group_panic event recording what enforcement found") + } +} + +// A member that fails without firing its edge ends the attempt at the join, +// as it would have sequentially: nothing after the group runs. +func TestAnUntriggeredFailingMemberEndsTheAttempt(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), nil) + failing := enginetest.Step{Text: envelope(map[string]any{ + "status": "fail", "summary": "could not finish the review", "approved": true, "blocking": []any{}})} + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, approve("correct", nil)) + f.fake.Route(securityRole, failing) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptFailed || !strings.Contains(outcome.Error, `"review-security"`) || + strings.Contains(outcome.Error, "acceptance") { + t.Fatalf("outcome = %s (%s), want failed on review-security itself", outcome.State, outcome.Error) + } + if calls := callsFor(f.fake, closerRole); len(calls) != 0 { + t.Fatal("the phase after the group ran after a member failed") + } +} + +// A member that fails on a runtime error takes the fail exit, which inside a +// group must not run a boundary check of its own: the member never took a +// snapshot, so it would read the builder's authorized work as the member's +// breach and roll it back. The group's one enforcement is the only check. +func TestAMemberRuntimeErrorFailsWithoutItsOwnBoundaryCheck(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, twoMemberGroup), nil) + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, approve("correct", nil)) + f.fake.Route(securityRole, enginetest.Step{IsError: true, ExitCode: 1, Text: "rate limit exceeded"}) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptFailed || !strings.Contains(outcome.Error, "rate limit exceeded") || + strings.Contains(outcome.Error, "allowlist") { + t.Fatalf("outcome = %s (%s), want failed on the runtime error alone", outcome.State, outcome.Error) + } + if body, err := os.ReadFile(filepath.Join(f.repo, "src", "app.txt")); err != nil || string(body) != "app" { + t.Fatalf("the builder's authorized work was rolled back: %q, %v", body, err) + } +} + +// A member's role may also own a phase outside the group. The two run in +// different HOMEs, so they must never share a session key — a runtime would +// try to resume a conversation that lives in the other HOME. +func TestAMemberRoleUsedOutsideTheGroupKeepsSeparateSessions(t *testing.T) { + snapshot := strings.Replace(panelSnapshot(nil, panelGroup), + "{name: wrap-up, kind: agent, owner: closer}", "{name: wrap-up, kind: agent, owner: maintainability}", 1) + f := newRepairFixture(t, snapshot, nil) + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, approve("correct", nil)) + f.fake.Route(securityRole, approve("secure", nil)) + f.fake.Route(maintainabilityRole, approve("maintainable", nil), success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + calls := callsFor(f.fake, maintainabilityRole) + if len(calls) != 2 || calls[0].SessionKey == calls[1].SessionKey { + t.Fatalf("maintainability calls = %d, keys %q — the member and the chain phase must not share a session", + len(calls), []string{calls[0].SessionKey, calls[len(calls)-1].SessionKey}) + } +} + +// A member hitting the send budget stops a sibling that is mid-send; the +// runner waits for it, enforces, and ends the attempt on the send budget — +// never one send past it. +func TestAMemberOnTheSendBudgetStopsItsSiblingAndEndsTheAttempt(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, twoMemberGroup), + func(config *engine.Config) { config.MaxAttemptSends = 3 }) + correctnessStarted := make(chan struct{}) + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, enginetest.Step{Hang: true, Started: correctnessStarted}) + f.fake.Route(securityRole, enginetest.Step{Text: "not an envelope", WaitFor: correctnessStarted}) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptFailed || !strings.Contains(outcome.Error, "send budget") { + t.Fatalf("outcome = %s (%s), want failed on the send budget", outcome.State, outcome.Error) + } + if got := len(f.fake.Calls()); got != 3 { + t.Fatalf("sends = %d, want exactly the budget (3)", got) + } + if got := f.fake.Kills(); got != 1 { + t.Fatalf("killed sends = %d when Execute returned, want the hanging sibling's 1", got) + } + if !f.sink.has(protocol.EventError, "attempt_send_budget_exhausted") { + t.Fatal("no attempt_send_budget_exhausted event") + } + // The stopped sibling still closes its agent_start, with its killed send + // counted as unmetered rather than free (R1). + var stopped struct { + Outcome string `json:"outcome"` + Unmetered int `json:"unmetered_sends"` + } + for _, event := range f.sink.all() { + if event.Type == protocol.EventAgentEnd && event.Phase == "review-correctness" { + if err := json.Unmarshal(event.Payload, &stopped); err != nil { + t.Fatal(err) + } + } + } + if stopped.Outcome != protocol.AgentStopped || stopped.Unmetered != 1 { + t.Fatalf("stopped sibling's agent_end = %+v, want outcome %q with 1 unmetered send", + stopped, protocol.AgentStopped) + } +} + +// Precedence: a breach outranks the send budget, which proves the one +// enforcement runs on a terminal exit too. +func TestABreachOutranksTheSendBudgetThatEndedTheGroup(t *testing.T) { + f := newRepairFixture(t, panelSnapshot(nil, twoMemberGroup), + func(config *engine.Config) { config.MaxAttemptSends = 3 }) + correctnessStarted := make(chan struct{}) + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, enginetest.Step{Hang: true, Started: correctnessStarted, + Files: map[string]string{"planted.txt": "while hanging"}}) + f.fake.Route(securityRole, enginetest.Step{Text: "not an envelope", WaitFor: correctnessStarted}) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptFailed || !strings.Contains(outcome.Error, "planted.txt") { + t.Fatalf("outcome = %s (%s), want failed on the breach, not the send budget", outcome.State, outcome.Error) + } + if f.sink.has(protocol.EventError, "attempt_send_budget_exhausted") { + t.Fatal("the send budget was reported over a breach") + } + if _, err := os.Stat(filepath.Join(f.repo, "planted.txt")); !os.IsNotExist(err) { + t.Fatalf("the breaching write survived: %v", err) + } +} + +// Every member runs in its own seeded HOME — not the chain's, not a +// sibling's — and every member HOME is wiped when the chain ends. A CI +// repair round runs the group through the same runner in fresh HOMEs. +func TestEachMemberRunsInItsOwnSeededHomeAndEveryHomeIsWiped(t *testing.T) { + var mutex sync.Mutex + seeded := map[string][]string{} + var inFlight, overlapped atomic.Int32 + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), func(config *engine.Config) { + config.SeedHome = func(home, role string, _ protocol.RoleSpec) error { + // Seeding is serial, before any member starts: a second seed + // arriving while this one sleeps is a concurrent seed. + if inFlight.Add(1) > 1 { + overlapped.Add(1) + } + defer inFlight.Add(-1) + time.Sleep(30 * time.Millisecond) + mutex.Lock() + seeded[role] = append(seeded[role], home) + mutex.Unlock() + return os.WriteFile(filepath.Join(home, "auth-"+role), []byte("token"), 0o600) + } + }) + scriptPanel(f.fake) + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished || outcome.Continuation == nil { + t.Fatalf("outcome = %s (%s), want accepted with a continuation", outcome.State, outcome.Error) + } + + if overlapped.Load() != 0 { + t.Fatalf("%d HOME seed(s) overlapped another: members must be seeded serially before fan-out", + overlapped.Load()) + } + chainHome := filepath.Join(f.scratch, "attempt-1", "home") + homes := map[string]string{} + for _, role := range []string{"correctness", "security", "maintainability"} { + if len(seeded[role]) != 1 { + t.Fatalf("%s seeded %d time(s), want once", role, len(seeded[role])) + } + home := seeded[role][0] + if home == chainHome { + t.Fatalf("member role %s was seeded into the chain's HOME", role) + } + if other, shared := homes[home]; shared { + t.Fatalf("members %s and %s share HOME %s", other, role, home) + } + homes[home] = role + for _, call := range callsFor(f.fake, enginetest.ForSystemPrompt( + map[string]string{"correctness": "Review correctness.", "security": "Review security.", + "maintainability": "Review maintainability."}[role])) { + if !contains(call.Options.Env, "HOME="+home) { + t.Fatalf("%s's send ran outside its seeded HOME %s: %v", role, home, call.Options.Env) + } + } + } + if seeded["builder"][0] != chainHome { + t.Fatalf("the builder was seeded into %s, want the chain's HOME", seeded["builder"][0]) + } + for home := range homes { + if _, err := os.Stat(home); !os.IsNotExist(err) { + t.Fatalf("member HOME %s survived the chain: %v", home, err) + } + } + + // A round re-runs the group, seeding every member again. + f.fake.Route(buildRole, writes(map[string]string{"src/fix.txt": "fixed"}, "fixed")) + f.fake.Route(correctnessRole, approve("correct", nil)) + f.fake.Route(securityRole, approve("secure", nil)) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + if round := outcome.Continuation.RepairCI(context.Background(), redLint); round.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("round = %s (%s), want accepted", round.State, round.Error) + } + if got := f.sink.count(protocol.EventLog, "parallel_group_start"); got != 2 { + t.Fatalf("group runs = %d, want the chain's and the round's", got) + } + for _, role := range []string{"correctness", "security", "maintainability"} { + if len(seeded[role]) != 2 { + t.Fatalf("%s seeded %d time(s), want again for the round", role, len(seeded[role])) + } + } +} + +func contains(items []string, want string) bool { + for _, item := range items { + if item == want { + return true + } + } + return false +} + +// A member HOME that cannot be wiped keeps the chain from being handed to +// a round, exactly as the chain's own HOME does. +func TestAnUnwipeableMemberHomeIsNeverHandedToARound(t *testing.T) { + if os.Geteuid() == 0 { + t.Skip("root ignores directory permissions, so the wipe cannot be made to fail") + } + f := newRepairFixture(t, panelSnapshot(nil, panelGroup), func(config *engine.Config) { + config.SeedHome = func(home, role string, _ protocol.RoleSpec) error { + if role == "security" { + return lockHome(t, home) + } + return nil + } + }) + scriptPanel(f.fake) + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished || outcome.Continuation != nil { + t.Fatalf("outcome = %s (%s) continuation=%v, want accepted with no continuation", + outcome.State, outcome.Error, outcome.Continuation != nil) + } +} + +// Both members' process groups are live at once and both are recorded; each +// is cleared when its send ends. +func TestEveryMembersProcessGroupIsRecordedWhileLive(t *testing.T) { + var mutex sync.Mutex + live := map[int64]bool{} + bothLive := make(chan struct{}) + var once sync.Once + f := newRepairFixture(t, panelSnapshot(nil, twoMemberGroup), func(config *engine.Config) { + config.RecordProcess = func(_ string, group int64, active bool) error { + mutex.Lock() + defer mutex.Unlock() + if active { + live[group] = true + } else { + delete(live, group) + } + if live[101] && live[102] { + once.Do(func() { close(bothLive) }) + } + return nil + } + }) + correctness := approve("correct", nil) + correctness.ProcessGroup, correctness.WaitFor = 101, bothLive + security := approve("secure", nil) + security.ProcessGroup, security.WaitFor = 102, bothLive + f.fake.Route(buildRole, writes(map[string]string{"src/app.txt": "app"}, "built the app")) + f.fake.Route(correctnessRole, correctness) + f.fake.Route(securityRole, security) + f.fake.Route(maintainabilityRole, approve("maintainable", nil)) + f.fake.Route(closerRole, success("wrapped up", nil)) + + outcome := executeWithin(t, f, 20*time.Second) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("outcome = %s (%s), want accepted", outcome.State, outcome.Error) + } + mutex.Lock() + defer mutex.Unlock() + if len(live) != 0 { + t.Fatalf("process groups still recorded live after the chain: %v", live) + } +} diff --git a/internal/engine/phase.go b/internal/engine/phase.go index 716b43e..804d5b4 100644 --- a/internal/engine/phase.go +++ b/internal/engine/phase.go @@ -29,6 +29,8 @@ import ( "reflect" "sort" "strings" + "sync" + "sync/atomic" "time" "github.com/StructuPath/jig/internal/protocol" @@ -207,7 +209,7 @@ func (r *Runner) Execute(ctx context.Context, attempt Attempt) worker.Outcome { reportedBy: make(map[string]string), gateReports: make(map[string]protocol.GateReport), seededRoles: make(map[string]bool), - phaseEntries: make(map[string]int), + counters: &attemptCounters{phaseEntries: make(map[string]int)}, touchedPaths: make(map[string]bool), sessionKey: attempt.Claim.Attempt.ID, } @@ -239,7 +241,6 @@ type execution struct { timeouts timeoutConfig deadline time.Time sessions map[string]*runtime.Session - sessionSeq int transcripts map[string][]exchange fieldView map[string]any // publishHeld is the hold reason once publishHold() fired, so the summary @@ -254,17 +255,78 @@ type execution struct { results []protocol.PhaseResult gateReports map[string]protocol.GateReport seededRoles map[string]bool - phaseEntries map[string]int touchedPaths map[string]bool // sessionKey prefixes every runtime session key. It is the attempt id // for the chain and gains a round suffix in each CI repair round, whose // sessions start in a fresh HOME and must never collide with a key the // runtime already saw. sessionKey string + // counters are the attempt-wide numbers a parallel group's members share + // with the chain: sends, the session-key sequence, and phase entries. + counters *attemptCounters + + // members are the parallel group's private member views, by phase name, + // kept across group runs within one chain or round (parallel.go). Nil + // until the group first runs; forgotten when the HOMEs are wiped. + members map[string]*execution + // grouped marks a member view: it runs inside a group, which owns the + // write boundary, so the view takes no snapshot, enforces nothing, and + // never rolls back — and defers its field merge to the join. + grouped bool + // stop is the group's stop signal, selected by every member send. Nil + // outside a group, where it never fires. + stop <-chan struct{} + // pendingFields are a member view's envelope merges, replayed into the + // chain's field view in declared order at the join. + pendingFields []fieldMerge + // handoffSeed fingerprints the pre-group notes copied into a member's + // private handoff directory, so the join publishes only what the member + // itself wrote or changed. + handoffSeed map[string]string +} + +// fieldMerge is one deferred mergeAgentFields call. +type fieldMerge struct { + phase, role string + fields map[string]any +} + +// attemptCounters are the attempt-wide counters. They are shared by pointer +// between the chain and a parallel group's member views, so every access is +// atomic or locked. +type attemptCounters struct { // sends counts every prompt send this attempt has issued, across phases, - // parse corrections, gate corrections, crash re-entries, and repair - // dispatches alike — the one number the whole ladder is bounded by. - sends int + // parse corrections, gate corrections, crash re-entries, repair + // dispatches, and group members alike — the one number the whole ladder + // is bounded by. + sends atomic.Int64 + // sessionSeq numbers fresh session keys so no two are ever alike. + sessionSeq atomic.Int64 + + mutex sync.Mutex + phaseEntries map[string]int +} + +// acquireSend takes one send from the budget, or refuses when the budget is +// spent. Check and increment are one atomic step: two members racing for +// the last send cannot both win it. +func (c *attemptCounters) acquireSend(limit int) bool { + for { + current := c.sends.Load() + if current >= int64(limit) { + return false + } + if c.sends.CompareAndSwap(current, current+1) { + return true + } + } +} + +func (c *attemptCounters) nextEntry(phase string) int { + c.mutex.Lock() + defer c.mutex.Unlock() + c.phaseEntries[phase]++ + return c.phaseEntries[phase] } // ---- attempt-level control ------------------------------------------------- @@ -319,7 +381,7 @@ func (e *execution) conclude(end chainEnd) worker.Outcome { // Its own terminal cause, like the ceiling: the correction ladder ran // out of budget, which is neither a phase that failed nor a clock. e.emit.emit(protocol.EventError, "", "attempt_send_budget_exhausted", - map[string]any{"max_sends": e.runner.config.MaxAttemptSends, "sends": e.sends}) + map[string]any{"max_sends": e.runner.config.MaxAttemptSends, "sends": e.counters.sends.Load()}) return worker.Outcome{State: protocol.AttemptFailed, Error: end.diagnostic, Result: e.summaryJSON(nil)} case endAborted, endFailed: @@ -526,9 +588,19 @@ func (e *execution) freshenLease(ctx context.Context) { // to the first phase run. A fresh edgeUses per call means each CI repair // round gets fresh per-edge budgets; the attempt-wide sends and ceiling are // what bound the rounds together. +// +// A declared parallel group runs as ONE step: the whole group through the +// group runner, then the chain resumes after its last member. func (e *execution) runChain(ctx context.Context, start int, previous *parsedEnvelope) chainEnd { edgeUses := make(map[string]int) - for _, phase := range e.spec.Phases[start:] { + groupStart, groupEnd, grouped := e.spec.ParallelRange() + if grouped && start > groupStart && start < groupEnd { + // Validation keeps every chain entry point outside the group. + return chainEnd{endFailed, fmt.Sprintf( + "chain cannot start at phase %q, inside the parallel group", e.spec.Phases[start].Name)} + } + for index := start; index < len(e.spec.Phases); index++ { + phase := e.spec.Phases[index] if e.cancelled() { return chainEnd{endCancelled, "cancelled between phases"} } @@ -536,6 +608,17 @@ func (e *execution) runChain(ctx context.Context, start int, previous *parsedEnv return chainEnd{endCeiling, ""} } e.freshenLease(ctx) + if grouped && index == groupStart { + end, envelope := e.runGroup(ctx, e.spec.Phases[groupStart:groupEnd], previous, edgeUses) + if end != nil { + return *end + } + if envelope != nil { + previous = envelope + } + index = groupEnd - 1 + continue + } if phase.If != "" && !guardHolds(phase.If, e.fieldView) { e.recordResult(protocol.PhaseResult{ Phase: phase.Name, Kind: phase.Kind, Status: phaseStatusSkipped, @@ -563,70 +646,93 @@ func (e *execution) runChain(ctx context.Context, start int, previous *parsedEnv func (e *execution) runPhaseWithEdge( ctx context.Context, phase protocol.PhaseSpec, previous *parsedEnvelope, edgeUses map[string]int, ) (*chainEnd, *parsedEnvelope) { - edge := phase.OnFail for { run := e.runPhaseOnce(ctx, phase, previous) if end := run.attemptEnd(); end != nil { return end, nil } - - triggered := false - if edge != nil && run.hasEnvelope { - switch phase.Kind { - case protocol.PhaseKindCode: - triggered = run.outcome == phaseFailed - case protocol.PhaseKindAgent: - predicate, err := protocol.ParsePredicate(edge.When) - triggered = err == nil && predicateHolds(predicate, run.envelope.Fields) - } - } - if !triggered { + if !edgeTriggered(phase, run) { if run.outcome == phaseFailed { return &chainEnd{endFailed, run.failure}, nil } return nil, run.envelopeRef() } - - if edgeUses[phase.Name] >= edge.Budget { - policy := protocol.RepairExhaustedProceed - if edgeExhaustionFailsJob(edge) { - policy = protocol.RepairExhaustedFailJob - } - e.emit.emit(protocol.EventLog, phase.Name, "repair_exhausted", map[string]any{ - "budget": edge.Budget, "policy": policy, - }) - if edgeExhaustionFailsJob(edge) { - return &chainEnd{endFailed, fmt.Sprintf( - "phase %q: repair budget (%d) exhausted", phase.Name, edge.Budget)}, nil - } - return nil, run.envelopeRef() + end, repaired, dispatched := e.followEdge(ctx, phase, run, edgeUses) + if end != nil { + return end, nil } - edgeUses[phase.Name]++ - e.emit.emit(protocol.EventLog, phase.Name, "repair_edge", map[string]any{ - "run": edge.Run, "use": edgeUses[phase.Name], "budget": edge.Budget, - }) - - // Dispatch the repair target with the FAILING envelope as its - // previous — a failing test suite and a rejecting review enter the - // repair loop through the same door (R8). The target's own repair - // edge does not fire on a dispatched run: budgets bound one edge, - // not a chain of them. - repairPhase, found := e.phaseByName(edge.Run) - if !found { - return &chainEnd{endFailed, fmt.Sprintf( - "phase %q: repair target %q missing from frozen snapshot", phase.Name, edge.Run)}, nil + if !dispatched { + return nil, run.envelopeRef() // exhausted under proceed } - repairRun := e.runPhaseOnce(ctx, repairPhase, run.envelopeRef()) - if end := repairRun.attemptEnd(); end != nil { - return end, nil + previous = repaired + // then: rerun-self. + } +} + +// edgeTriggered reports whether a finished run fires its phase's declared +// repair edge: nonzero exit for code, the envelope predicate for agent. +func edgeTriggered(phase protocol.PhaseSpec, run phaseRun) bool { + edge := phase.OnFail + if edge == nil || !run.hasEnvelope { + return false + } + switch phase.Kind { + case protocol.PhaseKindCode: + return run.outcome == phaseFailed + case protocol.PhaseKindAgent: + predicate, err := protocol.ParsePredicate(edge.When) + return err == nil && predicateHolds(predicate, run.envelope.Fields) + } + return false +} + +// followEdge acts on a triggered edge. With budget spent it applies the +// exhaustion policy: fail-job returns the attempt's end, proceed returns +// dispatched=false. With budget left it charges ONE use to this phase's edge +// and dispatches the repair target, returning the target's envelope — the +// previous for the rerun that follows. +func (e *execution) followEdge( + ctx context.Context, phase protocol.PhaseSpec, run phaseRun, edgeUses map[string]int, +) (end *chainEnd, repaired *parsedEnvelope, dispatched bool) { + edge := phase.OnFail + if edgeUses[phase.Name] >= edge.Budget { + policy := protocol.RepairExhaustedProceed + if edgeExhaustionFailsJob(edge) { + policy = protocol.RepairExhaustedFailJob } - if repairRun.outcome == phaseFailed { + e.emit.emit(protocol.EventLog, phase.Name, "repair_exhausted", map[string]any{ + "budget": edge.Budget, "policy": policy, + }) + if edgeExhaustionFailsJob(edge) { return &chainEnd{endFailed, fmt.Sprintf( - "repair phase %q: %s", repairPhase.Name, repairRun.failure)}, nil + "phase %q: repair budget (%d) exhausted", phase.Name, edge.Budget)}, nil, false } - previous = repairRun.envelopeRef() - // then: rerun-self. + return nil, nil, false } + edgeUses[phase.Name]++ + e.emit.emit(protocol.EventLog, phase.Name, "repair_edge", map[string]any{ + "run": edge.Run, "use": edgeUses[phase.Name], "budget": edge.Budget, + }) + + // Dispatch the repair target with the FAILING envelope as its + // previous — a failing test suite and a rejecting review enter the + // repair loop through the same door (R8). The target's own repair + // edge does not fire on a dispatched run: budgets bound one edge, + // not a chain of them. + repairPhase, found := e.phaseByName(edge.Run) + if !found { + return &chainEnd{endFailed, fmt.Sprintf( + "phase %q: repair target %q missing from frozen snapshot", phase.Name, edge.Run)}, nil, false + } + repairRun := e.runPhaseOnce(ctx, repairPhase, run.envelopeRef()) + if end := repairRun.attemptEnd(); end != nil { + return end, nil, false + } + if repairRun.outcome == phaseFailed { + return &chainEnd{endFailed, fmt.Sprintf( + "repair phase %q: %s", repairPhase.Name, repairRun.failure)}, nil, false + } + return nil, repairRun.envelopeRef(), true } func (e *execution) phaseByName(name string) (protocol.PhaseSpec, bool) { @@ -649,6 +755,9 @@ const ( phaseCancelled phaseCeiling phaseSendBudget // the attempt's send ladder hit its enforced bound + // phaseStopped: a group member stopped because a sibling ended the + // attempt. Only the group runner ever sees it, and it discards it. + phaseStopped ) type phaseRun struct { @@ -668,6 +777,9 @@ func (run phaseRun) attemptEnd() *chainEnd { return &chainEnd{endCeiling, run.failure} case phaseSendBudget: return &chainEnd{endSendBudget, run.failure} + case phaseStopped: + // Unreachable outside a group; never let it read as a success. + return &chainEnd{endFailed, run.failure} } return nil } @@ -684,6 +796,8 @@ func terminalSend(kind sendEnd, detail string) *phaseRun { return &phaseRun{outcome: phaseCeiling} case sendBudget: return &phaseRun{outcome: phaseSendBudget, failure: detail} + case sendStopped: + return &phaseRun{outcome: phaseStopped, failure: detail} } return nil } @@ -702,6 +816,8 @@ func (run phaseRun) agentOutcome() string { return protocol.AgentCeiling case phaseSendBudget: return protocol.AgentSendBudget + case phaseStopped: + return protocol.AgentStopped } return protocol.AgentFailed } @@ -714,10 +830,7 @@ func (run phaseRun) envelopeRef() *parsedEnvelope { return &envelope } -func (e *execution) nextEntry(phase string) int { - e.phaseEntries[phase]++ - return e.phaseEntries[phase] -} +func (e *execution) nextEntry(phase string) int { return e.counters.nextEntry(phase) } func (e *execution) recordResult(result protocol.PhaseResult) { now := time.Now().UTC() @@ -818,7 +931,14 @@ func (e *execution) mergeFields(fields map[string]any) { // mergeAgentFields merges an agent envelope into the field view, except for // fields a reports_fields code phase already set: those keep the reported // value, and the attempted overwrite is traced. +// +// A group member's view defers the merge: its fields join the chain's view +// at the join, in declared order, through this same function. func (e *execution) mergeAgentFields(phase, role string, fields map[string]any) { + if e.grouped { + e.pendingFields = append(e.pendingFields, fieldMerge{phase: phase, role: role, fields: fields}) + return + } for key, value := range fields { if reporter, reported := e.reportedBy[key]; reported { if !reflect.DeepEqual(e.fieldView[key], value) { @@ -972,6 +1092,11 @@ func (e *execution) enforceWriteBoundary( // the nested parse/gate correction loops, boundary enforcement, envelope // persistence. died=true means the subprocess was killed or crashed and the // worktree has been rolled back to the pre-phase snapshot. +// +// In a group member's view (e.grouped) the group owns the worktree: one +// snapshot before every member, one enforcement after, so this entry takes +// no snapshot, enforces nothing, and rolls nothing back on death — a dead +// member's leftover write is the group's breach to find. func (e *execution) runAgentPhaseAttempt( ctx context.Context, phase protocol.PhaseSpec, previous *parsedEnvelope, ) (run phaseRun, died bool) { @@ -982,9 +1107,12 @@ func (e *execution) runAgentPhaseAttempt( "kind": phase.Kind, "phase_attempt": entry, "owner": phase.Owner, }) - before, err := snapshotTree(ctx, e.attempt.WorktreePath) - if err != nil { - return e.failInfra(phase.Name, err.Error()), false + var before treeSnapshot + if !e.grouped { + var err error + if before, err = snapshotTree(ctx, e.attempt.WorktreePath); err != nil { + return e.failInfra(phase.Name, err.Error()), false + } } if err := e.seedRole(phase.Owner, role); err != nil { return e.failInfra(phase.Name, "seed ephemeral HOME: "+err.Error()), false @@ -1041,9 +1169,11 @@ func (e *execution) runAgentPhaseAttempt( e.emit.emit(protocol.EventPhaseDeath, phase.Name, phase.Owner, map[string]any{ "phase_attempt": entry, "error": detail, }) - if rollbackErr := restoreSnapshot(ctx, e.attempt.WorktreePath, before); rollbackErr != nil { - e.emit.emit(protocol.EventError, phase.Name, "rollback_failed", - map[string]string{"error": rollbackErr.Error()}) + if !e.grouped { + if rollbackErr := restoreSnapshot(ctx, e.attempt.WorktreePath, before); rollbackErr != nil { + e.emit.emit(protocol.EventError, phase.Name, "rollback_failed", + map[string]string{"error": rollbackErr.Error()}) + } } e.recordResult(protocol.PhaseResult{ Phase: phase.Name, Kind: phase.Kind, Status: protocol.EnvelopeFail, @@ -1057,11 +1187,13 @@ func (e *execution) runAgentPhaseAttempt( // retained worktree and the next phase — a failed phase is not a phase // that wrote nothing (R10). fail := func(detail string) (phaseRun, bool) { - if _, terminal := e.enforceWriteBoundary( - ctx, phase, entry, started, before, role.Writes); terminal != nil { - terminal.failure = detail + " — and " + terminal.failure - agentEnd(terminal.agentOutcome()) - return *terminal, false + if !e.grouped { + if _, terminal := e.enforceWriteBoundary( + ctx, phase, entry, started, before, role.Writes); terminal != nil { + terminal.failure = detail + " — and " + terminal.failure + agentEnd(terminal.agentOutcome()) + return *terminal, false + } } agentEnd(protocol.AgentFailed) e.recordResult(protocol.PhaseResult{ @@ -1085,7 +1217,7 @@ func (e *execution) runAgentPhaseAttempt( // authorized work the operator retained the worktree to look at, while // discarding nothing that could carry forward. terminalExit := func(run phaseRun, detail string) (phaseRun, bool) { - if sender.sendCount == 0 { + if sender.sendCount == 0 || e.grouped { // The budget check refuses before the subprocess starts: no agent // ran in this entry, so there is nothing of its to enforce. agentEnd(run.agentOutcome()) @@ -1164,17 +1296,19 @@ func (e *execution) runAgentPhaseAttempt( // Permission is checked after every send is done and before the envelope // is accepted: an agent does not get to report success on a phase in // which it wrote somewhere it was not allowed to (R10). - touched, terminal := e.enforceWriteBoundary(ctx, phase, entry, started, before, role.Writes) - if terminal != nil { - agentEnd(terminal.agentOutcome()) - return *terminal, false - } - for _, path := range touched { - e.touchedPaths[path] = true - } - if len(touched) > 0 { - e.emit.emit(protocol.EventLog, phase.Name, "paths_touched", - map[string]any{"role": phase.Owner, "paths": touched}) + if !e.grouped { + touched, terminal := e.enforceWriteBoundary(ctx, phase, entry, started, before, role.Writes) + if terminal != nil { + agentEnd(terminal.agentOutcome()) + return *terminal, false + } + for _, path := range touched { + e.touchedPaths[path] = true + } + if len(touched) > 0 { + e.emit.emit(protocol.EventLog, phase.Name, "paths_touched", + map[string]any{"role": phase.Owner, "paths": touched}) + } } e.mergeAgentFields(phase.Name, phase.Owner, envelope.Fields) @@ -1259,6 +1393,9 @@ const ( // edge against a rate limit is a retry storm with the wrong diagnosis, so // this ends the PHASE with no parse ladder at all. sendRuntimeError + // sendStopped: the parallel group's stop signal fired — a sibling ended + // the attempt — so this send never started, or was killed. + sendStopped ) // agentSender owns one phase attempt's sends: session continuity (or its @@ -1279,14 +1416,21 @@ type agentSender struct { } func (s *agentSender) send(ctx context.Context, prompt string) (runtime.Result, sendEnd, string) { + // A member whose group has been told to stop starts nothing: the sibling + // that stopped it already ended the attempt. + select { + case <-s.stop: + return runtime.Result{}, sendStopped, "stopped: a parallel group sibling ended the attempt" + default: + } // Checked BEFORE the subprocess starts: the send that would exceed the - // bound is the one that must not happen. - if s.execution.sends >= s.runner.config.MaxAttemptSends { + // bound is the one that must not happen. One atomic check-and-take, so + // group members racing for the last send cannot both have it. + if !s.counters.acquireSend(s.runner.config.MaxAttemptSends) { return runtime.Result{}, sendBudget, fmt.Sprintf( "attempt send budget (%d) exhausted in phase %q", s.runner.config.MaxAttemptSends, s.phase) } - s.execution.sends++ session := s.sessionFor(s.role) if !s.capability.CanResume && session.Sends > 0 { @@ -1338,6 +1482,10 @@ func (s *agentSender) send(ctx context.Context, prompt string) (runtime.Result, case <-s.attempt.Cancelled: s.abandon(handle) return runtime.Result{}, sendCancelled, "cancelled during phase " + s.phase + case <-s.stop: + s.abandon(handle) + return runtime.Result{}, sendStopped, "stopped during phase " + s.phase + + ": a parallel group sibling ended the attempt" } } result, err := handle.Result() @@ -1422,15 +1570,20 @@ func (e *execution) sessionFor(role string) *runtime.Session { if session := e.sessions[role]; session != nil { return session } + if e.grouped { + // A member's role may also run outside the group, in the chain's + // HOME under the plain key; the member's session lives in its own + // HOME and takes a sequenced key no other session can hold. + return e.freshSession(role) + } session := &runtime.Session{Key: e.sessionKey + "-" + role} e.sessions[role] = session return session } func (e *execution) freshSession(role string) *runtime.Session { - e.sessionSeq++ session := &runtime.Session{ - Key: fmt.Sprintf("%s-%s-r%d", e.sessionKey, role, e.sessionSeq), + Key: fmt.Sprintf("%s-%s-r%d", e.sessionKey, role, e.counters.sessionSeq.Add(1)), } e.sessions[role] = session return session diff --git a/internal/engine/resume.go b/internal/engine/resume.go index 4ac9a72..cd846b2 100644 --- a/internal/engine/resume.go +++ b/internal/engine/resume.go @@ -39,8 +39,16 @@ func (e *execution) repairable(outcome worker.Outcome) bool { // (an agent can make one fail, with a read-only directory) is an error the // caller must act on: a HOME that may still hold what an agent left there // is never handed to another round. +// +// Every parallel group member's HOME goes too, and with it the member views +// whose sessions lived there: the next group run starts them afresh. func (e *execution) wipeHome() error { e.seededRoles = make(map[string]bool) + e.members = nil + if err := os.RemoveAll(e.scratch.members); err != nil { + e.runner.config.Logger.Warn("attempt_member_homes_wipe_failed", "error", err) + return fmt.Errorf("wipe the parallel group members' ephemeral HOMEs: %w", err) + } if err := os.RemoveAll(e.scratch.home); err != nil { e.runner.config.Logger.Warn("attempt_home_wipe_failed", "error", err) return fmt.Errorf("wipe the ephemeral HOME: %w", err) diff --git a/internal/protocol/definition.go b/internal/protocol/definition.go index d749a06..c5ec808 100644 --- a/internal/protocol/definition.go +++ b/internal/protocol/definition.go @@ -63,10 +63,63 @@ type DefinitionSpec struct { Name string `yaml:"name"` Roster map[string]RoleSpec `yaml:"roster"` Phases []PhaseSpec `yaml:"phases"` + Parallel ParallelGroup `yaml:"parallel"` Acceptance []string `yaml:"acceptance"` Publish *PublishSpec `yaml:"publish"` } +// ParallelGroup is the one opt-in parallel construct: the names of two or +// more consecutive agent phases, in chain order, that run concurrently as +// one step of the chain. Every member is a read-only role (`writes: []`) +// with its own owner, and no member's `if:` guard reads a field a sibling +// reports, because every member is judged against the envelope and field +// view that preceded the group. Members' results merge in declared order, +// and the phase after the group receives the last member's envelope. +// +// It reopens v1's "no parallel phases" non-goal (V1:KTD2) for exactly this +// shape and no other: reviewers that read the same tree and report +// verdicts. One group per definition. +type ParallelGroup []string + +// UnmarshalYAML accepts a flat list of phase names. A list of lists is a +// second group, which is refused by name rather than as a type error. +func (group *ParallelGroup) UnmarshalYAML(node *yaml.Node) error { + if node.Kind != yaml.SequenceNode { + return fmt.Errorf("parallel: line %d: must be a list of phase names", node.Line) + } + names := make(ParallelGroup, 0, len(node.Content)) + for _, item := range node.Content { + if item.Kind == yaml.SequenceNode { + return fmt.Errorf("parallel: line %d: only one parallel group per definition — "+ + "declare it as a flat list of phase names", item.Line) + } + if item.Kind != yaml.ScalarNode { + return fmt.Errorf("parallel: line %d: a member must be a phase name", item.Line) + } + names = append(names, item.Value) + } + *group = names + return nil +} + +// ParallelRange locates the group in the chain: Phases[start:end] are its +// members. ok is false when the definition declares no group. Call it only +// on a spec that passed Validate. +func (spec *DefinitionSpec) ParallelRange() (start, end int, ok bool) { + if len(spec.Parallel) == 0 { + return 0, 0, false + } + for i, phase := range spec.Phases { + if phase.Name == spec.Parallel[0] { + return i, i + len(spec.Parallel), true + } + } + return 0, 0, false +} + +// envelopeBaseFields are reported by every agent phase: the base contract. +var envelopeBaseFields = []string{"status", "summary", "artifacts", "notes_for_next_agent"} + // PublishSpec is the definition's say over delivery. HoldWhen is a declared // envelope predicate (the `on_fail.when` language) evaluated against the // chain's merged field view once acceptance has passed: when it holds, the @@ -289,9 +342,97 @@ func (spec *DefinitionSpec) Validate() error { if err := spec.validateAcceptance(); err != nil { return err } + if err := spec.validateParallel(phasesByName); err != nil { + return err + } return spec.validatePublish(phasesByName) } +// validateParallel enforces R8. What a member "reports" is what the +// definition can see it report: the envelope base fields every agent phase +// emits and the field its own repair edge's `when` reads. An agent envelope +// may carry more, which no save-time check can know. +func (spec *DefinitionSpec) validateParallel(phases map[string]PhaseSpec) error { + group := spec.Parallel + if group == nil { + return nil + } + if len(group) < 2 { + return fmt.Errorf("parallel: a group needs two or more phases, got %d", len(group)) + } + owners := make(map[string]string, len(group)) + seen := make(map[string]bool, len(group)) + for _, name := range group { + phase, defined := phases[name] + if !defined { + return fmt.Errorf("parallel: member %q is not a phase in the chain", name) + } + if seen[name] { + return fmt.Errorf("parallel: member %q is listed twice", name) + } + seen[name] = true + if phase.Kind != PhaseKindAgent { + return fmt.Errorf("parallel: member %q is a %s phase; only agent phases run in a group", + name, phase.Kind) + } + if writes := spec.Roster[phase.Owner].Writes; writes == nil || len(writes) > 0 { + return fmt.Errorf("parallel: member %q runs role %q, which may write; "+ + "every member's role must declare writes: []", name, phase.Owner) + } + if other, shared := owners[phase.Owner]; shared { + return fmt.Errorf("parallel: members %q and %q share owner role %q; "+ + "every member needs its own role", other, name, phase.Owner) + } + owners[phase.Owner] = name + } + start, _, _ := spec.ParallelRange() + for i, name := range group { + if i > 0 && (start+i >= len(spec.Phases) || spec.Phases[start+i].Name != name) { + return fmt.Errorf("parallel: members must be consecutive phases listed in chain order; "+ + "%q does not directly follow %q in the chain", name, group[i-1]) + } + } + for _, name := range group { + guard := strings.TrimSpace(phases[name].If) + if guard == "" { + continue + } + field := guard + if IsPredicate(guard) { + predicate, _ := ParsePredicate(guard) + field = predicate.Field + } + for _, sibling := range group { + if sibling != name && reportsField(phases[sibling], field) { + return fmt.Errorf("parallel: member %q's if: guard reads %q, which member %q reports; "+ + "every member is judged before any member runs", name, field, sibling) + } + } + } + if spec.Publish != nil && spec.Publish.CI != nil && spec.Publish.CI.OnFail != nil && + seen[spec.Publish.CI.OnFail.Run] { + return fmt.Errorf("publish: ci: on_fail: run phase %q is a parallel group member; "+ + "a CI repair round cannot start inside the group", spec.Publish.CI.OnFail.Run) + } + return nil +} + +// reportsField reports whether an agent phase declares that it reports the +// field: a base envelope field, or the field its repair edge keys on. +func reportsField(phase PhaseSpec, field string) bool { + for _, base := range envelopeBaseFields { + if field == base { + return true + } + } + if phase.OnFail != nil && phase.OnFail.When != "" { + if predicate, err := ParsePredicate(phase.OnFail.When); err == nil && predicate.Field == field { + return true + } + } + return false +} + func (spec *DefinitionSpec) validatePublish(phases map[string]PhaseSpec) error { if spec.Publish == nil { return nil diff --git a/internal/protocol/definition_test.go b/internal/protocol/definition_test.go index 9b912ab..cabaf07 100644 --- a/internal/protocol/definition_test.go +++ b/internal/protocol/definition_test.go @@ -544,3 +544,97 @@ phases: - {name: build, kind: agent, owner: builder} `, `"builder"`, "budget_usd must not be negative") } + +// parallelPanel is a builder, three read-only reviewers, and a code check; +// the tests below splice roster entries, phases, and group lines into it. +func parallelPanel(extraRoster, phases, tail string) string { + return ` +name: panel +roster: + builder: {model: opus, system_prompt: s, user_prompt: u, writes: ["src/"]} + a: {model: opus, system_prompt: s, user_prompt: u, writes: []} + b: {model: opus, system_prompt: s, user_prompt: u, writes: []} + c: {model: opus, system_prompt: s, user_prompt: u, writes: []} +` + extraRoster + ` +phases: +` + phases + tail +} + +const panelPhases = ` + - {name: build, kind: agent, owner: builder} + - {name: review-a, kind: agent, owner: a, on_fail: {when: "approved == false", run: build, then: rerun-self, budget: 1}} + - {name: review-b, kind: agent, owner: b} + - {name: review-c, kind: agent, owner: c} + - {name: check, kind: code, command: "true"} +` + +func TestAParallelGroupOfConsecutiveReadOnlyReviewersValidates(t *testing.T) { + spec := mustParse(t, parallelPanel("", panelPhases, "parallel: [review-a, review-b, review-c]\n")) + start, end, ok := spec.ParallelRange() + if !ok || start != 1 || end != 4 { + t.Fatalf("ParallelRange = %d, %d, %v, want 1, 4, true", start, end, ok) + } + if _, _, ok := mustParse(t, parallelPanel("", panelPhases, "")).ParallelRange(); ok { + t.Fatal("a definition without a group reports one") + } +} + +// R8: every rule that keeps a group's members concurrent-safe is enforced at +// save time, naming what broke it. +func TestAParallelGroupIsRejectedUnlessItsMembersAreConcurrentSafe(t *testing.T) { + unrestricted := " d: {model: opus, system_prompt: s, user_prompt: u}\n" + cases := []struct { + name string + source string + want []string + }{ + {"a writing member", parallelPanel("", panelPhases, "parallel: [build, review-a]\n"), + []string{`"build"`, `"builder"`, "writes: []"}}, + {"a member whose role omits writes, which is unrestricted", parallelPanel(unrestricted, + panelPhases+" - {name: review-d, kind: agent, owner: d}\n - {name: review-e, kind: agent, owner: a}\n", + "parallel: [review-d, review-e]\n"), + []string{`"review-d"`, `"d"`, "writes: []"}}, + {"a code-phase member", parallelPanel("", panelPhases, "parallel: [review-c, check]\n"), + []string{`"check"`, "code phase"}}, + {"two members sharing a role", parallelPanel("", + panelPhases+" - {name: review-a2, kind: agent, owner: a}\n - {name: review-a3, kind: agent, owner: a}\n", + "parallel: [review-a2, review-a3]\n"), + []string{`"review-a2"`, `"review-a3"`, `"a"`, "own role"}}, + {"non-consecutive members", parallelPanel("", panelPhases, "parallel: [review-a, review-c]\n"), + []string{"consecutive", `"review-c"`, `"review-a"`}}, + {"members out of chain order", parallelPanel("", panelPhases, "parallel: [review-b, review-a]\n"), + []string{"consecutive"}}, + {"a member guarded on a sibling's repair field", parallelPanel("", + strings.Replace(panelPhases, "owner: b}", `owner: b, if: "approved == true"}`, 1), + "parallel: [review-a, review-b]\n"), + []string{`"review-b"`, `"approved"`, `"review-a"`}}, + {"a member guarded on a sibling's base envelope field", parallelPanel("", + strings.Replace(panelPhases, "owner: c}", "owner: c, if: summary}", 1), + "parallel: [review-b, review-c]\n"), + []string{`"review-c"`, `"summary"`}}, + {"a second group", parallelPanel("", panelPhases, "parallel: [[review-a, review-b], [review-c]]\n"), + []string{"only one parallel group"}}, + {"a second group key", parallelPanel("", panelPhases, + "parallel: [review-a, review-b]\nparallel: [review-c]\n"), + []string{"parallel"}}, + {"a group of one", parallelPanel("", panelPhases, "parallel: [review-a]\n"), + []string{"two or more"}}, + {"an undefined member", parallelPanel("", panelPhases, "parallel: [review-a, review-z]\n"), + []string{`"review-z"`}}, + {"a member listed twice", parallelPanel("", panelPhases, "parallel: [review-a, review-a]\n"), + []string{"twice"}}, + {"CI repair starting inside the group", parallelPanel("", panelPhases, + "parallel: [review-a, review-b]\npublish: {ci: {wait: true, on_fail: {run: review-a, budget: 1}}}\n"), + []string{"on_fail", `"review-a"`, "parallel group member"}}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + mustReject(t, tc.source, tc.want...) + }) + } + // A guard on a field no sibling declares is fine: it is judged against + // the view that preceded the group, exactly as the definition says. + mustParse(t, parallelPanel("", + strings.Replace(panelPhases, "owner: c}", `owner: c, if: "risk != low"}`, 1), + "parallel: [review-a, review-b, review-c]\n")) +} diff --git a/internal/protocol/types.go b/internal/protocol/types.go index 8be8c79..4e3ea4e 100644 --- a/internal/protocol/types.go +++ b/internal/protocol/types.go @@ -91,6 +91,9 @@ const ( AgentCancelled = "cancelled" AgentCeiling = "ceiling" AgentSendBudget = "send_budget" + // AgentStopped: a parallel group member stopped because a sibling ended + // the attempt. + AgentStopped = "stopped" ) // Envelope statuses. Success must be earned: everything that constructs a diff --git a/internal/worker/manifest.go b/internal/worker/manifest.go index ac0892e..383d2d3 100644 --- a/internal/worker/manifest.go +++ b/internal/worker/manifest.go @@ -71,8 +71,16 @@ type attemptManifest struct { WorktreePath string `json:"worktree_path"` Branch string `json:"branch"` - // ProcessGroupID and ProcessActive are recorded by U4's supervisor; - // U3 keeps them zero. + // ProcessGroups is the SET of agent process groups live in this attempt + // right now, written only while more than one is live. A parallel + // reviewer group runs several agent subprocesses at once, each in its + // own group, and start-time reconciliation must stop every one a crashed + // worker left behind. + ProcessGroups []int64 `json:"process_groups,omitempty"` + // ProcessGroupID and ProcessActive are the single-group record older jig + // versions read and wrote. They still name the lowest live group + // (writeProcessGroups), so a worker rolled back to an older binary stops + // at least that one. ProcessGroupID int64 `json:"process_group_id,omitempty"` ProcessActive bool `json:"process_active"` @@ -452,6 +460,16 @@ func (store *manifestStore) validate(manifest attemptManifest) error { if manifest.ProcessActive && !signallableProcessGroup(manifest.ProcessGroupID) { return errors.New("attempt manifest advertises a live process without a signallable group id") } + seenGroups := make(map[int64]bool, len(manifest.ProcessGroups)) + for _, groupID := range manifest.ProcessGroups { + if !signallableProcessGroup(groupID) { + return fmt.Errorf("attempt manifest process group %d can never name a real group", groupID) + } + if seenGroups[groupID] { + return fmt.Errorf("attempt manifest lists process group %d twice", groupID) + } + seenGroups[groupID] = true + } if !manifestLifecycles[manifest.Lifecycle] { return fmt.Errorf("attempt manifest lifecycle %q is invalid", manifest.Lifecycle) } diff --git a/internal/worker/reconcile.go b/internal/worker/reconcile.go index 37a912e..d7032ac 100644 --- a/internal/worker/reconcile.go +++ b/internal/worker/reconcile.go @@ -219,28 +219,32 @@ func (w *Worker) Reconcile(ctx context.Context) (ReconcileReport, error) { // ---- process-group reconciliation (U4's read side) ------------------------- // stopRecordedProcessGroups stops every agent process group a previous worker -// process recorded as live. Each candidate is identity-checked first: a -// recorded group id is just a number, and by the time we read it the pid may -// belong to someone else's shell — or to nobody at all, if the field was -// zeroed or corrupted. Only a group id that can name a real group AND whose -// leader still leads that exact group AND that started inside the manifest's -// own lifetime is ours to signal. Either way the flag is cleared, so one +// process recorded as live — every group in an attempt's set, since a +// parallel reviewer group runs several at once, plus the single group an +// older jig recorded. Each candidate is identity-checked first: a recorded +// group id is just a number, and by the time we read it the pid may belong +// to someone else's shell — or to nobody at all, if the field was zeroed or +// corrupted. Only a group id that can name a real group AND whose leader +// still leads that exact group AND that started inside the manifest's own +// lifetime is ours to signal. Either way the record is cleared, so one // unverifiable manifest cannot make every later reconciliation re-examine it // forever. func (w *Worker) stopRecordedProcessGroups(ctx context.Context, manifests []attemptManifest) []int64 { var stopped []int64 for _, manifest := range manifests { - if !manifest.ProcessActive { + groups := recordedProcessGroups(manifest) + if len(groups) == 0 { continue } - groupID := manifest.ProcessGroupID - ours, reason := processGroupIsOurs(ctx, manifest) - if ours { - stopProcessGroup(groupID) - stopped = append(stopped, groupID) - w.logger.Info("orphan_process_group_stopped", - "attempt_id", manifest.AttemptID, "process_group_id", groupID) - } else { + for _, groupID := range groups { + ours, reason := processGroupIsOurs(ctx, manifest, groupID) + if ours { + stopProcessGroup(groupID) + stopped = append(stopped, groupID) + w.logger.Info("orphan_process_group_stopped", + "attempt_id", manifest.AttemptID, "process_group_id", groupID) + continue + } // Not provably ours: never signalled. A recycled pid belongs to // someone else, and this is the line where that is decided. w.logger.Info("orphan_process_group_skipped", @@ -248,6 +252,7 @@ func (w *Worker) stopRecordedProcessGroups(ctx context.Context, manifests []atte } if _, err := w.manifests.update(manifest.AttemptID, func(value *attemptManifest) error { value.ProcessActive = false + value.ProcessGroups = nil return nil }); err != nil { w.logger.Warn("process_group_clear_failed", @@ -257,32 +262,45 @@ func (w *Worker) stopRecordedProcessGroups(ctx context.Context, manifests []atte return stopped } -// processGroupIsOurs answers whether the recorded group's leader is still the -// process this manifest recorded, and says why when it is not. -func processGroupIsOurs(ctx context.Context, manifest attemptManifest) (bool, string) { - if !signallableProcessGroup(manifest.ProcessGroupID) { +// recordedProcessGroups is every group a manifest records as live: its set, +// and the legacy single group when that is flagged active. +func recordedProcessGroups(manifest attemptManifest) []int64 { + groups := append([]int64(nil), manifest.ProcessGroups...) + if manifest.ProcessActive { + for _, groupID := range groups { + if groupID == manifest.ProcessGroupID { + return groups + } + } + groups = append(groups, manifest.ProcessGroupID) + } + return groups +} + +// processGroupIsOurs answers whether a group the manifest recorded is still +// led by the process this attempt started, and says why when it is not. +func processGroupIsOurs(ctx context.Context, manifest attemptManifest, recorded int64) (bool, string) { + if !signallableProcessGroup(recorded) { // -1 signals every process the operator's user may signal, 0 signals // jig's own group, 1 is init. None can be an attempt's agent group, so // there is nothing here to verify and nothing to signal. return false, fmt.Sprintf( - "recorded process group %d can never name an attempt's own group", - manifest.ProcessGroupID) + "recorded process group %d can never name an attempt's own group", recorded) } - groupID, started, found, err := inspectProcessGroupLeader(ctx, manifest.ProcessGroupID) + groupID, started, found, err := inspectProcessGroupLeader(ctx, recorded) switch { case err != nil: return false, "process identity could not be read: " + err.Error() case !found: return false, "the recorded group leader no longer exists" - case groupID != manifest.ProcessGroupID: - return false, fmt.Sprintf("pid %d now leads group %d, not %d", - manifest.ProcessGroupID, groupID, manifest.ProcessGroupID) + case groupID != recorded: + return false, fmt.Sprintf("pid %d now leads group %d, not %d", recorded, groupID, recorded) } earliest := manifest.CreatedAt.Add(-processIdentitySlack) latest := manifest.UpdatedAt.Add(processIdentitySlack) if started.Before(earliest) || started.After(latest) { return false, fmt.Sprintf("pid %d started at %s, outside this attempt's lifetime (%s..%s) — recycled", - manifest.ProcessGroupID, started.UTC().Format(time.RFC3339), + recorded, started.UTC().Format(time.RFC3339), earliest.UTC().Format(time.RFC3339), latest.UTC().Format(time.RFC3339)) } return true, "" diff --git a/internal/worker/reconcile_test.go b/internal/worker/reconcile_test.go index 0bfef1b..1278de1 100644 --- a/internal/worker/reconcile_test.go +++ b/internal/worker/reconcile_test.go @@ -105,11 +105,174 @@ func TestReconcileStopsTheProcessGroupsACrashedWorkerLeftRunning(t *testing.T) { if err != nil { t.Fatal(err) } - if manifest.ProcessActive { + if manifest.ProcessActive || len(manifest.ProcessGroups) != 0 { t.Fatal("the manifest still advertises a live process group after it was stopped") } } +// A parallel reviewer group runs several agents at once. The manifest keeps +// the SET of live groups — recording one never overwrites another, and +// clearing one leaves the rest — and reconciliation stops every group a +// crashed worker left in it. +func TestReconcileStopsEveryLiveGroupOfAParallelGroup(t *testing.T) { + h := newHarness(t) + dataDir := filepath.Join(t.TempDir(), "worker") + w, attemptID := claimOneAttempt(t, h, dataDir) + + first, second, finished := startGroupLeader(t), startGroupLeader(t), startGroupLeader(t) + for _, groupID := range []int{first, finished, second} { + if err := w.RecordProcessGroup(attemptID, int64(groupID), true); err != nil { + t.Fatalf("record process group %d: %v", groupID, err) + } + } + if err := w.RecordProcessGroup(attemptID, int64(finished), false); err != nil { + t.Fatalf("clear process group: %v", err) + } + manifest, err := w.manifests.load(attemptID) + if err != nil { + t.Fatal(err) + } + want := []int64{int64(first), int64(second)} + if want[0] > want[1] { + want[0], want[1] = want[1], want[0] + } + if len(manifest.ProcessGroups) != 2 || manifest.ProcessGroups[0] != want[0] || manifest.ProcessGroups[1] != want[1] { + t.Fatalf("recorded groups = %v, want exactly the two still live %v", manifest.ProcessGroups, want) + } + + restarted := newTestWorker(t, h, dataDir, 1, RunnerFunc( + func(context.Context, *PreparedAttempt) Outcome { + return Outcome{State: protocol.AttemptFailed, Error: "unused"} + })) + report, err := restarted.Reconcile(context.Background()) + if err != nil { + t.Fatalf("reconcile: %v", err) + } + if len(report.StoppedProcessGroups) != 2 { + t.Fatalf("reconcile stopped %v, want both live members' groups %v", report.StoppedProcessGroups, want) + } + waitForProcessExit(t, first, 15*time.Second) + waitForProcessExit(t, second, 15*time.Second) + if err := syscall.Kill(-finished, 0); err != nil { + t.Fatalf("reconcile stopped a group the manifest had already cleared: %v", err) + } + after, err := restarted.manifests.load(attemptID) + if err != nil { + t.Fatal(err) + } + if len(after.ProcessGroups) != 0 { + t.Fatalf("the manifest still records %v after reconciliation", after.ProcessGroups) + } +} + +// preSetManifest is the manifest shape jig wrote before the process-group +// set, decoded as that version did: strictly, unknown fields refused. +type preSetManifest struct { + SchemaVersion int `json:"schema_version"` + WorkerID string `json:"worker_id"` + JobID string `json:"job_id"` + AttemptID string `json:"attempt_id"` + AttemptNumber int `json:"attempt_number"` + Repository string `json:"repository"` + RepositoryDir string `json:"repository_dir"` + BaseSHA string `json:"base_sha"` + WorktreePath string `json:"worktree_path"` + Branch string `json:"branch"` + ProcessGroupID int64 `json:"process_group_id,omitempty"` + ProcessActive bool `json:"process_active"` + Lifecycle string `json:"lifecycle"` + TerminalState string `json:"terminal_state,omitempty"` + RetentionReason string `json:"retention_reason,omitempty"` + CleanupIntent string `json:"cleanup_intent,omitempty"` + CleanupResult string `json:"cleanup_result,omitempty"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +// A worker rolled back to a jig that predates the set still sees a live +// group: the legacy fields always name the lowest live one. With one group +// live the manifest is exactly the old shape; with several, the legacy +// fields still name the lowest. +func TestAnOlderReaderStillSeesALiveProcessGroup(t *testing.T) { + h := newHarness(t) + dataDir := filepath.Join(t.TempDir(), "worker") + w, attemptID := claimOneAttempt(t, h, dataDir) + path := filepath.Join(dataDir, "attempts", attemptID+".json") + readOld := func() (preSetManifest, error) { + body, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + var old preSetManifest + return old, decodeStrictJSON(body, &old) + } + + if err := w.RecordProcessGroup(attemptID, 5100, true); err != nil { + t.Fatal(err) + } + old, err := readOld() + if err != nil { + t.Fatalf("an older reader cannot read a one-group manifest: %v", err) + } + if !old.ProcessActive || old.ProcessGroupID != 5100 { + t.Fatalf("older reader sees active=%v group=%d, want the live group 5100", old.ProcessActive, old.ProcessGroupID) + } + + for _, groupID := range []int64{5300, 5200} { + if err := w.RecordProcessGroup(attemptID, groupID, true); err != nil { + t.Fatal(err) + } + } + if err := w.RecordProcessGroup(attemptID, 5100, false); err != nil { + t.Fatal(err) + } + current, err := w.manifests.load(attemptID) + if err != nil { + t.Fatal(err) + } + if !current.ProcessActive || current.ProcessGroupID != 5200 { + t.Fatalf("legacy fields = active %v group %d, want the lowest live group 5200", + current.ProcessActive, current.ProcessGroupID) + } + + if err := w.RecordProcessGroup(attemptID, 5200, false); err != nil { + t.Fatal(err) + } + if old, err = readOld(); err != nil || !old.ProcessActive || old.ProcessGroupID != 5300 { + t.Fatalf("older reader after the set shrank to one: %+v, %v; want active group 5300", old, err) + } + if err := w.RecordProcessGroup(attemptID, 5300, false); err != nil { + t.Fatal(err) + } + if old, err = readOld(); err != nil || old.ProcessActive { + t.Fatalf("older reader after every group ended: %+v, %v; want nothing live", old, err) + } +} + +// Concurrent records — members starting at once — all land. +func TestConcurrentProcessGroupRecordsAllLand(t *testing.T) { + h := newHarness(t) + dataDir := filepath.Join(t.TempDir(), "worker") + w, attemptID := claimOneAttempt(t, h, dataDir) + const members = 8 + errs := make(chan error, members) + for i := 0; i < members; i++ { + go func(groupID int64) { errs <- w.RecordProcessGroup(attemptID, groupID, true) }(int64(1000 + i)) + } + for i := 0; i < members; i++ { + if err := <-errs; err != nil { + t.Fatalf("record: %v", err) + } + } + manifest, err := w.manifests.load(attemptID) + if err != nil { + t.Fatal(err) + } + if len(manifest.ProcessGroups) != members { + t.Fatalf("recorded %v, want all %d concurrent records", manifest.ProcessGroups, members) + } +} + func TestReconcileNeverSignalsARecordedGroupWhosePidNoLongerMatches(t *testing.T) { h := newHarness(t) dataDir := filepath.Join(t.TempDir(), "worker") @@ -174,7 +337,7 @@ func TestAProcessGroupIdAtOrBelowOneIsNeverSignalled(t *testing.T) { "stopProcessGroup would negate it into a machine-wide kill", groupID) } manifest := attemptManifest{ProcessGroupID: groupID, ProcessActive: true} - ours, reason := processGroupIsOurs(context.Background(), manifest) + ours, reason := processGroupIsOurs(context.Background(), manifest, groupID) if ours { t.Errorf("process group %d was claimed as ours", groupID) } @@ -264,6 +427,19 @@ func TestAManifestCannotCarryAnUnsignallableProcessGroup(t *testing.T) { if err := store.validate(live); err != nil { t.Errorf("a real live process group must stay valid: %v", err) } + // The set is held to the same bar, and never lists a group twice. + for _, groups := range [][]int64{{4242, 1}, {0}, {-1}, {4242, 4242}} { + manifest := base + manifest.ProcessGroups = groups + if err := store.validate(manifest); err == nil { + t.Errorf("a manifest recording process groups %v validated", groups) + } + } + set := base + set.ProcessGroups = []int64{4242, 4343} + if err := store.validate(set); err != nil { + t.Errorf("a set of real live process groups must stay valid: %v", err) + } } // End to end: a manifest that reaches disk carrying process_group_id 1 — diff --git a/internal/worker/registration.go b/internal/worker/registration.go index 5d570a3..5994756 100644 --- a/internal/worker/registration.go +++ b/internal/worker/registration.go @@ -179,18 +179,51 @@ func (noEngineRunner) Run(context.Context, *PreparedAttempt) Outcome { // RecordProcessGroup records an attempt's live subprocess group into its // manifest (U4). The phase engine reports each agent process group as it -// starts and stops — the manifest's ProcessGroupID/ProcessActive fields are -// what start-time reconciliation uses to stop orphaned groups a crashed -// worker left behind. +// starts and stops; the manifest keeps the SET of groups live at once — a +// parallel reviewer group runs several — and start-time reconciliation stops +// every one a crashed worker left behind. Calls may arrive concurrently; the +// manifest store serializes them. func (w *Worker) RecordProcessGroup(attemptID string, processGroupID int64, active bool) error { + if active && !signallableProcessGroup(processGroupID) { + return fmt.Errorf("process group %d can never name a real group", processGroupID) + } _, err := w.manifests.update(attemptID, func(manifest *attemptManifest) error { - manifest.ProcessGroupID = processGroupID - manifest.ProcessActive = active + groups := make([]int64, 0, len(manifest.ProcessGroups)+2) + for _, groupID := range recordedProcessGroups(*manifest) { + if groupID != processGroupID { + groups = append(groups, groupID) + } + } + if active { + groups = append(groups, processGroupID) + } + sort.Slice(groups, func(i, j int) bool { return groups[i] < groups[j] }) + writeProcessGroups(manifest, groups) return nil }) return err } +// writeProcessGroups stores the live set so a downgraded worker still sees +// what it can. The legacy single-group fields always name the lowest live +// group, so an older jig's reconciliation stops at least that one; the full +// set is written only when more than one group is live. One live group +// therefore leaves a manifest an older binary reads exactly as before, and +// only the rare crash mid-group leaves one it refuses to read — and an +// older jig fails closed on an unreadable manifest, retaining the worktree +// rather than guessing. +func writeProcessGroups(manifest *attemptManifest, groups []int64) { + manifest.ProcessActive = len(groups) > 0 + manifest.ProcessGroupID = 0 + manifest.ProcessGroups = nil + if len(groups) > 0 { + manifest.ProcessGroupID = groups[0] + } + if len(groups) > 1 { + manifest.ProcessGroups = groups + } +} + // Config configures the single implicit worker. type Config struct { // ServerURL is the control plane's http://host:port address.