Standalone orchestration for software development in Claude Code: managers, workers in their own git worktrees and one reviewer run your tasks in parallel, and…

flow is a Claude Code plugin that works through a batch of software tasks in parallel, inside your session. Managers plan each task. Workers build it, each in its own git worktree, and open pull requests. A reviewer merges them, runs your checks, cuts the release and deploys. You only make the decisions.
you ── main session
├── manager: csv-export one per task: plans, reviews, hands PRs over
│ ├── worker: csv-endpoint one per piece of work, in its own worktree
│ └── worker: csv-button
├── manager: login-redirect
│ └── worker: redirect-fix
└── reviewer merges PRs, runs the full check, deploys
In a Claude Code session:
/plugin install flow --marketplace Ying-Kai-Liao/flow
Ask for work in your main session:
Add CSV export to the reports page, and fix the login redirect loop.
Here is what happens next:
Type /flow to open a pane that shows who is doing what. For a single change, "start a worker to fix X" is enough.
| Command | What it does | ||
|---|---|---|---|
/flow | Opens the pane: managers, workers, and what needs you first. | ||
/flow inbox | Lists what waits for you: a summary line, the questions (blocking first, deploy/push/env marked NEEDS YOU), the decisions agents made as one flat table (newest 15; /flow inbox decisions lists all, /flow inbox d12 one item in full) and what standing answers did. | ||
| `/flow ok [d10 d11 \ | owner \ | topic]` | Keeps decisions (no words: every open decision) or takes the default of an ordinary question. Never answers a deploy, push, env or guard item. |
/flow no d12 [what instead] | Undoes a decision; the agent is told what to do instead. | ||
/flow answer q3 <choice> | Answers any question addressed to main (option letter, number, text, or your own words), deploy, push and env items included. | ||
/flow checks | Lists the after-deploy checks that only a person can do as one table (id, PR, title, what it needs, age). /flow checks <id> shows one in full; pass, fail and skip <id...> <why> close them. | ||
/flow approve <n> | Approves PR n when it is waiting for you. | ||
/flow push | Pushes the batch the reviewer built and checked, when push_mode is confirm. /flow push back <pr> returns one PR, /flow push drop all of them. | ||
/flow hold <target> / /flow release <target> | Keeps a deploy target from deploying, or lets it go again. | ||
/flow resume | Picks up unfinished work after a restart. | ||
/flow clean | Lists leftover worktrees and branches. Add --yes to remove them. |
Questions, approvals and checks also show up in the pane and as notifications. With anything open, press i in the pane to read and answer the inbox there (j/k move, a digit or letter picks an option, y takes the default or keeps a decision, w keeps all decisions, n undoes, r types an answer); you may answer any question, a manager's too. The key row sits at the pane's bottom.
flow works with no setup. Most people set a few things in .claude/flow.json at the repo root:
{
"test_command": "pnpm test",
"full_check_command": "pnpm check",
"deploy_targets": [{ "name": "staging", "deploy": ["./deploy.sh staging"] }],
"merge_mode": "confirm",
"max_workers": 3
}
| Setting | Set it when |
|---|---|
test_command | your repo has a test command workers should run for the files they change |
full_check_command | you have a full check (lint, build, the whole suite) the reviewer should run once per batch |
deploy_targets | you want the reviewer to deploy after merging. For one target, deploy_command is the shortcut |
merge_mode | you want to approve each handed-over PR yourself (confirm) instead of merging automatically |
push_mode | you want to approve each batch before it is pushed (confirm) |
release | you want each merged batch to bump the version and cut a changelog section |
manager_model, worker_model, max_workers | you want other models for managers or workers, or more or fewer parallel workers |
Anything you leave out is skipped and reported. flow never guesses a command. Everything else, with defaults, is in the reference.
/flow resume picks them up in a new session.claude --plugin-dir . # a session with your working copy loaded
npm ci # install TypeScript
npm run typecheck # must report 0 errors
npm run validate # must pass
npm run check # typecheck, validate and the plugin tests
Add your changelog lines under ## [Unreleased] in CHANGELOG.md. Don't change the version; the reviewer does that at merge. More in the reference.
hooks/register.tsx 5278 lines1// This file is the wiring only: atoms, `on(...)` hooks, `$.tool.register` calls, and io builders whose
2// closures spell out each `$.noun.event(...)` call. New logic goes into a module under hooks/, either
3// pure or taking an io object of closures (see deliver.ts), because the engine refuses `$` across an
4// import and wants atoms declared in this file.
5import { atom, read, update } from 'claude-code'
6import type { AgentInfo, AgentSpawnInput, EngineInterface, Register } from 'claude-code'
7
8import type { Activity, AgentRow, Ledger, Role, EnvChange, Handover, HandoffRecord, Leftovers, LogEvent, OpenPr, PrCache, Session, SlotEntry, TestSlots } from '../types'
9import { absolutePath, parseAttachments, rewriteAttachments } from './attachments'
10import { grantFile, grantSlots, heldBy, reapSlots, slotLine, span, CLAIM_MS, LEASE_MS, WAIT_DEFAULT_S, WAIT_MAX_S } from './slots'
11import { checkEvidence, evidenceRefusal, evidenceSummary, evidenceText, type Evidence } from './evidence'
12import { noteWork, runLine, serialFastForward, type RunWork } from './mainff'
13import { ancestorPids, ancestryQueries, containedCandidates, dirtyFiles, isLive, leftoverLine, parsePorcelain, selectCleanup, sweepText, waitingPaths } from './clean'
14import { sweep as sweepTo } from './clean'
15import type { CleanInputs, CleanIo, Kept, PrRow, Sweep } from './clean'
16import { gatherLeftovers as gatherLeftoversTo, resumeInstructions } from './resume'
17import type { Gathered, Leftover, OwnerNotes, ResumeIo } from './resume'
18import { deliver as deliverTo, flushAgent as flushTo, holdForReviewer as holdTo, resetDelivery, ENDED, LIVE, type DeliverIo, type DeliverOptions } from './deliver'
19import { addStep, addTurn, baseName, costBlock, entriesOfBranch, mergeLedgers, normalizeLedger, pruneLedger, prCost, reportSuffix, setIdentity } from './cost'
20import { backNote, effectiveSize, floorFrom, generation, isSuccessorName, modelFor, parseSize } from './routing'
21import type { Size, SizeModels } from './routing'
22import { analyze, cleanDir, findRefs, render, UNSET_TEXT } from './migrations'
23import type { PrInput } from './migrations'
24import { addNodes, agentFor, asksQuestion, describe, noticeText, settle, waitsOnReport } from './dag'
25import type { AgentFact, Facts, Graph, Notice, Plan } from './dag'
26import {
27 addQuestions, answerMessage, askingNames, EMPTY_INBOX, fyiAsked, inboxHead, isFyi, parseFyi, openAll, renderInbox, renderItem, expandOk, findItem, stillOpen, markAnswered, overridesManager, clip, paneRows, paneRowText, protectedWhy, tagsOf, OVERTURN, needsMessage, normalizeInbox, notesOwner, openFor, parseAsk, parseChoice,
28} from './inbox'
29import type { PaneRow } from './inbox'
30import {
31 closeStale, denyText, dueRound, EMPTY_PREFLIGHT, followUp, FILE_HELP, gateOf, isSkip, markDelivered, normalizePreflight, parseFiling, phaseOf,
32 recordFiling, recordSpawn, renderFollowUp, renderRound, renderStatus,
33} from './preflight'
34import type { Preflight } from './preflight'
35import {
36 addChecks, anyMatch, CHECKS_USAGE, closeChecks, dueForPrompt, EMPTY_CHECKS, inboxChecksSection, markStarted, normalizeChecks, paneChecksLine, parseNeeds, renderCheckDetail, renderChecksForAgents, renderChecksTable,
37 resumeChecksLines, SKIP_NOTE, versionInSteps,
38} from './checks'
39import type { Check, Checks } from './checks'
40import { termWidth } from './table'
41import type { Inbox, Marked, Question } from './inbox'
42import { AUTO, escalation, matchRule, nextRuleId, removeRule, renderRules, renderSeeds, ruleFromQuestion, sameRule, SEEDS, seedIds, seedsToOffer, suggest, validateRule } from './standing'
43import type { Resolved, Rule } from './standing'
44import { graphNodes, layoutGraph, moveFocus } from './graph'
45import { shimmerParts } from './shimmer'
46import type { GNode, Seg } from './graph'
47import {
48 fill, MANAGER_PROMPT, NO_REVIEWER_RULE, REVIEWER_PROMPT, REVIEWER_RULE, WORKER_PROMPT,
49} from './prompts'
50import type { Settings } from './prompts'
51import { deployModeWarnings, deployTargetsOf, stateFileOf, targetsOf } from './prompts'
52import { isFable, mergeLayers, renameOptions, settingsOf } from './settings'
53import type { WarnSettings } from './settings'
54import { bumpVersion, changelogSection, cutChangelog, highestBump, isBump, labelBump, localDate, readVersion, setVersion } from './release'
55import type { Bump } from './release'
56import {
57 allowed, allowList, killRefusal, mainCheckoutRefusal, mainRelative, parseWorktrees, resolvePath, writeTargets,
58} from './guards'
59import type { WriteTarget } from './guards'
60import {
61 ago, claudeLimits, harnessesOf, runsOutIn, SESSION, sessionKey, sessionRow,
62} from './sessions'
63import type { HarnessSpec, Limit } from './sessions'
64import { sessionTool as sessionToolTo, watchSessions as watchSessionsTo } from './session-run'
65import type { Caller, Ran, SessionIo } from './session-run'
66import { branchOwners, buildDigest, findWorktree, noteKey, ownerFor } from './state'
67import { ADD_OPTION, addMapping, guardReport, guardTestsFor, parseGuardTests, pathsFromStatus, suggestionQuestion, suggestionsFor } from './guardtests'
68import type { GuardMap } from './guardtests'
69import { autoRefused, effectiveMode, labelSpec, parseMode, takeDecision } from './mergemode'
70import {
71 batchId, closePushItem, DROPPED_REASON, dropStep, EMPTY_PUSH, itemOf, newBatch, normalizePush, openPushItem, parsePushMode, parseVerdict,
72 PUSH_KIND, recordable, refOf, releaseStep, renderBatch, reviewerNote, sendBackStep, SENT_BACK_REASON, settlePrs,
73} from './pushgate'
74import type { PushState, ReadyBatch } from './pushgate'
75import {
76 applyAnswer, applyViewOf, approvalContext, behindLines, closeEnvItems, declinedEntries, decideEnv, decideGate, EMPTY_DEPLOYS, ENV_KIND, envCommandFor, envListLine,
77 envSummary, isApprove, normalizeDeploys, openApplyItem, openApprovalItem, reopenItem, openDeployIds, openEnvItems, parseEnvInput, pendingEnv, recordDeployed,
78 release, renderList, retargetItem, unknownTarget, itemViewOf, withApproval, withEnvDone, withHold,
79} from './deploy'
80import type { Deploys, DeployMode, Draft, EnvInput, Hold, TargetInfo, TargetState } from './deploy'
81
82// The orca-flow pattern inside one Claude Code session. The main session is the super manager
83// (the `dispatch` skill); it starts `flow:manager` agents, which start
84// `flow:worker` agents in worktrees of their own and hand approved PRs to the
85// `flow:reviewer` agent through this plugin's tools. The pane in main shows the tree.
86
87const PANE = 'flow'
88const POLL_MS = 3000
89const LOG_MAX = 40
90const MANAGER = 'flow:manager'
91const WORKER = 'flow:worker'
92// A worker continued in its predecessor's worktree: the plugin rewrites a spawn to it, the model never picks it.
93const CONTINUE = 'flow:continue'
94const WORKERS = new Set([WORKER, CONTINUE, SESSION])
95const REVIEWER = 'flow:reviewer'
96// The reviewer's old agent type: a reviewer started by an older version still runs under it.
97const QUEUE = 'flow:queue'
98const isReviewer = (type: string): boolean => type === REVIEWER || type === QUEUE
99const LIVE_STATUS = new Set(['running', 'pending'])
100// Display order: what may need a person first, finished agents last.
101const ORDER = ['waiting', 'idle', 'running', 'pending', 'failed', 'killed', 'completed']
102const GLYPH: Record<string, string> = {
103 pending: '○', running: '●', waiting: '◐', idle: '◌', completed: '✓', failed: '✗', killed: '■',
104}
105const COLOR: Record<string, string> = {
106 running: 'suggestion', waiting: 'warning', idle: 'warning', completed: 'success', failed: 'error', killed: 'error',
107}
108const ROLE: Record<string, string> = { [MANAGER]: 'manager', [WORKER]: 'worker', [CONTINUE]: 'worker', [SESSION]: 'worker', [REVIEWER]: 'reviewer', [QUEUE]: 'reviewer' }
109const ROOT_GLYPH = '◆'
110// Plan states, drawn like the agent statuses they turn into; a waiting node has no agent yet.
111const PLAN_GLYPH: Record<string, string> = { waiting: '○', ready: '◌', running: '●', done: '✓', blocked: '✗' }
112const PLAN_COLOR: Record<string, string | undefined> = { waiting: undefined, ready: 'warning', running: 'suggestion', done: 'success', blocked: 'error' }
113const HANDOVER_GLYPH: Record<Handover['status'], string> = {
114 pending: '…', awaiting: '⏸', taken: '●', ready: '⇪', done: '✓', returned: '↩',
115}
116
117const roster = atom({ plugin: 'flow', key: 'roster' } as const, [] as AgentRow[])
118const activity = atom({ plugin: 'flow', key: 'activity' } as const, {} as Record<string, Activity>)
119const selected = atom({ plugin: 'flow', key: 'selected' } as const, null as string | null)
120// The highlighted card of the tree (an agent id), and the collapse the person chose per card: true
121// is one row with its children hidden, false is expanded even when the tree is crowded. Absent =
122// the role's default (managers and the reviewer collapsed), else automatic. MERGE_QUEUE_KEY is the
123// Reviewer section's entry.
124const cursor = atom({ plugin: 'flow', key: 'cursor' } as const, null as string | null)
125const MERGE_QUEUE_KEY = '#merge-queue'
126const folded = atom({ plugin: 'flow', key: 'folded' } as const, {} as Record<string, boolean>)
127// Following the chat view: the view (an agent id, null for main) in which the person last acted in
128// the pane. While the view differs from it, the pane shows the viewed agent; written by handlers only.
129const overrideView = atom({ plugin: 'flow', key: 'overrideView' } as const, undefined as string | null | undefined)
130// Set once the first card click has told the person how to see that agent's chat.
131const hinted = atom({ plugin: 'flow', key: 'hinted' } as const, false)
132const now = atom({ plugin: 'flow', key: 'now' } as const, 0)
133const handovers = atom({ plugin: 'flow', key: 'handovers' } as const, {} as Record<string, Handover>)
134// The decision inbox, mirrored from <state dir>/inbox.json.
135// The last status and answer seen per agent name, so an agent the host drops from its list keeps its
136// last known state instead of reading as never started.
137const seenAgents = atom({ plugin: 'flow', key: 'seen-agents' } as const, {} as Record<string, AgentFact>)
138const inbox = atom({ plugin: 'flow', key: 'inbox' } as const, EMPTY_INBOX as Inbox)
139// Pre-flight records, mirrored from <state dir>/preflight.json.
140const preflight = atom({ plugin: 'flow', key: 'preflight' } as const, EMPTY_PREFLIGHT as Preflight)
141const checks = atom({ plugin: 'flow', key: 'checks' } as const, EMPTY_CHECKS as Checks)
142// Tokens spent per agent and model, mirrored (throttled) to <state dir>/ledger.json. Never forgets an ended agent.
143const ledger = atom({ plugin: 'flow', key: 'ledger' } as const, {} as Ledger)
144const handoffs = atom({ plugin: 'flow', key: 'handoffs' } as const, {} as Record<string, HandoffRecord>)
145const queueRuns = atom({ plugin: 'flow', key: 'queueRuns' } as const, 0)
146// The open PRs gh listed last, so the 3 s refresh never calls gh itself.
147const prCache = atom({ plugin: 'flow', key: 'prCache' } as const, { prs: [], fetchedAt: 0 } as PrCache)
148// The last dry cleanup sweep, refreshed with the PR list so a render never runs git.
149const sessions = atom({ plugin: 'flow', key: 'sessions' } as const, {} as Record<string, Session>)
150// What the reviewer routing already sent main, so a retried or repeated report reaches it once:
151// Scoped to the sending reviewer run (its agent id), so a later run's different outcome is never swallowed.
152// `keys` maps "reviewer|manager|PR numbers" to the needs-a-person / pending-decisions lines forwarded for it,
153// `lines` is, per reviewer id, every line forwarded (normalized). Persisted: plugins reload mid-session.
154type Forwarded = { keys: Record<string, string[]>; lines: Record<string, string[]> }
155const forwarded = atom({ plugin: 'flow', key: 'forwarded' } as const, { keys: {}, lines: {} } as Forwarded)
156const harnessLimits = atom({ plugin: 'flow', key: 'harnessLimits' } as const, [] as Limit[])
157const armed = atom({ plugin: 'flow', key: 'armed' } as const, null as { key: string; at: number } | null)
158const leftovers = atom({ plugin: 'flow', key: 'leftovers' } as const, { worktrees: 0, branches: 0, needsLook: 0 } as Leftovers)
159// Deploy gates, mirrored from <state dir>/deploys.json.
160const deploys = atom({ plugin: 'flow', key: 'deploys' } as const, EMPTY_DEPLOYS as Deploys)
161// Commits each target is behind the base, from git off the 3 s refresh path (refreshBehind).
162const behind = atom({ plugin: 'flow', key: 'behind' } as const, {} as Record<string, number>)
163// The ready batch awaiting the user's push, mirrored from <state dir>/push.json.
164const pushState = atom({ plugin: 'flow', key: 'push' } as const, EMPTY_PUSH as PushState)
165
166const PR_POLL_MS = 5 * 60_000
167const PR_MIN_GAP_MS = 60_000
168// A worker that just ended: its manager is probably reviewing the PR.
169const GRACE_MS = 20 * 60_000
170// Set by register(); refresh() and the pane flag nothing when there is no reviewer.
171let queueOn = true
172// The cleanup setting and the base it sweeps against; set by register() like queueOn.
173let cleanupMode: 'auto' | 'off' = 'auto'
174let cleanupBase = 'main'
175// The deploy targets and their modes; set by register() like queueOn.
176let deployInfos: TargetInfo[] = []
177const infosOf = (s: Parameters<typeof targetsOf>[0]): TargetInfo[] => targetsOf(s).map(t => ({ name: t.name, mode: t.mode, ...(t.envCommand !== undefined ? { envCommand: t.envCommand } : {}) }))
178
179export type Unhanded = { pr: number; title: string; branch: string; note: string }
180
181// Open flow/* PRs that nobody handed to the reviewer and nobody is working on any more.
182// A draft is a WIP handoff, a pending/taken/done handover is in hand, a returned one needs a
183// manager again. The worker is the agent named like the branch, or its continuation (-2, -3).
184export function unhandedPrs(
185 prs: OpenPr[], handovers: Record<string, Handover>, roster: AgentRow[],
186 activity: Record<string, Activity>, t: number, graceMs = GRACE_MS,
187): Unhanded[] {
188 const out: Unhanded[] = []
189 for (const p of prs) {
190 if (p.isDraft || !p.headRefName.startsWith('flow/')) continue
191 const h = handovers[String(p.number)]
192 if (h !== undefined && h.status !== 'returned') continue
193 const slug = p.headRefName.slice('flow/'.length)
194 const mine = roster.filter(a => a.name === slug || (a.name?.startsWith(`${slug}-`) === true && /^\d+$/.test(a.name.slice(slug.length + 1))))
195 if (mine.some(a => !ENDED.has(a.status))) continue
196 const endedAt = Math.max(...mine.map(a => activity[a.id]?.endedAt ?? activity[a.id]?.lastAt ?? 0))
197 if (mine.length > 0 && endedAt > 0 && t - endedAt < graceMs) continue
198 const last = mine.length > 0 ? mine.reduce((x, y) => (activity[y.id]?.endedAt ?? 0) >= (activity[x.id]?.endedAt ?? 0) ? y : x) : undefined
199 const note = h !== undefined ? `returned: ${h.reason ?? ''}`
200 : last === undefined ? 'no agent in this session'
201 : `no handover, worker ${labelOf(last)} ended${endedAt > 0 ? ` ${ago(t - endedAt)} ago` : ''}`
202 out.push({ pr: p.number, title: p.title, branch: p.headRefName, note })
203 }
204 return out
205}
206
207const unhandedLine = (u: Unhanded): string => `#${u.pr} ${u.title} (${u.branch}) — ${u.note}`
208
209async function currentUnhanded($: EngineInterface): Promise<Unhanded[]> {
210 if (!queueOn) return []
211 const [cache, hs, rows, acts, t] = await Promise.all([
212 read($, prCache), read($, handovers), read($, roster), read($, activity), $.clock.now(),
213 ])
214 return unhandedPrs(cache.prs, hs, rows, acts, t)
215}
216const plan = atom({ plugin: 'flow', key: 'plan' } as const, {} as Plan)
217// The pane draws the agents as a tree of cards or as the dependency graph; the graph's highlight is a node id.
218const viewMode = atom({ plugin: 'flow', key: 'viewMode' } as const, 'tree' as 'tree' | 'graph' | 'inbox')
219// The inbox view: the highlighted row's key, the expanded groups, the typed answer and the last result line.
220const inboxCursor = atom({ plugin: 'flow', key: 'inboxCursor' } as const, null as string | null)
221const inboxDraft = atom({ plugin: 'flow', key: 'inboxDraft' } as const, '')
222const inboxNote = atom({ plugin: 'flow', key: 'inboxNote' } as const, '')
223const graphFocus = atom({ plugin: 'flow', key: 'graphFocus' } as const, null as string | null)
224const testSlots = atom({ plugin: 'flow', key: 'testSlots' } as const, { holders: [], waiters: [] } as TestSlots)
225
226// <git-common-dir>/flow/test-slots, set at session start; undefined when not in a git repo.
227let slotDir: string | undefined
228// The test_slots setting, for refresh()'s status line (set where the settings are read).
229let slotLimit = 1
230// The decision_phrases setting, for the question check in refresh() and syncPlans().
231let decisionPhrases: string[] = []
232// Pre-flight setting and round wait (ms); set by register() like the others.
233let preflightOn = true
234let preflightWaitMs = 10 * 60_000
235
236// Applies `pre`, reaps, grants, and signals: a grant file and a message per new holder, the files
237// of everyone who left the line removed. Every slot change goes through here, inside one update().
238async function settleSlots($: EngineInterface, rows: AgentRow[], t: number, pre: (s: TestSlots) => TestSlots = s => s, me?: string): Promise<void> {
239 const notes: string[] = []
240 let granted: SlotEntry[] = []
241 let left: string[] = []
242 await update($, testSlots, st => {
243 notes.length = 0
244 const r = reapSlots(pre(st), rows, t)
245 const g = grantSlots(r.state, slotLimit, t)
246 notes.push(...r.notes, ...g.notes)
247 // The caller of acquire is here to hear it: its grant is confirmed on the spot, with no signal.
248 const state = me === undefined ? g.state
249 : { ...g.state, holders: g.state.holders.map(h => h.key === me ? { ...h, claimed: true } : h) }
250 granted = g.granted.filter(h => h.key !== me)
251 const stay = new Set([...state.holders, ...state.waiters].map(e => e.key))
252 left = [...st.holders, ...st.waiters].map(e => e.key).filter(k => !stay.has(k))
253 return state
254 })
255 for (const n of notes) void $.ui.toast(n)
256 await signalSlots($, granted, left)
257}
258
259const sh = ($: EngineInterface, script: string, ...args: string[]) =>
260 $.process.run(['sh', '-c', script, 'sh', ...args], { timeoutMs: 5_000 }).catch(() => undefined)
261
262async function signalSlots($: EngineInterface, granted: SlotEntry[], left: string[]): Promise<void> {
263 for (const key of left) {
264 const f = grantFile(slotDir, key)
265 if (f) await sh($, 'rm -f "$1"', f)
266 }
267 for (const g of granted) {
268 const f = grantFile(slotDir, g.key)
269 if (f) await sh($, 'mkdir -p "$(dirname "$1")" && : > "$1"', f)
270 if (g.key !== 'main') {
271 await deliver($, g.key, `flow: your test slot is granted (${g.label}). Call mcp__flow__test_slot acquire to confirm, run, then release.`, { onGone: () => undefined }).catch(() => undefined)
272 }
273 }
274}
275
276export type TreeItem = { a: AgentRow; depth: number; kids: number; collapsed: boolean }
277
278// Managers and the reviewer start collapsed; everything else has no default.
279export const foldDefault = (a: AgentRow): boolean | undefined => (a.type === MANAGER || isReviewer(a.type) ? true : undefined)
280
281// The rows of the tree in the order they are drawn and walked. A collapsed agent is one row with
282// its children hidden: the person's choice wins, then the role's default, then `auto`, which folds
283// every agent with children except the ones on the way to the highlight (a crowded tree).
284// `at` is the highlight, moved up to the nearest row that is drawn.
285export function treeItems(
286 list: AgentRow[], fold: Record<string, boolean>, cur: string | null | undefined, auto: boolean,
287 acts: Record<string, Activity> = {},
288): { items: TreeItem[]; at: string | undefined } {
289 const ids = new Set(list.map(a => a.id))
290 const kids = (id: string | undefined) => list
291 .filter(a => (id === undefined ? a.parentId === undefined || !ids.has(a.parentId) : a.parentId === id))
292 .sort((a, b) => rankOf(a, acts) - rankOf(b, acts))
293 const path = new Set<string>()
294 for (let a = list.find(x => x.id === cur); a !== undefined && !path.has(a.id); a = list.find(x => x.id === a!.parentId)) path.add(a.id)
295 const items: TreeItem[] = []
296 const walk = (a: AgentRow, depth: number) => {
297 const under = kids(a.id)
298 const collapsed = fold[a.id] ?? foldDefault(a) ?? (under.length > 0 && auto && !path.has(a.id))
299 items.push({ a, depth, kids: under.length, collapsed })
300 if (!collapsed) for (const c of under) walk(c, depth + 1)
301 }
302 for (const a of kids(undefined)) walk(a, 0)
303 const drawn = new Set(items.map(i => i.a.id))
304 let at: string | undefined = cur ?? undefined
305 while (at !== undefined && !drawn.has(at)) at = list.find(a => a.id === at)?.parentId
306 return { items, at }
307}
308
309// The window of rows to draw: all of them, or `size` rows with the highlight in the middle.
310export function viewOf<T>(items: T[], at: number, size: number): { top: number; rows: T[] } {
311 if (items.length <= size) return { top: 0, rows: items }
312 const top = Math.min(Math.max(0, at - Math.floor(size / 2)), items.length - size)
313 return { top, rows: items.slice(top, top + size) }
314}
315
316// A plan node with an agent is drawn by the agent's status; one without has only its plan state.
317const planState = (n: GNode): boolean => n.agentId === undefined && n.state in PLAN_GLYPH
318const nodeGlyph = (n: GNode): string => (planState(n) ? PLAN_GLYPH[n.state] ?? '?' : GLYPH[n.state] ?? PLAN_GLYPH[n.state] ?? '?')
319const nodeColor = (n: GNode): string | undefined => (planState(n) ? PLAN_COLOR[n.state] : COLOR[n.state] ?? PLAN_COLOR[n.state])
320
321// One line for a tool call: the tool and its most telling argument.
322function describeCall(e: Record<string, unknown>): string {
323 const tool = String(e.tool ?? '?').replace(/^mcp__flow__/, '')
324 const arg = [e.file_path, e.command, e.pattern, e.path, e.url, e.description, e.action, e.prompt]
325 .find(v => typeof v === 'string' && v.length > 0) as string | undefined
326 // A command is its first line, cut at a heredoc: its body is code, not news.
327 const text = arg === undefined ? '' : (e.command === arg ? arg.split('\n')[0]!.replace(/<<.*$/, '') : arg)
328 const short = text === '' ? '' : ' ' + text.replace(/\s+/g, ' ').trim().slice(0, 90)
329 return tool + short
330}
331
332const base = (path: unknown): string => String(path ?? '').split('/').pop() ?? ''
333
334// What a call is, in a few plain words for a card: never a command's body or a prompt.
335export function summarizeCall(e: Record<string, unknown>): string {
336 const tool = String(e.tool ?? '?').replace(/^mcp__flow__/, '')
337 const file = base(e.file_path)
338 switch (tool) {
339 case 'Edit': case 'MultiEdit': case 'NotebookEdit': return file ? `editing ${file}` : 'editing'
340 case 'Write': return file ? `writing ${file}` : 'writing'
341 case 'Read': return file ? `reading ${file}` : 'reading'
342 case 'Grep': case 'Glob': return 'searching'
343 case 'Agent': {
344 const role = String(e.subagent_type ?? e.subagentType ?? '').replace(/^flow:/, '')
345 const name = String(e.name ?? '')
346 return ['started', role === 'manager' || role === 'worker' ? role : 'agent', name].filter(Boolean).join(' ')
347 }
348 case 'Bash': {
349 const cmd = String(e.command ?? '').split('\n')[0]!.trim()
350 // git and gh first: a commit message may say "test".
351 if (/^git\s+\S+/.test(cmd)) return 'git ' + cmd.split(/\s+/)[1]
352 if (/^gh\s+pr\b/.test(cmd)) return 'gh pr ' + (cmd.split(/\s+/)[2] ?? '')
353 if (/\b(tsc|test|vitest|jest|pytest)\b/.test(cmd)) return 'running tests'
354 return 'running a command'
355 }
356 default: return tool
357 }
358}
359
360// Running time as the native subagent row has it: 12s, 1m 43s, 1h 5m.
361export function elapsed(ms: number): string {
362 const s = Math.max(0, Math.floor(ms / 1000))
363 if (s < 60) return `${s}s`
364 const m = Math.floor(s / 60)
365 return m < 60 ? `${m}m ${s % 60}s` : `${Math.floor(m / 60)}h ${m % 60}m`
366}
367
368// "↓ 84.0k tokens": the context the agent has taken in, one decimal like the native row.
369export function tokensDown(n: number): string {
370 return `↓ ${n >= 1_000_000 ? `${(n / 1_000_000).toFixed(1)}M` : n >= 1000 ? `${(n / 1000).toFixed(1)}k` : n} tokens`
371}
372
373function labelOf(a: AgentRow): string {
374 return a.name ?? a.description ?? a.id.slice(0, 8)
375}
376
377const DEFAULT_WINDOW = 200_000
378const METER_CELLS = 12
379const CARD_ROWS = 5
380const DANGER_PERCENT = 90
381
382const LARGE_WINDOW = 1_000_000
383
384// The model id with what differs between the spellings of one model dropped: `[1m]`, a date suffix, `-latest`.
385// `e.model` is the id the engine resolved (it may carry `[1m]`); `usage.model` is the id the API reports.
386function modelKey(model: string): string {
387 return model.toLowerCase().replace(/\[1m\]/g, '').replace(/-latest$/, '').replace(/-\d{8}$/, '')
388}
389
390// A subagent reports no window of its own: borrow the main session's when it runs the same model
391// (compared by modelKey), else 200k, or 1M for a `[1m]` model id. Tokens past the assumed window
392// prove the window is the 1M one, so a guess is never shown past 100%.
393// `spawned` is the window the agent was started with (see the agent.spawn hook); it wins over the guess.
394function windowOf(model: string, mainModel: string | undefined, mainWindow: number | undefined, tokens = 0, spawned?: number): number {
395 const window = spawned !== undefined ? spawned : mainWindow !== undefined && mainModel !== undefined && modelKey(model) === modelKey(mainModel) ? mainWindow
396 : model.includes('[1m]') ? LARGE_WINDOW : DEFAULT_WINDOW
397 return tokens > window ? Math.max(window, LARGE_WINDOW) : window
398}
399
400
401// The percent in force: a 1M window has its own, so the 200k percent never applies to it.
402function windowPercent(s: WarnSettings, window: number): number {
403 return window >= LARGE_WINDOW ? s.contextWarn1m : s.contextWarn
404}
405
406// The trigger: the percent of the window, capped by the token setting when that is set.
407function thresholdOf(s: WarnSettings, window: number): number {
408 const byPercent = Math.round(window * windowPercent(s, window) / 100)
409 return s.contextWarnTokens > 0 ? Math.min(s.contextWarnTokens, byPercent) : byPercent
410}
411
412// The same threshold as a percent of the window, for the meter's marker and colour.
413function warnPercent(s: WarnSettings, window: number): number {
414 return Math.min(100, Math.max(1, Math.round(thresholdOf(s, window) / window * 100)))
415}
416
417// What the threshold is called in a message: "40%" or "100k tokens".
418function limitLabel(s: WarnSettings, window: number): string {
419 return s.contextWarnTokens > 0 && s.contextWarnTokens < Math.round(window * windowPercent(s, window) / 100)
420 ? `${tokensLabel(s.contextWarnTokens)} tokens` : `${windowPercent(s, window)}%`
421}
422
423// The threshold in tokens for the meter, only when the token limit is the one in force.
424function limitTokens(s: WarnSettings, window: number): string | undefined {
425 return limitLabel(s, window).endsWith('tokens') ? tokensLabel(s.contextWarnTokens) : undefined
426}
427
428// The line an agent past the threshold reads. A manager keeps going until its workers are done.
429function wrapUpText(role: string, percent: number, limit: string): string {
430 const handoff = 'follow the Handoff section of your instructions now: commit WIP, push, update the draft PR\'s ## Handoff note, end with HANDOFF: <branch>.'
431 return role === 'manager'
432 ? `flow: WRAP UP: your context is at ${percent}% (past the ${limit} limit). Start no new work; once none of your workers is running, ${handoff}`
433 : `flow: WRAP UP: your context is at ${percent}% (past the ${limit} limit). Stop new work, ${handoff}`
434}
435
436function tokensLabel(n: number): string {
437 if (n >= 1_000_000) return `${+(n / 1_000_000).toFixed(1)}M`
438 return n >= 1000 ? `${Math.round(n / 1000)}k` : String(n)
439}
440
441// The bar with a marker cell at the warn threshold, so how close an agent is stays readable.
442function cells(percent: number, warn: number): { ch: string; kind: 'mark' | 'fill' | 'empty' }[] {
443 const filled = Math.round(Math.min(100, percent) / 100 * METER_CELLS)
444 const mark = Math.min(METER_CELLS - 1, Math.floor(warn / 100 * METER_CELLS))
445 return Array.from({ length: METER_CELLS }, (_, i) =>
446 i === mark ? { ch: '│', kind: 'mark' } : i < filled ? { ch: '█', kind: 'fill' } : { ch: '░', kind: 'empty' })
447}
448
449function meterColor(percent: number, warn: number): string | undefined {
450 return percent >= DANGER_PERCENT && percent >= warn ? 'error' : percent >= warn ? 'warning' : undefined
451}
452
453function rank(status: string): number {
454 const i = ORDER.indexOf(status)
455 return i === -1 ? ORDER.length : i
456}
457
458// An agent told to wrap up sorts with what needs a person, ahead of plain running.
459function rankOf(a: AgentRow, acts: Record<string, Activity>): number {
460 return handoffOf(a, acts[a.id])?.kind === 'wrapping' ? Math.min(rank(a.status), rank('idle')) : rank(a.status)
461}
462
463// Where an agent is in its handoff: told to wrap up and still live, or ended (idle counts) on a
464// report whose last line is `HANDOFF: <branch>` (a manager's is `HANDOFF: manager <name>`).
465// Old rows without the fields give undefined.
466export function handoffOf(a: AgentRow, act: Activity | undefined):
467 { kind: 'wrapping'; percent?: number; reminders: number } | { kind: 'done'; to: string } | undefined {
468 if (act === undefined) return undefined
469 if (ENDED.has(a.status) || a.status === 'idle') {
470 const last = (act.answer ?? '').trim().split('\n').pop()?.trim() ?? ''
471 if (last.startsWith('HANDOFF:')) return { kind: 'done', to: last.slice('HANDOFF:'.length).trim() }
472 }
473 if (!ENDED.has(a.status) && act.handoffNotifiedAt !== undefined) {
474 return { kind: 'wrapping', percent: act.handoffPercent, reminders: act.remindersSent ?? 0 }
475 }
476 return undefined
477}
478
479function handoffText(h: NonNullable<ReturnType<typeof handoffOf>>): string {
480 return h.kind === 'done'
481 ? `handed off${h.to ? ` → ${h.to}` : ''}`
482 : `wrapping up${h.percent === undefined ? '' : ` (told at ${h.percent}%)`}${h.reminders > 1 ? ` · ${h.reminders} reminders` : ''}`
483}
484
485// The last batch the release tool cut, so a retried push does not release twice.
486let lastRelease: { key: string; version: string } | undefined
487
488// Tag the release commit, push the tag and create the GitHub Release. Every step is idempotent and retried, and
489// a failure is a returned line, never a throw: the release commit is already pushed and the batch stays done.
490const PUBLISH_ATTEMPTS = 4
491async function publishRelease($: EngineInterface, settings: Settings, dir: string, wanted: string | undefined): Promise<string> {
492 if (!settings.release) return 'Refused: the release setting is off, so nothing is published.'
493 if (!settings.releaseGithub) return 'Refused: the release_github setting is off, so nothing is published.'
494 if (!dir.startsWith('/')) return 'Refused: dir must be the absolute path of your worktree.'
495 let version = wanted
496 if (version === undefined) {
497 let last = lastRelease
498 if (last === undefined) {
499 const stateD = await stateDir($)
500 const disk = stateD === undefined ? undefined : await readJson($, `${stateD}/release.json`) as { key?: string; version?: string } | undefined
501 if (typeof disk?.version === 'string') last = { key: String(disk.key ?? ''), version: disk.version }
502 }
503 version = last?.version
504 }
505 if (version === undefined) return 'Refused: no version to publish; pass version or cut a release first.'
506 const tag = `v${version}`
507 const git = (...a: string[]) => $.process.run(['git', '-C', dir, ...a])
508 const tail = (r: { stderr: string; stdout: string }) => (r.stderr.trim() || r.stdout.trim()).split('\n').slice(-3).join(' ').slice(0, 300)
509 const run = async (step: string, argv: string[]): Promise<{ ok: true } | { ok: false; line: string }> => {
510 let last = ''
511 for (let i = 0; i < PUBLISH_ATTEMPTS; i++) {
512 if (i > 0) await $.clock.sleep(2000 * 2 ** (i - 1))
513 try {
514 const r = await $.process.run(argv, { cwd: dir })
515 if (r.exitCode === 0) return { ok: true }
516 last = tail(r)
517 } catch (err) { last = err instanceof Error ? err.message : String(err) }
518 }
519 return { ok: false, line: `Not published: ${step} failed: ${last}. The release commit is pushed; report this, the batch stays done.` }
520 }
521 try {
522 // The release commit may sit below a merge commit (a rejected push is fetched, merged and pushed again).
523 const log = await git('log', '--first-parent', '-n', '50', '--format=%H %s')
524 if (log.exitCode !== 0) return `Not published: could not read the log of ${dir}: ${tail(log)}. The release commit is pushed; report this, the batch stays done.`
525 const sha = log.stdout.split('\n').find(l => l.slice(41) === `Release ${version}`)?.slice(0, 40)
526 if (sha === undefined) return `Refused: no commit "Release ${version}" in the last 50 first-parent commits of ${dir}. Publish only after the release commit is in HEAD.`
527 const anc = await git('merge-base', '--is-ancestor', sha, 'HEAD')
528 if (anc.exitCode !== 0) return `Refused: the commit "Release ${version}" is not an ancestor of HEAD in ${dir}.`
529 const existing = await git('rev-parse', '--verify', '--quiet', `refs/tags/${tag}^{commit}`)
530 if (existing.exitCode === 0 && existing.stdout.trim() !== sha) return `Not published: tag ${tag} already exists at ${existing.stdout.trim().slice(0, 8)}, not at the release commit ${sha.slice(0, 8)}; it was not moved. The release commit is pushed; report this, the batch stays done.`
531 if (existing.exitCode !== 0) {
532 const t = await git('tag', '-a', tag, '-m', `Release ${version}`, sha)
533 if (t.exitCode !== 0) return `Not published: tagging failed: ${tail(t)}. The release commit is pushed; report this, the batch stays done.`
534 }
535 const pushed = await run('pushing the tag', ['git', '-C', dir, 'push', 'origin', tag])
536 if (!pushed.ok) return pushed.line
537 const view = await $.process.run(['gh', 'release', 'view', tag], { cwd: dir })
538 if (view.exitCode === 0) return `Published ${tag}: tag pushed, GitHub Release already existed.`
539 const logName = settings.changelogFile || 'CHANGELOG.md'
540 const section = await $.fs.exists(`${dir}/${logName}`) ? changelogSection(await $.fs.read(`${dir}/${logName}`), version) : undefined
541 const stateD = await stateDir($)
542 const notes = `${stateD ?? '/tmp'}/release-notes-${tag}.md`
543 await $.fs.write(notes, section !== undefined && section !== '' ? section + '\n' : `Release ${version}\n`)
544 const made = await run('gh release create', ['gh', 'release', 'create', tag, '--title', tag, '--notes-file', notes, '--verify-tag'])
545 if (!made.ok) return made.line
546 const url = await $.process.run(['gh', 'release', 'view', tag, '--json', 'url', '--jq', '.url'], { cwd: dir })
547 return `Published ${tag}: tag pushed, GitHub Release created${url.exitCode === 0 && url.stdout.trim() !== '' ? ' ' + url.stdout.trim() : ''}.${section === undefined || section === '' ? ` The changelog has no ${version} section, so the notes are "Release ${version}".` : ''}`
548 } catch (err) {
549 return `Not published: ${err instanceof Error ? err.message : String(err)}. The release commit is pushed; report this, the batch stays done.`
550 }
551}
552
553// How many managers main runs at once; refresh needs it to tell main how many slots are free.
554let maxManagers = 20
555// Managers whose end already freed a slot (and woke main), so a poll does not wake it twice.
556const freedSeen = new Set<string>()
557
558const isManagerOf = (owner: string, name: string | undefined) =>
559 name !== undefined && (name === owner || (name.startsWith(owner + '-') && /^\d+$/.test(name.slice(owner.length + 1))))
560
561// Whose graph an agent works on: its own name, or the manager it continues (`<name>-N`); main has none.
562function planOwner(plans: Plan, name: string | undefined): string {
563 if (name === undefined) return 'main'
564 if (plans[name]) return name
565 const m = /^(.+)-\d+$/.exec(name)
566 return m && plans[m[1]!] ? m[1]! : name
567}
568
569const liveManagers = (rows: AgentRow[]) => rows.filter(a => a.type === MANAGER && !ENDED.has(a.status))
570
571// Settles every plan against the roster and handovers in one update, then tells each owner once.
572// `edit` changes the plan first (the plan tool); `quietOwner` is the caller, who sees the outcome
573// in its own tool result and needs no message.
574async function syncPlans(
575 $: EngineInterface, rows: AgentRow[], edit?: (p: Plan) => Plan, quietOwner?: string,
576): Promise<Plan> {
577 // Marked before the first await, so two overlapping polls count an ended manager once.
578 const freed = rows.filter(a => a.type === MANAGER && ENDED.has(a.status) && !freedSeen.has(a.id))
579 for (const a of freed) freedSeen.add(a.id)
580 if (edit === undefined && Object.keys(await read($, plan)).length === 0) return {}
581 const [acts, hs] = await Promise.all([read($, activity), read($, handovers)])
582 const agents: AgentFact[] = rows.map(a => ({
583 id: a.id, name: a.name, status: a.status, answer: acts[a.id]?.answer,
584 children: rows.filter(c => c.parentId === a.id && LIVE_STATUS.has(c.status)).length,
585 at: acts[a.id]?.lastAt,
586 childAt: Math.max(0, ...rows.filter(c => c.parentId === a.id).map(c => acts[c.id]?.lastAt ?? 0)),
587 }))
588 const known = { ...(await read($, seenAgents)) }
589 const present = new Set(rows.map(a => a.name))
590 for (const a of agents) if (a.name !== undefined) known[a.name] = a
591 for (const [name, a] of Object.entries(known)) if (!present.has(name)) agents.push(a)
592 if (JSON.stringify(known) !== JSON.stringify(await read($, seenAgents))) await update($, seenAgents, () => known)
593 const owners = branchOwners(await readLog($))
594 const facts: Facts = {
595 agents, owners,
596 handovers: Object.values(hs),
597 phrases: decisionPhrases,
598 asking: askingNames(await read($, inbox)),
599 }
600 const slots = Math.max(0, maxManagers - liveManagers(rows).length)
601 let notices: Notice[] = []
602 let result: Plan = {}
603 await update($, plan, p => {
604 const r = settle(edit ? edit(p) : p, facts, { slots })
605 notices = r.notices
606 result = r.plan
607 // A slot freed while ready tasks wait: tell main again, once per ended manager.
608 const waiting = Object.values(r.plan['main'] ?? {}).filter(n => n.state === 'ready')
609 if (freed.length > 0 && slots > 0 && waiting.length > 0 && !notices.some(n => n.owner === 'main')) {
610 notices = [...notices, { owner: 'main', ready: waiting, blocked: [], done: [], slots }]
611 }
612 return r.plan
613 })
614 for (const n of notices) {
615 if (n.owner === quietOwner) continue
616 const text = noticeText(n)
617 if (n.owner === 'main') {
618 $.clock.after(0, () => void $.prompt.submit({ text }).catch(() => undefined))
619 } else {
620 const agent = agentFor(n.owner, rows.filter(a => !ENDED.has(a.status)))
621 if (agent) await deliver($, agent.id, text).catch(() => undefined)
622 else {
623 // The owner ended its turn to wait (the host reports it completed or drops it): a ready node means it
624 // has work left, so the notice wakes it. Only a manager that cannot be resumed sends it to main.
625 const gone = agentFor(n.owner, agents.filter(a => a.status !== 'failed' && a.status !== 'killed'))
626 if (gone?.id !== undefined) {
627 resumable.add(gone.id)
628 await deliver($, gone.id, text, { onGone: t => tellMain($, n.owner, t) }).catch(() => undefined)
629 } else if (agentFor(n.owner, agents) !== undefined) await tellMain($, n.owner, text)
630 }
631 }
632 }
633 return result
634}
635
636function limitsLine(rows: AgentRow[], workers: number): string {
637 return `Limits: managers ${liveManagers(rows).length}/${maxManagers}, workers per manager ${workers}`
638}
639
640// The roster as the board shows it: every agent of the session, with its role.
641async function refresh($: EngineInterface): Promise<AgentRow[]> {
642 const [list, t, before, acts] = await Promise.all([
643 $.agent.list(), $.clock.now(), read($, roster), read($, activity),
644 ])
645 const rows: AgentRow[] = list.map(a => ({
646 id: a.id, name: a.name, description: a.description, type: a.type,
647 status: a.status, parentId: a.parentId,
648 }))
649 // Workers in other harnesses stand in the tree like agents, under the manager that started them.
650 rows.push(...Object.values(await read($, sessions)).map(sessionRow))
651 const was = new Map(before.map(a => [a.id, a.status]))
652 const ended: string[] = []
653 const askers = askingNames(await read($, inbox))
654 for (const a of rows) {
655 const prev = was.get(a.id)
656 // The reviewer's worktree (detached at the base) goes once the reviewer is gone, also when
657 // its ending was not seen as a transition (a reload left the roster record empty).
658 if (ENDED.has(a.status) && isReviewer(a.type) && (prev === undefined || (prev !== a.status && !ENDED.has(prev)))) autoSweep($)
659 if (prev === undefined || prev === a.status || ENDED.has(prev)) continue
660 if (ENDED.has(a.status) && acts[a.id] !== undefined) ended.push(a.id)
661 if (ENDED.has(a.status) || a.status === 'idle') {
662 const asks = asksQuestion(acts[a.id]?.answer, decisionPhrases) || askers.includes(a.name ?? '')
663 const role = ROLE[a.type] ? `${ROLE[a.type]} ` : ''
664 void $.ui.toast(`${role}${labelOf(a)}: ${asks ? 'asks a question' : a.status === 'idle' ? 'finished its turn' : a.status}`)
665 }
666 }
667 if (ended.length > 0) {
668 await update($, activity, all => {
669 const next = { ...all }
670 for (const id of ended) if (next[id] !== undefined && next[id]!.endedAt === undefined) next[id] = { ...next[id]!, endedAt: t }
671 return next
672 })
673 }
674 if (JSON.stringify(rows) !== JSON.stringify(before)) await update($, roster, () => rows)
675 await update($, now, () => t)
676 // Free the test slots of ended agents and pass them on.
677 await settleSlots($, rows, t)
678
679 const hs = Object.values(await read($, handovers))
680 const live = rows.filter(a => !ENDED.has(a.status))
681 const count = (type: string) => live.filter(a => a.type === type || (type === WORKER && WORKERS.has(a.type))).length
682 const queued = hs.filter(h => h.status === 'pending' || h.status === 'taken').length
683 const awaiting = hs.filter(h => h.status === 'awaiting').length
684 const parts = [
685 count(MANAGER) && `${count(MANAGER)} managers`,
686 count(WORKER) && `${count(WORKER)} workers`,
687 queued && `reviewer: ${queued} PR${queued > 1 ? 's' : ''}`,
688 awaiting && `${awaiting} awaiting approval`,
689 ].filter(Boolean)
690 const unhanded = (await currentUnhanded($)).length
691 if (unhanded) parts.push(`${unhanded} unhanded`)
692 const slots = await read($, testSlots)
693 const slotPart = slots.holders.length || slots.waiters.length ? ` · tests ${slots.holders.length}/${slotLimit}` : ''
694 $.ui.status(rows.length === 0 && hs.length === 0 && !unhanded ? undefined
695 : `flow: ${parts.length ? parts.join(' · ') : `${live.length} live`}${slotPart} · /flow`)
696 await syncPlans($, rows).catch(() => undefined)
697 await preflightTick($, rows).catch(() => undefined)
698 return rows
699}
700
701async function openPane($: EngineInterface): Promise<void> {
702 const r = await $.ui.open({ id: PANE, title: 'Flow' })
703 if (!r.isPlaced) void $.ui.toast('flow is running agents: type /flow to watch them')
704}
705
706// flow's state on disk: <git-common-dir>/flow/. Shared by every worktree of the repo, never part
707// of a working tree. Every helper is best-effort: outside a repo, or when a write fails, it
708// does nothing (a failure goes to $.ui.log) and never throws, so a tool call or turn never fails
709// over state. `config.json` in this dir belongs to the settings package: never touched here.
710
711const TEXT_MAX = 300
712const NOTES_MAX = 3000
713
714let cached: string | undefined
715
716// Forget the resolved dir; a new session resolves it again.
717function resetStateDir(): void {
718 cached = undefined
719}
720
721async function warn($: EngineInterface, what: string, err: unknown): Promise<void> {
722 try {
723 await $.ui.log(`flow state: ${what}: ${err instanceof Error ? err.message : String(err)}`)
724 } catch {
725 // No log to write to.
726 }
727}
728
729// The state dir, or undefined when this is not a repo (or the answer is no absolute path).
730async function stateDir($: EngineInterface): Promise<string | undefined> {
731 if (cached !== undefined) return cached
732 try {
733 const r = await $.process.run(['git', 'rev-parse', '--path-format=absolute', '--git-common-dir'])
734 const out = r.stdout.trim()
735 if (r.exitCode !== 0 || !out.startsWith('/') || out.includes('\n')) return undefined
736 cached = `${out.replace(/\/+$/, '')}/flow`
737 return cached
738 } catch {
739 return undefined
740 }
741}
742
743// Write a temp file next to the target, then rename it over: a reader sees the old file or the new one.
744async function writeJsonAtomic($: EngineInterface, path: string, obj: unknown): Promise<boolean> {
745 try {
746 const tmp = `${path}.${Date.now()}.tmp`
747 await $.fs.write(tmp, `${JSON.stringify(obj, null, 2)}\n`)
748 const r = await $.process.run(['mv', tmp, path])
749 if (r.exitCode !== 0) throw new Error(r.stderr.trim() || 'mv failed')
750 return true
751 } catch (err) {
752 await warn($, `writing ${path}`, err)
753 return false
754 }
755}
756
757// The parsed file, or undefined when it is missing, unreadable or not JSON.
758async function readJson($: EngineInterface, path: string): Promise<unknown> {
759 try {
760 return JSON.parse(await $.fs.read(path))
761 } catch {
762 return undefined
763 }
764}
765
766// --- Cost ledger ---
767// Loaded from disk once, before the first change, so what this run counts is added to the history.
768let ledgerLoad: Promise<void> | undefined
769let ledgerSaveDue = false
770
771function loadLedger($: EngineInterface): Promise<void> {
772 ledgerLoad ??= (async () => {
773 try {
774 const dir = await stateDir($)
775 if (dir === undefined) return
776 const disk = normalizeLedger(await readJson($, `${dir}/ledger.json`))
777 const now = await $.clock.now()
778 await update($, ledger, cur => pruneLedger(mergeLedgers(disk, cur), now))
779 } catch {
780 // No history: count from here.
781 }
782 })()
783 return ledgerLoad
784}
785
786async function saveLedger($: EngineInterface): Promise<void> {
787 try {
788 const dir = await stateDir($)
789 if (dir === undefined) return
790 await $.process.run(['mkdir', '-p', dir])
791 await writeJsonAtomic($, `${dir}/ledger.json`, pruneLedger(await read($, ledger), await $.clock.now()))
792 } catch {
793 // The meter is best-effort.
794 }
795}
796
797// Every step changes the ledger; the file is written at most every few seconds.
798async function changeLedger($: EngineInterface, fn: (l: Ledger, now: number) => Ledger): Promise<void> {
799 await loadLedger($)
800 const now = await $.clock.now()
801 await update($, ledger, l => fn(l, now))
802 if (ledgerSaveDue) return
803 ledgerSaveDue = true
804 $.clock.after(3000, () => { ledgerSaveDue = false; void saveLedger($) })
805}
806
807// Main's entry is per session: the ledger outlives the session, and its steps must not merge with an earlier one's.
808async function mainKey($: EngineInterface): Promise<string> {
809 const u = await $.session.usage().then(x => x, () => undefined)
810 return `main@${u?.startedAt ?? 0}`
811}
812
813// The cost block of status, or nothing before the first counted step.
814async function costLines($: EngineInterface, prs: { pr: number; branch: string }[], keep?: (e: Ledger[string]) => boolean): Promise<string[]> {
815 try {
816 await loadLedger($)
817 const start = (await $.session.usage().then(x => x, () => undefined))?.startedAt ?? 0
818 const live = new Set((await read($, roster)).map(a => a.id))
819 const all = await read($, ledger)
820 const led = keep === undefined ? all : Object.fromEntries(Object.entries(all).filter(([, e]) => keep(e)))
821 return costBlock(led, prs, Object.keys(await read($, sessions)).length > 0, { start, live })
822 } catch {
823 return []
824 }
825}
826
827// Inbox changes run one after another: two answers at once must not each start from the same file.
828let inboxChain: Promise<unknown> = Promise.resolve()
829
830// Read the inbox from disk, let fn change it, write it back and set the atom. fn returns the new inbox (the
831// same object when nothing changed) and whatever the caller wants back.
832function withInbox<T>($: EngineInterface, fn: (cur: Inbox) => { inbox: Inbox; out: T }): Promise<T> {
833 const run = async (): Promise<T> => {
834 const dir = await stateDir($)
835 const cur = dir === undefined ? await read($, inbox) : normalizeInbox(await readJson($, `${dir}/inbox.json`))
836 const { inbox: next, out } = fn(cur)
837 if (next !== cur) {
838 if (dir !== undefined) {
839 await $.process.run(['mkdir', '-p', dir])
840 await writeJsonAtomic($, `${dir}/inbox.json`, next)
841 }
842 await update($, inbox, () => next)
843 }
844 return out
845 }
846 const result = inboxChain.then(run, run)
847 inboxChain = result.catch(() => undefined)
848 return result
849}
850
851// Pre-flight changes run one after another, like the inbox's.
852let preflightChain: Promise<unknown> = Promise.resolve()
853
854function withPreflight<T>($: EngineInterface, fn: (cur: Preflight) => { state: Preflight; out: T }): Promise<T> {
855 const run = async (): Promise<T> => {
856 const dir = await stateDir($)
857 const cur = dir === undefined ? await read($, preflight) : normalizePreflight(await readJson($, `${dir}/preflight.json`))
858 const { state, out } = fn(cur)
859 if (state !== cur) {
860 if (dir !== undefined) {
861 await $.process.run(['mkdir', '-p', dir])
862 await writeJsonAtomic($, `${dir}/preflight.json`, state)
863 }
864 await update($, preflight, () => state)
865 }
866 return out
867 }
868 const result = preflightChain.then(run, run)
869 preflightChain = result.catch(() => undefined)
870 return result
871}
872
873// Check changes run one after another, like the inbox's.
874let checksChain: Promise<unknown> = Promise.resolve()
875
876function withChecks<T>($: EngineInterface, fn: (cur: Checks) => { checks: Checks; out: T }): Promise<T> {
877 const run = async (): Promise<T> => {
878 const dir = await stateDir($)
879 const cur = dir === undefined ? await read($, checks) : normalizeChecks(await readJson($, `${dir}/checks.json`))
880 const { checks: next, out } = fn(cur)
881 if (next !== cur) {
882 if (dir !== undefined) {
883 await $.process.run(['mkdir', '-p', dir])
884 await writeJsonAtomic($, `${dir}/checks.json`, next)
885 }
886 await update($, checks, () => next)
887 }
888 return out
889 }
890 const result = checksChain.then(run, run)
891 checksChain = result.catch(() => undefined)
892 return result
893}
894
895// The version the running plugin was installed as: its plugin.json next to the module. Undefined when unreadable.
896async function installedVersion($: EngineInterface): Promise<string | undefined> {
897 try {
898 const v = (JSON.parse(await $.fs.read(`${$.plugin.root}/.claude-plugin/plugin.json`)) as { version?: unknown }).version
899 return typeof v === 'string' && /^\d+\.\d+\.\d+/.test(v) ? v : undefined
900 } catch {
901 return undefined
902 }
903}
904
905// The version a PR's merge needs installed: plugin.json at the merged sha, else one named in the steps.
906async function versionAt($: EngineInterface, sha: string | undefined, steps: string): Promise<string | undefined> {
907 if (sha) {
908 try {
909 const r = await $.process.run(['git', 'show', `${sha}:.claude-plugin/plugin.json`])
910 const v = r.exitCode === 0 ? (JSON.parse(r.stdout) as { version?: unknown }).version : undefined
911 if (typeof v === 'string' && v !== '') return v
912 } catch { /* not a plugin repo, or the sha is unknown */ }
913 }
914 return versionInSteps(steps)
915}
916
917// Adds a check for each needs-a-person line in a done PR's report (the same PR and steps never twice). Live: the
918// handover's verify_command rides along, unless verify_paths says no changed file warrants it. Returns the new ones.
919async function captureChecks($: EngineInterface, h: Handover, verifyPaths: string[], live: boolean): Promise<Check[]> {
920 const needs = parseNeeds(h.report)
921 if (needs.length === 0) return []
922 let command = live && h.verifyCommand ? h.verifyCommand : undefined
923 let note: string | undefined
924 if (command !== undefined && verifyPaths.length > 0) {
925 // gh failing runs the command anyway: fail toward verifying.
926 const r = await $.process.run(['gh', 'pr', 'view', String(h.pr), '--json', 'files'])
927 let files: string[] | undefined
928 try { files = r.exitCode === 0 ? (JSON.parse(r.stdout) as { files: { path: string }[] }).files.map(f => f.path) : undefined } catch { files = undefined }
929 if (files !== undefined && !anyMatch(verifyPaths, files)) { command = undefined; note = SKIP_NOTE }
930 }
931 const t = live ? await $.clock.now() : h.at
932 const items = await Promise.all(needs.map(async n => ({
933 steps: n.steps, version: await versionAt($, h.sha, n.steps),
934 ...(n.pr === h.pr && command !== undefined ? { verifyCommand: command } : {}), ...(n.pr === h.pr && note !== undefined ? { note } : {}),
935 })))
936 return withChecks($, cur => {
937 const r = addChecks(cur, h.pr, h.title, h.sha, items, t)
938 return { checks: r.checks, out: r.added }
939 })
940}
941
942// Managers that ended with no live manager of the same notes key (a successor is the same manager).
943function endedManagers(rows: AgentRow[]): Set<string> {
944 const live = new Set(liveManagers(rows).map(a => noteKey(a.name ?? '')))
945 return new Set(rows.filter(a => a.type === MANAGER && ENDED.has(a.status) && !live.has(noteKey(a.name ?? ''))).map(a => noteKey(a.name ?? '')))
946}
947
948// Sends the round to main once it is due: everyone filed, skipped or ended, or the wait is over. The
949// delivered flag is written inside the update, so two callers deliver it once.
950async function preflightTick($: EngineInterface, rows: AgentRow[]): Promise<void> {
951 if (!(await read($, preflight)).rounds.some(r => !r.delivered)) return
952 const t = await $.clock.now()
953 const ended = endedManagers(rows)
954 const due = await withPreflight($, cur => {
955 const d = dueRound(cur, ended, t, preflightWaitMs)
956 if (d === undefined) return { state: cur, out: undefined }
957 const next = markDelivered(cur, d.round.id, t)
958 return { state: next, out: { ...d, state: next } }
959 })
960 if (due === undefined) return
961 const members = due.round.members.map(k => due.state.entries[k])
962 // A round of skipped managers only has nothing to say.
963 if (members.every(e => e === undefined || e.phase === 'skipped')) return
964 const text = renderRound(due.state, await read($, inbox), due.round, ended, t, due.timedOut)
965 $.clock.after(0, () => void $.prompt.submit({ text }).catch(() => undefined))
966}
967
968// Why this manager may not start workers yet, or undefined. Only a recorded manager is ever gated.
969async function preflightGate($: EngineInterface, agentId: string | undefined): Promise<string | undefined> {
970 if (!preflightOn || agentId === undefined) return undefined
971 const me = (await $.agent.list()).find(a => a.id === agentId)
972 if (me?.type !== MANAGER || me.name === undefined) return undefined
973 const g = gateOf(await read($, preflight), await read($, inbox), me.name)
974 return g === undefined ? undefined : denyText(g)
975}
976
977// A manager main just started joins the open round, or opens one with a timer for the wait.
978async function recordManager($: EngineInterface, name: string, prompt: string): Promise<void> {
979 const t = await $.clock.now()
980 const opened = await withPreflight($, cur => {
981 const next = recordSpawn(cur, name, isSkip(prompt), t, preflightWaitMs)
982 return { state: next, out: next.rounds.length > cur.rounds.length }
983 })
984 if (opened) $.clock.after(preflightWaitMs, () => void refresh($).catch(() => undefined))
985}
986
987const today = async ($: EngineInterface) => new Date(await $.clock.now()).toISOString().slice(0, 10)
988
989// The one way a question gets answered: marks it, tells the asker when that is needed, notes the decision.
990// choice null takes the question's default. Returns one result line.
991// user: the person answering from the pane or a /flow command, who may answer any question, a manager's included.
992async function answerQuestion($: EngineInterface, wanted: string, choice: string | null, by: string, options: Record<string, unknown>, user = false): Promise<string> {
993 const at = await $.clock.now()
994 // A decision filed before decisions had d-ids is also found by its old q-id.
995 const id = findItem(await read($, inbox), wanted)?.id ?? wanted.trim().toLowerCase()
996 const marked = await withInbox($, (cur): { inbox: Inbox; out: Marked | { kind: 'escalated'; q: Question } } => {
997 // A standing rule made this the user's decision: only main answers it.
998 const held = cur.items.find(x => x.id === id)
999 if (held !== undefined && held.state === 'open' && held.escalated !== undefined && by !== 'main') return { inbox: cur, out: { kind: 'escalated', q: held } }
1000 const m = markAnswered(cur, id, choice, by, at, undefined, user)
1001 return { inbox: m.kind === 'ok' ? m.inbox : cur, out: m }
1002 })
1003 if (marked.kind === 'escalated') return `${id}: refused, standing rule ${marked.q.escalated} makes this the user's decision. It is in main's inbox as ${id}; main answers it directly and the asker gets the answer. Don't answer or re-ask it.`
1004 if (marked.kind === 'unknown') return `${id}: no such question.`
1005 if (marked.kind === 'answered') return `${id}: already answered ("${marked.q.answer ?? ''}" by ${marked.q.answeredBy ?? '?'}).`
1006 if (marked.kind === 'refused') return `${id}: refused, it is addressed to ${marked.q.addressee}, not ${by}.`
1007 if (marked.kind === 'empty') return `${id}: refused, the choice is empty.`
1008 const { q, answer, isDefault } = marked
1009 // The user over a manager's head: the asker hears "the user", and the manager learns the question is closed.
1010 const over = user && overridesManager(q)
1011 const who = over ? 'the user' : by
1012 const deployNote = q.kind === 'deploy' ? await onDeployAnswer($, q, answer) : q.kind === ENV_KIND ? await onEnvAnswer($, q, answer) : q.kind === PUSH_KIND ? await onPushAnswer($, q, answer, by) : ''
1013 let delivered = true
1014 let hint = ''
1015 if (needsMessage(q, isDefault)) {
1016 const asker = q.askerId === undefined ? undefined : (await $.agent.list()).find(a => a.id === q.askerId)
1017 delivered = false
1018 const alive = (a: AgentInfo | undefined): a is AgentInfo => a !== undefined && (LIVE.has(a.status) || a.status === 'idle')
1019 if (alive(asker)) {
1020 delivered = await deliver($, asker.id, answerMessage(q, answer, who, isDefault), { urgent: q.blocking === true, onGone: () => undefined }).catch(() => false)
1021 } else if (isFyi(q) && q.addressee !== 'main') {
1022 // A finished owner cannot act on an undo: its manager (the notes owner) gets it, naming the owner.
1023 const mgr = (await $.agent.list()).find(a => a.type === MANAGER && a.name !== undefined && noteKey(a.name) === noteKey(q.addressee) && alive(a))
1024 if (mgr !== undefined) {
1025 delivered = await deliver($, mgr.id, `${answerMessage(q, answer, by, isDefault)} (This was ${q.owner}'s decision; ${q.owner} is no longer running, so act on it yourself or with a new worker.)`, { urgent: q.blocking === true, onGone: () => undefined }).catch(() => false)
1026 }
1027 }
1028 if (!delivered) hint = `; ${q.owner} is gone: main should relay it to ${noteKey(q.owner)}-2`
1029 }
1030 await withInbox($, cur => ({
1031 inbox: { ...cur, items: cur.items.map(x => (x.id === id ? { ...x, delivered } : x)) }, out: undefined,
1032 }))
1033 if (over) {
1034 const mgr = (await $.agent.list()).find(a => a.type === MANAGER && a.name !== undefined && noteKey(a.name) === noteKey(q.addressee) && (LIVE.has(a.status) || a.status === 'idle'))
1035 if (mgr !== undefined) {
1036 await deliver($, mgr.id, `flow: the user answered ${q.id}, which ${q.owner} asked you: ${answer}. It is closed; don't answer it.`, { urgent: false, onGone: () => undefined }).catch(() => false)
1037 }
1038 }
1039 if (!(isFyi(q) && isDefault)) await appendNote($, notesOwner(q), `- ${await today($)} decision: "${q.id} ${noteText(q)}: ${isFyi(q) ? `undone, ${answer}` : answer}"`)
1040 const guard = q.guard !== undefined && answer === ADD_OPTION ? ` ${await addGuardTest($, options, q.guard.glob, q.guard.test)}` : ''
1041 return `${id}: ${answer}${isDefault ? ' (default)' : ''}, ${delivered ? 'delivered' : 'undelivered'}${hint}.${guard}${deployNote}`
1042}
1043
1044// The pane's inbox rows, read fresh.
1045async function paneInboxRows($: EngineInterface): Promise<{ rows: PaneRow[]; now: number }> {
1046 const [box, now] = await Promise.all([read($, inbox), $.clock.now()])
1047 return { rows: paneRows(box), now }
1048}
1049
1050// Pane presses run one after another, so a second press sees the first one's answer and highlight.
1051let paneChain: Promise<unknown> = Promise.resolve()
1052const paneSerial = <T,>(fn: () => Promise<T>): Promise<T> => {
1053 const run = paneChain.then(fn, fn)
1054 paneChain = run.catch(() => undefined)
1055 return run
1056}
1057
1058// The question as a notes line: an env change names the variable, never its value.
1059const noteText = (q: Question): string => (q.env === undefined ? q.question : `env ${q.env.name} on ${q.env.target} (${q.env.role})`)
1060
1061// Deploy gate changes run one after another, like the inbox's.
1062let deploysChain: Promise<unknown> = Promise.resolve()
1063
1064// Read deploys.json, let fn change it, write it back atomically and set the atom. fn returns the new state
1065// (the same object when nothing changed) and whatever the caller wants back.
1066function withDeploys<T>($: EngineInterface, fn: (cur: Deploys) => { deploys: Deploys; out: T }): Promise<T> {
1067 const run = async (): Promise<T> => {
1068 const dir = await stateDir($)
1069 const cur = dir === undefined ? await read($, deploys) : normalizeDeploys(await readJson($, `${dir}/deploys.json`))
1070 const { deploys: next, out } = fn(cur)
1071 if (next !== cur) {
1072 if (dir !== undefined) {
1073 await $.process.run(['mkdir', '-p', dir])
1074 await writeJsonAtomic($, `${dir}/deploys.json`, next)
1075 }
1076 await update($, deploys, () => next)
1077 }
1078 return out
1079 }
1080 const result = deploysChain.then(run, run)
1081 deploysChain = result.catch(() => undefined)
1082 return result
1083}
1084
1085const setTarget = (d: Deploys, name: string, ts: TargetState): Deploys => ({ targets: { ...d.targets, [name]: ts } })
1086
1087// Commits behind the base for each target with a recorded deploy. One git call per target, never from a render.
1088let behindRun: Promise<void> | undefined
1089function refreshBehind($: EngineInterface): Promise<void> {
1090 behindRun ??= (async () => {
1091 const d = await read($, deploys)
1092 const out: Record<string, number> = {}
1093 for (const t of deployInfos) {
1094 const sha = d.targets[t.name]?.deployedSha
1095 if (sha === undefined) continue
1096 const r = await $.process.run(['git', 'rev-list', '--count', `${sha}..origin/${cleanupBase}`], { timeoutMs: 10_000 }).catch(() => undefined)
1097 const n = r !== undefined && r.exitCode === 0 ? Number(r.stdout.trim()) : NaN
1098 if (Number.isInteger(n)) out[t.name] = n
1099 }
1100 await update($, behind, () => out)
1101 })().catch(() => undefined).finally(() => { behindRun = undefined })
1102 return behindRun
1103}
1104
1105// The commits an approval would ship: one line each, at most 15.
1106async function commitsSince($: EngineInterface, last: string | undefined, sha: string): Promise<string[]> {
1107 if (last === undefined) return []
1108 const r = await $.process.run(['git', 'log', '--oneline', '-15', `${last}..${sha}`], { timeoutMs: 10_000 }).catch(() => undefined)
1109 return r !== undefined && r.exitCode === 0 ? r.stdout.split('\n').map(l => l.trim()).filter(Boolean) : []
1110}
1111
1112// The user's answer to a deploy approval item: record it on its target, and start a deploy-only run for "deploy".
1113async function onDeployAnswer($: EngineInterface, q: Question, answer: string): Promise<string> {
1114 const name = await withDeploys($, cur => {
1115 const hit = Object.entries(cur.targets).find(([, ts]) => ts.approval?.qid === q.id)
1116 if (hit === undefined) return { deploys: cur, out: undefined }
1117 return { deploys: setTarget(cur, hit[0], applyAnswer(hit[1], answer)), out: hit[0] }
1118 })
1119 if (name === undefined) return ''
1120 if (!isApprove(answer)) return ` ${name} stays behind; the next batch asks again.`
1121 return ` ${name} approved. ${await ensureQueue($)}`
1122}
1123
1124// The user's answer to an env item. An applied `apply` item records the change as done by the user. A
1125// positive answer that leaves a target with nothing open, on an auto target, starts a deploy-only run (the
1126// gate then applies what is approved); a negative one holds the target until main releases it.
1127async function onEnvAnswer($: EngineInterface, q: Question, answer: string): Promise<string> {
1128 const e = q.env
1129 if (e === undefined) return ''
1130 const yes = e.role === 'change' ? answer.trim().toLowerCase() === 'yes' : answer.trim().toLowerCase() === 'done'
1131 const at = await $.clock.now()
1132 const info = deployInfos.find(t => t.name === e.target)
1133 const hs = Object.values(await read($, handovers))
1134 const waiting = pendingEnv(hs, e.target, (await read($, deploys)).targets[e.target]).length > 0
1135 if (e.role === 'apply' && yes) {
1136 await withDeploys($, cur => ({ deploys: setTarget(cur, e.target, withEnvDone(cur.targets[e.target], [{ pr: e.pr, name: e.name, how: 'user', at }])), out: undefined }))
1137 }
1138 if (!yes) {
1139 return e.role === 'change'
1140 ? ` ${e.target} is held by this: the gate skips it until main runs release on ${e.target}, which drops the declined change.`
1141 : ` Not yet: ${e.target} keeps waiting; the next gate opens a fresh item for it.`
1142 }
1143 const ib = await read($, inbox)
1144 const stillOpen = ib.items.some(x => x.state === 'open' && x.kind === ENV_KIND && x.env?.target === e.target)
1145 if (!waiting || stillOpen || info === undefined || info.mode !== 'auto') return ''
1146 await withDeploys($, cur => ({ deploys: setTarget(cur, e.target, { ...(cur.targets[e.target] ?? {}), due: true }), out: undefined }))
1147 return ` ${e.target} has its env answers. ${await ensureQueue($)}`
1148}
1149
1150// A person's (or main's) hold on a target.
1151async function holdTarget($: EngineInterface, name: string, until: Hold['until'], by: string, reason: string | undefined): Promise<string> {
1152 const bad = unknownTarget(deployInfos, name)
1153 if (bad !== undefined) return bad
1154 const at = await $.clock.now()
1155 await withDeploys($, cur => {
1156 const hold: Hold = { until, by, at, ...(reason ? { reason } : {}) }
1157 return { deploys: setTarget(cur, name, withHold(cur.targets[name], hold)), out: undefined }
1158 })
1159 return until === 'batch'
1160 ? `${name} held for the next batch: the reviewer's next gate call for it answers Held and the hold ends. Release earlier with release.`
1161 : `${name} held until released: the reviewer skips it every batch. Release with release (/flow release ${name}).`
1162}
1163
1164async function releaseTarget($: EngineInterface, name: string): Promise<string> {
1165 const info = deployInfos.find(t => t.name === name)
1166 const bad = unknownTarget(deployInfos, name)
1167 if (bad !== undefined || info === undefined) return bad ?? ''
1168 const hs = Object.values(await read($, handovers))
1169 const ib = await read($, inbox)
1170 const at = await $.clock.now()
1171 // Releasing also drops the env changes the user declined: they are recorded as dropped and stop holding the target.
1172 const r = await withDeploys($, cur => {
1173 const released = release(cur.targets[name], info.mode)
1174 const pending = pendingEnv(hs, name, released.state)
1175 const dropped = declinedEntries({ pending, item: itemViewOf(ib) })
1176 if (!released.had && dropped.length === 0) return { deploys: cur, out: { had: false, dropped } }
1177 const state = withEnvDone(released.state, dropped.map(d => ({ pr: d.pr, name: d.change.name, how: 'dropped' as const, at })))
1178 return { deploys: setTarget(cur, name, dropped.length > 0 && info.mode === 'auto' ? { ...state, due: true } : state), out: { had: released.had, dropped } }
1179 })
1180 if (!r.had && r.dropped.length === 0) return `${name} has no hold.`
1181 if (r.dropped.length > 0) {
1182 const qids = r.dropped.flatMap(d => [d.change.qid, ...(d.change.loginQid === undefined ? [] : [d.change.loginQid])])
1183 await withInbox($, cur => ({ inbox: closeEnvItems(cur, qids, 'dropped: released without it', at), out: undefined }))
1184 }
1185 const dropNote = r.dropped.length === 0 ? '' : ` Dropped the declined env changes: ${[...new Set(r.dropped.map(d => d.change.name))].join(', ')}.`
1186 if (info.mode === 'confirm') return `${name} released.${dropNote} It is a confirm target: the next gate still asks for approval.`
1187 return `${name} released.${dropNote} ${await ensureQueue($)}`
1188}
1189
1190// ---- the push gate (pushgate.ts) ----
1191
1192// Push state changes run one after another, like the inbox's.
1193let pushChain: Promise<unknown> = Promise.resolve()
1194
1195// Read push.json, let fn change it, write it back atomically and set the atom. fn returns the new state (the same
1196// object when nothing changed) and whatever the caller wants back.
1197function withPush<T>($: EngineInterface, fn: (cur: PushState) => { push: PushState; out: T }): Promise<T> {
1198 const run = async (): Promise<T> => {
1199 const dir = await stateDir($)
1200 const cur = dir === undefined ? await read($, pushState) : normalizePush(await readJson($, `${dir}/push.json`))hooks/attachments.ts 74 lines1// Attachments in a brief: file paths (screenshots, mockups, logs) a worker opens at start. The
2// spawn hook checks that each exists and rewrites relative paths to absolute ones, because a
3// worker runs in another worktree where an untracked file of the manager's checkout is missing.
4// Pure text handling; the file checks are the hook's.
5
6export type Attachment = { line: number; path: string }
7
8const HEADING = /^#{1,6}\s+Attachments\s*:?\s*$/i
9const LABEL = /^Attachments:\s*(.*)$/i
10const ANY_HEADING = /^#{1,6}\s/
11
12// A bullet or bare line, with one layer of quotes or backticks removed.
13function entry(line: string): string | undefined {
14 let t = line.trim().replace(/^[-*]\s+/, '').trim()
15 if (t === '') return undefined
16 const q = /^(["'`])(.*)\1$/.exec(t)
17 if (q) t = q[2]!.trim()
18 // "none" and friends say there is nothing attached.
19 return t === '' || /^\(?\s*(none|n\/a|-)\s*\)?\.?$/i.test(t) ? undefined : t
20}
21
22// The paths listed under a `## Attachments` heading or an `Attachments:` line, with the index of
23// their line. The list ends at the next heading, or at a blank line followed by text that is not a
24// list item. A path on the `Attachments:` line itself counts.
25export function parseAttachments(prompt: string): Attachment[] {
26 const lines = prompt.split('\n')
27 const out: Attachment[] = []
28 for (let i = 0; i < lines.length; i++) {
29 const line = lines[i]!
30 const label = LABEL.exec(line)
31 if (!HEADING.test(line) && !label) continue
32 if (label && label[1]!.trim() !== '') {
33 const p = entry(label[1]!)
34 if (p !== undefined) out.push({ line: i, path: p })
35 }
36 let blank = false
37 for (let j = i + 1; j < lines.length; j++) {
38 const l = lines[j]!
39 if (ANY_HEADING.test(l)) break
40 if (l.trim() === '') { blank = true; continue }
41 // After a blank line only list items continue the section.
42 if (blank && !/^\s*[-*]\s/.test(l)) break
43 const p = entry(l)
44 if (p !== undefined) out.push({ line: j, path: p })
45 }
46 break
47 }
48 return out
49}
50
51// `~` and relative paths made absolute: `home` for `~`, `cwd` for the rest.
52export function absolutePath(path: string, cwd: string, home: string): string {
53 if (path === '~') return home
54 if (path.startsWith('~/')) return `${home.replace(/\/$/, '')}/${path.slice(2)}`
55 if (path.startsWith('/')) return path
56 return `${cwd.replace(/\/$/, '')}/${path.replace(/^\.\//, '')}`
57}
58
59// The prompt with each listed path replaced by its absolute form. Quotes around a path are
60// dropped; an absolute path written bare stays as it was.
61export function rewriteAttachments(prompt: string, items: Attachment[], abs: (path: string) => string): string {
62 if (items.length === 0) return prompt
63 const lines = prompt.split('\n')
64 for (const { line, path } of items) {
65 const l = lines[line]!
66 const at = l.lastIndexOf(path)
67 const quoted = at > 0 && /["'`]/.test(l[at - 1]!)
68 const start = quoted ? at - 1 : at
69 const end = at + path.length + (quoted ? 1 : 0)
70 lines[line] = l.slice(0, start) + abs(path) + l.slice(end)
71 }
72 return lines.join('\n')
73}
74hooks/slots.ts 64 lines1import type { AgentRow, SlotEntry, TestSlots } from '../types'
2import { ENDED } from './deliver'
3
4// A hook has a 10 s budget, a $.clock wait included, so acquire blocks only briefly. A free slot
5// is granted to the head waiter at once (claimed: false); it has CLAIM_MS to confirm with acquire,
6// else the grant passes on. The grant is signalled by a file and a message, so waiting costs no turns.
7export const LEASE_MS = 45 * 60_000
8export const CLAIM_MS = 2 * 60_000
9export const WAIT_MAX_S = 8
10export const WAIT_DEFAULT_S = 2
11
12// The grant file of a waiter, under the slot directory `dir`; undefined when there is none.
13export const grantFile = (dir: string | undefined, key: string): string | undefined => dir && `${dir}/${key.replace(/[^\w.-]/g, '_')}.granted`
14
15export const span = (ms: number): string => {
16 const m = Math.floor(ms / 60_000)
17 return m < 1 ? `${Math.max(0, Math.round(ms / 1000))}s` : `${m}m`
18}
19export const heldBy = (h: SlotEntry, t: number): string => `${h.name} (${h.label}, ${span(t - h.since)})`
20
21// Drops holders and waiters whose agent ended or is gone and holders past the lease. The main session (key "main") never ends. Returns the new state and what was dropped.
22export function reapSlots(state: TestSlots, rows: AgentRow[], t: number): { state: TestSlots; notes: string[] } {
23 const status = new Map(rows.map(a => [a.id, a.status]))
24 const gone = (key: string) => key !== 'main' && (status.get(key) === undefined || ENDED.has(status.get(key)!))
25 const notes: string[] = []
26 const holders = state.holders.filter(h => {
27 if (gone(h.key)) return false
28 if (t - h.since >= LEASE_MS) {
29 notes.push(`test slot of ${h.name} (${h.label}) released after ${span(LEASE_MS)}: its lease ran out`)
30 return false
31 }
32 return true
33 })
34 const waiters = state.waiters.filter(w => !gone(w.key))
35 const changed = holders.length !== state.holders.length || waiters.length !== state.waiters.length
36 return { state: changed ? { holders, waiters } : state, notes }
37}
38
39// Gives free slots to the head waiters, after dropping grants nobody confirmed within CLAIM_MS.
40// Pure: the caller does the signalling for `granted`.
41export function grantSlots(state: TestSlots, limit: number, t: number): { state: TestSlots; granted: SlotEntry[]; notes: string[] } {
42 const notes: string[] = []
43 const holders = state.holders.filter(h => {
44 if (h.claimed !== false || t - h.since < CLAIM_MS) return true
45 notes.push(`test slot offered to ${h.name} (${h.label}) passed on: not confirmed within ${span(CLAIM_MS)}`)
46 return false
47 })
48 const waiters = [...state.waiters]
49 const granted: SlotEntry[] = []
50 while (holders.length < limit && waiters.length > 0) {
51 const w = { ...waiters.shift()!, since: t, claimed: false }
52 holders.push(w)
53 granted.push(w)
54 }
55 const changed = holders.length !== state.holders.length || granted.length > 0
56 return { state: changed ? { holders, waiters } : state, granted, notes }
57}
58
59export function slotLine(state: TestSlots, limit: number, t: number): string {
60 if (state.holders.length === 0 && state.waiters.length === 0) return ''
61 const held = state.holders.length ? `held by ${state.holders.map(h => heldBy(h, t)).join(', ')}` : 'free'
62 return `Test slots: ${state.holders.length}/${limit} ${held}${state.waiters.length ? ` · ${state.waiters.length} waiting` : ''}`
63}
64hooks/evidence.ts 118 lines1// The `## Verification` section of a PR description: what the worker ran, how the change was exercised for
2// real, and what was not verified. `handover` refuses a PR without it; the parsed record travels with the
3// handover into the reviewer report and the status file. Pure, so it is tested without the plugin runtime.
4
5import type { GuardHit } from './guardtests'
6
7export type Evidence = { ran: string[]; exercised: string; notVerified: string[] }
8
9export const EVIDENCE_FORMAT = [
10 '## Verification',
11 'Ran:',
12 '- `<command>`: pass (<short result, e.g. 42 tests>)',
13 'Exercised: <how the change was run for real: app launched and what was seen, screenshot path, curl output, or "n/a: <reason>" for docs/prompt-only changes>',
14 'Not verified:',
15 '- <what you did not check> (or one line "Not verified: none, because <reason>")',
16].join('\n')
17
18const LABELS = { ran: /^ran$/i, exercised: /^exercised$/i, notVerified: /^not verified$/i }
19type Key = keyof typeof LABELS
20
21const bullet = /^\s*[-*]\s+/
22const label = /^\s*(ran|exercised|not verified)\s*:\s*(.*)$/i
23
24// The text of the section, or undefined when the body has none. Case-insensitive heading, ends at the next `## `.
25export function verificationSection(body: string): string | undefined {
26 const lines = body.replace(/\r\n?/g, '\n').split('\n')
27 const start = lines.findIndex(l => /^##\s+verification\s*$/i.test(l.trim()))
28 if (start < 0) return undefined
29 const rest = lines.slice(start + 1)
30 const end = rest.findIndex(l => /^##\s/.test(l))
31 return (end < 0 ? rest : rest.slice(0, end)).join('\n')
32}
33
34// Labels may stand alone with bullets below, or carry text after the colon. Bullets are entries; any other
35// text under Ran or Not verified is one entry per line; Exercised joins its lines.
36export function parseVerification(section: string): Partial<Evidence> {
37 const out: Evidence = { ran: [], exercised: '', notVerified: [] }
38 const seen = new Set<Key>()
39 let cur: Key | undefined
40 const add = (key: Key, text: string) => {
41 const t = text.trim()
42 if (!t) return
43 if (key === 'exercised') out.exercised = out.exercised ? `${out.exercised} ${t}` : t
44 else out[key].push(t)
45 }
46 for (const raw of section.split('\n')) {
47 const m = label.exec(raw)
48 if (m && !bullet.test(raw)) {
49 cur = (Object.keys(LABELS) as Key[]).find(k => LABELS[k].test(m[1]!.trim()))!
50 seen.add(cur)
51 add(cur, m[2]!)
52 } else if (cur) {
53 add(cur, raw.replace(bullet, ''))
54 }
55 }
56 return {
57 ...(seen.has('ran') ? { ran: out.ran } : {}),
58 ...(seen.has('exercised') ? { exercised: out.exercised } : {}),
59 ...(seen.has('notVerified') ? { notVerified: out.notVerified } : {}),
60 }
61}
62
63const unticked = (s: string) => s.replace(/`/g, '')
64
65// Evidence when the body qualifies, else every reason it does not. `required` are the settings' workerChecks
66// and alwaysTests: each must appear (backticks optional) in some Ran entry. `guard` are the guard tests the PR's
67// changed files require (guardtests.ts): same matching, and the refusal names the globs that triggered them.
68export function checkEvidence(body: string, required: string[], guard: GuardHit[] = []): { evidence: Evidence } | { problems: string[] } {
69 const section = verificationSection(body)
70 if (section === undefined) return { problems: ['the PR description has no `## Verification` section'] }
71 if (!section.trim()) return { problems: ['the `## Verification` section is empty'] }
72 const p = parseVerification(section)
73 const problems: string[] = []
74 if (!p.ran?.length) problems.push('`Ran:` has no entries')
75 if (!p.exercised) problems.push('`Exercised:` is missing or empty')
76 else if (/^n\/a:?$/i.test(p.exercised.trim())) problems.push('`Exercised: n/a` needs a reason ("n/a: <reason>")')
77 if (!p.notVerified?.length) problems.push('`Not verified:` is missing or empty')
78 else if (p.notVerified.length === 1 && /^(nothing|none|n\/a|-)\.?$/i.test(unticked(p.notVerified[0]!).trim())) {
79 problems.push('`Not verified:` must name at least one honest item, or say "none, because <reason>"')
80 }
81 if (p.ran?.length) {
82 const ran = p.ran.map(unticked)
83 for (const cmd of required) {
84 if (!ran.some(r => r.includes(unticked(cmd)))) problems.push(`required command \`${cmd}\` does not appear under \`Ran:\``)
85 }
86 for (const g of guard) {
87 if (required.includes(g.test) || ran.some(r => r.includes(unticked(g.test)))) continue
88 problems.push(`guard test \`${g.test}\` (for ${g.globs.map(x => `\`${x}\``).join(', ')}) does not appear under \`Ran:\``)
89 }
90 }
91 if (problems.length) return { problems }
92 return { evidence: { ran: p.ran!, exercised: p.exercised!, notVerified: p.notVerified! } }
93}
94
95export function evidenceRefusal(pr: number, problems: string[]): string {
96 return [
97 `Refused: PR #${pr} does not carry the proof handover requires:`,
98 ...problems.map(p => `- ${p}`),
99 'The PR description must contain this section:',
100 EVIDENCE_FORMAT,
101 `Have the worker fix the PR description (gh pr edit ${pr} --body-file <file>), then call handover again.`,
102 ].join('\n')
103}
104
105const clip = (s: string, n: number) => (s.length > n ? `${s.slice(0, n - 1)}…` : s)
106
107// Full text, for the reviewer's list.
108export function evidenceText(e: Evidence | undefined): string {
109 if (!e) return 'no evidence recorded'
110 return `ran: ${e.ran.join('; ')} | exercised: ${e.exercised} | not verified: ${e.notVerified.join('; ')}`
111}
112
113// Short, for the status line.
114export function evidenceSummary(e: Evidence | undefined): string {
115 if (!e) return 'no evidence recorded'
116 return `${e.ran.length} ran, exercised: ${clip(e.exercised, 60)}, ${e.notVerified.length} not verified`
117}
118hooks/mainff.ts 64 lines1// Fast-forward of the repo's main checkout after a reviewer batch. The reviewer runs worktree-isolated
2// and the host refuses its git calls on the main checkout, so the plugin (not subject to that
3// isolation) does it when the reviewer reports "done". Never stash, reset, checkout, clean or force.
4import { parsePorcelain } from './clean'
5
6export type Run = { exitCode: number; stdout: string; stderr: string }
7// Runs one argv, git first.
8export type GitRunner = (argv: string[]) => Promise<Run>
9
10const reason = (r: Run): string => (r.stderr.trim() || r.stdout.trim()).split('\n')[0]?.slice(0, 160) || `exit ${r.exitCode}`
11const notUpdated = (why: string): string => `main checkout not updated: ${why}`
12
13// One line: "main checkout fast-forwarded to <sha>" or "main checkout not updated: <reason>".
14export async function fastForwardMain(git: GitRunner, base: string): Promise<string> {
15 try {
16 const wl = await git(['git', 'worktree', 'list', '--porcelain'])
17 if (wl.exitCode !== 0) return notUpdated(`worktree list failed: ${reason(wl)}`)
18 const main = parsePorcelain(wl.stdout)[0]?.path
19 if (main === undefined) return notUpdated('no main checkout found')
20 // Untracked files do not count as dirty; an incoming file that would overwrite one fails the pull.
21 const st = await git(['git', '-C', main, 'status', '--porcelain', '--untracked-files=no'])
22 if (st.exitCode !== 0) return notUpdated(`status failed: ${reason(st)}`)
23 if (st.stdout.trim() !== '') return notUpdated('dirty')
24 const br = await git(['git', '-C', main, 'branch', '--show-current'])
25 if (br.exitCode !== 0) return notUpdated(`branch check failed: ${reason(br)}`)
26 const on = br.stdout.trim()
27 if (on === '') return notUpdated('on a detached HEAD')
28 if (on !== base) return notUpdated(`on branch ${on}`)
29 const pull = await git(['git', '-C', main, 'pull', '--ff-only', 'origin', base])
30 if (pull.exitCode !== 0) return notUpdated(`ff failed: ${reason(pull)}`)
31 const sha = await git(['git', '-C', main, 'rev-parse', '--short', 'HEAD'])
32 return `main checkout fast-forwarded to ${sha.exitCode === 0 ? sha.stdout.trim() : 'origin/' + base}`
33 } catch (err) {
34 return notUpdated(`ff failed: ${(err instanceof Error ? err.message : String(err)).split('\n')[0]?.slice(0, 160) ?? ''}`)
35 }
36}
37
38// One fast-forward at a time; several "done" calls in a batch queue behind each other. A repeat is
39// harmless: the second pull finds nothing to do and answers fast-forwarded to the same sha.
40let chain: Promise<unknown> = Promise.resolve()
41export function serialFastForward(git: GitRunner, base: string): Promise<string> {
42 const next = chain.then(() => fastForwardMain(git, base), () => fastForwardMain(git, base))
43 chain = next
44 return next
45}
46
47// What a reviewer run did with the queue, for the line of a run that never reached "done".
48export type RunWork = { touched: boolean; ready: boolean; line?: string }
49
50export const noteWork = (w: RunWork | undefined, action: string): RunWork => {
51 const cur = w ?? { touched: false, ready: false }
52 if (action === 'ready') return { ...cur, touched: true, ready: true }
53 if (action === 'take' || action === 'back') return { ...cur, touched: true }
54 return cur
55}
56
57// The main-checkout line for the final report of a run: the one "done" produced, else why none ran.
58// undefined for a run that did no batch work.
59export function runLine(w: RunWork | undefined): string | undefined {
60 if (w === undefined || (!w.touched && w.line === undefined)) return undefined
61 if (w.line !== undefined) return w.line
62 return notUpdated(w.ready ? 'batch awaits /flow push' : 'nothing pushed')
63}
64hooks/clean.ts 577 lines1import type { HandoffRecord, Leftovers, LogEvent, Session } from '../types'
2import { sessionKey } from './sessions'
3
4// Which leftover worktrees and local branches of finished agents can go. Pure: register.tsx runs
5// git and gh, hands the parsed answers here, and removes what comes back in `remove`.
6// Only clean work that is on the base, in a merged PR, or pushed with its PR closed is removed;
7// uncommitted, unpushed, locked and live work is listed for a person, never touched.
8
9export type Worktree = {
10 path: string
11 head?: string
12 // Short name, without refs/heads/; absent on a detached HEAD.
13 branch?: string
14 // The lock reason ('' when locked without one); absent when not locked.
15 locked?: string
16 // The directory is gone; `git worktree prune` drops the entry.
17 prunable: boolean
18}
19
20// `git worktree list --porcelain`, the main checkout first; bare entries are skipped.
21export function parsePorcelain(text: string): Worktree[] {
22 const out: Worktree[] = []
23 for (const block of text.split(/\n\n+/)) {
24 const lines = block.split('\n')
25 if (lines.includes('bare')) continue
26 const get = (k: string) => {
27 const l = lines.find(x => x === k || x.startsWith(`${k} `))
28 return l === undefined ? undefined : l.slice(k.length + 1)
29 }
30 const path = get('worktree')
31 if (path === undefined || path === '') continue
32 const branch = get('branch')?.replace(/^refs\/heads\//, '')
33 const locked = get('locked')
34 out.push({
35 path,
36 ...(get('HEAD') && { head: get('HEAD') }),
37 ...(branch && { branch }),
38 ...(locked !== undefined && { locked }),
39 prunable: get('prunable') !== undefined,
40 })
41 }
42 return out
43}
44
45// Untracked links to the plugin type definitions a worker makes for tsc; never work.
46const TYPE_LINKS = new Set(['types', 'types/', '.claude-plugin/types', '.claude-plugin/types/'])
47
48// The files `git status --porcelain` names, minus the type-definition links.
49export function dirtyFiles(status: string): string[] {
50 return status.split('\n').filter(l => l.trim() !== '')
51 .filter(l => !(l.startsWith('?? ') && TYPE_LINKS.has(l.slice(3).trim())))
52 .map(l => l.slice(3).trim())
53}
54
55// An agent that may still use its worktree: running, waiting, or idle between turns.
56export const isLive = (status: string): boolean => ['pending', 'running', 'waiting', 'idle'].includes(status)
57
58// The type-link paths `git status` lists as untracked: the only files apply deletes itself so a
59// plain `git worktree remove` goes through. Exactly the ignored ones, nothing else.
60export function typeLinks(status: string): string[] {
61 return status.split('\n').filter(l => l.startsWith('?? ') && TYPE_LINKS.has(l.slice(3).trim()))
62 .map(l => l.slice(3).trim().replace(/\/$/, ''))
63}
64
65export type PrRow = { number: number; headRefName: string; headRefOid: string; state: string }
66export type RosterEntry = { id: string; name?: string; live: boolean; cwd?: string }
67
68export type CleanInputs = {
69 main: string
70 base: string
71 worktrees: Worktree[]
72 // `git status --porcelain` per worktree path; absent when it could not be read.
73 status: Record<string, string>
74 // Local branches and their tips.
75 branches: Record<string, string>
76 // Shas known to be ancestors of origin/<base>.
77 onBase: Set<string>
78 // origin/<branch> tips.
79 remote: Record<string, string>
80 // undefined: gh failed, ancestry only.
81 prs?: PrRow[]
82 // "a b": a is an ancestor of b (b a sha or origin/<branch>), as asked by ancestryQueries.
83 ancestry: Set<string>
84 roster: RosterEntry[]
85 // Pids from lock reasons that no longer run.
86 deadPids?: Set<number>
87 // Worktree paths a handoff record keeps for a successor not spawned yet.
88 waiting?: Set<string>
89 // sha -> the commits (short sha + subject, origin/<base>..sha) whose whole change from the merge
90 // base is already on origin/<base>, file for file; see containedCandidates.
91 contained?: Record<string, string[]>
92 // The pid of the Claude process running this plugin; absent when it could not be told. A lock
93 // that names it was made in this session, so $.agent.list() knows its agent.
94 ownPid?: number
95 // The pid of the plugin's process and its ancestors, up to the Claude process: the plugin's
96 // process runner need not be a direct child of Claude, so a lock may name any of them.
97 ownPids?: Set<number>
98 // Worktrees whose lock file is older than 10 minutes; a younger or unaged lock may belong to an
99 // agent that is still starting and not yet in the roster.
100 oldLocks?: Set<string>
101}
102
103export type Kept = { kind: 'worktree' | 'branch'; name: string; reason: string; needsLook: boolean }
104export type Sweep = {
105 remove: { worktrees: string[]; branches: string[] }
106 // Locked worktrees whose lock belongs to an agent that ended: unlocked before removal.
107 unlock: string[]
108 keep: Kept[]
109 // Type links to delete before the plain remove, per worktree path that is in remove.
110 links: Record<string, string[]>
111 // Unpushed commits dropped because a merged PR of the branch family already holds their content.
112 dropped: { kind: 'worktree' | 'branch'; name: string; commits: string[] }[]
113}
114
115const linked = (i: Pick<CleanInputs, 'main'>, w: Worktree) =>
116 w.path !== i.main && w.path.startsWith(`${i.main.replace(/\/+$/, '')}/.claude/worktrees/`)
117
118// flow/x and flow/x-2 are one line of work: a continuation or rebase reuses the name with a suffix.
119const family = (b: string) => b.replace(/-\d+$/, '')
120
121const mergedFamily = (i: Pick<CleanInputs, 'prs'>, branch: string) =>
122 (i.prs ?? []).filter(p => p.state === 'MERGED' && family(p.headRefName) === family(branch))
123
124// Worktrees whose handoff record keeps them for a successor: not taken, not continued, and the
125// branch has no merged PR (the work was complete after all, so nobody will come).
126export function waitingPaths(
127 hs: { branch: string; worktree?: string; takenBy?: string; at: number }[],
128 continued: { branch?: string; ts: string }[],
129 prs: PrRow[] | undefined,
130): Set<string> {
131 return new Set(hs.filter(h => h.worktree !== undefined && h.takenBy === undefined
132 && !continued.some(l => l.branch === h.branch && Date.parse(l.ts) >= h.at)
133 && !(prs ?? []).some(p => p.state === 'MERGED' && p.headRefName === h.branch)).map(h => h.worktree!))
134}
135
136// The shas register.tsx checks for content equivalence: tips not on the base whose branch family
137// has a merged PR. It answers with `contained` (every file the commits changed since the merge
138// base has the same tree entry on origin/<base>, so nothing is lost), which landed() trusts.
139export function containedCandidates(i: Omit<CleanInputs, 'ancestry'>): string[] {
140 const out = new Set<string>()
141 const ask = (sha: string | undefined, branch: string | undefined) => {
142 if (sha === undefined || branch === undefined || i.onBase.has(sha)) return
143 if (mergedFamily(i, branch).length > 0) out.add(sha)
144 }
145 for (const w of i.worktrees) if (linked(i, w) && !w.prunable) ask(w.head, w.branch)
146 for (const [b, sha] of Object.entries(i.branches)) ask(sha, b)
147 return [...out]
148}
149
150const merged = (i: Pick<CleanInputs, 'prs'>, branch?: string) =>
151 (i.prs ?? []).filter(p => p.state === 'MERGED' && (branch === undefined || p.headRefName === branch))
152
153// The pairs ancestry needs beyond onBase, as "a b": a sha against a merged PR head of its branch,
154// and against origin/<branch>. register.tsx answers each with `git merge-base --is-ancestor`.
155export function ancestryQueries(i: Omit<CleanInputs, 'ancestry'>): string[] {
156 const out = new Set<string>()
157 const ask = (sha: string | undefined, branch: string | undefined) => {
158 if (sha === undefined || i.onBase.has(sha)) return
159 if (branch === undefined) return
160 for (const p of merged(i, branch)) if (p.headRefOid && p.headRefOid !== sha) out.add(`${sha} ${p.headRefOid}`)
161 if (i.remote[branch] !== undefined && i.remote[branch] !== sha) out.add(`${sha} origin/${branch}`)
162 }
163 for (const w of i.worktrees) if (linked(i, w) && !w.prunable) ask(w.head, w.branch)
164 for (const [b, sha] of Object.entries(i.branches)) ask(sha, b)
165 return [...out]
166}
167
168type Verdict = { ok: true; dropped?: string[] } | { ok: false; reason: string; needsLook: boolean }
169
170// Whether a commit's work is safe to drop: on the base, in a merged PR, or (worktrees only)
171// pushed with a closed PR.
172function landed(i: CleanInputs, sha: string, branch: string | undefined, closedOk: boolean): Verdict {
173 if (i.onBase.has(sha)) return { ok: true }
174 if (merged(i).some(p => p.headRefOid === sha)) return { ok: true }
175 if (branch !== undefined && merged(i, branch).some(p => i.ancestry.has(`${sha} ${p.headRefOid}`))) return { ok: true }
176 const same = i.contained?.[sha]
177 if (branch !== undefined && same !== undefined && mergedFamily(i, branch).length > 0) {
178 return same.length > 0 ? { ok: true, dropped: same } : { ok: true }
179 }
180 const pushed = branch !== undefined && (i.remote[branch] === sha || i.ancestry.has(`${sha} origin/${branch}`))
181 const prs = branch === undefined ? [] : (i.prs ?? []).filter(p => p.headRefName === branch)
182 const closed = prs.find(p => p.state === 'CLOSED')
183 if (closedOk && pushed && closed !== undefined) return { ok: true }
184 if (closed !== undefined) return { ok: false, reason: `PR #${closed.number} closed without merging${pushed ? ', pushed' : ''}`, needsLook: true }
185 if (branch === undefined) return { ok: false, reason: 'detached HEAD with commits not on the base', needsLook: true }
186 if (merged(i, branch).length > 0) return { ok: false, reason: 'unpushed commits after its PR was merged', needsLook: true }
187 if (pushed) return { ok: false, reason: 'pushed, not merged, no PR', needsLook: true }
188 return { ok: false, reason: 'unpushed commits', needsLook: true }
189}
190
191// The live agents' claim on a worktree: its path ends in agent-<id>, it is their cwd, or it is
192// on their branch flow/<name>.
193function liveOwner(i: CleanInputs, path: string | undefined, branch: string | undefined): RosterEntry | undefined {
194 return i.roster.find(a => a.live && (
195 (path !== undefined && (path.endsWith(`/agent-${a.id}`) || a.cwd === path || a.cwd?.startsWith(`${path}/`) === true)) ||
196 (branch !== undefined && a.name !== undefined && branch === `flow/${a.name}`)))
197}
198
199const LOCK_AGENT = /\bagent-([0-9a-zA-Z]+)\b/
200const LOCK_PID = /\bpid (\d+)\b/
201
202// A lock is stale when the agent it names ended in this session's roster, is unknown to it while
203// the lock is this session's own, or its process is gone.
204function staleLock(i: CleanInputs, reason: string, path: string): boolean {
205 const id = LOCK_AGENT.exec(reason)?.[1]
206 const pid = Number(LOCK_PID.exec(reason)?.[1] ?? NaN)
207 const agent = id === undefined ? undefined : i.roster.find(a => a.id === id || a.id === `agent-${id}`)
208 if (agent !== undefined) return !agent.live
209 // Same process, same session: the roster lists every agent it started, nested ones included
210 // (AgentInfo.parentId), so an agent missing from it is gone. Another pid proves nothing.
211 if (id !== undefined && Number.isInteger(pid) && (pid === i.ownPid || i.ownPids?.has(pid) === true)) return i.oldLocks?.has(path) === true
212 return id !== undefined && Number.isInteger(pid) && i.deadPids?.has(pid) === true
213}
214
215export function selectCleanup(i: CleanInputs): Sweep {
216 const sweep: Sweep = { remove: { worktrees: [], branches: [] }, unlock: [], keep: [], links: {}, dropped: [] }
217 const keep = (kind: Kept['kind'], name: string, reason: string, needsLook: boolean) =>
218 sweep.keep.push({ kind, name, reason, needsLook })
219 // Branches checked out in a worktree that stays (the main checkout included).
220 const held = new Map<string, string>()
221 for (const w of i.worktrees) {
222 if (!linked(i, w)) {
223 if (w.branch !== undefined) held.set(w.branch, w.path)
224 continue
225 }
226 const hold = () => { if (w.branch !== undefined) held.set(w.branch, w.path) }
227 if (w.prunable) {
228 sweep.remove.worktrees.push(w.path)
229 continue
230 }
231 const live = liveOwner(i, w.path, w.branch)
232 if (live !== undefined) {
233 keep('worktree', w.path, `in use by ${live.name ?? live.id}`, false)
234 hold()
235 continue
236 }
237 if (i.waiting?.has(w.path)) {
238 keep('worktree', w.path, 'waiting for its successor', false)
239 hold()
240 continue
241 }
242 const status = i.status[w.path]
243 if (status === undefined || w.head === undefined) {
244 keep('worktree', w.path, 'git status failed', true)
245 hold()
246 continue
247 }
248 const dirty = dirtyFiles(status)
249 if (dirty.length > 0) {
250 const named = dirty.slice(0, 3).join(', ') + (dirty.length > 3 ? ` and ${dirty.length - 3} more` : '')
251 keep('worktree', w.path, `uncommitted changes: ${named}`, true)
252 hold()
253 continue
254 }
255 const v = landed(i, w.head, w.branch, true)
256 if (!v.ok) {
257 keep('worktree', w.path, v.reason, v.needsLook)
258 hold()
259 continue
260 }
261 if (w.locked !== undefined) {
262 if (!staleLock(i, w.locked, w.path)) {
263 keep('worktree', w.path, 'locked', true)
264 hold()
265 continue
266 }
267 sweep.unlock.push(w.path)
268 }
269 if (v.dropped) sweep.dropped.push({ kind: 'worktree', name: w.path, commits: v.dropped })
270 const links = typeLinks(status)
271 if (links.length > 0) sweep.links[w.path] = links
272 sweep.remove.worktrees.push(w.path)
273 }
274
275 for (const [b, tip] of Object.entries(i.branches)) {
276 if (b === i.base) continue
277 const at = held.get(b)
278 if (at !== undefined) {
279 // The main checkout's branch is the person's own business.
280 if (at !== i.main) keep('branch', b, `checked out in ${at}`, false)
281 continue
282 }
283 const open = (i.prs ?? []).find(p => p.state === 'OPEN' && p.headRefName === b)
284 if (open !== undefined) {
285 keep('branch', b, `PR #${open.number} is open`, false)
286 continue
287 }
288 const live = liveOwner(i, undefined, b)
289 if (live !== undefined) {
290 keep('branch', b, `in use by ${live.name ?? live.id}`, false)
291 continue
292 }
293 const v = landed(i, tip, b, false)
294 if (v.ok) {
295 sweep.remove.branches.push(b)
296 if (v.dropped) sweep.dropped.push({ kind: 'branch', name: b, commits: v.dropped })
297 }
298 else keep('branch', b, v.reason, v.needsLook)
299 }
300 return sweep
301}
302
303// The listing /flow clean and mcp__flow__clean answer with.
304export function sweepText(s: Sweep, opts: { applied: boolean; removed?: Sweep['remove']; failed?: Kept[]; note?: string; dryHint?: string }): string {
305 const lines: string[] = []
306 if (opts.note) lines.push(opts.note)
307 const gone = opts.applied ? opts.removed ?? s.remove : s.remove
308 const verb = opts.applied ? 'Removed' : 'Would remove'
309 if (gone.worktrees.length) lines.push(`${verb} ${gone.worktrees.length} worktree${gone.worktrees.length > 1 ? 's' : ''}:`, ...gone.worktrees.map(p => ` ${p}`))
310 if (gone.branches.length) lines.push(`${verb} ${gone.branches.length} branch${gone.branches.length > 1 ? 'es' : ''}:`, ...gone.branches.map(b => ` ${b}`))
311 if (!gone.worktrees.length && !gone.branches.length) lines.push(opts.applied ? 'Removed nothing.' : 'Nothing to remove.')
312 const lost = s.dropped.filter(d => gone[d.kind === 'worktree' ? 'worktrees' : 'branches'].includes(d.name))
313 if (lost.length) {
314 lines.push('Dropping unpushed commits whose content a merged PR already holds:')
315 for (const d of lost) lines.push(` ${d.kind} ${d.name}:`, ...d.commits.map(c => ` ${c}`))
316 }
317 const kept = [...(opts.failed ?? []), ...s.keep]
318 const look = kept.filter(k => k.needsLook)
319 const busy = kept.filter(k => !k.needsLook)
320 if (look.length) lines.push('Kept for a person to decide:', ...look.map(k => ` ${k.kind} ${k.name}: ${k.reason}`))
321 if (busy.length) lines.push('Kept, in use:', ...busy.map(k => ` ${k.kind} ${k.name}: ${k.reason}`))
322 if (!opts.applied && opts.dryHint && (s.remove.worktrees.length || s.remove.branches.length)) lines.push(opts.dryHint)
323 return lines.join('\n')
324}
325
326// The pane's and status's one line; undefined when there is nothing to say.
327export function leftoverLine(c: { worktrees: number; branches: number; needsLook: number }): string | undefined {
328 if (c.worktrees === 0 && c.branches === 0 && c.needsLook === 0) return undefined
329 // worktrees and branches are what a sweep would remove; needsLook is what it keeps for a person to decide.
330 const gone = [
331 ...(c.worktrees ? [`${c.worktrees} worktree${c.worktrees === 1 ? '' : 's'}`] : []),
332 ...(c.branches ? [`${c.branches} branch${c.branches === 1 ? '' : 'es'}`] : []),
333 ]
334 const parts = [
335 ...(gone.length ? [`${gone.join(' and ')} can be removed`] : []),
336 ...(c.needsLook ? [`${c.needsLook} hold${c.needsLook === 1 ? 's' : ''} work that isn't merged`] : []),
337 ]
338 return `Cleanup: ${parts.join(', ')}. /flow clean lists them`
339}
340
341// This session's pids: `start` and its ancestors up to and including the nearest `claude` process,
342// walked with `step` (undefined when ps fails). A chain that meets no claude within `max` steps,
343// pid 1 or a repeat is not trusted: empty. Stopping at the nearest one keeps a nested claude from
344// claiming the locks of the session whose shell it runs in.
345export async function ancestorPids(start: number, step: (pid: number) => Promise<{ ppid: number; comm: string } | undefined>, max = 12): Promise<Set<number>> {
346 const out = new Set<number>()
347 let pid: number | undefined = start
348 for (let n = 0; n < max && pid !== undefined && Number.isInteger(pid) && pid > 1 && !out.has(pid); n++) {
349 out.add(pid)
350 const r = await step(pid)
351 if (r === undefined) break
352 if (r.comm.trim().split('/').pop() === 'claude') return out
353 pid = r.ppid
354 }
355 return new Set()
356}
357
358// --- Sweep: reading git and gh, and removing what is safe ----------------------------------------
359// The engine handle `$` is never passed across an import, so the caller hands in these calls.
360
361export type CleanIo = {
362 // A process run; a failure to start throws.
363 exec: (argv: string[], timeoutMs: number) => Promise<{ exitCode: number; stdout: string; stderr: string }>
364 agents: () => Promise<{ id: string; status: string; name?: string }[]>
365 cwd: () => Promise<string | undefined>
366 handoffs: () => Promise<HandoffRecord[]>
367 sessions: () => Promise<Session[]>
368 log: () => Promise<LogEvent[]>
369 append: (event: Omit<LogEvent, 'ts'>) => Promise<void>
370 setLeftovers: (counts: Leftovers) => Promise<void>
371}
372
373export type CleanGathered = { inputs?: CleanInputs; notes: string[]; error?: string }
374
375// Everything selectCleanup needs, read once: one fetch, one gh call, then local git.
376async function gatherClean(io: CleanIo, base: string): Promise<CleanGathered> {
377 const run = async (argv: string[]) => {
378 try {
379 return await io.exec(argv, 60_000)
380 } catch (err) {
381 return { exitCode: 1, stdout: '', stderr: err instanceof Error ? err.message : String(err) }
382 }
383 }
384 const first = (r: { stderr: string }) => r.stderr.trim().split('\n')[0]?.slice(0, 200) || 'no output'
385 const notes: string[] = []
386 const fetched = await run(['git', 'fetch', 'origin', '--prune'])
387 if (fetched.exitCode !== 0) notes.push(`git fetch origin failed (${first(fetched)}): judged against the refs as they were.`)
388 const wl = await run(['git', 'worktree', 'list', '--porcelain'])
389 if (wl.exitCode !== 0) return { notes, error: `git worktree list failed: ${first(wl)}` }
390 const worktrees = parsePorcelain(wl.stdout)
391 const main = worktrees[0]?.path
392 if (main === undefined) return { notes, error: 'git worktree list gave no checkout.' }
393 const refs = async (prefix: string) => {
394 const r = await run(['git', 'for-each-ref', '--format=%(refname:short) %(objectname)', prefix])
395 const out: Record<string, string> = {}
396 for (const l of r.stdout.split('\n')) {
397 const [name, sha] = l.trim().split(' ')
398 if (name && sha) out[name] = sha
399 }
400 return out
401 }
402 const branches = await refs('refs/heads')
403 const remote: Record<string, string> = {}
404 for (const [name, sha] of Object.entries(await refs('refs/remotes/origin'))) {
405 if (name.startsWith('origin/') && name !== 'origin/HEAD') remote[name.slice('origin/'.length)] = sha
406 }
407 if (remote[base] === undefined) return { notes, error: `origin/${base} not found: nothing is judged without the base.` }
408 const onBase = new Set((await run(['git', 'for-each-ref', '--merged', `origin/${base}`, '--format=%(objectname)', 'refs/heads'])).stdout
409 .split('\n').map(l => l.trim()).filter(Boolean))
410 const status: Record<string, string> = {}
411 const deadPids = new Set<number>()
412 for (const w of worktrees.slice(1)) {
413 if (w.prunable) continue
414 const st = await run(['git', '-C', w.path, 'status', '--porcelain'])
415 if (st.exitCode === 0) status[w.path] = st.stdout
416 if (w.head !== undefined && !onBase.has(w.head)
417 && (await run(['git', 'merge-base', '--is-ancestor', w.head, `origin/${base}`])).exitCode === 0) onBase.add(w.head)
418 const pid = /\bpid (\d+)\b/.exec(w.locked ?? '')?.[1]
419 if (pid !== undefined && (await run(['ps', '-p', pid])).exitCode !== 0) deadPids.add(Number(pid))
420 }
421 let prs: PrRow[] | undefined
422 const gh = await run(['gh', 'pr', 'list', '--state', 'all', '--limit', '300', '--json', 'number,headRefName,headRefOid,state'])
423 try {
424 if (gh.exitCode !== 0) throw new Error(first(gh))
425 const parsed = JSON.parse(gh.stdout || '[]') as unknown
426 if (!Array.isArray(parsed)) throw new Error('unexpected gh output')
427 prs = parsed as PrRow[]
428 } catch (err) {
429 notes.push(`gh pr list failed (${err instanceof Error ? err.message : String(err)}): squash-merged work is not recognised, only what is on origin/${base}.`)
430 }
431 // A successor continuing in a handed-off worktree runs there; one not spawned yet may still claim it.
432 const hs = await io.handoffs()
433 const continued = (await io.log()).filter(l => l.event === 'continue')
434 const cwdOf = new Map(hs.filter(h => h.takenBy !== undefined && h.worktree !== undefined).map(h => [h.takenBy!, h.worktree!]))
435 const waiting = waitingPaths(hs, continued, prs)
436 const roster = (await io.agents()).map(a => ({
437 id: a.id, live: isLive(a.status), ...(a.name !== undefined && { name: a.name }),
438 ...(a.name !== undefined && cwdOf.has(a.name) && { cwd: cwdOf.get(a.name) }),
439 }))
440 // The main session may itself run in a linked worktree: never pull the floor from under it.
441 const here = await io.cwd().catch(() => undefined)
442 if (here) roster.push({ id: 'main', live: true, name: 'the main session', cwd: here })
443 // A worker in another harness lives in its worktree until it is stopped or its terminal is gone.
444 for (const s of await io.sessions()) {
445 roster.push({ id: sessionKey(s.name), live: s.status === 'running' || s.status === 'reported', name: s.name, cwd: s.worktree })
446 }
447 const partial: Omit<CleanInputs, 'ancestry'> = { main, base, worktrees, status, branches, onBase, remote, prs, roster, deadPids, waiting }
448 // Content equivalence for tips a merged PR of the branch family may have replaced (a rebase
449 // leaves the old commits behind): contained when every file they changed since the merge base
450 // has the very same tree entry on origin/<base>. Anything odd (quoted names, many files) is not.
451 const contained: Record<string, string[]> = {}
452 for (const sha of containedCandidates(partial)) {
453 const mb = (await run(['git', 'merge-base', sha, `origin/${base}`])).stdout.trim()
454 const diff = await run(['git', 'diff', '--name-only', '--no-renames', mb, sha])
455 const files = diff.stdout.split('\n').filter(Boolean)
456 if (!mb || diff.exitCode !== 0 || files.length > 200 || files.some(f => f.startsWith('"'))) continue
457 let same = true
458 for (const f of files) {
459 const [a, b] = await Promise.all([sha, `origin/${base}`].map(r => run(['git', 'ls-tree', r, '--', f])))
460 if (a!.exitCode !== 0 || b!.exitCode !== 0 || a!.stdout.trim() !== b!.stdout.trim()) { same = false; break }
461 }
462 if (!same) continue
463 const log = await run(['git', 'log', '--format=%h %s', `origin/${base}..${sha}`])
464 if (log.exitCode === 0) contained[sha] = log.stdout.split('\n').filter(Boolean)
465 }
466 partial.contained = contained
467 // The lock names the Claude process; the plugin's runner may sit below it, so the chain up to the
468 // nearest claude process counts as this session. Works with an empty roster (after a reload): an agent absent
469 // from it is judged by its lock's age. An unreadable chain means no lock is broken.
470 const own = await ancestorPids(Number((await run(['sh', '-c', 'echo $PPID'])).stdout.trim()), async pid => {
471 const r = await run(['ps', '-o', 'ppid=,comm=', '-p', String(pid)])
472 const m = /^\s*(\d+)\s+(.*)$/.exec(r.stdout.trim())
473 return r.exitCode === 0 && m ? { ppid: Number(m[1]), comm: m[2]! } : undefined
474 })
475 if (own.size > 0) {
476 partial.ownPids = own
477 // A fresh agent's worktree is locked before the roster lists it; only an old lock is judged.
478 // Unknown age (no admin dir, no lock file) counts as young.
479 const oldLocks = new Set<string>()
480 for (const w of worktrees.slice(1)) {
481 const lp = Number(/\bpid (\d+)\b/.exec(w.locked ?? '')?.[1] ?? NaN)
482 if (!own.has(lp)) continue
483 const dir = await run(['git', '-C', w.path, 'rev-parse', '--absolute-git-dir'])
484 if (dir.exitCode !== 0) continue
485 const old = await run(['find', `${dir.stdout.trim()}/locked`, '-mmin', '+10'])
486 if (old.exitCode === 0 && old.stdout.trim() !== '') oldLocks.add(w.path)
487 }
488 partial.oldLocks = oldLocks
489 }
490 const ancestry = new Set<string>()
491 for (const q of ancestryQueries(partial)) {
492 const [a, b] = q.split(' ')
493 if ((await run(['git', 'merge-base', '--is-ancestor', a!, b!])).exitCode === 0) ancestry.add(q)
494 }
495 return { inputs: { ...partial, ancestry }, notes }
496}
497
498// One sweep at a time: the reviewer's `done`, a reviewer ending and /flow clean may meet.
499let sweepChain: Promise<unknown> = Promise.resolve()
500function exclusive<T>(fn: () => Promise<T>): Promise<T> {
501 const run = sweepChain.then(fn, fn)
502 sweepChain = run.catch(() => undefined)
503 return run
504}
505
506const leftoverCounts = (s: Sweep, failed: Kept[] = [], applied = false): Leftovers => ({
507 worktrees: applied ? 0 : s.remove.worktrees.length,
508 branches: applied ? 0 : s.remove.branches.length,
509 needsLook: s.keep.filter(k => k.needsLook).length + failed.length,
510})
511
512// A sweep: the listing, and with apply the removal of what selectCleanup marked safe. A removal
513// that fails is kept with git's message; nothing is forced. `started` runs once it holds the lock.
514export async function sweep(io: CleanIo, base: string, apply: boolean, dryHint?: string, started?: () => void): Promise<string> {
515 return exclusive(async () => {
516 started?.()
517 const g = await gatherClean(io, base)
518 if (g.inputs === undefined) return [...g.notes, g.error ?? 'Cleanup failed.'].join('\n')
519 const s = selectCleanup(g.inputs)
520 const note = g.notes.join('\n') || undefined
521 if (!apply) {
522 await io.setLeftovers(leftoverCounts(s))
523 return sweepText(s, { applied: false, ...(note && { note }), ...(dryHint && { dryHint }) })
524 }
525 const git = async (argv: string[]) => {
526 try {
527 return await io.exec(['git', ...argv], 60_000)
528 } catch (err) {
529 return { exitCode: 1, stdout: '', stderr: err instanceof Error ? err.message : String(err) }
530 }
531 }
532 const failed: Kept[] = []
533 const removed: Sweep['remove'] = { worktrees: [], branches: [] }
534 const stuck = new Set<string>()
535 const byPath = new Map(g.inputs.worktrees.map(w => [w.path, w]))
536 for (const path of s.remove.worktrees) {
537 const w = byPath.get(path)
538 // A missing directory: prune drops the entry.
539 if (w?.prunable) { removed.worktrees.push(path); continue }
540 if (s.unlock.includes(path)) await git(['worktree', 'unlock', path])
541 // Untracked type links block a plain remove: unlink the links themselves (never followed). A real directory under that name is somebody's files: left for git to refuse.
542 for (const l of s.links[path] ?? []) {
543 const at = `${path}/${l}`
544 // `test -L` then `rm` (no trailing slash, no -r): removes the link, never its target.
545 const isLink = await io.exec(['test', '-L', at], 10_000).then(r => r.exitCode === 0, () => false)
546 if (isLink) await io.exec(['rm', '-f', '--', at], 10_000).catch(() => undefined)
547 }
548 const r = await git(['worktree', 'remove', path])
549 if (r.exitCode === 0) removed.worktrees.push(path)
550 else {
551 failed.push({ kind: 'worktree', name: path, reason: `kept: ${r.stderr.trim().split('\n')[0] || 'git worktree remove failed'}`, needsLook: true })
552 if (w?.branch !== undefined) stuck.add(w.branch)
553 }
554 }
555 await git(['worktree', 'prune'])
556 for (const b of s.remove.branches) {
557 if (stuck.has(b)) {
558 failed.push({ kind: 'branch', name: b, reason: 'its worktree could not be removed', needsLook: true })
559 continue
560 }
561 const r = await git(['branch', '-D', b])
562 if (r.exitCode === 0) removed.branches.push(b)
563 else failed.push({ kind: 'branch', name: b, reason: `kept: ${r.stderr.trim().split('\n')[0] || 'git branch -D failed'}`, needsLook: true })
564 }
565 await io.setLeftovers(leftoverCounts(s, failed, true))
566 if (removed.worktrees.length || removed.branches.length) {
567 await io.append({
568 event: 'clean', owner: 'main',
569 text: `removed ${[...removed.worktrees.map(p => `worktree ${p.split('/').pop()}`), ...removed.branches.map(b => `branch ${b}`)].join(', ')}`
570 + s.dropped.filter(d => (d.kind === 'worktree' ? removed.worktrees : removed.branches).includes(d.name))
571 .map(d => `; dropped unpushed in ${d.kind} ${d.name.split('/').pop()}: ${d.commits.join(' | ')}`).join(''),
572 })
573 }
574 return sweepText(s, { applied: true, removed, failed, ...(note && { note }) })
575 })
576}
577hooks/resume.ts 204 lines1// The /flow resume leftovers: unfinished flow work an earlier session left behind, and the
2// instructions that restart it. The engine handle `$` is never passed across an import, so the
3// caller hands in these calls.
4
5import type { Handover, LogEvent, Session } from '../types'
6import { LIVE } from './deliver'
7import { parsePorcelain } from './clean'
8import { ownerFor } from './state'
9
10export type ResumeIo = {
11 // A process run; a failure to start throws.
12 exec: (argv: string[], timeoutMs: number) => Promise<{ exitCode: number; stdout: string; stderr: string }>
13 agents: () => Promise<{ id: string; status: string; name?: string }[]>
14 sessions: () => Promise<Session[]>
15 log: () => Promise<LogEvent[]>
16 loadHandovers: () => Promise<Record<string, Handover>>
17 saveHandover: (h: Handover) => Promise<void>
18 // State on disk is best-effort: a failure goes to the UI log.
19 best: (what: string, fn: () => Promise<void>) => Promise<void>
20 // Adds the handovers read from disk to the atom; what the atom already holds wins.
21 mergeHandovers: (all: Record<string, Handover>) => Promise<void>
22}
23
24const OWNED = new Set([...LIVE, 'idle'])
25
26const resumeHead = (limit: number): string => [
27 'The user ran /flow resume. Unfinished flow work was found (below). Start flow managers for it with the Agent tool, without asking:',
28 '- One flow:manager per task, named resume-<slug>, run_in_background true, at most ' + limit + ' at a time; start the rest as each finishes. Items whose branch names share a manager prefix (flow/csv-export-endpoint and flow/csv-export-button) are one task.',
29 '- Each manager\'s prompt carries, for every item of its task: the branch, the PR number and URL, the PR description (including any ## Handoff section), and for a worktree its path. It carries the work on from there and must not redo work already merged into the base branch.',
30].join('\n')
31
32// Unfinished flow work left behind by an earlier session, as `/flow resume` lists it.
33export type Leftover = { key: string; branch?: string; kind: 'pr' | 'branch' | 'worktree'; line: string; detail: string; owner?: string }
34
35// `error` with items means GitHub was unavailable and the items come from the state dir alone.
36export type Gathered = { items: Leftover[]; skipped: Leftover[]; error?: string }
37
38// Handovers on disk that are not finished, as resume items. A handover whose PR GitHub shows as
39// merged counts as done: it is marked so in the atom and not restarted.
40export async function diskItems(io: ResumeIo, items: Leftover[], merged: { prs: Set<number>; branches: Set<string> } | undefined): Promise<void> {
41 const all = await io.loadHandovers()
42 const events = await io.log()
43 const hs = Object.values(all)
44 if (hs.length === 0) return
45 const finished = (h: Handover) => h.status === 'done' || (merged !== undefined && (merged.prs.has(h.pr) || merged.branches.has(h.branch)))
46 await io.best('restoring handovers', async () => {
47 for (const h of hs) {
48 if (h.status !== 'done' && finished(h)) {
49 all[String(h.pr)] = { ...h, status: 'done' }
50 await io.saveHandover(all[String(h.pr)] as Handover)
51 }
52 }
53 await io.mergeHandovers(all)
54 })
55 for (const h of hs) {
56 if (finished(h)) continue
57 const note = `Handover #${h.pr} is ${h.status}${h.status === 'returned' ? ` (${h.reason ?? 'no reason'})` : ''}, reported to ${h.reportTo}. Verified: ${h.verified} Pending: ${h.pending}` +
58 (h.status === 'awaiting' ? ` It awaits the user's /flow approve ${h.pr}; do not restart work on it.` : '') +
59 (h.status === 'ready' ? ' It is in a batch the reviewer built and checked, which awaits the user\'s /flow push; do not restart work on it.' : '')
60 const mate = items.find(i => i.key === h.branch)
61 if (mate !== undefined) {
62 mate.line += ` | handover ${h.status}`
63 mate.detail += `\n${note}`
64 } else {
65 items.push({ key: h.branch, branch: h.branch, kind: 'pr', line: `#${h.pr} ${h.branch}: handover ${h.status} — ${h.title}`, detail: `Branch ${h.branch}, PR #${h.pr}, title "${h.title}".\n${note}` })
66 }
67 const at = items.find(i => i.key === h.branch)
68 if (at !== undefined) at.owner = ownerFor(events, { pr: h.pr, branch: h.branch }) ?? h.reportTo
69 }
70 for (const i of items) {
71 if (i.owner === undefined && i.branch !== undefined) i.owner = ownerFor(events, { branch: i.branch })
72 }
73}
74
75// Finds what a restart leaves behind: flow/* PRs and pushed branches, and worktrees with
76// uncommitted or unpushed work. Any git or gh failure becomes one line, never a throw.
77export async function gatherLeftovers(io: ResumeIo, base: string, resumed: Set<string>): Promise<Gathered> {
78 const run = async (argv: string[]) => {
79 try {
80 return await io.exec(argv, 60_000)
81 } catch (err) {
82 return { exitCode: 1, stdout: '', stderr: err instanceof Error ? err.message : String(err) }
83 }
84 }
85 const fail = async (what: string, r: { stderr: string }): Promise<Gathered> => {
86 const items: Leftover[] = []
87 await diskItems(io, items, undefined)
88 return {
89 items: items.filter(i => !resumed.has(i.key)), skipped: items.filter(i => resumed.has(i.key)),
90 error: `Cannot look for unfinished work: ${what} failed: ${r.stderr.trim().split('\n')[0]?.slice(0, 200) || 'no output'}${items.length ? '. GitHub was unavailable: these come from the state dir only.' : ''}`,
91 }
92 }
93
94 const fetched = await run(['git', 'fetch', 'origin', '--prune'])
95 if (fetched.exitCode !== 0) return await fail('git fetch origin', fetched)
96 const open = await run(['gh', 'pr', 'list', '--state', 'open', '--json', 'number,title,headRefName,isDraft,url,body', '--limit', '100'])
97 if (open.exitCode !== 0) return await fail('gh pr list', open)
98 const all = await run(['gh', 'pr', 'list', '--state', 'all', '--json', 'number,headRefName,headRefOid,state', '--limit', '200'])
99 if (all.exitCode !== 0) return await fail('gh pr list', all)
100 const refs = await run(['git', 'for-each-ref', '--format=%(refname:short)', 'refs/remotes/origin/flow/'])
101 if (refs.exitCode !== 0) return await fail('git for-each-ref', refs)
102
103 type Pr = { number: number; title: string; headRefName: string; isDraft: boolean; url: string; body: string }
104 const prs = (JSON.parse(open.stdout || '[]') as Pr[]).filter(p => p.headRefName.startsWith('flow/'))
105 const ended = (JSON.parse(all.stdout || '[]') as { number?: number; headRefName: string; headRefOid?: string; state: string }[])
106 .filter(p => p.state !== 'OPEN')
107 const closed = new Set(ended.map(p => p.headRefName))
108 // A squash-merged branch is deleted on the remote, so its worktree looks unpushed: match by head too.
109 const endedHeads = new Set(ended.map(p => p.headRefOid).filter(Boolean))
110 const withPr = new Set(prs.map(p => p.headRefName))
111 const branches = refs.stdout.split('\n').map(l => l.trim().replace(/^origin\//, ''))
112 .filter(b => b.startsWith('flow/') && !withPr.has(b) && !closed.has(b))
113
114 // A live agent named like the branch owns it.
115 const liveAgents = (await io.agents()).filter(a => OWNED.has(a.status))
116 const owners = new Set(liveAgents.filter(a => a.name !== undefined).map(a => `flow/${a.name}`))
117 for (const s of await io.sessions()) if (s.status === 'running' || s.status === 'reported') owners.add(s.branch)
118 // A worktree path ends in agent-<agentId>: that covers a worker that has not renamed its branch, and the reviewer.
119 const liveIds = new Set(liveAgents.map(a => `agent-${a.id}`))
120
121 const items: Leftover[] = []
122 for (const p of prs) {
123 items.push({
124 key: p.headRefName, kind: 'pr',
125 line: `#${p.number} ${p.headRefName}${p.isDraft ? ' (draft)' : ''}: ${p.title} ${p.url}`,
126 detail: `Branch ${p.headRefName}, PR #${p.number} ${p.url}${p.isDraft ? ' (draft)' : ''}, title "${p.title}".\nPR description:\n${p.body.trim() || '(empty)'}`,
127 })
128 }
129 for (const b of branches) {
130 items.push({ key: b, kind: 'branch', line: `${b} (pushed, no PR)`, detail: `Branch ${b}, pushed, no PR yet.` })
131 }
132
133 const wt = await run(['git', 'worktree', 'list', '--porcelain'])
134 if (wt.exitCode === 0) {
135 for (const { path, branch } of parsePorcelain(wt.stdout)) {
136 if (!path.includes('/.claude/worktrees/')) continue
137 if (liveIds.has(path.split('/').pop() ?? '') || (branch !== undefined && closed.has(branch))) continue
138 const dirty = (await run(['git', '-C', path, 'status', '--porcelain'])).stdout.trim() !== ''
139 if (branch === undefined && !dirty) continue
140 let unpushed = ''
141 if (!dirty) {
142 const up = await run(['git', '-C', path, 'rev-parse', '--abbrev-ref', '@{u}'])
143 unpushed = (await run(up.exitCode === 0
144 ? ['git', '-C', path, 'log', '--oneline', '@{u}..']
145 : ['git', '-C', path, 'log', '--oneline', `origin/${base}..HEAD`])).stdout.trim()
146 }
147 if (!dirty && unpushed === '') continue
148 const head = (await run(['git', '-C', path, 'rev-parse', 'HEAD'])).stdout.trim()
149 if (endedHeads.has(head)) continue
150 if (!dirty && (await run(['git', '-C', path, 'branch', '-r', '--contains', 'HEAD'])).stdout.trim() !== '') continue
151 const what = dirty ? 'uncommitted changes' : `${unpushed.split('\n').length} unpushed commit(s)`
152 const mate = branch === undefined ? undefined : items.find(i => i.key === branch)
153 if (mate !== undefined) {
154 mate.line += ` | worktree ${path} (${what})`
155 mate.detail += `\nIts worktree: ${path} (${what}).`
156 continue
157 }
158 items.push({
159 key: path, branch, kind: 'worktree', line: `${path} on ${branch ?? 'detached HEAD'}: ${what}`,
160 detail: `Worktree ${path}, branch ${branch ?? 'detached HEAD'}: ${what}.`,
161 })
162 }
163 }
164
165 await diskItems(io, items, {
166 prs: new Set(ended.filter(p => p.state === 'MERGED' && p.number !== undefined).map(p => p.number as number)),
167 branches: new Set(ended.filter(p => p.state === 'MERGED').map(p => p.headRefName)),
168 })
169
170 const live = (i: Leftover) => owners.has(i.branch ?? i.key)
171 return {
172 items: items.filter(i => !live(i) && !resumed.has(i.key)),
173 skipped: items.filter(i => !live(i) && resumed.has(i.key)),
174 }
175}
176
177export type OwnerNotes = { path: string; text: string }
178
179export function resumeInstructions(items: Leftover[], limit: number, notes: Map<string, OwnerNotes> = new Map()): string {
180 const owners = [...new Set(items.map(i => i.owner).filter((o): o is string => o !== undefined))]
181 // Recorded owners: the state on disk says who ran each item and what the user told them.
182 const grouped = owners.length === 0 ? [] : [
183 '',
184 'Flow state on disk names the manager that owned each item. Items of one owner go to ONE manager. Its prompt tells it to read its notes first (mcp__flow__note with manager = the owner name below and no text), and carries the notes quoted here. A manager with notes but no item below is not restarted.',
185 ...owners.flatMap(o => {
186 const n = notes.get(o)
187 return [
188 '',
189 `Owner ${o}:`,
190 ...items.filter(i => i.owner === o).map(i => `- ${i.detail.replaceAll('\n', '\n ')}`),
191 n?.text ? `Notes (${n.path}):\n${n.text}` : 'Notes: none.',
192 ]
193 }),
194 ...(items.some(i => i.owner === undefined) ? ['', 'No recorded owner:', ...items.filter(i => i.owner === undefined).map(i => `- ${i.detail.replaceAll('\n', '\n ')}`)] : []),
195 ]
196 if (grouped.length) return [resumeHead(limit), ...grouped].join('\n')
197 return [
198 resumeHead(limit),
199 '',
200 'Found:',
201 ...items.map(i => `- ${i.detail.replaceAll('\n', '\n ')}`),
202 ].join('\n')
203}
204hooks/deliver.ts 153 lines1// Message delivery to agents. A message to an idle agent starts a new turn, and when the prompt
2// cache (5 minutes) has expired that turn re-writes the whole context at full price. So a running
3// agent gets a message at once, an idle one gets messages that arrive close together as ONE wake-up,
4// and a message that needs no action yet (a non-blocking ask) waits for the agent's next turn.
5// The queue is in memory: a plugin reload loses it.
6
7// The engine handle `$` is never passed across an import, so the caller hands in these calls.
8export type DeliverIo = {
9 // The agent's status; undefined when it is not on the roster; throws when the roster cannot be read.
10 status: (id: string) => Promise<string | undefined>
11 // Throws when the agent is gone or the session refuses the message.
12 send: (id: string, text: string) => Promise<void>
13 after: (ms: number, fn: () => void) => void
14 // The name to show in a fallback text; the id is used when absent.
15 name?: (id: string) => Promise<string | undefined>
16}
17
18export const ENDED = new Set(['completed', 'failed', 'killed'])
19export const LIVE = new Set(['pending', 'running', 'waiting'])
20
21// How long an idle agent's first queued message waits for company.
22export const WINDOW_MS = 8_000
23// The longest a held message waits when nothing else wakes the agent.
24export const HOLD_CAP_MS = 600_000
25// How long a queue held for a reviewer's message waits if that message never lands.
26export const REVIEWER_HOLD_MS = 60_000
27
28export type DeliverOptions = {
29 // Flush the queue with this text at once, even if the agent is idle.
30 urgent?: boolean
31 // Do not wake an idle agent for this text: wait for the next delivery, its next turn or the cap.
32 held?: boolean
33 // Where the text goes when the agent has ended by the time it is flushed (default: main).
34 onGone?: (text: string) => void | Promise<void>
35}
36
37type Item = { text: string; held: boolean; onGone?: DeliverOptions['onGone'] }
38// refused: a send of this queue was refused while the agent was not ended; the next cap flush is final.
39type Queue = { items: Item[]; windowMs: number; generation: number; refused?: boolean }
40
41const queues = new Map<string, Queue>()
42let generation = 0
43let toMain: (text: string) => void = () => undefined
44
45export const resetDelivery = (main: (text: string) => void): void => {
46 queues.clear()
47 toMain = main
48}
49
50// The joined text of the queue: arrival order, identical texts once, a blank line between.
51const joined = (items: Item[]): string => [...new Set(items.map(i => i.text))].join('\n\n')
52
53// An unreadable roster is not an ended agent: 'unknown' is sent to now, and the send says if it is gone.
54const statusOf = (io: DeliverIo, id: string): Promise<string | undefined> => io.status(id).catch(() => 'unknown')
55
56// Send a text now; false when the agent is gone or the send fails.
57const sendNow = (io: DeliverIo, id: string, text: string): Promise<boolean> =>
58 io.send(id, text).then(() => true, () => false)
59
60// The agent is gone, or kept refusing: each item goes to its own fallback, else to main with the target named.
61async function gone(io: DeliverIo, id: string, items: Item[], why: 'ended' | 'refused' = 'ended'): Promise<void> {
62 const name = (await io.name?.(id).catch(() => undefined)) ?? id
63 for (const i of items) {
64 if (i.onGone !== undefined) await i.onGone(i.text)
65 else toMain(why === 'ended' ? `${i.text} (${name} ended before this was delivered.)` : `${i.text} (This could not be delivered to ${name}: it kept refusing the message.)`)
66 }
67}
68
69// Put refused items back, ahead of anything queued since, as held: they go out on the agent's next
70// turn or delivery, and at the cap timer at the latest.
71function requeue(io: DeliverIo, id: string, items: Item[]): void {
72 const q = queues.get(id) ?? { items: [], windowMs: 0, generation: 0 }
73 queues.set(id, q)
74 const seen = new Set<string>()
75 q.items = [...items, ...q.items].filter(i => !seen.has(i.text) && seen.add(i.text))
76 for (const i of q.items) i.held = true
77 q.refused = true
78 arm(io, id, q, HOLD_CAP_MS)
79}
80
81// Send items as one message. A refusal by an agent that has not ended is not a gone agent (the host
82// refuses one waiting between turns): the items are re-queued, or fall back when this was the final
83// try. True when sent or re-queued; false when the items went to their fallbacks.
84async function sendItems(io: DeliverIo, id: string, items: Item[], final = false): Promise<boolean> {
85 if (await sendNow(io, id, joined(items))) return true
86 const status = await statusOf(io, id)
87 if (status === undefined || ENDED.has(status)) await gone(io, id, items)
88 else if (final) await gone(io, id, items, 'refused')
89 else {
90 requeue(io, id, items)
91 return true
92 }
93 return false
94}
95
96// Send everything queued for an agent as one message. A no-op when nothing is queued.
97// The cap timer passes capTimer: a queue that was refused before and is refused again falls back.
98export async function flushAgent(io: DeliverIo, id: string, capTimer = false): Promise<void> {
99 const q = queues.get(id)
100 if (q === undefined || q.items.length === 0) return
101 queues.delete(id)
102 const status = await statusOf(io, id)
103 if (status === undefined || ENDED.has(status)) await gone(io, id, q.items)
104 else await sendItems(io, id, q.items, capTimer && q.refused === true)
105}
106
107// Arm the flush timer: the window when something is not held, else the cap. A shorter deadline
108// replaces a longer one; the stale timer sees its generation changed and does nothing.
109function arm(io: DeliverIo, id: string, q: Queue, ms: number): void {
110 const g = ++generation
111 q.generation = g
112 q.windowMs = ms
113 io.after(ms, () => {
114 if (queues.get(id)?.generation === g) void flushAgent(io, id, true)
115 })
116}
117
118// True when the text was sent or is queued, also re-queued after a refusal (it goes out later); false when the
119// agent is gone. A caller that must know about a refusal reads it from its io (see sendWrapUp).
120export async function deliver(io: DeliverIo, id: string, text: string, opts: DeliverOptions = {}): Promise<boolean> {
121 const status = await statusOf(io, id)
122 if (status === undefined || ENDED.has(status)) {
123 // Not a live target: whatever was queued for it and this text fall back together.
124 const q = queues.get(id)
125 queues.delete(id)
126 await gone(io, id, [...(q?.items ?? []), { text, held: false, onGone: opts.onGone }])
127 return false
128 }
129 const q = queues.get(id)
130 const item: Item = { text, held: opts.held === true, onGone: opts.onGone }
131 if (status !== 'idle' || opts.urgent === true) {
132 // Running (or urgent): the agent takes everything queued for it now, with this text last.
133 queues.delete(id)
134 return sendItems(io, id, [...(q?.items ?? []), item])
135 }
136 const queue: Queue = q ?? { items: [], windowMs: 0, generation: 0 }
137 queues.set(id, queue)
138 if (!queue.items.some(i => i.text === text)) queue.items.push(item)
139 // A held text must not shorten the wait of a queue that already has a window running.
140 if (queue.generation === 0 || (!item.held && queue.windowMs > WINDOW_MS)) arm(io, id, queue, item.held ? HOLD_CAP_MS : WINDOW_MS)
141 return true
142}
143
144// A reviewer's message to this agent is about to land and will start its turn: what is queued waits
145// for that turn (flushAgent on its first step) instead of waking the agent a second time. The short
146// timer sends it anyway if the message never lands.
147export function holdForReviewer(io: DeliverIo, id: string): void {
148 const q = queues.get(id)
149 if (q === undefined || q.items.length === 0) return
150 for (const i of q.items) i.held = true
151 arm(io, id, q, REVIEWER_HOLD_MS)
152}
153hooks/cost.ts 275 lines1// The cost meter's pure part: one price table, the ledger arithmetic and the text of the report.
2//
3// Prices are ESTIMATES at Anthropic first-party API list prices, taken 2026-10-10. Subscription
4// plans are billed differently, so read the dollars as "what this would cost at list price".
5// To update: edit PRICES below (USD per million tokens) and nothing else; the ledger keeps tokens,
6// not dollars, so a changed table reprices history.
7
8import { localDate } from './release'
9import type { Bucket, Ledger, LedgerEntry, Role, Tokens } from '../types'
10
11export type Rates = { input: number; write5m: number; write1h: number; read: number; output: number }
12
13// Cache write 5m = 1.25x input, 1h = 2x input, cache read = 0.1x input, unless a model lists its own read price.
14function rates(input: number, output: number, read = input * 0.1): Rates {
15 return { input, write5m: input * 1.25, write1h: input * 2, read, output }
16}
17
18type Price = { rates: Rates; long?: { over: number; rates: Rates } }
19
20const PRICES: Record<string, Price> = {
21 'claude-opus-5-5': { rates: rates(4, 20, 0.2) },
22 'claude-opus-5': { rates: rates(5, 25) },
23 'claude-opus-4-8': { rates: rates(5, 25) },
24 'claude-opus-4-7': { rates: rates(5, 25) },
25 'claude-opus-4-6': { rates: rates(5, 25) },
26 'claude-sonnet-5-5': { rates: rates(2, 10, 0.2) },
27 'claude-sonnet-5': { rates: rates(2, 10, 0.2) },
28 'claude-sonnet-4-6': { rates: rates(3, 15) },
29 // Prompts over 100K tokens (input + cache write + cache read of one step) cost more.
30 'claude-haiku-5-5': { rates: rates(0.1, 0.5), long: { over: 100_000, rates: rates(0.5, 2.5) } },
31 'claude-haiku-4-5': { rates: rates(1, 5) },
32 'claude-fable-5-1': { rates: rates(10, 50, 0.25) },
33 'claude-fable-5': { rates: rates(10, 50) },
34}
35
36// An alias, or a version not in the table, is priced as the newest of its family.
37const FAMILIES: Record<string, string> = {
38 opus: 'claude-opus-5-5',
39 sonnet: 'claude-sonnet-5-5',
40 haiku: 'claude-haiku-5-5',
41 fable: 'claude-fable-5-1',
42}
43
44// The same normalisation as modelKey in register.tsx: `[1m]`, a date suffix and `-latest` do not change the price.
45export function normalizeModel(model: string): string {
46 return model.toLowerCase().replace(/\[1m\]/g, '').replace(/-latest$/, '').replace(/-\d{8}$/, '')
47}
48
49// The price of a model id: the exact id, else the longest known id it extends, else its family, else undefined.
50export function priceOf(model: string): Price | undefined {
51 const key = normalizeModel(model)
52 const exact = PRICES[key]
53 if (exact) return exact
54 const known = Object.keys(PRICES).filter(k => key.startsWith(`${k}-`)).sort((a, b) => b.length - a.length)[0]
55 if (known) return PRICES[known]
56 const family = Object.keys(FAMILIES).find(f => key.includes(f))
57 return family === undefined ? undefined : PRICES[FAMILIES[family]!]
58}
59
60export const ZERO: Tokens = { input: 0, write5m: 0, write1h: 0, read: 0, output: 0 }
61
62const add = (a: Tokens, b: Tokens): Tokens => ({
63 input: a.input + b.input, write5m: a.write5m + b.write5m, write1h: a.write1h + b.write1h, read: a.read + b.read, output: a.output + b.output,
64})
65const sum = (list: Tokens[]): Tokens => list.reduce(add, ZERO)
66const all = (t: Tokens): number => t.input + t.write5m + t.write1h + t.read + t.output
67const prompt = (t: Tokens): number => t.input + t.write5m + t.write1h + t.read
68
69const dollars = (t: Tokens, r: Rates): number =>
70 (t.input * r.input + t.write5m * r.write5m + t.write1h * r.write1h + t.read * r.read + t.output * r.output) / 1_000_000
71
72// The USD estimate of a bucket, or undefined when the model has no price (shown as "?", never as 0).
73export function costOf(model: string, b: Bucket): number | undefined {
74 const p = priceOf(model)
75 if (p === undefined) return undefined
76 if (p.long === undefined || b.long === undefined) return dollars(b, p.rates)
77 const base: Tokens = {
78 input: b.input - b.long.input, write5m: b.write5m - b.long.write5m, write1h: b.write1h - b.long.write1h,
79 read: b.read - b.long.read, output: b.output - b.long.output,
80 }
81 return dollars(base, p.rates) + dollars(b.long, p.long.rates)
82}
83
84// cache read / everything the prompts carried; undefined before any prompt.
85export function hitRate(t: Tokens): number | undefined {
86 return prompt(t) === 0 ? undefined : t.read / prompt(t)
87}
88
89const num = (v: unknown): number => (typeof v === 'number' && Number.isFinite(v) && v > 0 ? v : 0)
90
91// One step's usage as tokens; the cache write splits 5m/1h when the usage carries the breakdown, else all is 5m.
92export function stepTokens(usage: unknown): Tokens | undefined {
93 if (typeof usage !== 'object' || usage === null) return undefined
94 const u = usage as Record<string, unknown>
95 const bd = typeof u.cache_creation === 'object' && u.cache_creation !== null ? u.cache_creation as Record<string, unknown> : undefined
96 const total = num(u.cache_creation_input_tokens)
97 const write1h = bd === undefined ? 0 : num(bd.ephemeral_1h_input_tokens)
98 const write5m = bd !== undefined && bd.ephemeral_5m_input_tokens !== undefined ? num(bd.ephemeral_5m_input_tokens) : Math.max(0, total - write1h)
99 return { input: num(u.input_tokens), write5m, write1h, read: num(u.cache_read_input_tokens), output: num(u.output_tokens) }
100}
101
102export type Identity = Partial<Omit<LedgerEntry, 'models' | 'firstAt' | 'lastAt' | 'turns'>>
103
104const definedOnly = <T extends object>(o: T): Partial<T> => Object.fromEntries(Object.entries(o).filter(([, v]) => v !== undefined)) as Partial<T>
105
106// The ledger with one step added (a new object; the old one is untouched). A step with no tokens at all is skipped.
107export function addStep(ledger: Ledger, key: string, model: string, usage: unknown, who: Identity | undefined, now: number): Ledger {
108 const t = stepTokens(usage)
109 if (t === undefined || all(t) === 0) return ledger
110 const isMain = key.startsWith('main@')
111 const cur: LedgerEntry = ledger[key] ?? { role: isMain ? 'main' : 'other', name: isMain ? 'main' : key, models: {} }
112 const entry: LedgerEntry = { ...cur, ...(who ? definedOnly(who) : {}), firstAt: cur.firstAt ?? now, lastAt: now }
113 const b: Bucket = entry.models[model] ?? { ...ZERO }
114 const over = priceOf(model)?.long?.over
115 const long = over !== undefined && prompt(t) > over
116 const next: Bucket = { ...add(b, t), ...(long || b.long ? { long: add(b.long ?? ZERO, long ? t : ZERO) } : {}) }
117 return { ...ledger, [key]: { ...entry, models: { ...entry.models, [model]: next } } }
118}
119
120// The ledger with one completed turn counted. An agent with no row yet gets an empty one (no models,
121// so the cost block hides it until a step arrives and its tokens are counted).
122export function addTurn(ledger: Ledger, key: string, now: number): Ledger {
123 const isMain = key.startsWith('main@')
124 const cur: LedgerEntry = ledger[key] ?? { role: isMain ? 'main' : 'other', name: isMain ? 'main' : key, models: {} }
125 return { ...ledger, [key]: { ...cur, turns: (cur.turns ?? 0) + 1, firstAt: cur.firstAt ?? now, lastAt: now } }
126}
127
128// Identity recorded at spawn (or filled from the roster): sets the fields, keeps the tokens.
129export function setIdentity(ledger: Ledger, key: string, who: Identity & { role: Role; name: string }, now: number): Ledger {
130 const cur = ledger[key]
131 return { ...ledger, [key]: { models: {}, ...cur, ...definedOnly(who), firstAt: cur?.firstAt ?? now, lastAt: now } as LedgerEntry }
132}
133
134// What came off disk: anything that is not a ledger entry is dropped.
135export function normalizeLedger(raw: unknown): Ledger {
136 if (typeof raw !== 'object' || raw === null || Array.isArray(raw)) return {}
137 const out: Ledger = {}
138 for (const [k, v] of Object.entries(raw)) {
139 if (typeof v !== 'object' || v === null) continue
140 const e = v as Record<string, unknown>
141 if (typeof e.name !== 'string' || typeof e.models !== 'object' || e.models === null) continue
142 out[k] = { ...(e as unknown as LedgerEntry), role: (['manager', 'worker', 'reviewer', 'main'] as const).find(r => r === e.role) ?? 'other' }
143 }
144 return out
145}
146
147// Both ledgers' tokens added (the file's history plus what this run counted before it was loaded).
148export function mergeLedgers(a: Ledger, b: Ledger): Ledger {
149 const out: Ledger = { ...a }
150 for (const [k, e] of Object.entries(b)) {
151 const o = out[k]
152 if (o === undefined) { out[k] = e; continue }
153 const models = { ...o.models }
154 for (const [m, bk] of Object.entries(e.models)) {
155 const x = models[m]
156 models[m] = x === undefined ? bk : { ...add(x, bk), ...(x.long || bk.long ? { long: add(x.long ?? ZERO, bk.long ?? ZERO) } : {}) }
157 }
158 out[k] = {
159 ...e, ...o, models,
160 ...((o.turns ?? 0) + (e.turns ?? 0) > 0 ? { turns: (o.turns ?? 0) + (e.turns ?? 0) } : {}),
161 ...(o.firstAt !== undefined || e.firstAt !== undefined ? { firstAt: Math.min(o.firstAt ?? Infinity, e.firstAt ?? Infinity) } : {}),
162 ...(o.lastAt !== undefined || e.lastAt !== undefined ? { lastAt: Math.max(o.lastAt ?? 0, e.lastAt ?? 0) } : {}),
163 }
164 }
165 return out
166}
167
168// The ledger file stays bounded: an entry not counted for 90 days goes, and so does one with no `lastAt`.
169export const KEEP_MS = 90 * 24 * 3600 * 1000
170export function pruneLedger(ledger: Ledger, now: number): Ledger {
171 const kept = Object.entries(ledger).filter(([, e]) => e.lastAt !== undefined && e.lastAt >= now - KEEP_MS)
172 return kept.length === Object.keys(ledger).length ? ledger : Object.fromEntries(kept)
173}
174
175// --- Totals ---
176
177export type Total = { tokens: Tokens; usd: number; unknown: boolean; turns: number }
178
179export function totalOf(entries: LedgerEntry[]): Total {
180 let usd = 0
181 let unknown = false
182 let turns = 0
183 const parts: Tokens[] = []
184 for (const e of entries) {
185 turns += e.turns ?? 0
186 for (const [m, b] of Object.entries(e.models)) {
187 parts.push(b)
188 const c = costOf(m, b)
189 if (c === undefined) unknown = true
190 else usd += c
191 }
192 }
193 return { tokens: sum(parts), usd, unknown, turns }
194}
195
196// A name without its successor suffix: `csv-export-2` is a successor of `csv-export`.
197export const baseName = (name: string): string => name.replace(/-\d+$/, '')
198
199// --- Text ---
200
201export function humanTokens(n: number): string {
202 if (n < 1000) return String(Math.round(n))
203 if (n < 1_000_000) return `${(n / 1000).toFixed(n < 100_000 ? 1 : 0).replace(/\.0$/, '')}k`
204 return `${(n / 1_000_000).toFixed(n < 100_000_000 ? 1 : 0).replace(/\.0$/, '')}M`
205}
206
207export function money(usd: number, unknown = false): string {
208 const s = usd < 0.1 && usd > 0 ? usd.toFixed(3) : usd.toFixed(2)
209 return unknown ? (usd > 0 ? `~$${s}+?` : '~$?') : `~$${s}`
210}
211
212const percent = (t: Tokens): string => {
213 const h = hitRate(t)
214 return h === undefined ? '-' : `${Math.round(h * 100)}%`
215}
216
217const tokenText = (t: Tokens): string =>
218 `${humanTokens(t.input)} in / ${humanTokens(t.write5m + t.write1h)} cache write / ${humanTokens(t.read)} cache read / ${humanTokens(t.output)} out, hit ${percent(t)}`
219
220// Turns and cache write per turn (the cost of each wake-up) show only once a turn was counted.
221const turnText = (t: Total): string =>
222 t.turns > 0 ? `, ${t.turns} turns, ~${humanTokens((t.tokens.write5m + t.tokens.write1h) / t.turns)} cache write/turn` : ''
223
224export const totalText = (t: Total): string => `${money(t.usd, t.unknown)} (tokens ${tokenText(t.tokens)}${turnText(t)})`
225
226// Appended to a reviewer's done report: the PR's workers' cost. Empty when they spent nothing.
227export function reportSuffix(t: Total): string {
228 if (all(t.tokens) === 0) return ''
229 return ` | cost: ${money(t.usd, t.unknown)} (tokens ${humanTokens(prompt(t.tokens))} in, ${humanTokens(t.tokens.output)} out, hit ${percent(t.tokens)})`
230}
231
232// The workers of a PR: those on its branch (a `-2` successor continues on the same branch).
233export const entriesOfBranch = (ledger: Ledger, branch: string): LedgerEntry[] =>
234 Object.values(ledger).filter(e => e.role === 'worker' && e.branch === branch)
235
236export const prCost = (ledger: Ledger, branch: string): Total => totalOf(entriesOfBranch(ledger, branch))
237
238// The block of `status`: one line per agent that took steps, then manager, PR, reviewer, main and session totals.
239// `scope`: this session started at `start` and has `live` agent ids in its roster. An entry is this
240// session's when its agent is live or was counted since the start; the older ones only reach the all-time line.
241export function costBlock(ledger: Ledger, prs: { pr: number; branch: string }[], otherSessions: boolean, scope: { start: number; live: Set<string> }): string[] {
242 const counted = Object.entries(ledger).filter(([, e]) => Object.keys(e.models).length > 0)
243 if (counted.length === 0) return []
244 const here = counted.filter(([k, e]) => scope.live.has(k) || (e.lastAt ?? 0) >= scope.start)
245 const entries = here.map(([, e]) => e)
246 const lines = ['Cost (estimates at API list prices):']
247 const order: Record<Role, number> = { main: 0, manager: 1, worker: 2, reviewer: 3, other: 4 }
248 for (const e of [...entries].sort((a, b) => order[a.role] - order[b.role] || a.name.localeCompare(b.name))) {
249 lines.push(` ${e.role} ${e.name} [${Object.keys(e.models).join(', ')}]: ${totalText(totalOf([e]))}`)
250 }
251 const ownerOf = (e: LedgerEntry): string | undefined => e.role === 'manager' ? baseName(e.name) : e.role === 'worker' && e.manager !== undefined ? baseName(e.manager) : undefined
252 for (const m of [...new Set(entries.map(ownerOf).filter((x): x is string => x !== undefined))].sort()) {
253 lines.push(` manager ${m} total (with its workers): ${totalText(totalOf(entries.filter(e => ownerOf(e) === m)))}`)
254 }
255 for (const p of [...prs].sort((a, b) => a.pr - b.pr)) {
256 const mine = entriesOfBranch(ledger, p.branch).filter(e => entries.includes(e) && Object.keys(e.models).length > 0)
257 if (mine.length > 0) lines.push(` PR #${p.pr} (${p.branch}) workers: ${totalText(totalOf(mine))}`)
258 }
259 const rev = entries.filter(e => e.role === 'reviewer')
260 if (rev.length > 0) lines.push(` reviewer total: ${totalText(totalOf(rev))}`)
261 const main = entries.filter(e => e.role === 'main')
262 if (main.length > 0) lines.push(` main total: ${totalText(totalOf(main))}`)
263 // Only once model routing has put a size on some entry; the rest is "unsized" (main, managers, the reviewer, older workers).
264 if (entries.some(e => e.size !== undefined)) {
265 const part = (label: string, of: LedgerEntry[]) => of.length === 0 ? [] : [`${label} ${money(totalOf(of).usd, totalOf(of).unknown)}`]
266 const by = [...['small', 'normal', 'large'].flatMap(s => part(s, entries.filter(e => e.size === s))), ...part('unsized', entries.filter(e => e.size === undefined))]
267 lines.push(` by size: ${by.join(', ')}`)
268 }
269 lines.push(` session total: ${totalText(totalOf(entries))}`)
270 const first = Math.min(...counted.map(([, e]) => e.firstAt ?? Infinity))
271 lines.push(` All time${Number.isFinite(first) ? ` (since ${localDate(first)})` : ''}: ${totalText(totalOf(counted.map(([, e]) => e)))}`)
272 if (otherSessions) lines.push(' Workers run in other harness sessions are not counted.')
273 return lines
274}
275hooks/routing.ts 57 lines1// Model routing, the pure part: a brief's `Size:` line picks the worker's model at spawn, and a
2// successor of a worker that failed runs one size up.
3
4export type Size = 'small' | 'normal' | 'large'
5export const SIZES: Size[] = ['small', 'normal', 'large']
6export type SizeModels = Record<Size, string>
7
8const SIZE_LINE = /^Size:\s*(small|normal|large)\b(.*)$/im
9
10// The brief's size and its reason (the rest of the line, minus the dash or parentheses around it).
11// An unknown word (`Size: tiny`) is no size line.
12export function parseSize(prompt: string): { size: Size; reason: string } | undefined {
13 const m = SIZE_LINE.exec(prompt)
14 if (m === null) return undefined
15 const reason = m[2]!.trim().replace(/^[\s\-–—(]+/, '').replace(/[\s)]+$/, '').trim()
16 return { size: m[1]!.toLowerCase() as Size, reason }
17}
18
19export const isSize = (v: unknown): v is Size => v === 'small' || v === 'normal' || v === 'large'
20
21export const modelFor = (size: Size, models: SizeModels): string => models[size]
22
23export const escalate = (size: Size): Size => SIZES[Math.min(SIZES.indexOf(size) + 1, SIZES.length - 1)]!
24
25export const biggerOf = (a: Size, b: Size): Size => (SIZES.indexOf(a) >= SIZES.indexOf(b) ? a : b)
26
27// The least size a successor may run, from its predecessors' recorded sizes: one up from the largest.
28// A predecessor with no recorded size (older data) says nothing; with none recorded there is no floor.
29export function floorFrom(predecessors: Array<Size | string | undefined>): Size | undefined {
30 const sizes = predecessors.filter(isSize)
31 if (sizes.length === 0) return undefined
32 return escalate(sizes.reduce(biggerOf))
33}
34
35// `x-2` and up is a successor of `x`; the number is its generation (a plain name is the first).
36export const generation = (name: string): number => {
37 const n = /-(\d+)$/.exec(name)
38 return n === null ? 1 : Number(n[1])
39}
40export const isSuccessorName = (name: string): boolean => generation(name) >= 2
41
42// The size a spawn ends up with: the brief's own, raised to the floor. Neither leaves it unset, which is
43// today's behaviour. A large floor is kept, so the next successor escalates from large, not from nothing.
44export function effectiveSize(declared: Size | undefined, floor: Size | undefined): Size | undefined {
45 if (declared === undefined || floor === undefined) return declared ?? floor
46 return biggerOf(declared, floor)
47}
48
49// Added to what a reviewer's "back" (or the user's send-back) tells the manager: when the PR's last
50// worker ran below large, the continuation is a fresh `-N` worker, which the plugin runs one size up.
51export function backNote(workers: Array<{ name: string; size?: string; spawnModel?: string }>, branch: string): string | undefined {
52 const last = [...workers].sort((a, b) => generation(a.name) - generation(b.name)).at(-1)
53 if (last === undefined || (last.size !== 'small' && last.size !== 'normal')) return undefined
54 const next = `${last.name.replace(/-\d+$/, '')}-${generation(last.name) + 1}`
55 return `Its worker ran ${last.size}${last.spawnModel === undefined ? '' : ` (${last.spawnModel})`}: continue with a fresh ${next} worker (Continue on branch: ${branch}); it runs one size up.`
56}
57hooks/migrations.ts 153 lines1// Migration numbers across parallel PRs: which ones clash or sit at or below the base, and the next free
2// number for each. Pure: the git and gh glue in register.tsx passes in path lists and file texts.
3
4export type Status = 'ok' | 'clash' | 'at-or-below'
5
6export type Ref = { file: string; line: number }
7
8// One migration of a PR: the files that share a number and a name (a paired up/down, or a folder).
9export type Migration = {
10 number: bigint
11 prefix: string
12 key: string
13 files: string[]
14 status: Status
15 clashWith?: string
16 newNumber?: string
17 moves?: Array<{ from: string; to: string }>
18 refs?: Ref[]
19}
20
21export type PrReport = { pr: number; error?: string; migrations: Migration[]; unnumbered: string[] }
22
23export type Report = { dir: string; baseHigh?: bigint; refHigh?: bigint; prs: PrReport[] }
24
25export type PrInput = { pr: number; error?: string; added: string[] }
26
27export type Inputs = { dir: string; baseFiles: string[]; refFiles: string[]; prs: PrInput[] }
28
29type Entry = { file: string; n: bigint; key: string }
30
31export function cleanDir(dir: string): string {
32 return dir.trim().replace(/^\.\//, '').replace(/\/+$/, '')
33}
34
35// The first path segment under dir, split into its leading digits and the rest. The same rule serves
36// `dir/0042_name.sql` and the folder style `dir/0042_name/up.sql`.
37export function parseEntry(dir: string, path: string): { prefix: string; segment: string; rest: string; key: string } | undefined {
38 const rel = path.startsWith(`${dir}/`) ? path.slice(dir.length + 1) : undefined
39 if (rel === undefined || rel === '') return undefined
40 const segment = rel.split('/')[0]!
41 const prefix = /^\d+/.exec(segment)?.[0]
42 if (prefix === undefined) return undefined
43 // A paired 0042_x.up.sql / 0042_x.down.sql is one migration: the key drops the extension and the up/down marker.
44 const stem = rel.includes('/') ? segment : segment.replace(/\.[^.]*$/, '').replace(/\.(up|down)$/, '')
45 return { prefix, segment, rest: rel.slice(segment.length), key: stem }
46}
47
48function entries(dir: string, files: string[]): Entry[] {
49 const out: Entry[] = []
50 for (const file of files) {
51 const e = parseEntry(dir, file)
52 if (e !== undefined) out.push({ file, n: BigInt(e.prefix), key: e.key })
53 }
54 return out
55}
56
57const top = (es: Entry[]): bigint | undefined => es.reduce<bigint | undefined>((a, e) => a === undefined || e.n > a ? e.n : a, undefined)
58
59export function analyze(input: Inputs): Report {
60 const dir = cleanDir(input.dir)
61 const base = entries(dir, input.baseFiles)
62 const ref = entries(dir, input.refFiles)
63 const baseHigh = top(base)
64 const refHigh = top(ref)
65 const landed = new Set([...input.baseFiles, ...input.refFiles])
66 // Everything a new migration must not collide with, in the order the batch lands.
67 const taken: Entry[] = [...base, ...ref]
68 let high = top(taken)
69
70 const prs: PrReport[] = []
71 for (const p of input.prs) {
72 const report: PrReport = { pr: p.pr, migrations: [], unnumbered: [] }
73 prs.push(report)
74 if (p.error !== undefined) { report.error = p.error; continue }
75 const groups = new Map<string, Migration>()
76 for (const file of [...p.added].sort()) {
77 const e = parseEntry(dir, file)
78 if (e === undefined) { report.unnumbered.push(file); continue }
79 const id = `${e.prefix}\0${e.key}`
80 const g = groups.get(id)
81 if (g !== undefined) g.files.push(file)
82 else groups.set(id, { number: BigInt(e.prefix), prefix: e.prefix, key: e.key, files: [file], status: 'ok' })
83 }
84 for (const m of groups.values()) {
85 report.migrations.push(m)
86 // Already on the base or at ref (the PR is part of the batch): nothing to renumber.
87 const present = m.files.every(f => landed.has(f))
88 const other = present ? undefined : taken.find(t => t.n === m.number && t.key !== m.key)
89 if (other !== undefined) { m.status = 'clash'; m.clashWith = other.file }
90 else if (!present && baseHigh !== undefined && m.number <= baseHigh) m.status = 'at-or-below'
91 if (m.status === 'ok') {
92 if (high === undefined || m.number > high) high = m.number
93 for (const file of m.files) taken.push({ file, n: m.number, key: m.key })
94 continue
95 }
96 // Highest + 1, never a gap fill; later flagged migrations get successive numbers.
97 const next = (high ?? 0n) + 1n
98 high = next
99 m.newNumber = String(next).padStart(m.prefix.length, '0')
100 m.moves = m.files.map(from => {
101 const e = parseEntry(dir, from)!
102 return { from, to: `${dir}/${m.newNumber}${e.segment.slice(e.prefix.length)}${e.rest}` }
103 })
104 for (const mv of m.moves) taken.push({ file: mv.to, n: next, key: m.key })
105 }
106 }
107 return { dir, ...(baseHigh !== undefined && { baseHigh }), ...(refHigh !== undefined && { refHigh }), prs }
108}
109
110// Lines in the PR's own changed files that mention a flagged migration's number (a whole token) or its
111// file or folder name. The migration's own files count by content (a migration that inserts its own
112// version into a schema-version table), never by path.
113export function findRefs(m: Migration, files: Array<{ path: string; text: string }>): Ref[] {
114 const names = new Set<string>()
115 for (const f of m.files) for (const s of f.split('/')) if (s.startsWith(m.prefix) && s.length > m.prefix.length) names.add(s)
116 const token = new RegExp(`(?<![0-9A-Za-z])${m.prefix}(?![0-9])`)
117 const out: Ref[] = []
118 for (const f of files) {
119 f.text.split('\n').forEach((line, i) => {
120 if (token.test(line) || [...names].some(n => line.includes(n))) out.push({ file: f.path, line: i + 1 })
121 })
122 }
123 return out
124}
125
126export function render(report: Report): string {
127 const show = (n: bigint | undefined) => n === undefined ? 'none' : String(n)
128 const lines = [
129 `Migrations in ${report.dir}`,
130 `Highest on the base: ${show(report.baseHigh)}`,
131 `Highest at ref: ${show(report.refHigh)}`,
132 ]
133 for (const p of report.prs) {
134 lines.push('', `PR #${p.pr}`)
135 if (p.error !== undefined) { lines.push(` not checked: ${p.error}`); continue }
136 if (p.migrations.length === 0 && p.unnumbered.length === 0) lines.push(' adds no migration')
137 for (const m of p.migrations) {
138 const label = m.files.join(', ')
139 if (m.status === 'ok') { lines.push(` ok: ${label} (number ${m.prefix})`); continue }
140 lines.push(m.status === 'clash'
141 ? ` clash: ${label} (number ${m.prefix}) has the same number as ${m.clashWith}`
142 : ` at-or-below: ${label} (number ${m.prefix}) is not above the highest on the base (${show(report.baseHigh)})`)
143 lines.push(` next free number: ${m.newNumber}`)
144 for (const mv of m.moves ?? []) lines.push(` git mv ${mv.from} ${mv.to}`)
145 for (const r of m.refs ?? []) lines.push(` references its number in ${r.file}:${r.line}`)
146 }
147 for (const u of p.unnumbered) lines.push(` unnumbered (ignored): ${u}`)
148 }
149 return lines.join('\n')
150}
151
152export const UNSET_TEXT = 'No migrations directory is set (the migrations_dir setting is empty): nothing checked.'
153hooks/dag.ts 321 lines1// The plan graph, pure: no `$` calls, so it is tested on its own. Each owner ("main" or a
2// manager's name) has a graph of nodes; an edge says "start B after A is done".
3import type { DagNode, Handover } from '../types'
4
5export type Graph = Record<string, DagNode>
6export type Plan = Record<string, Graph>
7
8// `children`: how many live agents work under this one (a manager with workers is not finished).
9// `at`: when this agent was last active; `childAt`: the newest activity among its children, live or not.
10export type AgentFact = { id?: string; name?: string; status: string; answer?: string; children?: number; at?: number; childAt?: number }
11// `owners`: branch -> the manager that started its worker, from the session log; it survives a wrong report_to.
12// `open`: owners whose own plan graph still has waiting, ready or running nodes (settle fills it in).
13export type Facts = { agents: AgentFact[]; handovers: Handover[]; phrases?: string[]; asking?: string[]; owners?: Record<string, string>; open?: string[] }
14
15export type NodeInput = { id: string; title?: string; after?: string[]; until?: 'merged' | 'reported' }
16
17export type Notice = {
18 owner: string
19 ready: DagNode[]
20 blocked: DagNode[]
21 // The finished dependencies behind the ready nodes, for the message.
22 done: DagNode[]
23 // Free manager slots; only set for owner "main".
24 slots?: number
25}
26
27// An agent with an open blocking ask in the decision inbox.
28const isAsking = (agent: AgentFact, facts: Facts): boolean => agent.name !== undefined && !!facts.asking?.includes(agent.name)
29
30// A phrase directly after one of these ("不需要你決定") says the opposite.
31const NEGATIONS = ['不', '不用', '不必', '無需', '毋需', '不需要', 'no ', 'not ', "don't ", 'no need to ']
32
33// The last line ends in a question mark, or the last paragraph holds one of the configured phrases
34// (for agents that write in a language where a question has no "?").
35export function asksQuestion(answer: string | undefined, phrases: string[] = []): boolean {
36 const text = (answer ?? '').trim()
37 const last = text.split('\n').pop() ?? ''
38 if (/[??][*_`'")\s]*$/.test(last)) return true
39 const paragraph = (text.split(/\n[ \t]*\n/).pop() ?? '').toLowerCase()
40 return phrases.some(p => {
41 const phrase = p.trim().toLowerCase()
42 if (phrase === '') return false
43 for (let i = paragraph.indexOf(phrase); i >= 0; i = paragraph.indexOf(phrase, i + 1)) {
44 const before = paragraph.slice(0, i)
45 if (!NEGATIONS.some(n => before.endsWith(n))) return true
46 }
47 return false
48 })
49}
50
51// The answer's last paragraph says the manager waits (for a merge or the reviewer's report) and that report
52// is the one it waits for: it names the same own PR, or names none and mentions merge/reviewer while the
53// report names one of the manager's PRs.
54export function waitsOnReport(answer: string | undefined, report: string, own: number[]): boolean {
55 const paragraph = ((answer ?? '').trim().split(/\n[ \t]*\n/).pop() ?? '')
56 if (!/\bwait/i.test(paragraph)) return false
57 const reported = new Set([...report.matchAll(/#(\d+)/g)].map(m => Number(m[1])))
58 const named = [...paragraph.matchAll(/#(\d+)/g)].map(m => Number(m[1]))
59 if (named.length > 0) return named.some(n => reported.has(n) && own.includes(n))
60 return /\bmerg|\breviewer/i.test(paragraph) && own.some(n => reported.has(n))
61}
62
63const lastLine = (answer: string | undefined) => (answer ?? '').trim().split('\n').pop() ?? ''
64const LIVE = new Set(['running', 'pending'])
65
66// Cycle path through the graph as "a -> b -> a", or undefined. `after` points at dependencies,
67// so the walk follows them; the path is printed in the order the dependencies run.
68function findCycle(graph: Graph): string | undefined {
69 const state: Record<string, 1 | 2> = {}
70 const stack: string[] = []
71 const visit = (id: string): string | undefined => {
72 if (state[id] === 2) return undefined
73 if (state[id] === 1) {
74 const from = stack.indexOf(id)
75 return [...stack.slice(from), id].reverse().join(' -> ')
76 }
77 state[id] = 1
78 stack.push(id)
79 for (const dep of graph[id]?.after ?? []) {
80 const found = visit(dep)
81 if (found) return found
82 }
83 stack.pop()
84 state[id] = 2
85 return undefined
86 }
87 for (const id of Object.keys(graph)) {
88 const found = visit(id)
89 if (found) return found
90 }
91 return undefined
92}
93
94// All or nothing: a refused call leaves the graph as it was.
95export function addNodes(graph: Graph, nodes: NodeInput[]): { graph: Graph } | { error: string } {
96 const next: Graph = { ...graph }
97 const added: string[] = []
98 for (const n of nodes) {
99 if (!n.id || !n.id.trim()) return { error: 'a node needs an id' }
100 if (next[n.id]) return { error: `duplicate id: ${n.id}` }
101 next[n.id] = {
102 id: n.id,
103 title: n.title ?? n.id,
104 after: [...new Set(n.after ?? [])],
105 until: n.until ?? 'merged',
106 state: 'waiting',
107 }
108 added.push(n.id)
109 }
110 for (const id of added) {
111 for (const dep of next[id]?.after ?? []) {
112 if (!next[dep]) return { error: `unknown dependency: ${id} comes after ${dep}, which is not in the plan` }
113 }
114 }
115 const cycle = findCycle(next)
116 if (cycle) return { error: `cycle: ${cycle}` }
117 return { graph: next }
118}
119
120// The agent named `id`, or its handoff continuations `<id>-N`; the newest has the highest N.
121// Among rows with the same name a live one beats a stale one.
122export function agentFor<T extends { name?: string; status?: string }>(id: string, agents: T[]): T | undefined {
123 let best: T | undefined
124 let bestN = -1
125 for (const a of agents) {
126 const m = a.name === id ? 1 : a.name?.startsWith(id + '-') && /^\d+$/.test(a.name.slice(id.length + 1)) ? Number(a.name.slice(id.length + 1)) : 0
127 if (m && (m > bestN || (m === bestN && best !== undefined && !LIVE.has(best.status ?? '') && LIVE.has(a.status ?? '')))) {
128 best = a
129 bestN = m
130 }
131 }
132 return best
133}
134
135const ownsBranch = (h: Handover, id: string) => h.branch === `flow/${id}`
136const reportsTo = (h: Handover, id: string) => h.reportTo === id || (h.reportTo.startsWith(id + '-') && /^\d+$/.test(h.reportTo.slice(id.length + 1)))
137
138const ownedBy = (h: Handover, id: string, facts: Facts) => {
139 const o = facts.owners?.[h.branch]
140 return o !== undefined && (o === id || (o.startsWith(id + '-') && /^\d+$/.test(o.slice(id.length + 1))))
141}
142
143// The newest handover per branch.
144function latestPerBranch(hs: Handover[]): Handover[] {
145 const by: Record<string, Handover> = {}
146 for (const h of hs) if (!by[h.branch] || h.at >= by[h.branch]!.at) by[h.branch] = h
147 return Object.values(by)
148}
149
150type Verdict = { state: 'done' | 'blocked' | 'running'; info?: string } | undefined
151
152function judge(owner: string, node: DagNode, facts: Facts): Verdict {
153 const agent = agentFor(node.id, facts.agents)
154 const live = agent && LIVE.has(agent.status)
155 const failed = agent && (agent.status === 'failed' || agent.status === 'killed')
156
157 if (node.until === 'reported') {
158 if (!agent) return undefined
159 if (agent.status === 'idle' || agent.status === 'completed') {
160 const last = lastLine(agent.answer)
161 if (last.startsWith('BLOCKED:')) return { state: 'blocked', info: last }
162 if (!asksQuestion(agent.answer, facts.phrases) && !isAsking(agent, facts) && !last.startsWith('HANDOFF:')) return { state: 'done', info: 'reported' }
163 }
164 if (failed) return { state: 'blocked', info: `agent ${agent.status}` }
165 return { state: 'running' }
166 }
167
168 if (owner !== 'main') {
169 // A worker's package: its PR merged.
170 const h = facts.handovers.filter(x => ownsBranch(x, node.id)).sort((a, b) => b.at - a.at)[0]
171 if (h?.status === 'done') return { state: 'done', info: `merged ${h.sha ?? h.head.slice(0, 7)}` }
172 if (h?.status === 'returned' && !live) return { state: 'blocked', info: `PR returned: ${h.reason ?? 'no reason given'}` }
173 if (failed) return { state: 'blocked', info: `agent ${agent.status}` }
174 return agent || h ? { state: 'running' } : undefined
175 }
176
177 // A manager's task: it finished cleanly and everything it handed over is merged.
178 if (!agent) return undefined
179 if (failed) return { state: 'blocked', info: `agent ${agent.status}` }
180 // A background manager that is done with its turn sits 'idle' (the host does not complete it), so
181 // idle counts as finished unless it has live workers, work still planned in its own graph, or a
182 // worker that reported after the manager's last turn (its notification has not reached the manager yet).
183 // Without a handover it is an investigation that never opened a PR.
184 const owned = latestPerBranch(facts.handovers.filter(x => reportsTo(x, node.id) || ownedBy(x, node.id, facts)))
185 if (agent.status === 'idle' ? (agent.children ?? 0) > 0 || (agent.childAt ?? 0) > (agent.at ?? 0) || (facts.open ?? []).some(o => o === node.id || (o.startsWith(node.id + '-') && /^\d+$/.test(o.slice(node.id.length + 1)))) : agent.status !== 'completed') return { state: 'running' }
186 const last = lastLine(agent.answer)
187 if (last.startsWith('BLOCKED:')) return { state: 'blocked', info: last }
188 if (asksQuestion(agent.answer, facts.phrases) || isAsking(agent, facts) || last.startsWith('HANDOFF:')) return { state: 'running' }
189 const hs = owned
190 const returned = hs.find(x => x.status === 'returned')
191 if (returned) return { state: 'blocked', info: `PR #${returned.pr} returned: ${returned.reason ?? 'no reason given'}` }
192 if (hs.some(x => x.status !== 'done')) return { state: 'running' }
193 return { state: 'done', info: hs.length ? `merged ${hs.map(x => x.sha ?? x.head.slice(0, 7)).join(', ')}` : 'completed' }
194}
195
196// One owner's graph against what the session shows now. Done is sticky; `manual` wins.
197export function evaluate(owner: string, graph: Graph, facts: Facts): Graph {
198 const out: Graph = {}
199 const visit = (id: string): DagNode => {
200 const seen = out[id]
201 if (seen) return seen
202 const prev = graph[id]
203 if (!prev) throw new Error(`unknown node: ${id}`)
204 // Marks the node while it is being worked out, so a (refused) cycle cannot recurse forever.
205 out[id] = prev
206 const deps = prev.after.map(visit)
207 let state = prev.state
208 let info = prev.info
209 if (prev.state === 'done') {
210 // sticky
211 } else if (prev.manual) {
212 state = prev.manual
213 info = info ?? 'set by hand'
214 } else {
215 const v = judge(owner, prev, facts)
216 if (v) {
217 state = v.state
218 info = v.info
219 } else {
220 state = deps.every(d => d.state === 'done') ? 'ready' : 'waiting'
221 info = undefined
222 }
223 }
224 const node: DagNode = { ...prev, state }
225 if (info === undefined) delete node.info
226 else node.info = info
227 out[id] = node
228 return node
229 }
230 for (const id of Object.keys(graph)) visit(id)
231 return out
232}
233
234// One pass over every owner: evaluate, and say once what newly became ready or blocked.
235export function settle(plan: Plan, facts: Facts, opts: { slots?: number } = {}): { plan: Plan; notices: Notice[] } {
236 const next: Plan = {}
237 const notices: Notice[] = []
238 const open = Object.entries(plan).filter(([, g]) => Object.values(g).some(n => n.state === 'waiting' || n.state === 'ready' || n.state === 'running')).map(([o]) => o)
239 facts = { ...facts, open: facts.open ?? open }
240 for (const [owner, graph] of Object.entries(plan)) {
241 const evaluated = evaluate(owner, graph, facts)
242 const ready: DagNode[] = []
243 const blocked: DagNode[] = []
244 const marked: Graph = {}
245 for (const n of Object.values(evaluated)) {
246 let m = n
247 // A node that left a state may enter it again later and deserves a new notice.
248 if (n.state === 'waiting' && n.readyNotified) m = { ...m, readyNotified: false }
249 if (n.state !== 'blocked' && n.blockedNotified) m = { ...m, blockedNotified: false }
250 if (n.state === 'ready' && !n.readyNotified) {
251 m = { ...m, readyNotified: true }
252 ready.push(m)
253 }
254 if (n.state === 'blocked' && !n.blockedNotified) {
255 m = { ...m, blockedNotified: true }
256 blocked.push(m)
257 }
258 marked[n.id] = m
259 }
260 next[owner] = marked
261 if (!ready.length && !blocked.length) continue
262 const doneIds = new Set(ready.flatMap(n => n.after))
263 const done = Object.values(marked).filter(n => doneIds.has(n.id) && n.state === 'done')
264 notices.push({ owner, ready, blocked, done, ...(owner === 'main' && opts.slots !== undefined ? { slots: opts.slots } : {}) })
265 }
266 return { plan: next, notices }
267}
268
269const label = (n: DagNode) => `${n.id} (${n.title})`
270
271export function noticeText(notice: Notice, slots: number | undefined = notice.slots): string {
272 const parts: string[] = []
273 if (notice.done.length) parts.push(notice.done.map(n => `${n.id} is done${n.info ? ` (${n.info})` : ''}`).join('; ') + '.')
274 if (notice.ready.length) {
275 parts.push(`Ready to start: ${notice.ready.map(label).join(', ')}.`)
276 const n = notice.ready.length
277 if (notice.owner === 'main' && slots !== undefined) {
278 const now = Math.max(0, Math.min(n, slots))
279 parts.push(
280 now === 0
281 ? `No manager slot is free; start them as slots free up.`
282 : now === n
283 ? `Start ${n === 1 ? 'it' : 'them'} now; build each brief on the merged code (git fetch origin first).`
284 : `Start ${now} now and let the other ${n - now} wait for a free manager slot; build each brief on the merged code (git fetch origin first).`,
285 )
286 } else {
287 parts.push(`Start ${n === 1 ? 'it' : 'them'} now; build each brief on the merged code (git fetch origin first).`)
288 }
289 }
290 if (notice.blocked.length) {
291 parts.push(`Blocked: ${notice.blocked.map(n => `${n.id}${n.info ? ` (${n.info})` : ''}`).join(', ')}. Their dependents stay waiting until you fix or lift them.`)
292 }
293 return `flow plan: ${parts.join(' ')}`
294}
295
296// Depth per node: the longest path from a root (a node without dependencies is 0).
297export function layers(graph: Graph): Record<string, number> {
298 const depth: Record<string, number> = {}
299 const visit = (id: string): number => {
300 const known = depth[id]
301 if (known !== undefined) return known
302 depth[id] = 0
303 const d = Math.max(-1, ...(graph[id]?.after ?? []).filter(d => graph[d]).map(visit)) + 1
304 depth[id] = d
305 return d
306 }
307 for (const id of Object.keys(graph)) visit(id)
308 return depth
309}
310
311// One line per node, shallowest first: "- login-redirect: waiting | after: csv-export (running)".
312export function describe(graph: Graph): string[] {
313 const depth = layers(graph)
314 return Object.values(graph)
315 .sort((a, b) => (depth[a.id] ?? 0) - (depth[b.id] ?? 0))
316 .map(n => {
317 const after = n.after.length ? ` | after: ${n.after.map(d => `${d} (${graph[d]?.state ?? '?'})`).join(', ')}` : ''
318 return `- ${n.id}: ${n.state}${n.info ? ` (${n.info})` : ''}${after}`
319 })
320}
321