diff --git a/.abcd/development/intents/planned/itd-2609201925079472-an-autonomous-implementation-run-paces-itself-by-default-and.md b/.abcd/development/intents/planned/itd-2609201925079472-an-autonomous-implementation-run-paces-itself-by-default-and.md deleted file mode 100644 index 5cc4c5aa6..000000000 --- a/.abcd/development/intents/planned/itd-2609201925079472-an-autonomous-implementation-run-paces-itself-by-default-and.md +++ /dev/null @@ -1,87 +0,0 @@ ---- -id: itd-2609201925079472 -slug: an-autonomous-implementation-run-paces-itself-by-default-and -spec_id: spc-2609202134341288 -kind: standalone -suggested_kind: null -reclassification_history: [] -builds_on: [itd-2609201916151817] -related_intents: [itd-29] -severity: minor -impact: additive -origin: researcher-authored -production_mode: hand-written ---- - -# An autonomous run paces itself by default, and one flag sets the pace - -## Press Release - -> **A run started with `abcd build` (or `abcd drain`, or the `abcd implement` loop they hand over to) follows a pace it was never told: a working window, a pause after it, and a ceiling on lanes alive at once.** The three numbers come from one place, layered: a flag on this run wins over the repository's abcd configuration, which wins over the machine's, which wins over the bundled default; the run record names which layer applied. `--pace /` and `--sub-agents ` set them for one run. At the window's end the loop starts no new lane, lets a running lane finish its step and checkpoint to its branch, writes `next_eligible_at` into the state file and exits; an invocation before that time refuses and names it, so the pause survives a killed process and needs no sleeping model. Today every autonomous run carries its pace in its prompt as prose, and two runs on one machine cannot share it. -> -> "Two pilots, two prompts, two paces typed by hand, and no way to know the other one was keeping to its ceiling," said a technical facilitator who ran both. "Now the pace is a line in the config and a timestamp in the state file, and a run that starts early is refused." - -## Why This Matters - -Two pilots on one machine in the week of 2026-09-15 ran on two paces written into their prompts by hand, two hours of work then five off with two lanes here, two hours then four off with three lanes there, and neither could be sure the other honoured its ceiling. The pace exists because the model budget resets on a window; a run that ignores it runs into the reset mid-lane and loses the lane. - -## Mechanism - -We expect a pace read from layered configuration and enforced by a timestamp in the state file to be honoured by every run without being told, because the loop that reads the state is the only thing that starts a lane; shown wrong if a run starts a lane inside a pause or above the ceiling, which the run record would show, or if the budget window the pause exists for moves in a way minutes cannot express. - -## Scope Conditions - -- Holds for a run driven by the implement verb's loop; a session pacing itself by prompt is outside it. -- Holds per run: The ceiling counts this run's lanes and validators. -- The bundled default is a choice, not a measurement: The two pilots ran 120/300 with two lanes and 120/240 with three, and a repository overrides it in its configuration. - -## What's In Scope - -- The three numbers, their four layers, the flags, the window clock and `next_eligible_at` in the state file, the ceiling on this run's lanes and validators, the budget check before a run starts, and the checkpoint on a rate-limit response. - -## What's Out of Scope - -- A ceiling across runs on one machine: The register's (`itd-2609150819440345`), where a lane registry can live. -- Mid-run telemetry and an operator's hand verbs over the state file: Dropped with `itd-29` until wanted. - -## Decisions - -Settled on 2026-09-20 with one defensible answer each, on the product thinker's ruling that such a decision is recorded rather than asked (`iss-2609202055103741`): - -1. **Re-entrant pause.** The pause is `next_eligible_at` in the state file; an invocation before it refuses and names the time. No process sleeps. -2. **Units are minutes**, as the product thinker specified the flag; a quota-window signal, where a harness exposes one, is a later refinement and is recorded as such. -3. **Cross-run enforcement is the register's.** This intent bounds one run; the design review found a repository file cannot bound two runs across repositories. -4. **Layering is flag, then repository, then machine, then bundled**, the order the model-tier intent already uses. -5. **The bundled default is 120 minutes of work, 300 of pause, two lanes**, the product thinker's numbers for this repository's runs on 2026-09-20; a repository that measured otherwise writes its own. -6. **The budget check and the rate-limit checkpoint come here from `itd-29`**, superseded on 2026-09-20 by the implement verb: A run refuses to start when the estimated cost exceeds the remaining quota where the runner reports one, and a rate-limit response checkpoints the lane and ends the window early. -7. **A run works in parallel up to its ceiling** (ruling DR6, the product thinker, 2026-09-29, verbatim: "per-run agent limit: WORK IN PARALLEL — a run may build several pieces and run reviewers concurrently up to its limit (new build-loop work; then AC6 is testable)."). The validators of a round run side by side, the lanes of steps that do not need each other run side by side, and implementers and reviewers share the ceiling; criterion 6 is tested through the concurrent loop the spec's piece 6 designs. -8. **Steps run one after another unless their plan says otherwise, and a hand-back holds the siblings' landings** (rulings DR6b and DR6c, the product thinker, 2026-09-30). DR6b, verbatim: "(a) ONE AFTER ANOTHER BY DEFAULT: a step runs alongside earlier ones only if its plan says so; nothing already planned changes; reviews run side by side; the 2026-09-21 wording stands." DR6c, verbatim: "(c) FINISH, BUT HOLD THEM: when one piece is handed back, pieces in flight finish but nothing merges until the person re-plans; the person then decides whether the held pieces land as they are." A step's default `needs` is every earlier step, opted out of per step; after a hand-back the sibling lanes run to completion and each passing one is `held`, never armed, until the person releases or discards it with `implement step --release ` or `--discard ` (the spec's piece 6, criteria C6, C11 and C13). - -## Open Questions - -_None open._ - -## Acceptance Criteria - -- **Given** no flag and no configuration, **when** a run starts, **then** it runs on 120/300 with two lanes and the run record names the bundled layer. -- **Given** a repository configuration and a machine configuration that disagree, **when** a run starts without a flag, **then** the repository's values apply and the record says so. -- **Given** `--pace 90/240 --sub-agents 3`, **when** a run starts, **then** those values apply over every configured layer and the record names the flag. -- **Given** a window that has elapsed, **when** the loop is invoked, **then** it starts no lane, lets a running lane finish its current step and checkpoint to its branch, writes `next_eligible_at`, and exits 0 naming the time. -- **Given** `next_eligible_at` in the future, **when** the loop is invoked, **then** it refuses naming the time and changes no state. -- **Given** the ceiling reached, **when** the loop would start a lane or a validator, **then** it starts nothing, names the lanes alive, and exits; the next invocation fills the slot, and the record counts the minutes the slot was waited for. -- **Given** a runner that reports remaining quota and an estimate that exceeds it, **when** a run starts, **then** it refuses naming both numbers and writes no state; a runner that reports no quota is named and the check is skipped out loud. -- **Given** a rate-limit response from a runner mid-lane, **when** the loop reads it, **then** the lane is checkpointed to its branch, the window ends early with `next_eligible_at` set, and the record names the response. -- **Given** a malformed pace or ceiling, **when** a run starts, **then** it refuses naming the value and the accepted form, and writes no state. - -## Typed Links - -- **builds_on `itd-2609201916151817`** (the implement verb): The loop this pace bounds. -- **refines `itd-2609170822093401`** (the model tier): The same configuration layering. - -## Audit Notes - -_Empty. Populated by intent-auditor when intent moves to shipped/._ - -## Grounds - -- pursued: we expect a pace read from layered configuration and enforced by a timestamp in the state file to be honoured by every run without being told, because the loop that reads the state is the only thing that starts a lane; shown wrong if a run starts a lane inside a pause or above the ceiling, which the run record would show diff --git a/.abcd/development/intents/shipped/itd-2609201925079472-an-autonomous-implementation-run-paces-itself-by-default-and.md b/.abcd/development/intents/shipped/itd-2609201925079472-an-autonomous-implementation-run-paces-itself-by-default-and.md new file mode 100644 index 000000000..8d33a4b57 --- /dev/null +++ b/.abcd/development/intents/shipped/itd-2609201925079472-an-autonomous-implementation-run-paces-itself-by-default-and.md @@ -0,0 +1,192 @@ +--- +id: itd-2609201925079472 +slug: an-autonomous-implementation-run-paces-itself-by-default-and +spec_id: spc-2609202134341288 +kind: standalone +suggested_kind: null +reclassification_history: [] +builds_on: [itd-2609201916151817] +related_intents: [itd-29] +severity: minor +impact: additive +origin: researcher-authored +production_mode: hand-written +--- + +# An autonomous run paces itself by default, and one flag sets the pace + +## Press Release + +> **A run started with `abcd build` (or `abcd drain`, or the `abcd implement` loop they hand over to) follows a pace it was never told: a working window, a pause after it, and a ceiling on lanes alive at once.** The three numbers come from one place, layered: a flag on this run wins over the repository's abcd configuration, which wins over the machine's, which wins over the bundled default; the run record names which layer applied. `--pace /` and `--sub-agents ` set them for one run. At the window's end the loop starts no new lane, lets a running lane finish its step and checkpoint to its branch, writes `next_eligible_at` into the state file and exits; an invocation before that time refuses and names it, so the pause survives a killed process and needs no sleeping model. Today every autonomous run carries its pace in its prompt as prose, and two runs on one machine cannot share it. +> +> "Two pilots, two prompts, two paces typed by hand, and no way to know the other one was keeping to its ceiling," said a technical facilitator who ran both. "Now the pace is a line in the config and a timestamp in the state file, and a run that starts early is refused." + +## Why This Matters + +Two pilots on one machine in the week of 2026-09-15 ran on two paces written into their prompts by hand, two hours of work then five off with two lanes here, two hours then four off with three lanes there, and neither could be sure the other honoured its ceiling. The pace exists because the model budget resets on a window; a run that ignores it runs into the reset mid-lane and loses the lane. + +## Mechanism + +We expect a pace read from layered configuration and enforced by a timestamp in the state file to be honoured by every run without being told, because the loop that reads the state is the only thing that starts a lane; shown wrong if a run starts a lane inside a pause or above the ceiling, which the run record would show, or if the budget window the pause exists for moves in a way minutes cannot express. + +## Scope Conditions + +- Holds for a run driven by the implement verb's loop; a session pacing itself by prompt is outside it. +- Holds per run: The ceiling counts this run's lanes and validators. +- The bundled default is a choice, not a measurement: The two pilots ran 120/300 with two lanes and 120/240 with three, and a repository overrides it in its configuration. + +## What's In Scope + +- The three numbers, their four layers, the flags, the window clock and `next_eligible_at` in the state file, the ceiling on this run's lanes and validators, the budget check before a run starts, and the checkpoint on a rate-limit response. + +## What's Out of Scope + +- A ceiling across runs on one machine: The register's (`itd-2609150819440345`), where a lane registry can live. +- Mid-run telemetry and an operator's hand verbs over the state file: Dropped with `itd-29` until wanted. + +## Decisions + +Settled on 2026-09-20 with one defensible answer each, on the product thinker's ruling that such a decision is recorded rather than asked (`iss-2609202055103741`): + +1. **Re-entrant pause.** The pause is `next_eligible_at` in the state file; an invocation before it refuses and names the time. No process sleeps. +2. **Units are minutes**, as the product thinker specified the flag; a quota-window signal, where a harness exposes one, is a later refinement and is recorded as such. +3. **Cross-run enforcement is the register's.** This intent bounds one run; the design review found a repository file cannot bound two runs across repositories. +4. **Layering is flag, then repository, then machine, then bundled**, the order the model-tier intent already uses. +5. **The bundled default is 120 minutes of work, 300 of pause, two lanes**, the product thinker's numbers for this repository's runs on 2026-09-20; a repository that measured otherwise writes its own. +6. **The budget check and the rate-limit checkpoint come here from `itd-29`**, superseded on 2026-09-20 by the implement verb: A run refuses to start when the estimated cost exceeds the remaining quota where the runner reports one, and a rate-limit response checkpoints the lane and ends the window early. +7. **A run works in parallel up to its ceiling** (ruling DR6, the product thinker, 2026-09-29, verbatim: "per-run agent limit: WORK IN PARALLEL — a run may build several pieces and run reviewers concurrently up to its limit (new build-loop work; then AC6 is testable)."). The validators of a round run side by side, the lanes of steps that do not need each other run side by side, and implementers and reviewers share the ceiling; criterion 6 is tested through the concurrent loop the spec's piece 6 designs. +8. **Steps run one after another unless their plan says otherwise, and a hand-back holds the siblings' landings** (rulings DR6b and DR6c, the product thinker, 2026-09-30). DR6b, verbatim: "(a) ONE AFTER ANOTHER BY DEFAULT: a step runs alongside earlier ones only if its plan says so; nothing already planned changes; reviews run side by side; the 2026-09-21 wording stands." DR6c, verbatim: "(c) FINISH, BUT HOLD THEM: when one piece is handed back, pieces in flight finish but nothing merges until the person re-plans; the person then decides whether the held pieces land as they are." A step's default `needs` is every earlier step, opted out of per step; after a hand-back the sibling lanes run to completion and each passing one is `held`, never armed, until the person releases or discards it with `implement step --release ` or `--discard ` (the spec's piece 6, criteria C6, C11 and C13). + +## Open Questions + +_None open._ + +## Acceptance Criteria + +- **Given** no flag and no configuration, **when** a run starts, **then** it runs on 120/300 with two lanes and the run record names the bundled layer. +- **Given** a repository configuration and a machine configuration that disagree, **when** a run starts without a flag, **then** the repository's values apply and the record says so. +- **Given** `--pace 90/240 --sub-agents 3`, **when** a run starts, **then** those values apply over every configured layer and the record names the flag. +- **Given** a window that has elapsed, **when** the loop is invoked, **then** it starts no lane, lets a running lane finish its current step and checkpoint to its branch, writes `next_eligible_at`, and exits 0 naming the time. +- **Given** `next_eligible_at` in the future, **when** the loop is invoked, **then** it refuses naming the time and changes no state. +- **Given** the ceiling reached, **when** the loop would start a lane or a validator, **then** it starts nothing, names the lanes alive, and exits; the next invocation fills the slot, and the record counts the minutes the slot was waited for. +- **Given** a runner that reports remaining quota and an estimate that exceeds it, **when** a run starts, **then** it refuses naming both numbers and writes no state; a runner that reports no quota is named and the check is skipped out loud. +- **Given** a rate-limit response from a runner mid-lane, **when** the loop reads it, **then** the lane is checkpointed to its branch, the window ends early with `next_eligible_at` set, and the record names the response. +- **Given** a malformed pace or ceiling, **when** a run starts, **then** it refuses naming the value and the accepted form, and writes no state. + +## Typed Links + +- **builds_on `itd-2609201916151817`** (the implement verb): The loop this pace bounds. +- **refines `itd-2609170822093401`** (the model tier): The same configuration layering. + +## Audit Notes + + +Fidelity review — receipt rcp-c411be255805 (verifier intent-auditor claude-fable-5-1). + +Provenance: intent-auditor@claude-fable-5-1 · rubric_hash sha256:effa65b3e9e88ff29433b443ec2be159522a8b0b71cf1434526514aa61edb13e · prompt_hash sha256:922986f5804db2851ba4918cd59959f541b9b4df9c2aa7a626ce91f58ab15c81 +Input attestations: diff:3f21f58b5d869ea6ed7299b6ad5c38a32a468f6f..8a08de576cdaa182b11719116cb64e43f5bcbc7a on build/run-2610100550297632-lane-1@sha256:73677a02544e3a04a092b71888d0014e06226aefae31de36621b3d27d78320d0; + +Acceptance rollup: MET 7 · MET_WITH_CONCERNS 2 · NOT_MET 0 · INCONCLUSIVE 0 + +Per-criterion verdicts: +- ac-1 — MET: The bundled constants are 120/300/2 (pace.go:28-30); TestARunWithNoConfigurationRunsOnTheBundledPace asserts the state and the run record name 120/300, 2 sub-agents and the bundled layer, and passes on a scratch copy of 8a08de576; the CLI test asserts the same line in implement status. + evidence: internal/core/implement/loop/pace.go:28 — "BundledWorkMinutes = 120" + evidence: internal/core/implement/loop/pace_test.go:82 — "the record names the pace and the bundled layer" + evidence: internal/surface/cli/build_pace_surface_test.go:61 — "pace: 120/300 minutes, 2 sub-agents (bundled)" +- ac-2 — MET: resolvePace reads each key through layered.Load (flag, repo, machine, bundled); TestTheRepositoryPaceWinsOverTheMachines sets disagreeing machine and repo files, asserts 100/200/4 from layer repo with origin .abcd/config.json, and that the record names the repository's file and not the machine's; it passes. + evidence: internal/core/implement/loop/pace.go:200 — "func resolvePace(roots layered.Roots, f paceFlags) (Pace, error)" + evidence: internal/core/implement/loop/pace_test.go:105 — "wantPace(t, st.Pace, 100, 200, 4, "repo")" + evidence: internal/core/implement/loop/pace_test.go:110 — "the record names the repository's file and not the machine's" +- ac-3 — MET: TestThePaceFlagsWinOverEveryLayer starts with both files set and --pace 90/240 --sub-agents 3, asserts 90/240/3 from layer flag with origins '--pace 90/240' and '--sub-agents 3' and the record naming both flags; the CLI test asserts the rendered line names the flag layer for each value; both pass. + evidence: internal/core/implement/loop/pace_test.go:143 — "wantPace(t, st.Pace, 90, 240, 3, "flag")" + evidence: internal/core/implement/loop/pace_test.go:144 — "st.Pace.WorkMinutes.Origin != "--pace 90/240" || st.Pace.SubAgents.Origin != "--sub-agents 3"" + evidence: internal/surface/cli/build_pace_surface_test.go:85 — "pace: 90/240 minutes, 3 sub-agents (work from --pace 90/240, the flag layer" +- ac-4 — MET: advance closes an elapsed window before any move (windowElapsed/closeWindow write next_eligible_at and the pause entry, the call returns nil error with Next naming the time); the pace test verifies the implementer's receipt is still taken inside the pause and the lane advances, and C11 verifies two lanes' receipts both verify with next_eligible_at written once. + evidence: internal/core/implement/loop/loop.go:647 — "if until, ok := windowElapsed(*st, now); ok {" + evidence: internal/core/implement/loop/loop.go:708 — "st.NextEligibleAt = &until" + evidence: internal/core/implement/loop/pace_test.go:293 — "the running lane checkpoints inside the pause" + evidence: internal/core/implement/loop/parallel_test.go:757 — "both receipts are verified during the pause" +- ac-5 — MET: advance refuses with a 'pause' contention naming next_eligible_at in RFC3339 before any state change; TestAPauseRefusesUntilNextEligibleAt asserts the refusal names 2026-09-25T13:00:00Z and no stage ran, and the pace test asserts the state bytes are identical after the refusal. + evidence: internal/core/implement/loop/loop.go:632 — "return false, contend("pause", "", "", "the run is paused until "+st.NextEligibleAt.UTC().Format(time.RFC3339)," + evidence: internal/core/implement/loop/loop_test.go:673 — "if r.Stage != "pause" || !r.Contention || !strings.Contains(r.Reason, "2026-09-25T13:00:00Z")" + evidence: internal/core/implement/loop/pace_test.go:304 — "a refused step inside the pause changes no state and performs nothing" +- ac-6 — MET: moveAgent at slotsInUse >= ceiling hands out nothing, sets CeilingReached and names the lanes alive; tookSlot records the whole minutes waited when the next call serves the item. C1-C3 (TestTheCeilingCountsValidatorsAndHoldsTheWaitingWork) assert the names, the unchanged waiting time, the slot filled by the next call and the record line '14 minute(s)'; it passes. + evidence: internal/core/implement/loop/schedule.go:538 — "if st.slotsInUse() >= st.ceiling() {" + evidence: internal/core/implement/loop/schedule.go:686 — "took a freed slot after waiting %d minute(s) at the run's ceiling of %d" + evidence: internal/core/implement/loop/parallel_test.go:356 — "the intent-auditor of lane-1 took a freed slot after waiting 14 minute(s)" +- ac-7 — MET_WITH_CONCERNS: budgetCheck runs after every other check and before the run directory is made; an estimate over a reported quota fails the row naming both numbers (TestAQuotaUnderTheEstimateRefusesTheStartNamingBoth asserts '5 agent run(s) left', 'estimated to start 6' and zero runs on disk); a runner reporting none is named and the row says 'skipped' in the checks and the run record. Concerns: (1) the comparison is reached only through the Options.Quota test seam, because QuotaReporter is implemented by no shipped runner (runner.go:135-137), so today every real run skips the check; (2) the estimate is a lower bound that omits fix rounds and syncs (budget.go:13). + evidence: internal/core/implement/loop/loop.go:332 — "budget := budgetCheck(chk, o)" + evidence: internal/core/implement/loop/budget.go:133 — "Detail: "the run's estimate exceeds the quota its runner reports: "" + evidence: internal/core/implement/loop/budget_test.go:79 — "a refused budget writes no state" + evidence: internal/core/implement/loop/budget_test.go:112 — "func TestARunnerThatReportsNoQuotaIsNamedAndTheCheckSkipped" + evidence: internal/core/runner/runner.go:135 — "QuotaReporter is a Runner that reports its remaining quota. Neither shipped" +- ac-8 — MET_WITH_CONCERNS: The claude runner reads a rejected rate_limit_event or an assistant rate_limit error as ReasonRateLimited; Dispatch returns it without fallback; Drive routes it to rateLimitWindow, which writes next_eligible_at once (a running pause kept), records the response with the lane and runner, and checkpoints the hit lane; the end-to-end test through Drive and the fake claude asserts all of it and passes. Concerns: (1) 'checkpointed to its branch' is delivered as save-aside plus reset to the branch's last commit, so the agent's uncommitted work lives in the run's aside directory, not on the branch; (2) only the claude CLI's response shapes are recognised; an opencode limit still falls back as an ordinary refusal (DECISIONS.md:2680); (3) a limit met by a host-run agent never reaches the loop (drive.go:152-158 sees only Dispatch failures). + evidence: internal/core/runner/claude.go:101 — "func claudeRateLimit(out []byte) string {" + evidence: internal/core/runner/dispatch.go:115 — "if fl.Reason == ReasonRateLimited {" + evidence: internal/core/implement/loop/drive.go:158 — "return rateLimitWindow(repoRoot, runID, res.Lane, aw, fl, o)" + evidence: internal/core/implement/loop/ratelimit.go:123 — "st.NextEligibleAt = until" + evidence: internal/core/implement/loop/ratelimit_test.go:93 — "the record names the response and the lane it came from" + evidence: internal/core/implement/loop/ratelimit.go:15 — "uncommitted work is kept for review and never built on), and its await" +- ac-9 — MET: parsePaceFlags and the resolver's range check refuse at stage pace naming the value and paceForm (the accepted form); TestAMalformedPaceIsRefusedAndWritesNoState covers 16 malformed flag and file cases and asserts the value, the form and an absent run tier; the CLI test asserts exit 2 with the form named; both pass. + evidence: internal/core/implement/loop/pace.go:62 — "var paceForm = "--pace takes < work-minutes>/< pause-minutes> in whole minutes" + evidence: internal/core/implement/loop/pace_test.go:202 — "runTierAbsent(t, repo.Root())" + evidence: internal/surface/cli/build_pace_surface_test.go:111 — "if code != 2 || !strings.Contains(errOut, "refused at pace") || !strings.Contains(errOut, "< work-minutes>/< pause-minutes>")" + +Gap audit: +- honoured: + - The three numbers come from one place, layered flag over repository over machine over bundled, and the run record names which layer applied + evidence: internal/core/implement/loop/pace.go:200 — "func resolvePace(roots layered.Roots, f paceFlags) (Pace, error)" + evidence: internal/core/implement/loop/pace_test.go:147 — "the record names the flags" + - At the window's end the loop starts no new lane, writes next_eligible_at and exits; an invocation before that time refuses and names it, so the pause survives a killed process + evidence: internal/core/implement/loop/loop.go:631 — "if st.NextEligibleAt != nil && now.Before(*st.NextEligibleAt) {" + evidence: internal/core/implement/loop/loop.go:707 — "func closeWindow(st *State, now, until time.Time) {" + - A ceiling on this run's lanes and validators, with the minutes a slot was waited for counted in the record + evidence: internal/core/implement/loop/schedule.go:671 — "func tookSlot(st *State, key, role string, now time.Time) {" + evidence: internal/core/implement/loop/parallel_test.go:321 — "if st.SlotsInUse() != 2 || st.Ceiling() != 2 {" + - A run refuses to start when the estimated cost exceeds the remaining quota where the runner reports one, naming both numbers and writing no state + evidence: internal/core/implement/loop/budget.go:126 — "if s.runs > q.Remaining {" + evidence: internal/core/implement/loop/loop.go:334 — "if !budget.OK {" + - A rate-limit response is never fallen back on; it ends the window early with next_eligible_at written once and the record naming the response, the runner and the lane + evidence: internal/core/runner/dispatch.go:115 — "if fl.Reason == ReasonRateLimited {" + evidence: internal/core/implement/loop/ratelimit.go:116 — "if until == nil || !now.Before(*until) {" + evidence: internal/core/implement/loop/ratelimit_test.go:185 — "next_eligible_at is written once" + - The budget line and the rate-limit block are wired on the CLI (build, build next, drain) and the plugin surface + evidence: internal/surface/cli/drain.go:132 — "cfg, err := loadRunners(cmd, roots)" + evidence: internal/surface/cli/build_pace_surface_test.go:190 — "func TestBuildNamesTheBudgetCheckAndTheRunnersItSkipped" + evidence: commands/build.md:180 — "Last, once every other check passes, the budget check asks each runner a role" + evidence: commands/implement.md:288 — "A runner's rate-limit response ends the window early in the same way" +- diverged: + - The rate-limited lane is 'checkpointed to its branch': delivered as the agent's uncommitted work and partial receipt saved aside under the run directory and the worktree reset to the branch's last commit, so nothing new reaches the branch + evidence: internal/core/implement/loop/ratelimit.go:157 — "aside, err := saveAside(repoRoot, root, st.RunID, *lane, lane.Awaits[k], why, aw.Role == RoleImplementer, now)" + evidence: internal/core/implement/loop/ratelimit_test.go:108 — "lane-2's worktree is reset to its last commit" + evidence: .abcd/work/DECISIONS.md:2680 — ""Checkpointed to its branch" reads as the restart of 2026-10-09 reads a gone agent" + - The budget comparison exists only behind an optional QuotaReporter that no shipped runner implements, so every real run today skips the check out loud; the comparison path is exercised through the Options.Quota test seam + evidence: internal/core/runner/runner.go:139 — "type QuotaReporter interface {" + evidence: internal/core/implement/loop/budget.go:139 — "detail = "skipped, as no runner of the run reports a quota: " + detail" + evidence: internal/core/implement/loop/loop.go:57 — "Quota func(name string) (runner.Quota, bool, error)" + - The estimate from the spec's size is a lower bound: fix rounds and syncs are not counted + evidence: internal/core/implement/loop/budget.go:13 — "Fix rounds and syncs" +- missing: + - Rate-limit detection for the opencode runner: an opencode limit is still an ordinary refusal and falls back as before + evidence: .abcd/work/DECISIONS.md:2680 — "an opencode limit stays a refusal and falls back as before" + evidence: internal/core/runner/claude.go:80 — "if what := claudeRateLimit(res.stdout); what != "" {" + - A rate limit met by a host-run sub-agent: the loop reads only a Dispatch failure, and no verb lets the host report one + evidence: internal/core/implement/loop/drive.go:152 — "out, err := d.Dispatch(ctx, req)" + +Scope-condition dispositions: +- cond-2609202134346930 — survived: Every enforcement point is inside the implement verb's loop: the pause and window close in advance, the ceiling in moveAgent, the budget check in start, and the rate-limit checkpoint reached only from Drive's Dispatch; a session pacing itself by prompt, and a host-run agent's limit, are outside what the loop can see, as the condition assumed. + evidence: internal/core/implement/loop/loop.go:639 — "The window clock (itd-2609201925079472) is the run's, one for every" + evidence: internal/core/implement/loop/drive.go:158 — "return rateLimitWindow(repoRoot, runID, res.Lane, aw, fl, o)" +- cond-2609202134341920 — survived: The ceiling is this run's pace.sub_agents counted against the awaits in this run's state file (implementers and validators alike), the budget estimate is per run, and a rate limit pauses only the run it hit (the implementer's report notes a drain's own window does not end early), so the per-run assumption held. + evidence: internal/core/implement/loop/schedule.go:538 — "if st.slotsInUse() >= st.ceiling() {" + evidence: internal/core/implement/loop/parallel_test.go:314 — "r.Slots != 2 || r.Ceiling != 2" + evidence: internal/core/implement/loop/ratelimit.go:81 — "func rateLimitWindow(repoRoot, runID, laneID string, aw Await, fl *runner.Failure, o Options) (StepResult, error) {" +- cond-2609202134341625 — survived: The bundled 120/300/2 is three named constants with the comment that a repository that measured otherwise writes its own under pace in .abcd/config.json; the repository override is tested, and the rate-limit pause falls to BundledPauseMinutes only when the run carries no pace. + evidence: internal/core/implement/loop/pace.go:23 — "The bundled pace: 120 minutes of work, 300 of pause, two lanes (decision 5," + evidence: internal/core/implement/loop/pace_test.go:99 — "repoConfig(t, repo, `{"pace": {"work_minutes": 100, "pause_minutes": 200, "sub_agents": 4}}`)" + evidence: internal/core/implement/loop/ratelimit.go:117 — "pause := BundledPauseMinutes" + + +## Grounds + +- pursued: we expect a pace read from layered configuration and enforced by a timestamp in the state file to be honoured by every run without being told, because the loop that reads the state is the only thing that starts a lane; shown wrong if a run starts a lane inside a pause or above the ceiling, which the run record would show diff --git a/.abcd/development/specs/open/spc-2609301921521360-the-budget-check-and-the-rate-limit-checkpoint.md b/.abcd/development/specs/closed/spc-2609301921521360-the-budget-check-and-the-rate-limit-checkpoint.md similarity index 100% rename from .abcd/development/specs/open/spc-2609301921521360-the-budget-check-and-the-rate-limit-checkpoint.md rename to .abcd/development/specs/closed/spc-2609301921521360-the-budget-check-and-the-rate-limit-checkpoint.md diff --git a/.abcd/work/DECISIONS.md b/.abcd/work/DECISIONS.md index feade8374..2c0ed7b5e 100644 --- a/.abcd/work/DECISIONS.md +++ b/.abcd/work/DECISIONS.md @@ -2677,3 +2677,4 @@ together (the script's header says why there is no escape hatch). - 2026-10-05 — Decided without asking (facilitator default, recorded by abcd-e8 at abcd-d7's direction while the person is away): the dashboard refuses a node shared into the tailnet from another account, and a request relayed by another node's Serve or Funnel, until the person decides; decision 4's 'anyone on your Tailscale' did not settle either, a shared-in node belongs to someone else and a Funnel relay is the open internet, so refusing is the safe default. A raw TCP relay cannot be told apart and stays a residual. - 2026-10-07 — Ruling DR3 of the person (2026-09-29, in the human-step interview run for autonomous run A and relayed by its interviewer abcd-23 [b1e81b]), restated for iss-2610020728128805 (recorded by lane-1 of run-2610071733272557, records only). The answers of that interview are kept in the local-tier file `.abcd/.work.local/scratch/reports/rulings-answered-2609-29-b.md`, line 31, and are restated here because that file is not committed. DR3, verbatim: "itd-60 docs-fidelity verdict: FOUND AUTOMATICALLY by spec close and launch ship (a saved verdict labelled with the reviewed code's fingerprint; match -> proceed; none/stale -> refuse 'run the docs review first'), like the preflight receipt." It is the ruling the doc-fidelity gate's shape rests on (itd-60, ac-2: the semantic layer is a saved review for HEAD that `spec close` and `launch ship` find automatically, rather than a review run at the close), cited by `internal/core/docfidelity/store.go`, `internal/core/docfidelity/docfidelity.go` and `internal/core/lint/gatereceipt.go`. The issue was captured on 2026-10-02 against the 2026-09-30 entry above, which restates DR1, DR2, DR4, DR5 and DR6 and skips DR3; that entry is left as written (DA002). The 2026-10-02 entry of lane rd2 above already restated DR3 later the same day, after the capture, together with CB1, CC1, CD1 and CD2, and DQ2b and CJ1b have their own 2026-09-30 entries, so none of the rulings the issue names as owed is still missing; this entry adds the issue's own pointer to the ruling and changes no earlier line. - 2026-10-09 — Rulings of the product thinker on the three hand-backs of the 2026-10-07/08 autonomous drain (asked by abcd-0b, one question at a time, answers verbatim): (1) iss-2609251618079479, whether the build-sequence chapter's build-milestone sense retires with the phase: "Rewrite in today's terms" — 06-delivery/01-build-sequence.md stays a current chapter, its build order restated as dependencies with no milestones (neither retired as history nor kept as a second sense of the word). (2) iss-2610040758569116, whether an agent may run 'abcd decide accept', which stamps the person's git name as accepted_by: "Agents too, disclosed" — no guard block as for 'source ledger --flip'; the commit's Assisted-by trailer is the disclosure. (3) iss-2610030956156354's open design point, whether a step still falls back to the agent harness when the size count already sent the brief to the person's own model server and the chat call then fails: "Move to the harness" — no answer received is the trigger, not nothing sent; recorded as iss-2610090701491123, since the resolved fix stops instead. +- 2026-10-10 — Four points spc-2609301921521360 left open, settled by lane-1 of run-2610100550297632 in building it (itd-2609201925079472 criteria 7 and 8). (1) A runner's quota is counted in agent runs (`runner.Quota.Remaining`, reported through the optional `runner.QuotaReporter`), and the estimate from the spec's size is the least the run starts: per step to build, one implementer and one round of the two reviewers, and the fidelity audit once for an intent; fix rounds and syncs are not foreseen. Neither shipped runner reports a quota (the claude CLI and opencode expose none a launch can read before it starts), so today every run names its runners and skips the check out loud; the comparison is reached by a runner that reports one. (2) The rate-limit response is read from the claude CLI's own events, a `rate_limit_event` whose `rate_limit_info.status` is `rejected` or an assistant message whose `error` is `rate_limit`, on a run that did not finish; opencode's error events carry no rate-limit shape this build has verified, so an opencode limit stays a refusal and falls back as before. A rate limit is never fallen back on: every route spends the budget the window paces. (3) "Checkpointed to its branch" reads as the restart of 2026-10-09 reads a gone agent (iss-2610080620372731): the limited agent's uncommitted work and partial receipt are saved aside for review, never built on, its worktree reset to the branch's last commit, and its await dropped so the first step after the pause hands the same work to a fresh agent; a limited validator has only its partial return saved aside, since it edits nothing and may share the worktree with its round. Every other lane with work in flight is checkpointed at its branch's head in the record and left running, its agents' receipts still taken inside the pause, because their agents are not this process's to stop. (4) The pause after a rate limit is the run's own pause minutes, not the reset time a harness reports (decision 2 of the intent), and a pause already running is kept. diff --git a/commands/build.md b/commands/build.md index 68a90e89b..150cb2560 100644 --- a/commands/build.md +++ b/commands/build.md @@ -177,6 +177,17 @@ stage with exit 2, naming the value and the accepted form, and nothing is written. Starting again keeps the run's pace: a flag naming another pace is refused, and one naming the same pace resumes. +Last, once every other check passes, the budget check asks each runner a role +of the run is routed to (below) for its remaining quota, in agent runs, and +compares the run's estimate on it: per step to build, one implementer and one +round of the two reviewers, and the fidelity audit once for an intent. An +estimate over a runner's quota is refused with exit 2 at the check `budget`, +naming both numbers, and nothing is written. A runner that reports no quota, +the host included, is named and the check is skipped out loud; neither shipped +runner reports one. The payload's `checks` carries the `budget` row, the text +output its `budget:` line, and the run record its `budget` line: tell the user +which runners were compared and which were skipped. + The window and the pause bind through `implement step` (below), and so does the ceiling: a run hands work to several agents at once, up to `sub_agents`, the validators of one round side by side and the lanes of steps that do not need @@ -241,6 +252,18 @@ that runs it): start the agent yourself as above. Every fallback is recorded; `implement status` and `implement record` count them per runner and per role. Tell the user each fallback's reason. +A runner that answers with a rate-limit response is the one failure not handed +to you: every lane spends the same budget, so the run's window ends early. The +step exits 0 with `next_eligible_at` (now plus the run's pause) and +`rate_limit`, naming the `lane`, the `role`, the `runner`, the `response` and +the `checkpoints`. The lane it came from is checkpointed to its branch: its +agent's uncommitted work and partial receipt are saved aside for review +(`aside`), never built on, and its worktree is reset to its last commit; the +first step after the pause hands the same work to a fresh agent. Every other +lane with work in flight is checkpointed at its branch's `head` and left +running: hand back the receipts of the agents still `out` as they return. Tell +the user the response, the lane and the time the run resumes. + The run's window opens when the run starts. Once its working minutes have elapsed, `implement step` starts nothing: it writes `next_eligible_at` (now plus the run's pause) into the state, records the pause, and exits 0 with diff --git a/commands/implement.md b/commands/implement.md index c4c4da310..b56e47d98 100644 --- a/commands/implement.md +++ b/commands/implement.md @@ -285,7 +285,12 @@ pause, and exits 0 with `next_eligible_at` in the result and `next` naming the time; an agent already started may still hand its receipt back. Before the run's `next_eligible_at` a `step` is refused as a pause (exit 3) naming the time, and nothing changes; at or after it a new window opens and the stage -proceeds. +proceeds. A runner's rate-limit response ends the window early in the same way +for every lane (`rate_limit` in the result, a pause already running kept): the +lane it came from is checkpointed to its last commit, its agent's uncommitted +work saved aside, and is handed to a fresh agent after the pause; every other +lane in flight is checkpointed at its branch's head and its agents may still +hand their receipts back. Without `--run`, both act on the one run in progress in this checkout, and are refused naming the runs when there are several. A refusal exits 2 (3 on a pause diff --git a/docs/reference/cli/commands.md b/docs/reference/cli/commands.md index 6b75dda4f..f2ce189ca 100644 --- a/docs/reference/cli/commands.md +++ b/docs/reference/cli/commands.md @@ -290,6 +290,15 @@ a runner is skipped with a warning on stderr, and the role runs on the host as i One it sets to host keeps the role on the host over a runner route in ~/.abcd.noindex/config.json, since that spends nothing of the person's, with a warning naming both routes. +The budget check runs last, once every other check passes: each runner a role of the +run is routed to is asked for its remaining quota, in agent runs, and compared with the +run's estimate on it (per step to build, one implementer and one round of the two +reviewers, and the fidelity audit once for an intent). An estimate over a runner's +quota is refused at the check budget, naming both numbers, and nothing is written. A +runner that reports no quota, the host included, is named and the check is skipped out +loud; neither shipped runner reports one. The result's budget line and the run record +say which. + An issue id (iss-N, validated by shape) is built as one lane. Its checks are the repository's own drain rule, read as `abcd drain` reads it (the issue is open, nothing open blocks it, its category and severity are ones the rule takes, it carries a remedy a @@ -2226,6 +2235,16 @@ nothing, writes next_eligible_at (now plus the run's pause) and exits 0 naming i agent already started may still hand back its receipt. Before next_eligible_at the call is refused as a pause and nothing changes; at or after it, a new window opens. +A runner that answers with a rate-limit response is not fallen back on, since every +lane spends the same budget: the run's window ends early, next_eligible_at is written +(now plus the run's pause; a pause already running is kept), and the call exits 0 +naming the response and the lane it came from. That lane is checkpointed to its branch: +its agent's uncommitted work and partial receipt are saved aside for review, never built +on, its worktree is reset to its last commit, and the first call after the pause hands +the same work to a fresh agent. Every other lane with work in flight is checkpointed at +its branch's head and left running; its agents may hand back their receipts inside the +pause. The run record names each checkpoint. + The run's lost connection (`abcd implement outage`) is read before every move. While the network is down, a lane whose move reaches the remote or the forge (a landing's push, pull request, arming or merged check; a hold's disarm) waits on the shared probe, diff --git a/internal/core/implement/loop/budget.go b/internal/core/implement/loop/budget.go new file mode 100644 index 000000000..690279056 --- /dev/null +++ b/internal/core/implement/loop/budget.go @@ -0,0 +1,142 @@ +package loop + +// budget.go is the budget check (itd-2609201925079472 criterion 7, from +// itd-29; spc-2609301921521360): before a run is created, each runner a role +// of the run is routed to is asked for its remaining quota, and the run's +// estimate on that runner is compared with it. A run the estimate exceeds is +// refused naming both numbers, and nothing is written. A runner that reports +// no quota, the host included, is named and the check skipped for it out loud, +// in the start's checks and in the run record. +// +// The estimate is the least the spec's size makes the run start: per step to +// build, one implementer and one round of the two reviewers, and the +// fidelity audit once for an intent (an issue has none). Fix rounds and syncs +// are not foreseen, so an estimate under a quota is no promise the quota +// holds; one over it is a run that cannot finish in the window. Quota and +// estimate are counted in agent runs, the unit runner.Quota names. + +import ( + "context" + "fmt" + "strings" + + "github.com/intentdriven/abcd/internal/abcdhome" + "github.com/intentdriven/abcd/internal/core/runner" + "github.com/intentdriven/abcd/internal/termsafe" +) + +// CheckBudget is the budget check's row among the start's checks, and the +// run record's stage for it. +const CheckBudget = "budget" + +// budgetShare is the part of the estimate one route carries. +type budgetShare struct { + route string + roles []string + runs int +} + +// estimateShares is the run's estimate per route, the routes in the order the +// roles first name them: steps implementers and as many rounds of the two +// reviewers, and, for an intent, the auditor once. +func estimateShares(cfg *runner.Config, steps int, audits bool) []budgetShare { + roles := []struct { + role string + runs int + }{{RoleImplementer, steps}, {RoleRuthless, steps}, {RoleSecurity, steps}} + if audits { + roles = append(roles, struct { + role string + runs int + }{RoleAuditor, 1}) + } + var out []budgetShare + for _, r := range roles { + route := runner.Host + if cfg != nil { + route = cfg.RouteFor(r.role).Runner + } + k := -1 + for i := range out { + if out[i].route == route { + k = i + } + } + if k < 0 { + out = append(out, budgetShare{route: route}) + k = len(out) - 1 + } + out[k].roles = append(out[k].roles, r.role) + out[k].runs += r.runs + } + return out +} + +// quota asks the named runner for its remaining quota: through the Quota seam +// when a test sets it, else the runner configuration. +func (o Options) quota(name string) (runner.Quota, bool, error) { + if o.Quota != nil { + return o.Quota(name) + } + if o.Runners == nil { + return runner.Quota{}, false, nil + } + return o.Runners.Quota(context.Background(), name) +} + +// budgetCheck is the budget row for a run of chk: the estimate on each route +// against the quota its runner reports. It fails only when a runner reports a +// quota the estimate on it exceeds. +func budgetCheck(chk CheckResult, o Options) CheckRow { + steps := len(chk.steps) + audits := chk.Intent != "" + shares := estimateShares(o.Runners, steps, audits) + total := 0 + for _, s := range shares { + total += s.runs + } + each := "an implementer, a ruthless-reviewer and a security-reviewer" + if audits { + each += ", and the intent-auditor once" + } + head := fmt.Sprintf("an estimate of %d agent run(s) from %d step(s), each %s", total, steps, each) + + var parts, over []string + compared := false + for _, s := range shares { + who := "the " + s.route + " runner" + if s.route == runner.Host { + parts = append(parts, fmt.Sprintf("the host reports no quota (%s: an estimated %d agent run(s)), so the check is skipped for it", + strings.Join(s.roles, ", "), s.runs)) + continue + } + q, ok, err := o.quota(s.route) + switch { + case err != nil: + parts = append(parts, fmt.Sprintf("%s did not report its quota (%s), so the check is skipped for it (%s: an estimated %d agent run(s))", + who, termsafe.Sanitize(err.Error()), strings.Join(s.roles, ", "), s.runs)) + case !ok: + parts = append(parts, fmt.Sprintf("%s reports no quota (%s: an estimated %d agent run(s)), so the check is skipped for it", + who, strings.Join(s.roles, ", "), s.runs)) + default: + compared = true + line := fmt.Sprintf("%s reports %d agent run(s) left, and the run is estimated to start %d on it (%s)", + who, q.Remaining, s.runs, strings.Join(s.roles, ", ")) + parts = append(parts, line) + if s.runs > q.Remaining { + over = append(over, line) + } + } + } + if len(over) > 0 { + return CheckRow{Name: CheckBudget, OK: false, + Detail: "the run's estimate exceeds the quota its runner reports: " + strings.Join(over, "; ") + " (" + head + ")", + Remedy: "start the run once the runner's quota window has reset, or route the roles it names to another runner in " + + abcdhome.Display("config.json") + "; nothing was written"} + } + detail := head + ": " + strings.Join(parts, "; ") + if !compared { + detail = "skipped, as no runner of the run reports a quota: " + detail + } + return CheckRow{Name: CheckBudget, OK: true, Detail: detail} +} diff --git a/internal/core/implement/loop/budget_test.go b/internal/core/implement/loop/budget_test.go new file mode 100644 index 000000000..fb07a2b48 --- /dev/null +++ b/internal/core/implement/loop/budget_test.go @@ -0,0 +1,151 @@ +package loop + +// budget_test.go is criterion 7 of itd-2609201925079472 (spc-2609301921521360, +// "The budget check"): where a runner a role of the run is routed to reports +// its remaining quota, the run's estimate from the spec's size is compared +// with it before the run starts; a run it exceeds is refused naming both +// numbers with no state written, and a runner that reports none is named and +// the check skipped out loud. + +import ( + "errors" + "slices" + "strings" + "testing" + + "github.com/intentdriven/abcd/internal/core/runner" +) + +// threeSteps is a spec of three steps, one after another. +const threeSteps = "1. One\n2. Two\n3. Three\n" + +// claudeImplements routes the implementer and the ruthless reviewer to the +// claude runner on the machine; the security reviewer and the auditor stay on +// the host. +const claudeImplements = `{"roles":{"implementer":{"runner":"claude"},"ruthless-reviewer":{"runner":"claude"}},"runner":{"claude":{}}}` + +// quotaOf is a quota seam: the runners it names report what it maps them to, +// and every other reports none. +func quotaOf(m map[string]int) func(string) (runner.Quota, bool, error) { + return func(name string) (runner.Quota, bool, error) { + n, ok := m[name] + return runner.Quota{Remaining: n}, ok, nil + } +} + +// budgetRow is the start's budget row. +func budgetRow(t *testing.T, rows []CheckRow) CheckRow { + t.Helper() + k := slices.IndexFunc(rows, func(r CheckRow) bool { return r.Name == CheckBudget }) + if k < 0 { + t.Fatalf("the start's checks carry no %s row: %+v", CheckBudget, rows) + } + return rows[k] +} + +// budgetEntry is the run record's budget line. +func budgetEntry(t *testing.T, st State) string { + t.Helper() + for _, e := range st.Record { + if e.Stage == CheckBudget { + return e.Note + } + } + t.Fatalf("the run record carries no budget line: %+v", st.Record) + return "" +} + +// The run of three steps on the claude runner is estimated at six agent runs +// there: an implementer and a ruthless reviewer per step. A quota of five is +// under it: the start is refused naming both numbers, and nothing is written. +func TestAQuotaUnderTheEstimateRefusesTheStartNamingBoth(t *testing.T) { + repo := loopRepo(t, readyIntent("", settledQuestions), specWithSteps(threeSteps)) + cfg := runnerConfig(t, claudeImplements, "") + _, err := Start(repo.Root(), "itd-10", Options{Runners: cfg, Quota: quotaOf(map[string]int{runner.Claude: 5})}) + r, ok := AsRefusal(err) + if !ok || r.Stage != "check" || r.Check != CheckBudget { + t.Fatalf("err = %v, want the budget check's refusal", err) + } + for _, want := range []string{"claude", "5 agent run(s) left", "estimated to start 6"} { + if !strings.Contains(r.Reason, want) { + t.Fatalf("the refusal names %q: %s", want, r.Reason) + } + } + if budgetRow(t, r.Checks).OK { + t.Fatalf("the budget row fails: %+v", r.Checks) + } + runs, err := Runs(repo.Root()) + if err != nil || len(runs) != 0 { + t.Fatalf("a refused budget writes no state: %d run(s), %v", len(runs), err) + } +} + +// A quota over the estimate starts the run, and the start and the record name +// both numbers; the roles on the host are named as reporting none. +func TestAQuotaOverTheEstimateStartsTheRun(t *testing.T) { + repo := loopRepo(t, readyIntent("", settledQuestions), specWithSteps(threeSteps)) + cfg := runnerConfig(t, claudeImplements, "") + res, err := Start(repo.Root(), "itd-10", Options{Runners: cfg, Quota: quotaOf(map[string]int{runner.Claude: 40})}) + if err != nil { + t.Fatal(err) + } + row := budgetRow(t, res.Checks) + if !row.OK { + t.Fatalf("the budget row passes: %+v", row) + } + st, err := ReadState(repo.Root(), res.RunID) + if err != nil { + t.Fatal(err) + } + for _, text := range []string{row.Detail, budgetEntry(t, st)} { + for _, want := range []string{"claude runner reports 40 agent run(s) left", "estimated to start 6", "the host reports no quota"} { + if !strings.Contains(text, want) { + t.Fatalf("the budget names %q: %s", want, text) + } + } + } +} + +// A runner that reports no quota is named and the check is skipped out loud: +// in the start's row, and in the run record. With no runner configured, every +// role is on the host, which reports none. +func TestARunnerThatReportsNoQuotaIsNamedAndTheCheckSkipped(t *testing.T) { + for _, tc := range []struct { + name string + o func(t *testing.T) Options + want []string + }{ + {"claude", func(t *testing.T) Options { return Options{Runners: runnerConfig(t, claudeImplements, "")} }, + []string{"skipped", "the claude runner reports no quota", "the host reports no quota"}}, + {"host", func(*testing.T) Options { return Options{} }, + []string{"skipped", "the host reports no quota", "implementer"}}, + {"failing", func(t *testing.T) Options { + return Options{Runners: runnerConfig(t, claudeImplements, ""), Quota: func(string) (runner.Quota, bool, error) { + return runner.Quota{}, true, errors.New("the quota endpoint\x1b[2J is down") + }} + }, []string{"skipped", "the claude runner did not report its quota", "is down"}}, + } { + t.Run(tc.name, func(t *testing.T) { + repo := loopRepo(t, readyIntent("", settledQuestions), specWithSteps(threeSteps)) + res, err := Start(repo.Root(), "itd-10", tc.o(t)) + if err != nil { + t.Fatal(err) + } + row := budgetRow(t, res.Checks) + st, err := ReadState(repo.Root(), res.RunID) + if err != nil { + t.Fatal(err) + } + for _, text := range []string{row.Detail, budgetEntry(t, st)} { + if !row.OK || strings.ContainsRune(text, '\x1b') { + t.Fatalf("a skipped check passes, its text plain: %+v", row) + } + for _, want := range tc.want { + if !strings.Contains(text, want) { + t.Fatalf("the skip names %q: %s", want, text) + } + } + } + }) + } +} diff --git a/internal/core/implement/loop/drive.go b/internal/core/implement/loop/drive.go index 5f286ec10..6408efe9c 100644 --- a/internal/core/implement/loop/drive.go +++ b/internal/core/implement/loop/drive.go @@ -21,7 +21,11 @@ package loop // receipt the verifier refuses leaves the lane awaiting, writes one // fallback receipt into the run's state (Fallbacks) and the record, and // hands the role to the host with the reason (criterion 3), which the run -// record counts per runner and per role (criterion 4). +// record counts per runner and per role (criterion 4); +// - a runner that answers with a rate-limit response is the one failure not +// handed to the host: every route spends the same budget, so the run's +// window ends early and every lane with work in flight is checkpointed to +// its branch (ratelimit.go; itd-2609201925079472 criterion 8). // // The runner is started outside the run's lock, which is held only for the // advance before it and the Receipt or the fallback write after it, so a @@ -34,6 +38,7 @@ package loop import ( "context" + "errors" "fmt" "os" "path/filepath" @@ -146,6 +151,12 @@ func Drive(ctx context.Context, repoRoot, runID string, steps Stages, o Options, } out, err := d.Dispatch(ctx, req) if err != nil { + var fl *runner.Failure + if errors.As(err, &fl) && fl.Reason == runner.ReasonRateLimited { + // Not a fallback and not a refusal: the run's window ends + // early and every lane in flight is checkpointed. + return rateLimitWindow(repoRoot, runID, res.Lane, aw, fl, o) + } return res, refuse(StageRunner, "", res.Lane, err.Error(), "the lane still awaits its receipt: correct what the reason names and step again, or start the agent the step names by hand") } diff --git a/internal/core/implement/loop/drive_test.go b/internal/core/implement/loop/drive_test.go index 806c68c0f..a4b37b90a 100644 --- a/internal/core/implement/loop/drive_test.go +++ b/internal/core/implement/loop/drive_test.go @@ -74,6 +74,17 @@ func loopFakeHarness(mode string) int { } return "" } + if mode == "ratelimit" { + // The claude CLI cut off by a rate limit mid-work: an edit left + // uncommitted in the worktree it runs in, a partial receipt, and the + // rejected rate_limit_event and error result its stream carries. + _ = os.WriteFile("halfway.txt", []byte("half done\n"), 0o600) + _ = os.WriteFile(field("Receipt: "), []byte(`{"schema_version": 1, "run_id": "`), 0o600) + fmt.Println(`{"type":"system","subtype":"init","session_id":"fake-session-3","model":"fake-model"}`) + fmt.Println(`{"type":"rate_limit_event","rate_limit_info":{"status":"rejected","resetsAt":1791000000,"rateLimitType":"five_hour"},"session_id":"fake-session-3"}`) + fmt.Println(`{"type":"result","subtype":"success","is_error":true,"result":"limit reached","session_id":"fake-session-3"}`) + return 1 + } if mode == "ok" || strings.HasPrefix(mode, "model-") { body := "{}\n" if role := field("You are the "); strings.HasPrefix(role, RoleRuthless) { diff --git a/internal/core/implement/loop/loop.go b/internal/core/implement/loop/loop.go index cc1178848..a34a909e4 100644 --- a/internal/core/implement/loop/loop.go +++ b/internal/core/implement/loop/loop.go @@ -48,6 +48,13 @@ type Options struct { // NetworkProbe proves the network back for the shared outage probe a // step runs when it is due; nil is implement.RemoteProbe at the checkout. NetworkProbe func() (bool, string) + // Runners is the runner configuration read before a run is created + // (LoadRunners). A new run's budget check (budget.go) asks each runner a + // role of the run is routed to for its remaining quota; nil leaves every + // role on the host, which reports none. + Runners *runner.Config + // Quota asks the named runner for its remaining quota; nil asks Runners. + Quota func(name string) (runner.Quota, bool, error) } // StageClaim is the refusal stage of a start whose shared-run claim is refused @@ -238,7 +245,11 @@ type StepResult struct { // reach the network or an agent wait on its shared probe (Blocked names // them) while the others move. Outage *OutageInfo `json:"outage,omitempty"` - Next string `json:"next"` + // RateLimit is the runner's rate-limit response that ended the run's + // window early in this call, with the lanes it checkpointed + // (ratelimit.go); NextEligibleAt is then the pause's end. + RateLimit *RateLimit `json:"rate_limit,omitempty"` + Next string `json:"next"` // handed is true when this call handed the lane to an agent, false when // it re-told an await an earlier call began: only the call that hands the // work out may start a runner for it. @@ -246,7 +257,8 @@ type StepResult struct { } // Start resumes the live run for key, or runs the checks and, when every one -// passes, creates the run: the state file with one lane at the sequence's first +// passes and the budget check (budget.go) finds no runner's quota under the +// run's estimate, creates the run: the state file with one lane at the sequence's first // stage, the spec's other unlanded steps pending, and the record's first line. A // refused check writes nothing. // @@ -315,6 +327,14 @@ func start(repoRoot, key string, o Options, pick *RunPick) (StartResult, error) if !chk.OK { return StartResult{}, chk.refusal() } + // The budget check runs once the record may start, and writes nothing: + // a run its estimate exceeds is refused before the run directory exists. + budget := budgetCheck(chk, o) + chk.Checks = append(chk.Checks, budget) + if !budget.OK { + chk.OK = false + return StartResult{}, chk.refusal() + } if err := fsutil.EnsureRealDirAll(repoRoot, RunRelDir, dirPerm); err != nil { return StartResult{}, fmt.Errorf("creating %s: %w", RunRelDir, err) } @@ -374,6 +394,7 @@ func start(repoRoot, key string, o Options, pick *RunPick) (StartResult, error) Note: "checks passed; " + st.Lanes[0].ID + " opened for " + laneWork(st, st.Lanes[0])}) st.Record = append(st.Record, Entry{At: now, Stage: StagePace, Note: "pace " + pace.String() + "; the first window opens now"}) + st.Record = append(st.Record, Entry{At: now, Stage: CheckBudget, Note: budget.Detail}) if pick != nil { rp := *pick rp.Lane = st.Lanes[0].ID diff --git a/internal/core/implement/loop/loop_test.go b/internal/core/implement/loop/loop_test.go index b289ab294..f320c6689 100644 --- a/internal/core/implement/loop/loop_test.go +++ b/internal/core/implement/loop/loop_test.go @@ -395,8 +395,8 @@ func TestStartCreatesOneLaneAndAStartAgainResumesIt(t *testing.T) { if len(st.Pending) != 1 || st.Pending[0].Number != 3 { t.Fatalf("the unlanded steps after the first wait as pending: %+v", st.Pending) } - if len(st.Record) != 2 || st.Record[0].Stage != "start" || st.Record[1].Stage != StagePace { - t.Fatalf("the record opens with the start, then names the pace: %+v", st.Record) + if len(st.Record) != 3 || st.Record[0].Stage != "start" || st.Record[1].Stage != StagePace || st.Record[2].Stage != CheckBudget { + t.Fatalf("the record opens with the start, then names the pace and the budget check: %+v", st.Record) } fi, err := os.Stat(filepath.Join(repo.Root(), filepath.FromSlash(res.State))) if err != nil || fi.Mode().Perm() != filePerm { diff --git a/internal/core/implement/loop/ratelimit.go b/internal/core/implement/loop/ratelimit.go new file mode 100644 index 000000000..cdbe35f31 --- /dev/null +++ b/internal/core/implement/loop/ratelimit.go @@ -0,0 +1,211 @@ +package loop + +// ratelimit.go is the rate-limit checkpoint (itd-2609201925079472 criterion 8, +// from itd-29; spc-2609301921521360, and spc-2609202134341288's "The pace and +// the fix rounds, per lane"). A runner's rate-limit response (runner +// ReasonRateLimited, which the dispatcher never falls back on) ends the window +// early for the whole run, since every lane spends the same budget: +// +// - next_eligible_at is written once: a response met while the run is +// already paused leaves the pause's end as it was; +// - the lane the response came from is checkpointed to its branch: its +// agent is gone, so what it left uncommitted is saved aside with its +// partial receipt and the worktree reset to the branch's last commit, as a +// restart does (restart.go; ruling of 2026-10-09: a gone agent's +// uncommitted work is kept for review and never built on), and its await +// is dropped, so the first step after the pause hands the same work to a +// fresh agent. A validator edits nothing and may share the worktree with +// its round, so only its partial return is saved aside; +// - every other lane with work in flight is checkpointed at its branch's +// head, which the record names: its agents are not this process's and are +// not stopped, its worktree is not touched, and each may still hand back +// its receipt inside the pause, as at a window's end. One that meets the +// same limit comes back through here and is checkpointed as the first; +// - the record names the response, the runner and the lane it came from. +// +// The pause is the run's own pause minutes (decision 2: minutes, a +// quota-window signal being a later refinement), not the reset time a harness +// may report. + +import ( + "fmt" + "os" + "slices" + "strings" + "time" + + "github.com/intentdriven/abcd/internal/core/runner" + "github.com/intentdriven/abcd/internal/fsutil" + "github.com/intentdriven/abcd/internal/gitutil" +) + +// StageRateLimit is the run record's stage for a rate-limit response, and +// StageCheckpoint for a lane it checkpointed. +const ( + StageRateLimit = "rate-limit" + StageCheckpoint = "checkpoint" +) + +// RateLimit is a runner's rate-limit response as a step reports it: the lane +// and the agent it came from, the runner, the response in abcd's words, and +// every lane it checkpointed. +type RateLimit struct { + Lane string `json:"lane"` + Role string `json:"role"` + Runner string `json:"runner"` + Response string `json:"response"` + Checkpoints []Checkpoint `json:"checkpoints"` +} + +// Checkpoint is one lane with work in flight, as a rate limit left it. +type Checkpoint struct { + Lane string `json:"lane"` + Branch string `json:"branch"` + // Head is the branch's last commit, "" when git cannot name it. + Head string `json:"head"` + // Aside is where the rate-limited agent's uncommitted work and partial + // receipt were saved, relative to the checkout root; "" for a lane whose + // agents were not the ones limited. + Aside string `json:"aside,omitempty"` + // Out are the roles of the lane's agents still out after the checkpoint, + // each of which may hand back its receipt inside the pause. + Out []string `json:"out"` + // Kept is why the limited agent's await was kept and its worktree left as + // it was (the save aside was refused); "" when it was not. + Kept string `json:"kept,omitempty"` +} + +// rateLimitWindow ends the run's window on the rate-limit response fl, which +// lane laneID's agent awaited at aw met, and checkpoints every lane with work +// in flight. +func rateLimitWindow(repoRoot, runID, laneID string, aw Await, fl *runner.Failure, o Options) (StepResult, error) { + var res StepResult + err := mutate(repoRoot, runID, func(root *os.Root, st *State) (bool, error) { + now := o.now() + hit := slices.IndexFunc(st.Lanes, func(l Lane) bool { return l.ID == laneID }) + if hit < 0 { + return false, refuse(StageRateLimit, "", "", fmt.Sprintf("%s has no lane %q to checkpoint", st.RunID, laneID), + "restore the run's state file") + } + rl := &RateLimit{Lane: laneID, Role: aw.Role, Runner: fl.Runner, Response: fl.Detail, Checkpoints: []Checkpoint{}} + st.Record = append(st.Record, Entry{At: now, Lane: laneID, Stage: StageRateLimit, + Note: fmt.Sprintf("the %s runner answered %s's %s with a rate-limit response (%s); every lane spends the same budget, so the run's window ends early", + fl.Runner, laneID, aw.Role, fl.Detail)}) + + for _, i := range st.laneOrder() { + lane := st.Lanes[i] + if i != hit && len(lane.Awaits) == 0 { + continue + } + cp := Checkpoint{Lane: lane.ID, Branch: lane.Branch, Head: branchHead(repoRoot, lane.Branch)} + if i == hit { + checkpointLimited(repoRoot, root, st, &lane, aw, fl, &cp, now) + st.Lanes[i] = lane + } + for _, a := range lane.Awaits { + cp.Out = append(cp.Out, a.Role) + } + if cp.Out == nil { + cp.Out = []string{} + } + st.Record = append(st.Record, Entry{At: now, Lane: lane.ID, Stage: StageCheckpoint, Note: checkpointNote(cp, i == hit, aw.Role)}) + rl.Checkpoints = append(rl.Checkpoints, cp) + } + + until := st.NextEligibleAt + if until == nil || !now.Before(*until) { + pause := BundledPauseMinutes + if st.Pace != nil { + pause = st.Pace.PauseMinutes.Value + } + end := now.Add(time.Duration(pause) * time.Minute) + until = &end + st.NextEligibleAt = until + opened := "the run's window" + if st.WindowStartedAt != nil { + opened = "the window opened at " + st.WindowStartedAt.UTC().Format(time.RFC3339) + } + st.Record = append(st.Record, Entry{At: now, Stage: "pause", + Note: fmt.Sprintf("a rate-limit response on %s ends %s early; no stage is taken before %s (a %d-minute pause)", + laneID, opened, end.UTC().Format(time.RFC3339), pause)}) + } + st.UpdatedAt = now + + res = idleResult(*st) + at := *until + res.NextEligibleAt = &at + res.RateLimit = rl + res.Next = fmt.Sprintf("the %s runner's rate limit on %s's %s ended the run's window early. %s", + fl.Runner, laneID, aw.Role, pausedMove(*st, at)) + return true, nil + }) + return res, err +} + +// checkpointLimited checkpoints the lane whose agent met the rate limit: its +// uncommitted work (an implementer's) and partial receipt saved aside, its +// worktree reset to the branch's last commit, and its await dropped. A save +// the loop refuses keeps the await and leaves the worktree, and cp says why. +func checkpointLimited(repoRoot string, root *os.Root, st *State, lane *Lane, aw Await, fl *runner.Failure, cp *Checkpoint, now time.Time) { + k := slices.IndexFunc(lane.Awaits, func(a Await) bool { return samePath(repoRoot, a.Receipt, aw.Receipt) }) + if k < 0 { + // The await is gone already: its receipt was handed back, or a + // restart re-told it, before the response was read. + return + } + why := fmt.Sprintf("rate limit: the %s runner answered the %s with one", fl.Runner, aw.Role) + aside, err := saveAside(repoRoot, root, st.RunID, *lane, lane.Awaits[k], why, aw.Role == RoleImplementer, now) + if err != nil { + reason := fsutil.RedactHome(err.Error()) + if r, ok := AsRefusal(err); ok { + reason = r.Reason + } + cp.Kept = reason + return + } + cp.Aside = aside.Path + if aw.Role == RoleImplementer { + cp.Head = aside.Head + } + lane.Awaits = slices.Delete(slices.Clone(lane.Awaits), k, k+1) + if len(lane.Awaits) == 0 { + lane.Awaits = nil + } +} + +// checkpointNote is the run record's line for one lane's checkpoint. +func checkpointNote(cp Checkpoint, limited bool, role string) string { + head := "an unnamed head" + if cp.Head != "" { + head = shortSHA(cp.Head) + } + at := fmt.Sprintf("%s is checkpointed at %s, the head of its branch %s", cp.Lane, head, cp.Branch) + out := "" + if len(cp.Out) > 0 { + out = "; still out, and free to hand back its receipt inside the pause: " + strings.Join(cp.Out, ", ") + } + switch { + case !limited: + return at + out + case cp.Kept != "": + return fmt.Sprintf("%s; the %s's await is kept and its worktree left as it was, since its work could not be saved aside (%s): settle it, then `abcd implement step --restart %s` once the pause ends%s", + at, role, cp.Kept, cp.Lane, out) + case cp.Aside == "": + return fmt.Sprintf("%s; the %s's await was already settled%s", at, role, out) + } + return fmt.Sprintf("%s; what the %s left (uncommitted work and partial receipt) is saved aside at %s for review, never built on, and a fresh %s takes the work once the pause ends%s", + at, role, cp.Aside, role, out) +} + +// branchHead is the last commit of branch in the checkout, "" when git cannot +// name it. +func branchHead(repoRoot, branch string) string { + if branch == "" { + return "" + } + sha, err := gitutil.Run(repoRoot, "rev-parse", "--verify", "--quiet", "refs/heads/"+branch+"^{commit}", "--") + if err != nil || !gitutil.IsFullSHA(sha) { + return "" + } + return sha +} diff --git a/internal/core/implement/loop/ratelimit_test.go b/internal/core/implement/loop/ratelimit_test.go new file mode 100644 index 000000000..8ec6a6b51 --- /dev/null +++ b/internal/core/implement/loop/ratelimit_test.go @@ -0,0 +1,246 @@ +package loop + +// ratelimit_test.go is criterion 8 of itd-2609201925079472 (spc-2609301921521360, +// "The rate-limit checkpoint"): a runner's rate-limit response ends the +// window early for the whole run, since every lane spends the same budget; +// every lane with work in flight is checkpointed to its own branch, +// next_eligible_at is written once, and the record names the lane the +// response came from. + +import ( + "bytes" + "context" + "errors" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/intentdriven/abcd/internal/core/runner" +) + +// claudeImplementer routes the implementer alone to the claude runner. +const claudeImplementer = `{"roles":{"implementer":{"runner":"claude"}},"runner":{"claude":{}}}` + +// rateLimitedPair is a run of two lanes side by side whose second +// implementer, started through the claude runner, meets a rate limit +// mid-work: lane-1's implementer, on the host, is out with a commit made and +// an edit left uncommitted. It returns the fixture, the driven step's result +// and lane-1's commit. +func rateLimitedPair(t *testing.T) (*parFixture, StepResult, string) { + t.Helper() + f := newParFixture(t, "1. One\n2. Two\n - needs: none\n", Options{SubAgents: strp("3")}) + f.stepUntil(t, "lane-1's implementer is out and lane-2 is at its implement stage", func(st State) bool { + return len(st.Lanes) == 2 && len(st.Lanes[0].Awaits) == 1 && st.Lanes[1].Stage == StageImplement && len(st.Lanes[1].Awaits) == 0 + }) + l1 := f.lane(t, "lane-1") + head1 := laneCommit(t, f.repo, l1, "one.txt") + if err := os.WriteFile(filepath.Join(l1.Worktree, "wip.txt"), []byte("lane-1 at work\n"), 0o600); err != nil { + t.Fatal(err) + } + + self, err := os.Executable() + if err != nil { + t.Fatal(err) + } + bin := t.TempDir() + if err := os.Symlink(self, filepath.Join(bin, "claude")); err != nil { + t.Fatal(err) + } + // The fake claude comes first on PATH, so no real harness is reached; + // git and the forge stub stay where the fixture put them. + t.Setenv("PATH", bin+string(os.PathListSeparator)+os.Getenv("PATH")) + t.Setenv(loopFakeEnv, "ratelimit") + t.Setenv(loopFakeLogEnv, "") + res, err := Drive(context.Background(), f.repo.Root(), f.runID, f.stages, f.opts(), + Runners{Config: runnerConfig(t, claudeImplementer, ""), Transcripts: &memTranscripts{}}) + if err != nil { + t.Fatalf("a rate-limit response is the step's answer, not its refusal: %v", err) + } + return f, res, head1 +} + +// entries are the run record's notes at stage, lane by lane. +func entries(st State, stage string) map[string][]string { + out := map[string][]string{} + for _, e := range st.Record { + if e.Stage == stage { + out[e.Lane] = append(out[e.Lane], e.Note) + } + } + return out +} + +func TestARateLimitMidLaneCheckpointsEveryLaneAndEndsTheWindow(t *testing.T) { + f, res, head1 := rateLimitedPair(t) + until := f.now.Add(BundledPauseMinutes * time.Minute) + + // The window ends early, for the whole run, with the response named. + if res.NextEligibleAt == nil || !res.NextEligibleAt.Equal(until) { + t.Fatalf("next_eligible_at = %v, want %v", res.NextEligibleAt, until) + } + if rl := res.RateLimit; rl == nil || rl.Lane != "lane-2" || rl.Role != RoleImplementer || rl.Runner != runner.Claude || + !strings.Contains(rl.Response, "rejected rate_limit_event") { + t.Fatalf("the step names the rate limit: %+v", res.RateLimit) + } + st := f.state(t) + if st.NextEligibleAt == nil || !st.NextEligibleAt.Equal(until) { + t.Fatalf("the state's next_eligible_at = %v, want %v", st.NextEligibleAt, until) + } + limits := entries(st, StageRateLimit) + if len(limits["lane-2"]) != 1 || !strings.Contains(limits["lane-2"][0], "claude") || !strings.Contains(limits["lane-2"][0], "rejected rate_limit_event") { + t.Fatalf("the record names the response and the lane it came from: %+v", limits) + } + if p := entries(st, "pause")[""]; len(p) != 1 || !strings.Contains(p[0], until.Format(time.RFC3339)) { + t.Fatalf("the record names the pause once: %+v", p) + } + + // The lane the response came from is checkpointed to its branch: its + // agent's uncommitted work and partial receipt are saved aside, its + // worktree is back at its last commit, and its implementer is out no more. + l2 := f.lane(t, "lane-2") + head2 := strings.TrimSpace(f.repo.Git("-C", l2.Worktree, "rev-parse", "HEAD")) + if len(l2.Awaits) != 0 { + t.Fatalf("the rate-limited agent's slot is freed: %+v", l2.Awaits) + } + if _, err := os.Stat(filepath.Join(l2.Worktree, "halfway.txt")); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("lane-2's worktree is reset to its last commit (%v)", err) + } + if s := f.repo.Git("-C", l2.Worktree, "status", "--porcelain"); s != "" { + t.Fatalf("lane-2's worktree is clean at its branch: %q", s) + } + dir := f.abs(RunRelDir + "/" + f.runID + "/lane-2") + at, meta := asideOf(t, dir) + if meta["head"] != head2 || !strings.Contains(meta["why"].(string), "rate limit") { + t.Fatalf("aside.json names the head and the rate limit: %v", meta) + } + if patch, err := os.ReadFile(filepath.Join(at, AsidePatchName)); err != nil || !bytes.Contains(patch, []byte("half done")) { + t.Fatalf("the uncommitted work is saved aside: %v", err) + } + cps := entries(st, StageCheckpoint) + if len(cps["lane-2"]) != 1 || !strings.Contains(cps["lane-2"][0], shortSHA(head2)) || !strings.Contains(cps["lane-2"][0], at[len(f.repo.Root())+1:]) { + t.Fatalf("the record names lane-2's checkpoint, its head and its aside: %+v", cps) + } + + // Every other lane with work in flight is checkpointed at its branch's + // head: its agent is not stopped and its worktree is not touched. + l1 := f.lane(t, "lane-1") + if len(l1.Awaits) != 1 || l1.Awaits[0].Role != RoleImplementer { + t.Fatalf("lane-1's implementer is still out: %+v", l1.Awaits) + } + if _, err := os.Stat(filepath.Join(l1.Worktree, "wip.txt")); err != nil { + t.Fatalf("lane-1's worktree is untouched: %v", err) + } + if len(cps["lane-1"]) != 1 || !strings.Contains(cps["lane-1"][0], shortSHA(head1)) || !strings.Contains(cps["lane-1"][0], RoleImplementer) { + t.Fatalf("the record names lane-1's checkpoint at its head %s: %+v", shortSHA(head1), cps) + } + + // Inside the pause a step refuses naming the time and changes nothing, + // while an agent already out may still hand back its receipt. + before := stateBytes(t, f.repo.Root(), f.runID) + _, err := advance(f.repo.Root(), f.runID, f.stages, f.opts()) + if r, ok := AsRefusal(err); !ok || r.Stage != "pause" || !strings.Contains(r.Reason, until.Format(time.RFC3339)) { + t.Fatalf("a step inside the pause refuses naming the time: %v", err) + } + if !bytes.Equal(before, stateBytes(t, f.repo.Root(), f.runID)) { + t.Fatal("a refused step inside the pause changes no state") + } + if err := os.Remove(filepath.Join(l1.Worktree, "wip.txt")); err != nil { + t.Fatal(err) + } + f.now = f.now.Add(10 * time.Minute) + f.receipt(t, "lane-1", f.await(t, "lane-1", RoleImplementer), head1) + + // After the pause, a fresh implementer takes lane-2's brief again, from + // the checkpoint. + f.now = until + f.stepUntil(t, "lane-2's implementer is handed out again", func(st State) bool { + for _, l := range st.Lanes { + if l.ID == "lane-2" && l.awaitsRole(RoleImplementer) { + return true + } + } + return false + }) + if a := f.await(t, "lane-2", RoleImplementer); a.Brief != l2.Brief { + t.Fatalf("the fresh implementer takes the same brief: %+v", a) + } +} + +// A second lane's rate-limit response inside the pause checkpoints that lane +// too, and next_eligible_at stays as it was written. +func TestASecondRateLimitInsideThePauseKeepsNextEligibleAt(t *testing.T) { + f, res, head1 := rateLimitedPair(t) + until := *res.NextEligibleAt + l1 := f.lane(t, "lane-1") + f.now = f.now.Add(20 * time.Minute) + again, err := rateLimitWindow(f.repo.Root(), f.runID, "lane-1", l1.Awaits[0], + &runner.Failure{Runner: runner.Claude, Reason: runner.ReasonRateLimited, Detail: "its event stream reports a rejected rate_limit_event and the run did not finish"}, f.opts()) + if err != nil { + t.Fatal(err) + } + st := f.state(t) + if again.NextEligibleAt == nil || !again.NextEligibleAt.Equal(until) || !st.NextEligibleAt.Equal(until) { + t.Fatalf("next_eligible_at is written once: %v, state %v, want %v", again.NextEligibleAt, st.NextEligibleAt, until) + } + if p := entries(st, "pause")[""]; len(p) != 1 { + t.Fatalf("one pause is recorded: %+v", p) + } + limits := entries(st, StageRateLimit) + if len(limits["lane-1"]) != 1 || len(limits["lane-2"]) != 1 { + t.Fatalf("the record names each response with its lane: %+v", limits) + } + l1 = f.lane(t, "lane-1") + if len(l1.Awaits) != 0 { + t.Fatalf("lane-1's rate-limited agent's slot is freed: %+v", l1.Awaits) + } + if _, err := os.Stat(filepath.Join(l1.Worktree, "wip.txt")); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("lane-1's uncommitted work is saved aside and its worktree reset (%v)", err) + } + if got := strings.TrimSpace(f.repo.Git("-C", l1.Worktree, "rev-parse", "HEAD")); got != head1 { + t.Fatalf("lane-1 keeps its committed work: HEAD %s, want %s", got, head1) + } +} + +// A validator that meets the rate limit edits nothing, and may share the +// worktree with its round: only its partial return is saved aside, the +// worktree is left as it is, and the first step after the pause hands its +// review to a fresh validator. +func TestARateLimitedValidatorLeavesTheWorktreeAndIsHandedOutAgain(t *testing.T) { + f := newParFixture(t, "1. One\n", Options{}) + f.stepUntil(t, "the implementer is out", func(st State) bool { return len(st.Lanes) == 1 && len(st.Lanes[0].Awaits) == 1 }) + f.implement(t, "lane-1", "one.txt") + f.stepUntil(t, "the ruthless reviewer is out", func(st State) bool { return st.Lanes[0].awaitsRole(RoleRuthless) }) + a := f.await(t, "lane-1", RoleRuthless) + if err := os.WriteFile(f.abs(a.Receipt), []byte("# Review\n\nhalf"), 0o600); err != nil { + t.Fatal(err) + } + l := f.lane(t, "lane-1") + if err := os.WriteFile(filepath.Join(l.Worktree, "note.txt"), []byte("left by no one in particular\n"), 0o600); err != nil { + t.Fatal(err) + } + res, err := rateLimitWindow(f.repo.Root(), f.runID, "lane-1", a, + &runner.Failure{Runner: runner.Claude, Reason: runner.ReasonRateLimited, Detail: "its event stream reports a rejected rate_limit_event and the run did not finish"}, f.opts()) + if err != nil { + t.Fatal(err) + } + if res.NextEligibleAt == nil || res.RateLimit == nil || res.RateLimit.Role != RoleRuthless { + t.Fatalf("the validator's rate limit ends the window: %+v", res) + } + if f.lane(t, "lane-1").awaitsRole(RoleRuthless) { + t.Fatal("the rate-limited validator's slot is freed") + } + if _, err := os.Stat(f.abs(a.Receipt)); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("the partial return leaves its path (%v)", err) + } + at, _ := asideOf(t, f.abs(RunRelDir+"/"+f.runID+"/lane-1")) + if got, err := os.ReadFile(filepath.Join(at, filepath.Base(a.Receipt))); err != nil || string(got) != "# Review\n\nhalf" { + t.Fatalf("the partial return is saved aside: %q %v", got, err) + } + if _, err := os.Stat(filepath.Join(l.Worktree, "note.txt")); err != nil { + t.Fatalf("a validator's rate limit leaves the worktree as it is: %v", err) + } + f.now = *res.NextEligibleAt + f.stepUntil(t, "the ruthless reviewer is handed out again", func(st State) bool { return st.Lanes[0].awaitsRole(RoleRuthless) }) +} diff --git a/internal/core/implement/loop/restart.go b/internal/core/implement/loop/restart.go index a98167146..a5cb1c31e 100644 --- a/internal/core/implement/loop/restart.go +++ b/internal/core/implement/loop/restart.go @@ -234,115 +234,144 @@ func Restart(repoRoot, runID, laneID, yielded string, o Options) (RestartResult, "run `abcd implement step` to take the lane's next stage; nothing was changed") } aw := lane.Awaits[k] - lw, err := restartWorktree(repoRoot, st.RunID, lane) + now := o.now() + aside, err := saveAside(repoRoot, root, st.RunID, lane, aw, why, true, now) if err != nil { return false, err } - head, err := gitutil.Run(lw.Path, "rev-parse", "--verify", "--quiet", "HEAD^{commit}") - if err != nil || !gitutil.IsFullSHA(head) { - return false, fmt.Errorf("resolving %s's last commit: %v", laneID, err) + dir, files, receiptName := aside.Path, aside.Files, aside.Receipt + head := aside.Head + + aw.Since = now + lane.Awaits[k] = aw + st.Lanes[i] = lane + saved := "nothing was left uncommitted" + if len(files) > 0 { + saved = fmt.Sprintf("%d uncommitted file(s) saved aside for review at %s", len(files), dir) + } + if receiptName != "" { + saved += ", with its partial receipt" } + st.Record = append(st.Record, Entry{At: now, Lane: laneID, Stage: stageRestart, + Note: fmt.Sprintf("restarted %s from its last commit %s (%s): %s (aside %s); a fresh implementer takes the brief again", + laneID, shortSHA(head), why, saved, dir)}) + st.UpdatedAt = now + res = RestartResult{StepResult: laneResult(*st, lane, "", &aw), Aside: aside} + return true, nil + }) + return res, err +} +// saveAside saves aside what the gone agent of aw left: with reset, everything +// uncommitted in the lane's worktree, proved to apply to the lane's last +// commit, after which the worktree is reset to that commit and cleaned; and, +// either way, the partial receipt at the await's path, which leaves it. The +// aside is a directory under the lane's directory of the run, named by the +// UTC second, its aside.json naming the head, the files and why. Without +// reset (a validator, who edits nothing and may share the worktree with the +// round's others) the worktree is not read or touched, and the head is the +// one the lane's state names. A refusal changes nothing, but for a reset git +// refuses after the save, which names the aside. +func saveAside(repoRoot string, root *os.Root, runID string, lane Lane, aw Await, why string, reset bool, now time.Time) (RestartAside, error) { + laneID := lane.ID + head := lane.HeadSHA + var lw LaneWorktree + var patch []byte + files := []string{} + if reset { + var err error + if lw, err = restartWorktree(repoRoot, runID, lane); err != nil { + return RestartAside{}, err + } + head, err = gitutil.Run(lw.Path, "rev-parse", "--verify", "--quiet", "HEAD^{commit}") + if err != nil || !gitutil.IsFullSHA(head) { + return RestartAside{}, fmt.Errorf("resolving %s's last commit: %v", laneID, err) + } tmp, err := os.MkdirTemp("", "abcd-restart-") if err != nil { - return false, err + return RestartAside{}, err } defer os.RemoveAll(tmp) - patch, files, err := snapshot(lw.Path, head, tmp) + patch, files, err = snapshot(lw.Path, head, tmp) if err != nil { - return false, refuse(stageRestart, "", laneID, "git could not capture the lane's uncommitted work: "+fsutil.RedactHome(err.Error()), + return RestartAside{}, refuse(stageRestart, "", laneID, "git could not capture the lane's uncommitted work: "+fsutil.RedactHome(err.Error()), "settle what git reports in the lane's worktree, then run `abcd implement step --restart "+laneID+"` again; nothing was changed") } if len(patch) > 0 { tmpPatch := filepath.Join(tmp, AsidePatchName) if err := os.WriteFile(tmpPatch, patch, filePerm); err != nil { - return false, err + return RestartAside{}, err } if err := checkAsidePatch(lw.Path, tmpPatch); err != nil { - return false, refuse(stageRestart, "", laneID, + return RestartAside{}, refuse(stageRestart, "", laneID, "the uncommitted work's patch does not apply to the lane's last commit "+shortSHA(head)+", so it is not saved and the worktree is not reset: "+fsutil.RedactHome(err.Error()), "save the lane's uncommitted work by hand, then run `abcd implement step --restart "+laneID+"` again; nothing was changed") } } + } - var partial []byte - receiptName := "" - data, err := fsutil.ReadGuardedInRoot(root, aw.Receipt, maxPartialReceiptBytes) - switch { - case err == nil: - partial, receiptName = data, filepath.Base(aw.Receipt) - case !errors.Is(err, os.ErrNotExist): - return false, refuse(stageRestart, "", laneID, "the partial receipt at "+aw.Receipt+" cannot be read: "+fsutil.RedactHome(err.Error()), - "move it out of the lane's directory yourself, then run `abcd implement step --restart "+laneID+"` again; nothing was changed") - } + var partial []byte + receiptName := "" + data, err := fsutil.ReadGuardedInRoot(root, aw.Receipt, maxPartialReceiptBytes) + switch { + case err == nil: + partial, receiptName = data, filepath.Base(aw.Receipt) + case !errors.Is(err, os.ErrNotExist): + return RestartAside{}, refuse(stageRestart, "", laneID, "the partial receipt at "+aw.Receipt+" cannot be read: "+fsutil.RedactHome(err.Error()), + "move it out of the lane's directory yourself, then run `abcd implement step --restart "+laneID+"` again; nothing was changed") + } - now := o.now() - laneDir, err := laneRel(st.RunID, laneID, stageRestart) - if err != nil { - return false, err - } - parent := laneDir + "/" + AsideDirName - if err := fsutil.EnsureRealDirAll(repoRoot, parent, dirPerm); err != nil { - return false, err - } - dir, err := asideDir(root, parent, now) - if err != nil { - return false, err - } - aside := RestartAside{RunID: st.RunID, Lane: laneID, At: now, Path: dir, Head: head, Files: files, Why: why, Receipt: receiptName} - if len(patch) > 0 { - aside.Patch = AsidePatchName - if err := fsutil.WriteFileAtomicInRoot(root, dir+"/"+AsidePatchName, patch, filePerm); err != nil { - return false, err - } - } - if partial != nil { - if err := fsutil.WriteFileAtomicInRoot(root, dir+"/"+receiptName, partial, filePerm); err != nil { - return false, err - } - } - meta, err := json.MarshalIndent(aside, "", " ") - if err != nil { - return false, err - } - if err := fsutil.WriteFileAtomicInRoot(root, dir+"/"+AsideFileName, append(meta, '\n'), filePerm); err != nil { - return false, err + laneDir, err := laneRel(runID, laneID, stageRestart) + if err != nil { + return RestartAside{}, err + } + parent := laneDir + "/" + AsideDirName + if err := fsutil.EnsureRealDirAll(repoRoot, parent, dirPerm); err != nil { + return RestartAside{}, err + } + dir, err := asideDir(root, parent, now) + if err != nil { + return RestartAside{}, err + } + aside := RestartAside{RunID: runID, Lane: laneID, At: now, Path: dir, Head: head, Files: files, Why: why, Receipt: receiptName} + if len(patch) > 0 { + aside.Patch = AsidePatchName + if err := fsutil.WriteFileAtomicInRoot(root, dir+"/"+AsidePatchName, patch, filePerm); err != nil { + return RestartAside{}, err } - // The partial receipt leaves the await's path, so the fresh agent - // finds none there and nothing of the gone one's is handed back. - if partial != nil { - if err := root.Remove(aw.Receipt); err != nil && !errors.Is(err, os.ErrNotExist) { - return false, err - } + } + if partial != nil { + if err := fsutil.WriteFileAtomicInRoot(root, dir+"/"+receiptName, partial, filePerm); err != nil { + return RestartAside{}, err } - - // Saved and proved: only now is the lane's own worktree reset. - for _, args := range [][]string{{"reset", "--hard", "--quiet", head}, {"clean", "-f", "-d", "--quiet", "--"}} { - if _, err := gitutil.Run(lw.Path, args...); err != nil { - return false, refuse(stageRestart, "", laneID, - "the uncommitted work is saved aside at "+dir+", but git could not reset the lane's worktree: "+fsutil.RedactHome(err.Error()), - "settle what git reports in the lane's worktree, then run `abcd implement step --restart "+laneID+"` again") - } + } + meta, err := json.MarshalIndent(aside, "", " ") + if err != nil { + return RestartAside{}, err + } + if err := fsutil.WriteFileAtomicInRoot(root, dir+"/"+AsideFileName, append(meta, '\n'), filePerm); err != nil { + return RestartAside{}, err + } + // The partial receipt leaves the await's path, so the fresh agent + // finds none there and nothing of the gone one's is handed back. + if partial != nil { + if err := root.Remove(aw.Receipt); err != nil && !errors.Is(err, os.ErrNotExist) { + return RestartAside{}, err } + } + if !reset { + return aside, nil + } - aw.Since = now - lane.Awaits[k] = aw - st.Lanes[i] = lane - saved := "nothing was left uncommitted" - if len(files) > 0 { - saved = fmt.Sprintf("%d uncommitted file(s) saved aside for review at %s", len(files), dir) - } - if receiptName != "" { - saved += ", with its partial receipt" + // Saved and proved: only now is the lane's own worktree reset. + for _, args := range [][]string{{"reset", "--hard", "--quiet", head}, {"clean", "-f", "-d", "--quiet", "--"}} { + if _, err := gitutil.Run(lw.Path, args...); err != nil { + return RestartAside{}, refuse(stageRestart, "", laneID, + "the uncommitted work is saved aside at "+dir+", but git could not reset the lane's worktree: "+fsutil.RedactHome(err.Error()), + "settle what git reports in the lane's worktree, then run `abcd implement step --restart "+laneID+"` again") } - st.Record = append(st.Record, Entry{At: now, Lane: laneID, Stage: stageRestart, - Note: fmt.Sprintf("restarted %s from its last commit %s (%s): %s (aside %s); a fresh implementer takes the brief again", - laneID, shortSHA(head), why, saved, dir)}) - st.UpdatedAt = now - res = RestartResult{StepResult: laneResult(*st, lane, "", &aw), Aside: aside} - return true, nil - }) - return res, err + } + return aside, nil } // asideDir makes the aside's directory under parent, named by the UTC second, diff --git a/internal/core/runner/claude.go b/internal/core/runner/claude.go index 9eb1e6352..890e6b1ab 100644 --- a/internal/core/runner/claude.go +++ b/internal/core/runner/claude.go @@ -75,12 +75,57 @@ func (c *ClaudeCLI) Run(ctx context.Context, req Request) (Answer, []byte, error res, err := c.launch.run(ctx, Claude, bin, c.args(req), nil, req.Dir, req.timeout()) transcript := res.transcript() if err != nil { + // A harness that stopped on a rate limit exits non-zero: what its + // events say is the reason, not the exit. + if what := claudeRateLimit(res.stdout); what != "" { + return Answer{}, transcript, rateLimited(what) + } return Answer{}, transcript, err } ans, err := parseClaude(res.stdout) return ans, transcript, err } +// rateLimited is the claude runner's rate-limit failure, what naming the +// response in abcd's own words. +func rateLimited(what string) *Failure { + return fail(Claude, ReasonRateLimited, "its event stream reports %s and the run did not finish", what) +} + +// claudeRateLimit names the rate-limit response the event stream carries, "" +// when it carries none. The claude CLI (2.1) reports one in two shapes: a +// rate_limit_event whose rate_limit_info.status is rejected (allowed and +// allowed_warning are a run going on), and an assistant message whose error is +// rate_limit. Each line is read leniently and on its own, so a field this +// reader does not know never turns a limit into an unparsable run. +func claudeRateLimit(out []byte) string { + for _, line := range bytes.Split(out, []byte("\n")) { + var ev struct { + Type string `json:"type"` + Error json.RawMessage `json:"error"` + Info json.RawMessage `json:"rate_limit_info"` + } + if json.Unmarshal(bytes.TrimSpace(line), &ev) != nil { + continue + } + switch ev.Type { + case "rate_limit_event": + var info struct { + Status string `json:"status"` + } + if json.Unmarshal(ev.Info, &info) == nil && info.Status == "rejected" { + return "a rejected rate_limit_event" + } + case "assistant": + var e string + if json.Unmarshal(ev.Error, &e) == nil && e == "rate_limit" { + return "an assistant message whose error is rate_limit" + } + } + } + return "" +} + // claudeEvent is the part of a stream-json event the runner reads. type claudeEvent struct { Type string `json:"type"` @@ -123,6 +168,13 @@ func parseClaude(out []byte) (Answer, error) { result = &e } } + if result == nil || result.IsError || result.Subtype != "success" { + // A run that did not finish on a rate limit is that limit, whatever + // its result says. + if what := claudeRateLimit(out); what != "" { + return Answer{}, rateLimited(what) + } + } if result == nil { return Answer{}, fail(Claude, ReasonUnparsable, "its output carries no result event") } diff --git a/internal/core/runner/config.go b/internal/core/runner/config.go index 1bdf5b88c..ac9b28f00 100644 --- a/internal/core/runner/config.go +++ b/internal/core/runner/config.go @@ -35,6 +35,7 @@ package runner // absent runner and falls back. import ( + "context" "fmt" "sort" "strings" @@ -314,6 +315,16 @@ func (c *Config) Runner(name string) (RunnerConfig, bool) { // "" when none is configured. func (c *Config) FallbackHost() string { return c.fallback } +// Quota asks the runner the machine enabled under name for its remaining +// quota. The host, and a name the machine has not enabled, report none. +func (c *Config) Quota(ctx context.Context, name string) (Quota, bool, error) { + rc, ok := c.runners[name] + if !ok { + return Quota{}, false, nil + } + return quotaOf(ctx, c.adapter(rc)) +} + // adapter builds the enabled runner's adapter. func (c *Config) adapter(rc RunnerConfig) Runner { if rc.Name == OpenCode { diff --git a/internal/core/runner/dispatch.go b/internal/core/runner/dispatch.go index 143df4249..2cf95aee7 100644 --- a/internal/core/runner/dispatch.go +++ b/internal/core/runner/dispatch.go @@ -6,7 +6,9 @@ package runner // contract's validator, the role goes to the host session, or with no host // session to the host the operator configured, and exactly one receipt is // written for the event. Because the branch and the receipt writer are one, -// the count the run's summary reports is every fallback there was. +// the count the run's summary reports is every fallback there was. A runner +// that reports a rate limit is the one failure not fallen back on: it reaches +// the caller as itself, with no receipt. import ( "context" @@ -110,6 +112,13 @@ func (d *Dispatcher) Dispatch(ctx context.Context, req Request) (Outcome, error) if !errors.As(err, &fl) { return Outcome{}, err } + if fl.Reason == ReasonRateLimited { + // Not a fallback: every route spends the budget the run's window + // paces, so the role is handed to no one and the caller ends the + // window (itd-2609201925079472 criterion 8). No receipt is written, + // since nothing fell back. + return Outcome{}, err + } fb := FallbackReceipt{At: d.now(), Role: req.Role, Asked: route.Runner, Reason: fl.Reason, Detail: fl.Detail, Ran: landing} if landing == route.Runner || landing == "" { fb.Ran = none diff --git a/internal/core/runner/main_test.go b/internal/core/runner/main_test.go index f2fde9b31..5b725edeb 100644 --- a/internal/core/runner/main_test.go +++ b/internal/core/runner/main_test.go @@ -92,6 +92,30 @@ func fakeHarness(mode string) int { fmt.Fprintln(os.Stdout, `{"type":"result","subtype":"error_during_execution","is_error":true,"result":"refused","session_id":"fake-session-1"}`) } return 0 + case "ratelimit-exit", "ratelimit-result", "ratelimit-assistant", "ratelimit-warning": + // A claude run that meets a rate limit, in the shapes its event + // stream carries one (the claude CLI 2.1's rate_limit_event, whose + // rate_limit_info.status is allowed, allowed_warning or rejected, and + // an assistant message whose error is rate_limit): rejected and the + // process exits non-zero; rejected and the result reports the error + // with a zero exit; the assistant's error alone; and a warning the + // run finishes past, which is no limit at all. + fmt.Fprintln(os.Stdout, `{"type":"system","subtype":"init","session_id":"fake-session-1","model":"fake-model"}`) + switch mode { + case "ratelimit-warning": + fmt.Fprintln(os.Stdout, `{"type":"rate_limit_event","rate_limit_info":{"status":"allowed_warning","resetsAt":1791000000,"rateLimitType":"five_hour","utilization":0.91},"session_id":"fake-session-1"}`) + fmt.Fprintln(os.Stdout, `{"type":"result","subtype":"success","is_error":false,"result":"done","session_id":"fake-session-1"}`) + return 0 + case "ratelimit-assistant": + fmt.Fprintln(os.Stdout, `{"type":"assistant","message":{"content":[{"type":"text","text":"limited"}]},"error":"rate_limit","session_id":"fake-session-1"}`) + default: + fmt.Fprintln(os.Stdout, `{"type":"rate_limit_event","rate_limit_info":{"status":"rejected","resetsAt":1791000000,"rateLimitType":"five_hour"},"session_id":"fake-session-1"}`) + } + fmt.Fprintln(os.Stdout, `{"type":"result","subtype":"success","is_error":true,"result":"limit reached","session_id":"fake-session-1"}`) + if mode == "ratelimit-result" { + return 0 + } + return 1 case "hugemodel", "ctrlmodel", "suffixmodel": // A harness whose init event reports a model: one past any bound, // one carrying control bytes, and a real id's bracketed suffix. diff --git a/internal/core/runner/ratelimit_test.go b/internal/core/runner/ratelimit_test.go new file mode 100644 index 000000000..dcde54e9d --- /dev/null +++ b/internal/core/runner/ratelimit_test.go @@ -0,0 +1,106 @@ +package runner + +// ratelimit_test.go is the runner's half of itd-2609201925079472's criteria 7 +// and 8 (spc-2609301921521360): a run that meets a rate limit says so as its +// own failure, never as a fallback, and a runner's remaining quota is asked +// for before a run starts, a runner that reports none reporting none. + +import ( + "context" + "errors" + "strings" + "testing" +) + +// claudeRouted routes the implementer to the claude runner, with no fallback +// host: a host session takes what the runner does not run. +const claudeRouted = `{"roles":{"implementer":{"runner":"claude"}},"runner":{"claude":{}}}` + +// TestARateLimitIsItsOwnFailure: a claude run whose event stream carries a +// rate-limit response fails as rate-limited, whether the harness exits +// non-zero or reports the error in its result, and whether the response is a +// rejected rate_limit_event or an assistant message whose error is +// rate_limit. A warning the run finishes past is an answer, not a limit. +func TestARateLimitIsItsOwnFailure(t *testing.T) { + for _, mode := range []string{"ratelimit-exit", "ratelimit-result", "ratelimit-assistant"} { + t.Run(mode, func(t *testing.T) { + f := newFake(t, mode, Claude) + _, transcript, err := newClaude("").Run(context.Background(), f.request("implementer")) + var fl *Failure + if !errors.As(err, &fl) || fl.Reason != ReasonRateLimited { + t.Fatalf("err = %v, want a %s failure", err, ReasonRateLimited) + } + if strings.Contains(fl.Detail, "limit reached") { + t.Fatalf("the detail is abcd's, never the harness's text: %q", fl.Detail) + } + if len(transcript) == 0 { + t.Fatal("the rate-limited run's transcript is returned") + } + }) + } + t.Run("ratelimit-warning", func(t *testing.T) { + f := newFake(t, "ratelimit-warning", Claude) + ans, _, err := newClaude("").Run(context.Background(), f.request("implementer")) + if err != nil || ans.Text != "done" { + t.Fatalf("a run past a rate-limit warning answers: %+v %v", ans, err) + } + }) +} + +// TestARateLimitIsNotFallenBackOn: a rate-limited runner hands the role to no +// one, since every route spends the budget the run's window paces; no +// fallback receipt is written, and the failure reaches the caller as itself. +func TestARateLimitIsNotFallenBackOn(t *testing.T) { + for _, hostSession := range []bool{true, false} { + f := newFake(t, "ratelimit-exit", Claude) + machine := claudeRouted + if !hostSession { + machine = `{"roles":{"implementer":{"runner":"claude"}},"runner":{"fallback_host":"claude","claude":{}}}` + } + h := newHarness(t, mustLoad(t, machine, ""), hostSession) + out, err := h.d.Dispatch(context.Background(), f.request("implementer")) + var fl *Failure + if !errors.As(err, &fl) || fl.Reason != ReasonRateLimited || fl.Runner != Claude { + t.Fatalf("host session %v: err = %v, want the claude runner's rate limit", hostSession, err) + } + if out.Handoff || out.Fallback != nil || len(h.receipts) != 0 { + t.Fatalf("host session %v: a rate limit records no fallback and hands nothing on: %+v %v", hostSession, out, h.receipts) + } + if len(h.store.got) != 1 { + t.Fatalf("host session %v: the rate-limited run's transcript is stored once: %d", hostSession, len(h.store.got)) + } + } +} + +// fakeQuota is a runner that reports its remaining quota. +type fakeQuota struct { + q Quota + err error +} + +func (fakeQuota) Name() string { return "fake" } +func (fakeQuota) Run(context.Context, Request) (Answer, []byte, error) { + return Answer{}, nil, errors.New("not run") +} +func (f fakeQuota) Quota(context.Context) (Quota, error) { return f.q, f.err } + +// TestQuotaIsAskedOfTheRunnerThatReportsIt: the host, a runner this machine +// has not enabled and the shipped runners report no quota; a runner that +// reports one is asked for it, and its failure is returned as it is. +func TestQuotaIsAskedOfTheRunnerThatReportsIt(t *testing.T) { + c := mustLoad(t, `{"runner":{"claude":{},"opencode":{"model":"local/qwen3-coder"}},`+localProvider+`}`, "") + for _, name := range []string{Host, Claude, OpenCode, "unknown"} { + if q, ok, err := c.Quota(context.Background(), name); ok || err != nil { + t.Fatalf("%s reports %+v %v %v, want none", name, q, ok, err) + } + } + if q, ok, err := quotaOf(context.Background(), fakeQuota{q: Quota{Remaining: 7}}); !ok || err != nil || q.Remaining != 7 { + t.Fatalf("a reporting runner's quota = %+v %v %v", q, ok, err) + } + if _, ok, err := quotaOf(context.Background(), fakeQuota{err: errors.New("down")}); !ok || err == nil { + t.Fatalf("a reporting runner's failure is returned: %v %v", ok, err) + } + if _, ok, err := quotaOf(context.Background(), newClaude("")); ok || err != nil { + t.Fatalf("the claude CLI reports no quota: %v %v", ok, err) + } +} diff --git a/internal/core/runner/runner.go b/internal/core/runner/runner.go index 36f1f2d41..bb98c86a1 100644 --- a/internal/core/runner/runner.go +++ b/internal/core/runner/runner.go @@ -117,8 +117,40 @@ const ( ReasonUnparsable Reason = "unparsable" // ReasonInvalid: the answer failed the contract's validator. ReasonInvalid Reason = "invalid" + // ReasonRateLimited: the harness reported a rate-limit response and did + // not finish. It is never fallen back on (dispatch.go): every route spends + // the budget the run's window paces, so the caller ends the window + // (itd-2609201925079472 criterion 8). + ReasonRateLimited Reason = "rate-limited" ) +// Quota is what a runner reports of its budget before a run starts +// (itd-2609201925079472 criterion 7): the agent runs it can still start in its +// current window. A runner whose harness counts its budget in another unit +// converts it, or reports none. +type Quota struct { + Remaining int `json:"remaining"` +} + +// QuotaReporter is a Runner that reports its remaining quota. Neither shipped +// runner is one: the claude CLI and opencode expose no remaining quota a +// launch can read before it starts, so a run routed to either names it and +// skips the budget check out loud. +type QuotaReporter interface { + Quota(ctx context.Context) (Quota, error) +} + +// quotaOf asks r for its remaining quota. reported is false when r reports +// none; an error is r's failure to report the one it keeps. +func quotaOf(ctx context.Context, r Runner) (q Quota, reported bool, err error) { + qr, ok := r.(QuotaReporter) + if !ok { + return Quota{}, false, nil + } + q, err = qr.Quota(ctx) + return q, true, err +} + // Failure is a runner's failure. Detail is written by abcd, never copied from // the harness's output, so it cannot carry anything the harness printed, a // credential it echoed included. diff --git a/internal/surface/cli/build.go b/internal/surface/cli/build.go index ae6060dc3..6b5caa30e 100644 --- a/internal/surface/cli/build.go +++ b/internal/surface/cli/build.go @@ -141,6 +141,14 @@ func newBuildCommand(asJSON *bool) *cobra.Command { "a runner is skipped with a warning on stderr, and the role runs on the host as if unrouted.\n" + "One it sets to host keeps the role on the host over a runner route in " + abcdhome.Display("config.json") + ",\n" + "since that spends nothing of the person's, with a warning naming both routes.\n\n" + + "The budget check runs last, once every other check passes: each runner a role of the\n" + + "run is routed to is asked for its remaining quota, in agent runs, and compared with the\n" + + "run's estimate on it (per step to build, one implementer and one round of the two\n" + + "reviewers, and the fidelity audit once for an intent). An estimate over a runner's\n" + + "quota is refused at the check budget, naming both numbers, and nothing is written. A\n" + + "runner that reports no quota, the host included, is named and the check is skipped out\n" + + "loud; neither shipped runner reports one. The result's budget line and the run record\n" + + "say which.\n\n" + "An issue id (iss-N, validated by shape) is built as one lane. Its checks are the\n" + "repository's own drain rule, read as `abcd drain` reads it (the issue is open, nothing\n" + "open blocks it, its category and severity are ones the rule takes, it carries a remedy a\n" + @@ -164,10 +172,11 @@ func newBuildCommand(asJSON *bool) *cobra.Command { for _, n := range notes { fmt.Fprintln(cmd.ErrOrStderr(), termsafe.Sanitize(n)) } - if _, err := loadRunners(cmd, roots); err != nil { + cfg, err := loadRunners(cmd, roots) + if err != nil { return loopFail(cmd.OutOrStdout(), *asJSON, prefix, err) } - o := loop.Options{Session: session, Roots: &roots} + o := loop.Options{Session: session, Roots: &roots, Runners: cfg} if cmd.Flags().Changed("pace") { o.Pace = &pace } @@ -189,6 +198,7 @@ func newBuildCommand(asJSON *bool) *cobra.Command { fmt.Fprintf(w, "build %s: %s run %s\n", termsafe.Sanitize(args[0]), verb, res.RunID) fmt.Fprintf(w, " state: %s\n", res.State) renderPace(w, res.Pace) + renderBudget(w, res.Checks) renderLaneLine(w, res.Lane) renderPending(w, res.Pending) switch { @@ -256,10 +266,11 @@ func newBuildNextCommand(asJSON *bool) *cobra.Command { for _, n := range notes { fmt.Fprintln(cmd.ErrOrStderr(), termsafe.Sanitize(n)) } - if _, err := loadRunners(cmd, roots); err != nil { + cfg, err := loadRunners(cmd, roots) + if err != nil { return loopFail(cmd.OutOrStdout(), *asJSON, prefix, err) } - o := loop.Options{Session: session, Roots: &roots} + o := loop.Options{Session: session, Roots: &roots, Runners: cfg} if cmd.Flags().Changed("pace") { o.Pace = &pace } @@ -322,6 +333,7 @@ func renderNext(w io.Writer, res loop.NextResult) { fmt.Fprintf(w, " (the lane's worktree stage commits it as %s's first commit)\n", res.Start.Lane.ID) fmt.Fprintf(w, " state: %s\n", res.Start.State) renderPace(w, res.Start.Pace) + renderBudget(w, res.Start.Checks) renderLaneLine(w, res.Start.Lane) renderPending(w, res.Start.Pending) fmt.Fprintf(w, "next: %s\n", termsafe.Sanitize(fsutil.RedactHome(res.Start.Next))) @@ -583,12 +595,50 @@ func renderSlots(w io.Writer, st loop.State) { } } +// renderBudget renders a new run's budget check (itd-2609201925079472 +// criterion 7): each runner's quota against the run's estimate, or the skip +// with the runners that report none. A resumed start runs no check and +// renders none. +func renderBudget(w io.Writer, checks []loop.CheckRow) { + for _, c := range checks { + if c.Name == loop.CheckBudget { + fmt.Fprintf(w, " budget: %s\n", termsafe.Sanitize(c.Detail)) + } + } +} + +// shortOrNone is a commit's short form, or "an unnamed head". +func shortOrNone(sha string) string { + if len(sha) < 12 { + return "an unnamed head" + } + return sha[:12] +} + // renderStepResult is the text form of an `implement step` or `implement receipt` result. func renderStepResult(w io.Writer, verb string, res loop.StepResult) { switch { case res.HandBack != nil: fmt.Fprintf(w, "%s: HANDED BACK: %s's %s stopped as %s after %d fix round(s); the run starts nothing further for it\n", verb, res.RunID, res.Lane, res.HandBack.Verdict, res.HandBack.FixRounds) + case res.RateLimit != nil: + rl := res.RateLimit + fmt.Fprintf(w, "%s: the %s runner's rate limit on %s's %s ended %s's window early; paused until %s\n", verb, + termsafe.Sanitize(rl.Runner), rl.Lane, termsafe.Sanitize(rl.Role), res.RunID, res.NextEligibleAt.UTC().Format(time.RFC3339)) + fmt.Fprintf(w, " response: %s\n", termsafe.Sanitize(rl.Response)) + for _, cp := range rl.Checkpoints { + line := fmt.Sprintf(" checkpoint: %s at %s on %s", cp.Lane, shortOrNone(cp.Head), termsafe.Sanitize(cp.Branch)) + if cp.Aside != "" { + line += "; the limited agent's work saved aside at " + termsafe.Sanitize(cp.Aside) + } + if cp.Kept != "" { + line += "; its await kept: " + termsafe.Sanitize(fsutil.RedactHome(cp.Kept)) + } + if len(cp.Out) > 0 { + line += "; still out: " + termsafe.Sanitize(strings.Join(cp.Out, ", ")) + } + fmt.Fprintln(w, line) + } case res.NextEligibleAt != nil: fmt.Fprintf(w, "%s: %s's window has elapsed; paused until %s\n", verb, res.RunID, res.NextEligibleAt.UTC().Format(time.RFC3339)) case res.CeilingReached: @@ -734,6 +784,15 @@ func newImplementStepCommand(asJSON *bool) *cobra.Command { "nothing, writes next_eligible_at (now plus the run's pause) and exits 0 naming it; an\n" + "agent already started may still hand back its receipt. Before next_eligible_at the call\n" + "is refused as a pause and nothing changes; at or after it, a new window opens.\n\n" + + "A runner that answers with a rate-limit response is not fallen back on, since every\n" + + "lane spends the same budget: the run's window ends early, next_eligible_at is written\n" + + "(now plus the run's pause; a pause already running is kept), and the call exits 0\n" + + "naming the response and the lane it came from. That lane is checkpointed to its branch:\n" + + "its agent's uncommitted work and partial receipt are saved aside for review, never built\n" + + "on, its worktree is reset to its last commit, and the first call after the pause hands\n" + + "the same work to a fresh agent. Every other lane with work in flight is checkpointed at\n" + + "its branch's head and left running; its agents may hand back their receipts inside the\n" + + "pause. The run record names each checkpoint.\n\n" + "The run's lost connection (`abcd implement outage`) is read before every move. While\n" + "the network is down, a lane whose move reaches the remote or the forge (a landing's\n" + "push, pull request, arming or merged check; a hold's disarm) waits on the shared probe,\n" + diff --git a/internal/surface/cli/build_pace_surface_test.go b/internal/surface/cli/build_pace_surface_test.go index 7e19c16b1..e194387da 100644 --- a/internal/surface/cli/build_pace_surface_test.go +++ b/internal/surface/cli/build_pace_surface_test.go @@ -182,3 +182,25 @@ func TestBuildHelpSaysTheCeilingBinds(t *testing.T) { t.Errorf("build --help still says the ceiling is not counted against") } } + +// TestBuildNamesTheBudgetCheckAndTheRunnersItSkipped is criterion 7 at the +// surface: the runner configuration the build reads reaches the budget check, +// whose line names the runner a role is routed to as reporting no quota, and +// the host for the rest, with the check skipped out loud. +func TestBuildNamesTheBudgetCheckAndTheRunnersItSkipped(t *testing.T) { + buildRepo(t) + home := os.Getenv("HOME") + if err := os.MkdirAll(abcdhome.Path(home), 0o700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(abcdhome.Path(home, "config.json"), + []byte(`{"roles": {"implementer": {"runner": "claude"}}, "runner": {"claude": {}}}`), 0o600); err != nil { + t.Fatal(err) + } + out := mustImplement(t, "build", "itd-10") + for _, want := range []string{" budget: skipped", "the claude runner reports no quota (implementer", "the host reports no quota (ruthless-reviewer"} { + if !strings.Contains(out, want) { + t.Fatalf("the build's budget line names %q:\n%s", want, out) + } + } +} diff --git a/internal/surface/cli/drain.go b/internal/surface/cli/drain.go index 2b13d0b8c..24142aae0 100644 --- a/internal/surface/cli/drain.go +++ b/internal/surface/cli/drain.go @@ -129,7 +129,11 @@ func runDrain(cmd *cobra.Command, asJSON bool, repoRoot string, maxLanes int, pa for _, n := range notes { fmt.Fprintln(cmd.ErrOrStderr(), termsafe.Sanitize(n)) } - o := loop.Options{Roots: &roots} + cfg, err := loadRunners(cmd, roots) + if err != nil { + return loopFail(cmd.OutOrStdout(), asJSON, prefix, err) + } + o := loop.Options{Roots: &roots, Runners: cfg} if cmd.Flags().Changed("pace") { o.Pace = &pace } @@ -178,6 +182,9 @@ func renderDrainRun(w io.Writer, res loop.DrainResult) { fmt.Fprintf(w, " cap: --max %d lane(s); %d opened\n", res.Max, len(res.Lanes)) } renderPace(w, res.Pace) + if res.Start != nil { + renderBudget(w, res.Start.Checks) + } fmt.Fprintf(w, " lanes: %d opened, one at a time\n", len(res.Lanes)) for _, l := range res.Lanes { line := fmt.Sprintf(" %s %s %s", l.Issue, l.RunID, l.Outcome)