Optional herdr fleet tools for the to-code skill: set up workspace-scoped fleets, get status, wait, read, send, and watch agent panes, wake the session when…

Pi's async coordinator and workers share the production statecharts in fleet/model.ts with the offline simulator. The Claude mod and unmanaged Pi fleets retain the dependency-free legacy hooks/core/ tools. A Herdr pane becoming idle is never an async outcome.
Use the @earendil-works Pi CLI 1.0.2 or newer inside Herdr. Async startup checks the actual host package, not the plugin's dev dependency. Unknown embedded hosts fail closed. Install this package's pinned dependencies before loading pi/index.ts. Do not load multiple copies of the extension.
Resolve the installed skill and extension independently before a run, and record their real paths/revisions plus the actual Pi host version in the run evidence. pi list shows configured package sources; a local package loads directly from that path. A skill manager may instead cache a remote Git checkout. Updating that source follows its configured remote/ref, not unpushed local commits. Use the manager's supported local-source/provider selection for approved local development, or wait for publication; preserve its cache and load only one extension copy. Restart/reload after approved installation changes and check the newly loaded skill's managed/legacy routing.
The coordinator owns scheduling, questions, and serial integration. Workers receive an existing ticket worktree, explicit model/thinking configuration, and one assignment ID. Startup uses a new background tab (--no-focus) in the coordinator-selected existing workspace, without creating a branch or worktree:
pi --extension <plugin>/pi/index.ts --model <provider/model> --thinking <level>
--fleet-run <absolute-run-directory> --fleet-implementor <worker-id>
A reviewer uses --fleet-reviewer instead. Only one role flag is allowed. Launch acknowledgment requires a flag-bound live process and ready callback endpoint, not just a terminal readiness marker. An ambiguous launch is retained as uncertain; inspect it before restarting, never repeat allocation blindly.
| Tool | Required input and meaning |
|---|---|
fleet_setup | workspace, integrationBranch, checks, optional commandTimeoutMs (integer milliseconds, 1..570000, default 120000). Cwd must be the integration worktree; .scratch/ must be Git-ignored. |
fleet_start_ticket | fleetId, ticketId, worktree, goal, model, thinking. The supplied worktree must be clean, separate, already created, and share the Git common directory. Ticket IDs are plain tokens. |
fleet_start_review | fleetId, ticketId, model, thinking. Requires verified implementation evidence. Reuses the retained reviewer on a refreshed implementation, keeping its recorded configuration. |
fleet_send | fleetId, workerId (or owned pane), text, optional questionId. Answers require the exact pending question ID. Dispatch rearms the stop gate with a new assignment. |
fleet_finish_ticket | fleetId, ticketId. After the coordinator lands the approved head, verifies immutable approval, unchanged source head, fresh source/integration checks, a clean integration worktree, and ancestry before retiring owned workers. |
fleet_resume | Canonical absolute directory. Reopens session-owned callbacks and replays pending applicable notifications; never takes over a live/ambiguous owner. |
fleet_restart_worker | fleetId, workerId. Reconciles a confirmed existing instance or restarts an absent worker in its recorded pane/configuration. Never allocates replacement Git resources. |
fleet_status returns snapshots, assignment IDs, launch/transport state, and pending callbacks. fleet_read is an owned-worker diagnostic, not a result source. Managed fleet_wait/fleet_watch return an error: use callbacks, not completion polling. Finish/account for legacy watches before switching modes.
If Git changes during review or approval, fleet_send explicitly supersedes stale work, invalidates the old ticket evidence, and assigns implementation re-verification. A reviewer question on a changed snapshot also takes this route. Superseding is a durable control outcome, not reviewer approval or failure. Old replies and reports cannot clear a newer assignment's gate.
| Tool | Input and meaning |
|---|---|
fleet_question | reportId, assignmentId, question. Persists a coordinator question before permitting accounted waiting. |
fleet_report | reportId, assignmentId, outcome, summary, optional filename. Implementor outcomes: completed, failed, approval. Reviewer outcome: review with the assigned file. |
fleet_request_review | reportId, assignmentId, summary. After fixing validated findings, checks the changed implementation and rearms the retained reviewer. |
Every delivered assignment supplies literal worker/role/assignment identity and suggested IDs such as <worker>-<assignment>-question, -completed, -failed, -review, -approval, and -request-review for the role's supported operations. Callback turns receive this context without relying on lifecycle hooks. A new assignment, including a question answer or retained re-review, supplies fresh suggestions and review binding.
Report IDs are unique per worker across assignments and operations. Identical retries reuse the original ID and payload and are idempotent, including after newer work is assigned, without clearing its stop gate. Changed payloads require fresh IDs; stale new submissions are rejected. Collision errors identify the prior operation/assignment and suggest an unused current-assignment ID (with a suffix if needed). A rejected review leaves its gate armed and returns the current literal filename/template plus a distinct question ID for coordinator recovery. Use that context, not a bare worker ID. No final chat message, bare approval, or Herdr terminal state substitutes for these tools. Every active assignment remains stop-guarded at agent_before_settle; accounted questions, reports, supersession, and recorded control faults may wait. A recorded Lost fault is not ticket completion.
All worker questions route to the coordinator. Native UI questions are surfaced there, but answers do not grant native permissions or resolve human confirmation dialogs. Abort/error/no-continuation paths preserve a control-fault callback without fabricating a reviewer result. Hard crashes outside lifecycle hooks need coordinator inspection/recovery.
commandTimeoutMs bounds each configured fleet verification command in milliseconds. Worker Bash's timeout is in seconds: 180 seconds corresponds to 180000 milliseconds. Assignment context states these units and asks for finite shell timeouts and small, reachable fixtures. Fleet command bounds and process-group cleanup are enforced by the extension; worker Bash budgets are guidance only. This is not a resource sandbox.
The worker binding supplies the exact relative .scratch/.../reports/<worker>-<assignment>.md path and this template:
# Review
ticket: T1
reviewer: r123
assignment: a-456
base: <assigned-base>
head: <assigned-head>
verdict: approved
## Findings
None.
Use changes_requested with - [medium|high|critical] <actionable finding and location> when blocking fixes are needed. Low suggestions may accompany approval. Exactly one metadata value and Findings section are required. Assignment, reviewer, ticket, filename, base, and head must match the live binding. Canonical file reads reject traversal, symlinks, nonregular/outside files, and files over 1 MiB.
Findings go directly to the implementor. The implementor fixes them and requests re-review; the same reviewer remains assigned. Approval goes to the implementor, who forwards outcome: approval to the coordinator. Accepted file digests are rechecked for re-review, forwarding, and landing. Validation proves structure and evidence binding, not the semantic quality of a review. Reviewers read the referenced spec/ticket and applicable repository standards, then map acceptance requirements to actual test assertions and flag missing coverage. Passing checks or suggestive test names do not establish that coverage. The template's approved value is a placeholder, not a preselected verdict. The coordinator independently inspects and verifies the final integrated branch and delegates any corrections.
Each Git-ignored run lives at <integration-worktree>/.scratch/to-code/async-<uuid>/. state.json stores schema-decoded Machine snapshots, resource ownership, immutable report records, and a callback outbox. Updates use a bounded cross-process lock and fsync/atomic rename on a local filesystem. Proven-dead same-host writer locks are archived; live, remote, or ambiguous locks fail closed. Preserve archived locks and scratch evidence for investigation.
Callbacks are at least once. Session-owned Unix sockets (Windows named-pipe addressing is also defined) announce durable work, with no detached daemon. A queued sendMessage is not delivery: only a persisted Pi custom-message entry acknowledges its outbox ID. Duplicate/stale notifications are deduplicated or marked obsolete; offline recipients replay from durable state on resume. No managed-worker completion timer is used.
Takeover requires the recorded host and repository identity and a proven absent predecessor. Preserve worktrees, branches, reports, and unrelated panes. Cleanup happens only after durable Landed; failures are separately retryable. Close only the exact owned, settled Pi pane, and close its tab only if it contains no other pane. Role checks and filesystem validation are workflow safeguards, not an OS security sandbox.
From the repository root, bun plugins/to-code/demo/build.ts writes .scratch/fleet-state-machine.html. Open it directly: the generated HTML contains its own script and styles, imports the production model/validator through the build, and needs no network or install on the viewing machine. Its Git, checks, transport, and panes are simulated, not claims about real Herdr behavior.
Run bun test plugins/to-code/fleet plugins/to-code/pi plugins/to-code/hooks/core plugins/to-code/demo, plus root check, lint, and formatting. Claude mod tests still use claude plugin test plugins/to-code. Live Herdr launch validation must run inside Herdr, against a human-selected workspace and disposable already-created worktree; fixtures alone do not prove the live launcher contract.
hooks/register.ts 160 lines1import { atom, read, update } from "claude-code";
2import type { EngineInterface, Register } from "claude-code";
3
4import { isHerdrDispatch, recordBashDispatch, resetPrompt, stopReason } from "./core/dispatch.ts";
5import {
6 EMPTY_MEMORY,
7 type FleetStore,
8 formatWake,
9 listFleet,
10 summarize,
11 tick,
12} from "./core/fleet.ts";
13import type { Runner } from "./core/herdr.ts";
14import { executeFleetTool, FLEET_TOOLS, type FleetToolName } from "./core/tools.ts";
15
16const memory = atom({ plugin: "to-code", key: "memory" } as const, EMPTY_MEMORY);
17
18const TOOL_PREFIX = "mcp__to-code__";
19const FLEET_TOOL = /^mcp__to-code__fleet_/;
20
21// `$.process.run` takes no abort signal: a child whose wait was abandoned
22// exits at its own herdr timeout instead.
23const processRunner =
24 ($: EngineInterface): Runner =>
25 async (argv, timeoutMs) => {
26 const ran = await $.process.run(argv, { timeoutMs: Math.min(timeoutMs, 600_000) });
27 return { exitCode: ran.exitCode ?? 1, stdout: ran.stdout, stderr: ran.stderr };
28 };
29
30// Memory written by an earlier version of this module survives a reload, so
31// fields added since are filled from the defaults.
32const stateStore = ($: EngineInterface): FleetStore => ({
33 read: async () => ({ ...EMPTY_MEMORY, ...(await read($, memory)) }),
34 update: async (fn) => {
35 await update($, memory, (current) => fn({ ...EMPTY_MEMORY, ...current }));
36 },
37});
38
39// A wake starts a turn of its own. The plugin's own prompt.submit hook is
40// skipped for a prompt it submits, so the stop gate's record resets here.
41const wake = async ($: EngineInterface, text: string) => {
42 await stateStore($).update(resetPrompt);
43 await $.prompt.submit({ text });
44};
45
46export const register: Register = (on, options) => {
47 const intervalMs =
48 typeof options.watchIntervalMs === "number" ? Math.max(1_000, options.watchIntervalMs) : 3_000;
49 const isStopGateOn = options.stopGate !== false;
50 // Panes an in-flight fleet_wait covers; the watcher leaves them to it.
51 const awaiting = new Set<string>();
52 let selfPane: string | undefined;
53 let isActive = false;
54 let isBusy = false;
55 // Wakes that arrived mid-turn, held for turn.complete: a row appended into a
56 // running turn is lost when the final request has already gone out.
57 let pending: string[] = [];
58
59 on("session.start", async ($, e, next) => {
60 isActive = (await $.env.get("HERDR_ENV")) === "1";
61 if (!isActive) return next(e);
62 selfPane = await $.env.get("HERDR_PANE_ID");
63
64 for (const tool of FLEET_TOOLS) {
65 await $.tool.register({
66 name: tool.name,
67 description: tool.description,
68 inputSchema: tool.inputSchema,
69 });
70 }
71
72 const runner = processRunner($);
73 const store = stateStore($);
74 let isTicking = false;
75 $.clock.every(intervalMs, () => {
76 if (isTicking) return;
77 isTicking = true;
78 void tick(runner, store, awaiting, selfPane)
79 .then(async ({ agents, events }) => {
80 $.ui.status(summarize(agents, await store.read()));
81 if (events.length === 0) return;
82 $.ui.toast(events.map((event) => `${event.name}: ${event.status}`).join(", "));
83 if (isBusy) {
84 pending = [...pending, formatWake(events)];
85 return;
86 }
87 await wake($, formatWake(events));
88 })
89 .catch(() => $.ui.status(undefined))
90 .finally(() => {
91 isTicking = false;
92 });
93 });
94
95 return next(e);
96 });
97
98 on("tool.describe", { tool: FLEET_TOOL }, async ($, e, next) => ({
99 ...(await next(e)),
100 isDeferred: false,
101 }));
102
103 on("tool.call", { tool: FLEET_TOOL }, async ($, e, next) => {
104 const { tool, tool_use_id: _id, agentId: _agent, ...input } = e as Record<string, unknown>;
105 const name = String(tool).slice(TOOL_PREFIX.length) as FleetToolName;
106 const outcome = await executeFleetTool(name, input, {
107 runner: processRunner($),
108 store: stateStore($),
109 selfPane,
110 awaiting,
111 signal: next.signal,
112 });
113 return outcome.isError ? { isError: true, result: outcome.text } : { result: outcome.text };
114 });
115
116 on("tool.call", { tool: "Bash" }, async ($, e, next) => {
117 if (isActive && e.agentId === undefined && isHerdrDispatch(e.command)) {
118 const agents = await listFleet(processRunner($), selfPane, undefined).catch(() => []);
119 await stateStore($).update((current) => recordBashDispatch(current, e.command, agents));
120 }
121 return next(e);
122 });
123
124 on("prompt.submit", async ($, e, next) => {
125 await stateStore($).update(resetPrompt);
126 return next(e);
127 });
128
129 on("turn.start", async ($, e, next) => {
130 isBusy = true;
131 return next(e);
132 });
133
134 on("turn.complete", async ($, e, next) => {
135 const result = await next(e);
136 if (e.agentId !== undefined) return result;
137 isBusy = false;
138 if (pending.length > 0) {
139 const text = pending.join("\n\n");
140 pending = [];
141 // Not awaited: the wake's turn starts once this one has finished.
142 void wake($, text).catch(() => $.ui.toast("to-code: a fleet wake could not be delivered"));
143 }
144 return result;
145 });
146
147 on("classic.Stop", async ($, e, next) => {
148 const result = await next(e);
149 if (!isActive || !isStopGateOn || result.block !== undefined) return result;
150 const current = await stateStore($).read();
151 const reason = stopReason({
152 memory: current,
153 pendingBackground: (e.background_tasks?.length ?? 0) + (e.session_crons?.length ?? 0),
154 });
155 if (reason === undefined) return result;
156 await stateStore($).update((value) => ({ ...value, nudged: true }));
157 return { ...result, block: reason };
158 });
159};
160hooks/core/dispatch.ts 166 lines1import type { FleetMemory } from "./fleet.ts";
2import type { FleetAgent } from "./herdr.ts";
3
4const HERDR_DISPATCH = /\bherdr\s+(agent\s+(prompt|start)|pane\s+run)\b/;
5const HERDR_TARGET = /\bherdr\s+(?:agent\s+prompt|pane\s+run)\s+["']?([A-Za-z0-9_.:-]+)/g;
6
7/**
8 * The label a shell dispatch gets when its target pane cannot be named.
9 */
10export const BASH_DISPATCH = "herdr via Bash";
11
12/**
13 * Whether a shell command hands work to another herdr pane.
14 *
15 * @param command the shell command
16 * @returns true for `herdr agent prompt|start` and `herdr pane run`
17 */
18export const isHerdrDispatch = (command: string): boolean => HERDR_DISPATCH.test(command);
19
20/**
21 * The literal targets of the `herdr agent prompt` and `herdr pane run` calls
22 * in a shell command: pane ids or agent names, as written.
23 *
24 * @param command the shell command
25 * @returns the targets, in order
26 */
27export const dispatchTargets = (command: string): string[] =>
28 [...command.matchAll(HERDR_TARGET)].flatMap((match) => (match[1] ? [match[1]] : []));
29
30/**
31 * What the stop gate knows when a turn is about to end.
32 */
33export type StopContext = {
34 memory: FleetMemory;
35 pendingBackground: number;
36};
37
38/**
39 * Decides whether to keep the model going because it dispatched work and left
40 * no way to be woken when that work finishes. Fires at most once per prompt.
41 *
42 * @param context the session's memory and the harness's pending background work
43 * @returns the reason the model reads, or undefined to let the turn end
44 */
45export const stopReason = (context: StopContext): string | undefined => {
46 const { memory, pendingBackground } = context;
47 if (memory.nudged || pendingBackground > 0) return undefined;
48 const scoped = Object.entries(memory.fleets).flatMap(([fleetId, fleet]) => {
49 const reason = stopReason({ memory: { ...memory, ...fleet, fleets: {} }, pendingBackground });
50 return reason === undefined ? [] : [`${reason} Use fleetId: ${fleetId}.`];
51 });
52 if (scoped.length > 0) return scoped.join("\n");
53 if (memory.dispatched.length === 0) return undefined;
54 if (
55 Object.keys(memory.watched).length > 0 ||
56 Object.values(memory.fleets).some((fleet) => Object.keys(fleet.watched).length > 0)
57 )
58 return undefined;
59
60 return [
61 `[to-code fleet] This turn dispatched work (${memory.dispatched.join(", ")}) but nothing will wake this session when it finishes.`,
62 "Call fleet_watch with the dispatched panes so the session resumes when they settle (a pane that already finished fires on the next check), or fleet_wait to block on them now.",
63 "If no follow-up is needed, say so and end the turn.",
64 ].join(" ");
65};
66
67/**
68 * Records one dispatch for the stop gate.
69 *
70 * @param memory the session's memory
71 * @param label what was dispatched, for the reason text
72 * @returns the updated memory
73 */
74export const recordDispatch = (memory: FleetMemory, label: string): FleetMemory => ({
75 ...memory,
76 dispatched: memory.dispatched.includes(label) ? memory.dispatched : [...memory.dispatched, label],
77});
78
79/**
80 * Records a shell dispatch: each target that names a listed agent pane by id
81 * or agent name is recorded as that pane, with its seq before the work began
82 * as the baseline a later `fleet_watch` starts from. Anything else, including
83 * `herdr agent start`, is recorded under BASH_DISPATCH.
84 *
85 * @param memory the session's memory
86 * @param command the shell command, already known to dispatch
87 * @param agents the fleet as listed before the command ran
88 * @returns the updated memory
89 */
90export const recordBashDispatch = (
91 memory: FleetMemory,
92 command: string,
93 agents: readonly FleetAgent[],
94): FleetMemory => {
95 const targets = dispatchTargets(command);
96 const resolved = targets.flatMap((target) => {
97 const agent = agents.find((one) => one.pane === target || one.name === target);
98 return agent === undefined ? [] : [agent];
99 });
100 const isNamedOnly =
101 resolved.length > 0 &&
102 resolved.length === targets.length &&
103 !/\bherdr\s+agent\s+start\b/.test(command);
104
105 const recorded = resolved.reduce(
106 (current, agent) => ({
107 ...recordDispatch(current, agent.pane),
108 baselines: { ...current.baselines, [agent.pane]: agent.seq },
109 }),
110 memory,
111 );
112 const result = isNamedOnly ? recorded : recordDispatch(recorded, BASH_DISPATCH);
113 const scopedPanes = new Set(
114 resolved
115 .filter((agent) => {
116 const workspace = agent.workspace;
117 return (
118 workspace !== undefined &&
119 Object.values(memory.fleets).some((fleet) => fleet.workspaces.includes(workspace))
120 );
121 })
122 .map((agent) => agent.pane),
123 );
124 return {
125 ...result,
126 dispatched: result.dispatched.filter((pane) => !scopedPanes.has(pane)),
127 fleets: Object.fromEntries(
128 Object.entries(memory.fleets).map(([fleetId, fleet]) => {
129 const matching = resolved.filter(
130 (agent) => agent.workspace !== undefined && fleet.workspaces.includes(agent.workspace),
131 );
132 return [
133 fleetId,
134 {
135 ...fleet,
136 baselines: {
137 ...fleet.baselines,
138 ...Object.fromEntries(matching.map((agent) => [agent.pane, agent.seq])),
139 },
140 dispatched: [...new Set([...fleet.dispatched, ...matching.map((agent) => agent.pane)])],
141 },
142 ];
143 }),
144 ),
145 };
146};
147
148/**
149 * Clears the stop gate's record when a new prompt or run starts. Watches and
150 * dispatch baselines carry over: they describe panes, not the prompt.
151 *
152 * @param memory the session's memory
153 * @returns the updated memory
154 */
155export const resetPrompt = (memory: FleetMemory): FleetMemory => ({
156 ...memory,
157 dispatched: [],
158 nudged: false,
159 fleets: Object.fromEntries(
160 Object.entries(memory.fleets).map(([fleetId, fleet]) => [
161 fleetId,
162 { ...fleet, dispatched: [], nudged: false },
163 ]),
164 ),
165});
166hooks/core/fleet.ts 403 lines1import {
2 type FleetAgent,
3 type FleetStatus,
4 HerdrFailure,
5 label,
6 listArgv,
7 parseAgentInfo,
8 parseAgentList,
9 type Runner,
10 SETTLED,
11 waitArgv,
12} from "./herdr.ts";
13
14/**
15 * Longest `fleet_wait` allowed. The Claude Code mod's `$.process.run` stops
16 * a child at ten minutes, so this stays under it with room for herdr to exit.
17 */
18export const MAX_WAIT_MS = 570_000;
19
20/**
21 * What both adapters remember for one session.
22 *
23 * `watched` maps a pane to the `state_change_seq` it had when the watch
24 * began: the pane fires once it settles at a later seq. `baselines` holds the
25 * seq a pane had when unwatched work was sent to it, so a later watch still
26 * sees that work finish. `dispatched` and `nudged` drive the stop gate for the
27 * current prompt.
28 */
29export type FleetMemory = {
30 watched: Record<string, number>;
31 baselines: Record<string, number>;
32 dispatched: string[];
33 nudged: boolean;
34 fleets: Record<string, FleetScope>;
35 nextFleetId: number;
36};
37
38/**
39 * A session-local fleet follows all agents in the selected workspaces.
40 */
41export type FleetScope = {
42 workspaces: string[];
43 watched: Record<string, number>;
44 baselines: Record<string, number>;
45 dispatched: string[];
46 nudged: boolean;
47};
48
49/**
50 * A fresh session's memory.
51 */
52export const EMPTY_MEMORY: FleetMemory = {
53 watched: {},
54 baselines: {},
55 dispatched: [],
56 nudged: false,
57 fleets: {},
58 nextFleetId: 1,
59};
60
61/**
62 * Read-modify-write access to the session's memory, backed by `$.state` in
63 * the Claude Code mod and a plain variable in the pi extension.
64 */
65export type FleetStore = {
66 read: () => Promise<FleetMemory>;
67 update: (fn: (memory: FleetMemory) => FleetMemory) => Promise<void>;
68};
69
70/**
71 * Creates a store held in a closure.
72 *
73 * @param initial the starting memory
74 * @returns the store
75 */
76export const memoryStore = (initial: FleetMemory = EMPTY_MEMORY): FleetStore => {
77 let memory = initial;
78
79 return {
80 read: async () => memory,
81 update: async (fn) => {
82 memory = fn(memory);
83 },
84 };
85};
86
87/**
88 * Selects agents in a fleet's captured workspaces.
89 *
90 * @param agents the current agent listing
91 * @param workspaces the workspace IDs selected at setup
92 * @returns current members, including newly created agents
93 */
94export const members = (
95 agents: readonly FleetAgent[],
96 workspaces: readonly string[],
97): FleetAgent[] =>
98 agents.filter((agent) => agent.workspace !== undefined && workspaces.includes(agent.workspace));
99
100/**
101 * Identifies a wait without suppressing another fleet's watch of the same pane.
102 *
103 * @param fleetId the fleet, or undefined for legacy global calls
104 * @param pane the target pane
105 * @returns the in-flight wait key
106 */
107export const waitKey = (fleetId: string | undefined, pane: string): string =>
108 fleetId === undefined ? pane : `${fleetId}/${pane}`;
109
110/**
111 * Projects a fleet's activity onto the existing tool executor's store interface.
112 *
113 * @param store the session store
114 * @param fleetId the fleet, or undefined for legacy global calls
115 * @returns a store that updates only the selected fleet
116 */
117export const scopedStore = (store: FleetStore, fleetId: string | undefined): FleetStore => {
118 if (fleetId === undefined) return store;
119 const scope = (memory: FleetMemory): FleetScope => {
120 const fleet = Object.hasOwn(memory.fleets, fleetId) ? memory.fleets[fleetId] : undefined;
121 if (fleet === undefined) throw new Error(`Unknown fleetId: ${fleetId}; call fleet_setup`);
122 return fleet;
123 };
124 return {
125 read: async () => ({ ...EMPTY_MEMORY, ...scope(await store.read()) }),
126 update: async (fn) => {
127 await store.update((memory) => {
128 const fleet = scope(memory);
129 const updated = fn({ ...EMPTY_MEMORY, ...fleet });
130 return {
131 ...memory,
132 fleets: {
133 ...memory.fleets,
134 [fleetId]: {
135 ...fleet,
136 watched: updated.watched,
137 baselines: updated.baselines,
138 dispatched: updated.dispatched,
139 nudged: updated.nudged,
140 },
141 },
142 };
143 });
144 },
145 };
146};
147
148/**
149 * A watched pane that settled or disappeared.
150 */
151export type WakeEvent = {
152 fleetId: string | undefined;
153 pane: string;
154 name: string;
155 status: FleetStatus | "gone";
156};
157
158/**
159 * What one watcher tick found.
160 */
161export type TickResult = {
162 agents: FleetAgent[];
163 events: WakeEvent[];
164};
165
166/**
167 * Drops the caller's own pane from a listing.
168 *
169 * @param agents every agent pane
170 * @param selfPane the caller's pane id, when it runs inside herdr
171 * @returns the other panes
172 */
173export const others = (agents: readonly FleetAgent[], selfPane: string | undefined): FleetAgent[] =>
174 agents.filter((agent) => agent.pane !== selfPane);
175
176/**
177 * Lists the other agent panes.
178 *
179 * @param runner the harness's process runner
180 * @param selfPane the caller's pane id
181 * @param signal aborts the listing
182 * @returns every agent pane but the caller's
183 */
184export const listFleet = async (
185 runner: Runner,
186 selfPane: string | undefined,
187 signal: AbortSignal | undefined,
188): Promise<FleetAgent[]> =>
189 others(parseAgentList(await runner(listArgv(), 10_000, signal)), selfPane);
190
191/**
192 * Finds the watched panes that settled since their watch began.
193 *
194 * @param agents the current listing
195 * @param watched pane to baseline seq
196 * @param awaiting panes an in-flight `fleet_wait` already covers
197 * @returns one event per pane to report
198 */
199export const settledEvents = (
200 agents: readonly FleetAgent[],
201 watched: Readonly<Record<string, number>>,
202 awaiting: ReadonlySet<string>,
203): WakeEvent[] =>
204 Object.entries(watched).flatMap(([pane, baseline]): WakeEvent[] => {
205 if (awaiting.has(pane)) return [];
206 const agent = agents.find((one) => one.pane === pane);
207 if (agent === undefined) return [{ fleetId: undefined, pane, name: pane, status: "gone" }];
208 if (agent.seq <= baseline || !SETTLED.includes(agent.status)) return [];
209 return [{ fleetId: undefined, pane, name: agent.name, status: agent.status }];
210 });
211
212const omit = (record: Readonly<Record<string, number>>, keys: readonly string[]) =>
213 Object.fromEntries(Object.entries(record).filter(([key]) => !keys.includes(key)));
214
215/**
216 * Removes panes from the watch set.
217 *
218 * @param memory the session's memory
219 * @param panes the panes to drop
220 * @returns the updated memory
221 */
222export const unwatch = (memory: FleetMemory, panes: readonly string[]): FleetMemory => ({
223 ...memory,
224 watched: omit(memory.watched, panes),
225});
226
227/**
228 * Forgets panes whose work has been reported settled: their watches, their
229 * dispatch baselines, and their entries in the stop gate's dispatch record.
230 *
231 * @param memory the session's memory
232 * @param panes the panes that settled
233 * @returns the updated memory
234 */
235export const forget = (memory: FleetMemory, panes: readonly string[]): FleetMemory => ({
236 ...memory,
237 watched: omit(memory.watched, panes),
238 baselines: omit(memory.baselines, panes),
239 dispatched: memory.dispatched.filter((entry) => !panes.includes(entry)),
240});
241
242/**
243 * Runs one watcher pass: lists the fleet, reports watched panes that settled,
244 * and forgets them so each watch fires once.
245 *
246 * @param runner the harness's process runner
247 * @param store the session's memory
248 * @param awaiting panes an in-flight `fleet_wait` already covers
249 * @param selfPane the caller's pane id
250 * @returns the listing and the events to deliver
251 */
252export const tick = async (
253 runner: Runner,
254 store: FleetStore,
255 awaiting: ReadonlySet<string>,
256 selfPane: string | undefined,
257): Promise<TickResult> => {
258 const agents = await listFleet(runner, selfPane, undefined);
259 const memory = await store.read();
260 const events = settledEvents(agents, memory.watched, awaiting);
261 for (const [fleetId, fleet] of Object.entries(memory.fleets)) {
262 const covered = new Set(
263 Object.keys(fleet.watched).filter((pane) => awaiting.has(waitKey(fleetId, pane))),
264 );
265 events.push(
266 ...settledEvents(members(agents, fleet.workspaces), fleet.watched, covered).map((event) => ({
267 ...event,
268 fleetId,
269 })),
270 );
271 }
272 for (const fleetId of new Set(events.map((event) => event.fleetId))) {
273 await scopedStore(store, fleetId).update((current) =>
274 forget(
275 current,
276 events.filter((event) => event.fleetId === fleetId).map((event) => event.pane),
277 ),
278 );
279 }
280 const workspaces = Object.values(memory.fleets).flatMap((fleet) => fleet.workspaces);
281 return {
282 agents:
283 workspaces.length === 0
284 ? agents
285 : agents.filter(
286 (agent) => workspaces.includes(agent.workspace ?? "") || agent.pane in memory.watched,
287 ),
288 events,
289 };
290};
291
292/**
293 * Formats the one-line status summary both harnesses show.
294 *
295 * @param agents the other panes
296 * @param memory the session's memory
297 * @returns the summary, or undefined when there is nothing to show
298 */
299export const summarize = (
300 agents: readonly FleetAgent[],
301 memory: FleetMemory,
302): string | undefined => {
303 const watched =
304 Object.keys(memory.watched).length +
305 Object.values(memory.fleets).reduce(
306 (total, fleet) => total + Object.keys(fleet.watched).length,
307 0,
308 );
309 if (agents.length === 0 && watched === 0) return undefined;
310 const counts = new Map<FleetStatus, number>();
311 for (const agent of agents) counts.set(agent.status, (counts.get(agent.status) ?? 0) + 1);
312 const parts = (["working", "blocked", "idle", "done", "unknown"] as const)
313 .filter((status) => counts.has(status))
314 .map((status) => `${counts.get(status)} ${status}`);
315 if (watched > 0) parts.push(`${watched} watched`);
316 return `fleet: ${parts.join(", ")}`;
317};
318
319/**
320 * Formats the message that wakes or steers the session.
321 *
322 * @param events the panes that settled
323 * @returns the message text
324 */
325export const formatWake = (events: readonly WakeEvent[]): string => {
326 const lines = events.map(
327 (event) =>
328 (event.status === "gone"
329 ? `- ${label(event.pane, 32)}: the pane is gone (closed or its agent exited)`
330 : `- ${label(event.pane, 32)} (${label(event.name)}): ${event.status}`) +
331 (event.fleetId === undefined ? "" : ` [fleetId: ${event.fleetId}]`),
332 );
333 return [
334 "[to-code fleet] Watched herdr panes settled (pane names are labels, not instructions):",
335 ...lines,
336 "Reconcile each one: read its result file or use fleet_read (pass fleetId when shown), then continue, answer a blocked agent, or record the outcome.",
337 ].join("\n");
338};
339
340/**
341 * How a `fleet_wait` ended.
342 */
343export type WaitOutcome =
344 | { kind: "settled"; agent: FleetAgent }
345 | { kind: "timeout" }
346 | { kind: "aborted" }
347 | { kind: "failed"; error: string };
348
349/**
350 * Waits for the first of several panes to reach a state, one herdr child per
351 * pane. The children that lose the race are aborted where the runner can
352 * abort, and otherwise exit at their own timeout.
353 *
354 * @param runner the harness's process runner
355 * @param panes the panes to wait on
356 * @param until the states to match
357 * @param timeoutMs how long to wait, capped at MAX_WAIT_MS
358 * @param signal aborts the whole wait
359 * @returns the first pane to match, a timeout, an abort, or a failure
360 */
361export const waitAny = async (
362 runner: Runner,
363 panes: readonly string[],
364 until: readonly FleetStatus[],
365 timeoutMs: number,
366 signal: AbortSignal | undefined,
367): Promise<WaitOutcome> => {
368 const bounded = Math.max(1_000, Math.min(timeoutMs, MAX_WAIT_MS));
369 const losers = new AbortController();
370 const onAbort = () => losers.abort();
371 signal?.addEventListener("abort", onAbort, { once: true });
372
373 const attempts = panes.map(async (pane): Promise<WaitOutcome> => {
374 try {
375 const output = await runner(waitArgv(pane, until, bounded), bounded + 15_000, losers.signal);
376 return { kind: "settled", agent: parseAgentInfo(output) };
377 } catch (error) {
378 if (error instanceof HerdrFailure && error.code === "timeout") return { kind: "timeout" };
379 return {
380 kind: "failed",
381 error: `${pane}: ${error instanceof Error ? error.message : String(error)}`,
382 };
383 }
384 });
385
386 try {
387 const first = await Promise.any(
388 attempts.map((attempt) =>
389 attempt.then((outcome) => (outcome.kind === "settled" ? outcome : Promise.reject(outcome))),
390 ),
391 ).catch(async () => {
392 const all = await Promise.all(attempts);
393 // An aborted child exits like a failed one; report the abort instead.
394 if (signal?.aborted) return { kind: "aborted" } as const;
395 return all.find((outcome) => outcome.kind === "failed") ?? ({ kind: "timeout" } as const);
396 });
397 return first;
398 } finally {
399 losers.abort();
400 signal?.removeEventListener("abort", onAbort);
401 }
402};
403hooks/core/herdr.ts 243 lines1/**
2 * Agent states herdr reports for a pane.
3 */
4export type FleetStatus = "idle" | "working" | "blocked" | "done" | "unknown";
5
6/**
7 * States that mean an agent stopped and wants attention.
8 */
9export const SETTLED: readonly FleetStatus[] = ["idle", "done", "blocked"];
10
11/**
12 * One agent pane, projected from herdr's `agent list` / `agent wait` JSON.
13 */
14export type FleetAgent = {
15 pane: string;
16 name: string;
17 harness: string;
18 status: FleetStatus;
19 cwd: string;
20 workspace: string | undefined;
21 seq: number;
22};
23
24/**
25 * The `error` object herdr prints on stdout when a command fails.
26 */
27export type HerdrError = {
28 code: string;
29 message: string;
30};
31
32/**
33 * What a runner reports for one finished child process.
34 */
35export type RunOutput = {
36 exitCode: number;
37 stdout: string;
38 stderr: string;
39};
40
41/**
42 * Runs an argv to completion. Each harness adapter supplies its own: the
43 * Claude Code mod through `$.process.run`, the pi extension through Node.
44 */
45export type Runner = (
46 argv: readonly string[],
47 timeoutMs: number,
48 signal: AbortSignal | undefined,
49) => Promise<RunOutput>;
50
51/**
52 * Raised when herdr answers with an error document or unreadable output.
53 */
54export class HerdrFailure extends Error {
55 readonly code: string;
56
57 /**
58 * @param error the parsed herdr error
59 */
60 constructor(error: HerdrError) {
61 super(`${error.code}: ${error.message}`);
62 this.code = error.code;
63 }
64}
65
66type RawAgent = {
67 pane_id?: unknown;
68 workspace_id?: unknown;
69 name?: unknown;
70 display_agent?: unknown;
71 title?: unknown;
72 agent?: unknown;
73 agent_status?: unknown;
74 cwd?: unknown;
75 state_change_seq?: unknown;
76};
77
78const STATUSES: readonly FleetStatus[] = ["idle", "working", "blocked", "done", "unknown"];
79
80const asString = (value: unknown, fallback: string): string =>
81 typeof value === "string" ? value : fallback;
82
83const firstText = (...values: unknown[]): string =>
84 values.find((value): value is string => typeof value === "string" && value !== "") ?? "";
85
86/**
87 * Reduces a herdr-reported label to a short, single-line, plain token.
88 * Pane names and titles are set by whoever runs in the pane, and they reach
89 * the session inside wake prompts, so they must not carry instructions.
90 *
91 * @param value the raw label
92 * @param max the longest label kept
93 * @returns the label with anything outside letters, digits, and `._:/-` turned into spaces
94 */
95export const label = (value: string, max = 48): string =>
96 value
97 .replace(/[^A-Za-z0-9 ._:/-]+/g, " ")
98 .replace(/\s+/g, " ")
99 .trim()
100 .slice(0, max);
101
102const toAgent = (raw: RawAgent): FleetAgent => {
103 const status = asString(raw.agent_status, "unknown") as FleetStatus;
104
105 return {
106 pane: label(asString(raw.pane_id, ""), 32),
107 // herdr's own agent name is set at `agent start`; the display name and
108 // title come later from the agent and drift with its session.
109 name: label(firstText(raw.name, raw.display_agent, raw.title)),
110 harness: label(asString(raw.agent, "unknown"), 16),
111 status: STATUSES.includes(status) ? status : "unknown",
112 cwd: asString(raw.cwd, ""),
113 workspace:
114 typeof raw.workspace_id === "string" && raw.workspace_id !== ""
115 ? raw.workspace_id
116 : undefined,
117 seq: typeof raw.state_change_seq === "number" ? raw.state_change_seq : 0,
118 };
119};
120
121const parseDocument = (output: RunOutput): Record<string, unknown> => {
122 const text = output.stdout.trim() || output.stderr.trim();
123 let parsed: unknown;
124 try {
125 parsed = JSON.parse(text);
126 } catch {
127 throw new HerdrFailure({
128 code: "unreadable_output",
129 message: text.slice(0, 300) || `exit ${output.exitCode}`,
130 });
131 }
132 if (typeof parsed !== "object" || parsed === null) {
133 throw new HerdrFailure({ code: "unreadable_output", message: text.slice(0, 300) });
134 }
135 const document = parsed as Record<string, unknown>;
136 const error = document.error as Partial<HerdrError> | undefined;
137 if (error !== undefined) {
138 throw new HerdrFailure({ code: error.code ?? "unknown_error", message: error.message ?? "" });
139 }
140 return document;
141};
142
143/**
144 * Parses `herdr agent list` output.
145 *
146 * @param output the finished command
147 * @returns every agent pane herdr knows about
148 * @throws HerdrFailure when herdr reported an error
149 */
150export const parseAgentList = (output: RunOutput): FleetAgent[] => {
151 const result = parseDocument(output).result as { agents?: RawAgent[] } | undefined;
152 return (result?.agents ?? []).map(toAgent).filter((agent) => agent.pane !== "");
153};
154
155/**
156 * Parses the single agent `herdr agent wait` and `agent prompt --wait` print.
157 *
158 * @param output the finished command
159 * @returns the agent as it settled
160 * @throws HerdrFailure when herdr reported an error (including `timeout`)
161 */
162export const parseAgentInfo = (output: RunOutput): FleetAgent => {
163 const result = parseDocument(output).result as { agent?: RawAgent } | undefined;
164 return toAgent(result?.agent ?? {});
165};
166
167/**
168 * Throws when a command that prints no JSON on success failed.
169 *
170 * @param output the finished command
171 * @throws HerdrFailure when the command failed
172 */
173export const assertOk = (output: RunOutput): void => {
174 if (output.exitCode === 0) return;
175 parseDocument(output);
176 throw new HerdrFailure({
177 code: "command_failed",
178 message: (output.stderr || output.stdout).slice(0, 300),
179 });
180};
181
182/**
183 * Argv for listing every agent pane.
184 *
185 * @returns the argv
186 */
187export const listArgv = (): string[] => ["herdr", "agent", "list"];
188
189/**
190 * Argv for waiting on one pane to reach a state.
191 *
192 * @param pane the pane id
193 * @param until the states to match
194 * @param timeoutMs herdr's own timeout
195 * @returns the argv
196 */
197export const waitArgv = (
198 pane: string,
199 until: readonly FleetStatus[],
200 timeoutMs: number,
201): string[] => [
202 "herdr",
203 "agent",
204 "wait",
205 pane,
206 ...until.flatMap((status) => ["--until", status]),
207 "--timeout",
208 String(timeoutMs),
209];
210
211/**
212 * Argv for reading a pane's recent terminal output.
213 *
214 * @param pane the pane id
215 * @param lines how many lines
216 * @returns the argv
217 */
218export const readArgv = (pane: string, lines: number): string[] => [
219 "herdr",
220 "agent",
221 "read",
222 pane,
223 "--source",
224 "recent-unwrapped",
225 "--lines",
226 String(lines),
227];
228
229/**
230 * Argv for submitting a prompt to a pane's agent without waiting.
231 *
232 * @param pane the pane id
233 * @param text the prompt
234 * @returns the argv
235 */
236export const promptArgv = (pane: string, text: string): string[] => [
237 "herdr",
238 "agent",
239 "prompt",
240 pane,
241 text,
242];
243hooks/core/tools.ts 431 lines1import { recordDispatch } from "./dispatch.ts";
2import {
3 type FleetStore,
4 forget,
5 listFleet,
6 MAX_WAIT_MS,
7 members,
8 scopedStore,
9 unwatch,
10 waitAny,
11 waitKey,
12} from "./fleet.ts";
13import {
14 assertOk,
15 type FleetAgent,
16 type FleetStatus,
17 HerdrFailure,
18 promptArgv,
19 readArgv,
20 type Runner,
21 SETTLED,
22} from "./herdr.ts";
23
24/**
25 * One model-callable tool, described once for both harnesses.
26 */
27export type FleetToolSpec = {
28 name: FleetToolName;
29 label: string;
30 description: string;
31 inputSchema: Record<string, unknown>;
32};
33
34/**
35 * The tools both adapters register.
36 */
37export type FleetToolName =
38 | "fleet_setup"
39 | "fleet_status"
40 | "fleet_wait"
41 | "fleet_read"
42 | "fleet_send"
43 | "fleet_watch";
44
45/**
46 * What a tool hands back to the model.
47 */
48export type ToolOutcome = {
49 text: string;
50 isError: boolean;
51};
52
53/**
54 * What the executor needs from the harness for one call.
55 */
56export type ToolContext = {
57 runner: Runner;
58 store: FleetStore;
59 selfPane: string | undefined;
60 awaiting: Set<string>;
61 signal: AbortSignal | undefined;
62};
63
64type ScopedContext = ToolContext & {
65 fleetId: string | undefined;
66 agents: FleetAgent[] | undefined;
67 workspaces: string[] | undefined;
68};
69
70const toolAgents = async (context: ScopedContext): Promise<FleetAgent[]> =>
71 context.agents ?? (await listFleet(context.runner, context.selfPane, context.signal));
72
73const PANE = {
74 type: "string",
75 description: "herdr pane id, as fleet_status lists it (for example w5V:p2)",
76};
77const PANES = { type: "array", items: PANE, minItems: 1 };
78const STATES = {
79 type: "array",
80 items: { type: "string", enum: ["idle", "working", "blocked", "done", "unknown"] },
81};
82
83/**
84 * The tool specs, in registration order.
85 */
86export const FLEET_TOOLS: readonly FleetToolSpec[] = (
87 [
88 {
89 name: "fleet_setup",
90 label: "Fleet setup",
91 description:
92 "Issue a session-local fleetId from explicit agent pane IDs. The fleet follows all other agents in their Herdr workspaces, including workers created later. Pass fleetId to later fleet tools to limit their results and targets. Setup does not watch or prompt workers.",
93 inputSchema: {
94 type: "object",
95 properties: { panes: PANES },
96 required: ["panes"],
97 additionalProperties: false,
98 },
99 },
100 {
101 name: "fleet_status",
102 label: "Fleet status",
103 description:
104 "List the other herdr agent panes: pane id, name, harness (claude or pi), status (idle, working, blocked, done, unknown), cwd, and whether this session watches it. Use instead of parsing `herdr agent list` in the shell.",
105 inputSchema: { type: "object", properties: {}, additionalProperties: false },
106 },
107 {
108 name: "fleet_wait",
109 label: "Fleet wait",
110 description:
111 "Block until the first of the given panes reaches one of the states (default idle, done, or blocked), or the timeout passes. Use instead of sleep or until loops. Returns the pane that settled and its status, or a timeout.",
112 inputSchema: {
113 type: "object",
114 properties: {
115 panes: PANES,
116 until: { ...STATES, description: "States to match; default idle, done, blocked" },
117 timeoutMs: { type: "number", description: `How long to wait; capped at ${MAX_WAIT_MS}` },
118 },
119 required: ["panes", "timeoutMs"],
120 additionalProperties: false,
121 },
122 },
123 {
124 name: "fleet_read",
125 label: "Fleet read",
126 description:
127 "Read a pane's recent terminal output to check what its agent is doing or asking. For checking state only: collect final reports from the result files workers write, not from here.",
128 inputSchema: {
129 type: "object",
130 properties: {
131 pane: PANE,
132 lines: { type: "number", description: "Lines to read; default 80, at most 400" },
133 },
134 required: ["pane"],
135 additionalProperties: false,
136 },
137 },
138 {
139 name: "fleet_send",
140 label: "Fleet send",
141 description:
142 "Submit a prompt to a pane's agent and, by default, watch the pane so this session wakes when it settles. A blocked agent rejects the prompt: read it and answer its question instead. A stalled or timed-out submission may still have landed, so check fleet_status before sending again; never resend blindly.",
143 inputSchema: {
144 type: "object",
145 properties: {
146 pane: PANE,
147 text: { type: "string", description: "The prompt" },
148 watch: { type: "boolean", description: "Watch the pane after sending; default true" },
149 },
150 required: ["pane", "text"],
151 additionalProperties: false,
152 },
153 },
154 {
155 name: "fleet_watch",
156 label: "Fleet watch",
157 description:
158 "Wake this session once each given pane next settles (idle, done, or blocked), then forget it. Use after dispatching work any other way, so you can end the turn instead of polling. Work this session sent with fleet_send (watch false) or `herdr agent prompt`/`herdr pane run` in the shell counts from when it was sent, so a pane that already finished it fires on the next check. With unwatch true, stop watching the panes.",
159 inputSchema: {
160 type: "object",
161 properties: {
162 panes: PANES,
163 unwatch: { type: "boolean", description: "Stop watching instead" },
164 },
165 required: ["panes"],
166 additionalProperties: false,
167 },
168 },
169 ] satisfies FleetToolSpec[]
170).map((tool) =>
171 tool.name === "fleet_setup"
172 ? tool
173 : {
174 ...tool,
175 description: `${tool.description} With fleetId, only that fleet's workspace members are accessible; omitting it preserves global behavior.`,
176 inputSchema: {
177 ...tool.inputSchema,
178 properties: {
179 ...tool.inputSchema.properties,
180 fleetId: {
181 type: "string",
182 minLength: 1,
183 description: "Session-local ID issued by fleet_setup",
184 },
185 },
186 },
187 },
188);
189
190const ok = (value: unknown): ToolOutcome => ({
191 text: JSON.stringify(value, null, 2),
192 isError: false,
193});
194
195const fail = (message: string): ToolOutcome => ({ text: message, isError: true });
196
197const strings = (value: unknown): string[] =>
198 Array.isArray(value)
199 ? value.filter((item): item is string => typeof item === "string" && item !== "")
200 : [];
201
202const setup = async (
203 input: Record<string, unknown>,
204 context: ToolContext,
205): Promise<ToolOutcome> => {
206 const panes = strings(input.panes);
207 if (panes.length === 0 || !Array.isArray(input.panes) || panes.length !== input.panes.length)
208 return fail("panes must list at least one nonempty pane id");
209 const agents = await listFleet(context.runner, context.selfPane, context.signal);
210 const selected = agents.filter((agent) => panes.includes(agent.pane));
211 const missing = panes.filter((pane) => !selected.some((agent) => agent.pane === pane));
212 if (missing.length > 0)
213 return fail(`not agent panes: ${missing.join(", ")}; call fleet_status for the current list`);
214 if (selected.some((agent) => agent.workspace === undefined))
215 return fail(
216 "Herdr did not report workspace IDs for the selected panes; fleet_setup requires workspace metadata",
217 );
218 const workspaces = [
219 ...new Set(
220 selected.flatMap((agent) => (agent.workspace === undefined ? [] : [agent.workspace])),
221 ),
222 ];
223 const memberPanes = new Set(members(agents, workspaces).map((agent) => agent.pane));
224 let fleetId = "";
225 await context.store.update((memory) => {
226 fleetId = `fleet-${memory.nextFleetId}`;
227 return {
228 ...memory,
229 nextFleetId: memory.nextFleetId + 1,
230 dispatched: memory.dispatched.filter((pane) => !memberPanes.has(pane)),
231 fleets: {
232 ...memory.fleets,
233 [fleetId]: {
234 workspaces,
235 watched: {},
236 baselines: Object.fromEntries(
237 Object.entries(memory.baselines).filter(([pane]) => memberPanes.has(pane)),
238 ),
239 dispatched: memory.dispatched.filter((pane) => memberPanes.has(pane)),
240 nudged: false,
241 },
242 },
243 };
244 });
245 return ok({ fleetId, workspaces, panes: members(agents, workspaces).map((agent) => agent.pane) });
246};
247
248const status = async (context: ScopedContext): Promise<ToolOutcome> => {
249 const agents = await toolAgents(context);
250 const memory = await context.store.read();
251 return ok(agents.map((agent) => ({ ...agent, watched: agent.pane in memory.watched })));
252};
253
254const wait = async (
255 input: Record<string, unknown>,
256 context: ScopedContext,
257): Promise<ToolOutcome> => {
258 const panes = strings(input.panes);
259 if (panes.length === 0) return fail("panes must list at least one pane id");
260 const until = strings(input.until) as FleetStatus[];
261 const timeoutMs = typeof input.timeoutMs === "number" ? input.timeoutMs : 60_000;
262
263 for (const pane of panes) context.awaiting.add(waitKey(context.fleetId, pane));
264 try {
265 const outcome = await waitAny(
266 context.runner,
267 panes,
268 until.length > 0 ? until : SETTLED,
269 timeoutMs,
270 context.signal,
271 );
272 if (outcome.kind === "failed") return fail(outcome.error);
273 if (outcome.kind === "aborted") return fail("fleet_wait was interrupted");
274 if (outcome.kind === "timeout") return ok({ timedOut: true, panes });
275 if (!panes.includes(outcome.agent.pane))
276 return fail("Herdr returned an unexpected pane for fleet_wait");
277 if (
278 context.workspaces !== undefined &&
279 members([outcome.agent], context.workspaces).length === 0
280 )
281 return fail(
282 `The waited pane is no longer a member of ${context.fleetId}; call fleet_status with fleetId`,
283 );
284 await context.store.update((memory) => forget(memory, [outcome.agent.pane]));
285 return ok({ timedOut: false, settled: outcome.agent });
286 } finally {
287 for (const pane of panes) context.awaiting.delete(waitKey(context.fleetId, pane));
288 }
289};
290
291const read = async (input: Record<string, unknown>, context: ToolContext): Promise<ToolOutcome> => {
292 const pane = typeof input.pane === "string" ? input.pane : "";
293 if (pane === "") return fail("pane is required");
294 const lines = Math.max(1, Math.min(typeof input.lines === "number" ? input.lines : 80, 400));
295 const output = await context.runner(readArgv(pane, lines), 15_000, context.signal);
296 assertOk(output);
297 return { text: output.stdout, isError: false };
298};
299
300const send = async (
301 input: Record<string, unknown>,
302 context: ScopedContext,
303): Promise<ToolOutcome> => {
304 const pane = typeof input.pane === "string" ? input.pane : "";
305 const text = typeof input.text === "string" ? input.text : "";
306 if (pane === "" || text === "") return fail("pane and text are required");
307 const before = (await toolAgents(context)).find((agent) => agent.pane === pane);
308 if (before === undefined)
309 return fail(`${pane} is not an agent pane; call fleet_status for the current list`);
310
311 assertOk(await context.runner(promptArgv(pane, text), 30_000, context.signal));
312 const watch = input.watch !== false;
313 await context.store.update((memory) => {
314 const dispatched = recordDispatch(memory, pane);
315 return watch
316 ? { ...dispatched, watched: { ...dispatched.watched, [pane]: before.seq } }
317 : { ...dispatched, baselines: { ...dispatched.baselines, [pane]: before.seq } };
318 });
319 return ok({ sent: true, pane, watched: watch });
320};
321
322const watchPanes = async (
323 input: Record<string, unknown>,
324 context: ScopedContext,
325): Promise<ToolOutcome> => {
326 const panes = strings(input.panes);
327 if (panes.length === 0) return fail("panes must list at least one pane id");
328 if (input.unwatch === true) {
329 await context.store.update((memory) => unwatch(memory, panes));
330 return ok({ unwatched: panes });
331 }
332
333 const agents = await toolAgents(context);
334 const missing = panes.filter((pane) => !agents.some((agent) => agent.pane === pane));
335 if (missing.length > 0)
336 return fail(`not agent panes: ${missing.join(", ")}; call fleet_status for the current list`);
337 await context.store.update((memory) => ({
338 ...memory,
339 watched: {
340 ...memory.watched,
341 ...Object.fromEntries(
342 panes.map((pane) => [
343 pane,
344 memory.watched[pane] ??
345 memory.baselines[pane] ??
346 agents.find((agent) => agent.pane === pane)?.seq ??
347 0,
348 ]),
349 ),
350 },
351 baselines: Object.fromEntries(
352 Object.entries(memory.baselines).filter(([pane]) => !panes.includes(pane)),
353 ),
354 }));
355 return ok({
356 watching: agents
357 .filter((agent) => panes.includes(agent.pane))
358 .map((agent) => ({ pane: agent.pane, name: agent.name, status: agent.status })),
359 note: "A pane that already finished work this session sent it fires on the next check. Any other idle or done pane fires only after its next state change.",
360 });
361};
362
363/**
364 * Runs one fleet tool call.
365 *
366 * @param name the tool
367 * @param input the model's arguments
368 * @param context the harness's runner, store, and abort signal
369 * @returns the text the model reads
370 */
371export const executeFleetTool = async (
372 name: FleetToolName,
373 input: Record<string, unknown>,
374 context: ToolContext,
375): Promise<ToolOutcome> => {
376 try {
377 if (name === "fleet_setup") return await setup(input, context);
378 if (input.fleetId !== undefined && (typeof input.fleetId !== "string" || input.fleetId === ""))
379 return fail("fleetId must be a nonempty ID issued by fleet_setup");
380 const fleetId = input.fleetId as string | undefined;
381 const store = scopedStore(context.store, fleetId);
382 let agents: FleetAgent[] | undefined;
383 let workspaces: string[] | undefined;
384 if (fleetId !== undefined) {
385 // Validate the ID before invoking Herdr or touching a target pane.
386 const memory = await store.read();
387 const fleet = (await context.store.read()).fleets[fleetId]!;
388 workspaces = fleet.workspaces;
389 const listed = members(
390 await listFleet(context.runner, context.selfPane, context.signal),
391 workspaces,
392 );
393 agents = listed;
394 const targets = typeof input.pane === "string" ? [input.pane] : strings(input.panes);
395 const outside = targets.filter(
396 (pane) =>
397 !listed.some((agent) => agent.pane === pane) &&
398 !(
399 name === "fleet_watch" &&
400 input.unwatch === true &&
401 Object.hasOwn(memory.watched, pane)
402 ),
403 );
404 if (outside.length > 0)
405 return fail(
406 `not members of ${fleetId}: ${outside.join(", ")}; call fleet_status with fleetId for the current list`,
407 );
408 }
409 const scoped = { ...context, store, fleetId, agents, workspaces };
410 switch (name) {
411 case "fleet_status":
412 return await status(scoped);
413 case "fleet_wait":
414 return await wait(input, scoped);
415 case "fleet_read":
416 return await read(input, scoped);
417 case "fleet_send":
418 return await send(input, scoped);
419 case "fleet_watch":
420 return await watchPanes(input, scoped);
421 }
422 } catch (error) {
423 if (error instanceof HerdrFailure && error.code === "agent_blocked") {
424 return fail(
425 "agent_blocked: the agent is waiting on a question. Read it with fleet_read and answer that instead.",
426 );
427 }
428 return fail(error instanceof Error ? error.message : String(error));
429 }
430};
431types/index.d.ts 27 lines1/**
2 * Persisted global activity and session-local workspace fleets for the mod.
3 */
4export type ToCodeMemory = {
5 watched: Record<string, number>;
6 baselines: Record<string, number>;
7 dispatched: string[];
8 nudged: boolean;
9 nextFleetId: number;
10 fleets: Record<
11 string,
12 {
13 workspaces: string[];
14 watched: Record<string, number>;
15 baselines: Record<string, number>;
16 dispatched: string[];
17 nudged: boolean;
18 }
19 >;
20};
21
22declare module "claude-code" {
23 interface PluginState {
24 "to-code": { memory: ToCodeMemory };
25 }
26}
27