pi agent for claude: a `pi` subagent type whose model steps run the pi CLI, so pi subagents and agent-team teammates live in Claude Code's own agent UI…

Vendored from FazalAAli/pi-agent-for-claude (MIT, by Fazal Ali). Install with
npx claude-code-templates@latest --mod integrations/pi-agent-for-claude, or from the author's marketplace as described upstream. Report issues upstream.
Runs pi as a native Claude Code subagent. It shows in the agent list, streams live, takes follow-ups, and opens in a tmux pane as a teammate, while the model behind it is anything pi can reach.
> Have pi run with model kimi k2 and summarize this repo.
⏺ pi-agent-for-claude:pi(Summarize repo)
⎿ ▸ bash: {"command":"ls -la && git log --oneline -10"}
▸ bash: {"command":"cat package.json"}
This repo is an Expo/React Native calling app with a Go admin service...
— answered by pi, openrouter/moonshotai/kimi-k2
Claude Code starts a real subagent. The mod replaces only its model requests (turn.step) with a detached pi --mode json run and streams pi's text, thinking and tool calls. Hook calls get 10 s, so long runs span several steps chained through a no-op pi_progress tool. pi's final answer goes back as text, SubagentHandback, or SendMessage for teammates, with pi's real token usage.
agent.spawn and turn.step)Ask Claude for pi in plain language. It picks the pi-agent-for-claude:pi agent type.
pi-model: <pattern> line, passed to pi --model. Without it pi uses its default. pi --list-models <search> lists models.Answers end with — answered by pi, <provider/model>, read from pi's own events. Trust it over what the model says about itself.
pi runs its own read/bash/edit/write tools, outside Claude Code's permission prompts, tool policy and safety classifier. A pi agent can change files without asking. Restrict pi in its own config if that matters.
env.~/.claude/projects/, exact engine message text, and ps to detect teammates.npx claude-code-templates@latest --mod integrations/pi-agent-for-claude
Then add the flags to ~/.claude/settings.json, since shell variables alone don't reach teammate panes, and restart (inside tmux for teammate panes):
{
"env": {
"CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS": "1"
},
"teammateMode": "tmux"
}
It is written to .claude/skills/pi-agent-for-claude/, which Claude Code auto-loads as pi-agent-for-claude@skills-dir once the workspace is trusted. For one session with hot reload: claude --plugin-dir .claude/skills/pi-agent-for-claude. claude plugin validate .claude/skills/pi-agent-for-claude prints every event it hooks and every $ call it makes.
Requirements. Mods are on by default in Claude Code 2.1.287+, and this one needs 2.1.275+. Typed against Anthropic's declarations: https://github.com/anthropics/claude-code/tree/main/mods
hooks/register.ts 608 lines1import type { EngineInterface, On, TurnStepChunk, TurnStepResult, TurnUsage } from 'claude-code'
2
3/** How often a running pi's event stream is read for new events. */
4const STREAM_POLL_MS = 150
5
6/**
7 * How long one step streams pi before handing off to the next. The engine
8 * gives a hook call 10 s and, past it, finishes the step with the real model;
9 * pi runs detached, so a longer run spans several steps, chained through
10 * PROGRESS_TOOL.
11 */
12const STEP_BUDGET_MS = 7_000
13
14/**
15 * The no-op tool a step ends on while pi is still running: its call makes the
16 * engine run another step, which carries on streaming the same pi run.
17 */
18const PROGRESS_TOOL = 'pi_progress'
19
20/** How long a follow-up step waits for its message to reach the transcript. */
21const FOLLOW_UP_WAIT_MS = 5_000
22const POLL_MS = 200
23
24/**
25 * The tool a subagent delivers its final report through, where the session
26 * runs the hand-back contract (`CLAUDE_CODE_SENDMESSAGE_HANDBACK`). A report
27 * left as plain text is bounced back with `[handback-send-enforce]`.
28 */
29const HANDBACK_TOOL = 'SubagentHandback'
30
31/** The tool a teammate answers another agent with. */
32const SEND_TOOL = 'SendMessage'
33
34/**
35 * User messages the engine writes into a subagent's transcript itself. They
36 * are not the person's follow-ups, so pi is never sent them.
37 */
38const ENGINE_PREFIXES = [
39 '<system-reminder>',
40 '[handback-send-enforce]',
41 '[Your previous response had no visible output',
42]
43
44/**
45 * Registers the `pi` subagent type.
46 *
47 * The spawn runs as it always does, so the engine starts a real subagent loop
48 * with its own id, transcript and entry in `$.agent.list()`; only the model
49 * request inside that loop is replaced, by a run of the pi CLI. Each loop gets
50 * a pi session named after its agent id, so a follow-up continues the same pi
51 * conversation.
52 *
53 * Which model pi runs and which tools it may use are pi's own settings, not
54 * this mod's: it shells out to `pi` and lets pi's config decide.
55 *
56 * @param on the engine's registrar
57 */
58export function register(on: On) {
59 const seeds = new Map<string, string>()
60 const consumed = new Map<string, number>()
61 // Loops whose last step delivered pi's answer through a tool call, with the
62 // text the step after it ends on.
63 const delivered = new Map<string, string>()
64 // This process's teammate key when it runs as a pi teammate (`claude
65 // --agent-type pi-agent-for-claude:pi` in its own pane): then its main loop is pi's.
66 let teammate: string | undefined
67 // PROGRESS_TOOL's full name as the engine registered it, `mcp__<plugin>__…`.
68 let progressTool = ''
69 const lastAnswers = new Map<string, string>()
70 const runs = new Map<string, Run>()
71 // The pi model each loop asked for with a `pi-model:` line; it holds for
72 // that loop's later messages until another such line names a new one.
73 const models = new Map<string, string>()
74
75 on('session.start', async ($, e, next) => {
76 const started = await next(e)
77 teammate = await teammateOf($)
78 ;({ tool: progressTool } = await $.tool.register({
79 name: PROGRESS_TOOL,
80 description: 'Internal to pi subagents: marks a pi run still in progress. Never call it.',
81 }))
82 return started
83 })
84
85 on('tool.call', async ($, e, next) => {
86 if (e.tool !== progressTool) return next(e)
87 const key = e.agentId ?? teammate
88 if (key === undefined || !runs.has(key)) {
89 return { deny: `${PROGRESS_TOOL} is internal to pi subagents.` }
90 }
91 return { result: 'pi is still working.' }
92 })
93
94 on('agent.spawn', async ($, e, next) => {
95 const started = await next(e)
96 if (isPi(e.subagentType) && started.agentId) {
97 seeds.set(started.agentId, e.prompt)
98 }
99 return started
100 })
101
102 on('turn.step', async function* ($, e, next) {
103 // A pi subagent's loop, keyed by its id, or this process's main loop when
104 // it is a pi teammate; any other loop is left to the model.
105 const seed = e.agentId === undefined ? undefined : seeds.get(e.agentId)
106 const agentId = e.agentId === undefined ? teammate : seed === undefined ? undefined : e.agentId
107 if (agentId === undefined) {
108 return yield* next(e)
109 }
110 const deadline = (await $.clock.now()) + STEP_BUDGET_MS
111
112 // The step after a delivery is that tool's result coming back: the answer
113 // is delivered, so the loop ends here without asking pi again. It ends on
114 // visible text, since an empty response is bounced back as one.
115 const done = delivered.get(agentId)
116 if (done !== undefined) {
117 delivered.delete(agentId)
118 return yield* respond(e, [
119 { kind: 'text', index: 0, text: done },
120 { kind: 'stop', stopReason: 'end_turn', usage: null },
121 ], done, [])
122 }
123
124 let run = runs.get(agentId)
125 if (run === undefined) {
126 const already = consumed.get(agentId) ?? 0
127 const { prompt, count, handback, bounced, replyTo } = seed === undefined
128 ? { ...(await teammatePromptOf($, already)), handback: false, bounced: false }
129 : { ...(await promptOf($, agentId, seed, already)), replyTo: undefined }
130 consumed.set(agentId, count)
131 if (bounced || prompt === undefined) {
132 // A bounce is the engine refusing a plain-text report: hand back the
133 // answer pi already gave rather than asking it again.
134 const text = bounced
135 ? lastAnswers.get(agentId) ?? ''
136 : "pi agent: no new message reached this agent's transcript."
137 return yield* finish(e, agentId, text, text, 1, handback || bounced, undefined, [{ kind: 'text', index: 0, text }])
138 }
139 const asked = modelOf(prompt)
140 if (asked.model) models.set(agentId, asked.model)
141 run = await start($, agentId, asked.prompt, count, handback, models.get(agentId))
142 run.replyTo = replyTo
143 runs.set(agentId, run)
144 }
145
146 const { shown, blocks, lastKind } = yield* pump($, run, deadline, next.signal)
147 if (!run.ended) {
148 const input = {}
149 return yield* respond(e, [
150 { kind: 'tool', index: blocks, id: `toolu_pi_${agentId}_${run.id}_${run.steps++}`, name: progressTool },
151 { kind: 'input', index: blocks, json: JSON.stringify(input) },
152 { kind: 'stop', stopReason: 'tool_use', usage: takeUsage(run) },
153 ], shown, [{ name: progressTool, input }])
154 }
155
156 runs.delete(agentId)
157 // The hand-back reminder can reach the transcript after the run's first
158 // step read it; checking again now saves a bounce.
159 const handback = seed !== undefined && (run.handback || (await transcriptOf($, agentId)).handback)
160 const answer = run.report.trim() || run.shown
161 // Named from pi's own events: a model asked what it is often answers
162 // wrongly, and the parent has nothing else to check it by.
163 const report = run.model ? `${answer}\n\n— answered by pi, ${run.model}` : answer
164 // A run spanning several steps streamed its answer across them; the last
165 // step carries the whole of it, since that step is what the parent reads.
166 // The parent reads the last text block, so it must hold the whole answer.
167 // A run this step only finished gets the answer again as its own block; one
168 // whose answer streamed whole here gets just the model line, onto that
169 // same block, not alone in a new one.
170 const whole = run.steps > 0 && !shown.includes(answer)
171 const extra = whole ? report : report.slice(answer.length)
172 const into = !whole && lastKind === 'text' ? blocks - 1 : blocks
173 const tail: TurnStepChunk[] = extra ? [{ kind: 'text', index: into, text: extra }] : []
174 return yield* finish(e, agentId, shown + extra, report, into + (extra ? 1 : 0), handback, run.replyTo, tail, takeUsage(run))
175 })
176
177 /**
178 * Ends a step on pi's answer: plain text, or text then the tool call that
179 * delivers it - HANDBACK_TOOL where the session runs the hand-back contract,
180 * or, in a pi teammate, SendMessage to whoever sent the message pi answered.
181 */
182 async function* finish(
183 e: { turnId: string; index: number },
184 agentId: string,
185 shown: string,
186 report: string,
187 blocks: number,
188 handback: boolean,
189 replyTo: string | undefined,
190 chunks: TurnStepChunk[],
191 usage: TurnUsage | null = null,
192 ) {
193 lastAnswers.set(agentId, report)
194 const delivery = handback
195 ? { name: HANDBACK_TOOL, input: { message: report }, done: 'Report delivered.' }
196 : replyTo !== undefined
197 ? { name: SEND_TOOL, input: { to: replyTo, message: report }, done: `Sent to ${replyTo}.` }
198 : undefined
199 if (delivery === undefined) {
200 return yield* respond(e, [...chunks, { kind: 'stop', stopReason: 'end_turn', usage }], shown, [])
201 }
202 delivered.set(agentId, delivery.done)
203 const id = `toolu_pi_${agentId.replace(/[^A-Za-z0-9_]/g, '_')}_${delivery.name}_${consumed.get(agentId) ?? 0}`
204 return yield* respond(e, [
205 ...chunks,
206 { kind: 'tool', index: blocks, id, name: delivery.name },
207 { kind: 'input', index: blocks, json: JSON.stringify(delivery.input) },
208 { kind: 'stop', stopReason: 'tool_use', usage },
209 ], shown, [{ name: delivery.name, input: delivery.input }])
210 }
211}
212
213/**
214 * The tokens a run used since the last step reported them, as this step's
215 * usage under the model pi named; null before pi named one. Resets the count.
216 */
217export function takeUsage(run: Run | undefined): TurnUsage | null {
218 if (!run?.model) return null
219 const { i, o, cr, cw } = run.usage
220 run.usage = { i: 0, o: 0, cr: 0, cw: 0 }
221 return {
222 model: run.model,
223 input_tokens: i,
224 output_tokens: o,
225 cache_read_input_tokens: cr,
226 cache_creation_input_tokens: cw,
227 }
228}
229
230/**
231 * This process's teammate key, when it is a pi teammate: a teammate runs as
232 * its own `claude --agent-id <id> --agent-type <type>` process, and a hooks
233 * module has no API for its process's arguments, so this reads them with `ps`.
234 * The key names the teammate's pi session; undefined in any other process.
235 */
236export async function teammateOf($: EngineInterface) {
237 const { stdout } = await $.process.run(['sh', '-c', 'ps -o command= -p "$PPID"'])
238 const type = stdout.match(/--agent-type\s+(\S+)/)?.[1]
239 const id = stdout.match(/--agent-id\s+(\S+)/)?.[1]
240 if (type === undefined || id === undefined || !isPi(type)) return undefined
241 return id.replace(/[^A-Za-z0-9_-]/g, '-')
242}
243
244/**
245 * Waits for a message past the first `already` in a pi teammate's own
246 * transcript - its main loop, so `$.session.messages()` reads it - and returns
247 * the newest as pi's prompt, with who sent it when another agent did.
248 *
249 * A message from another agent arrives wrapped as `<teammate-message
250 * teammate_id="team-lead">…</teammate-message>`; one typed into the pane is
251 * plain, and is answered in the pane rather than sent anywhere.
252 */
253export async function teammatePromptOf($: EngineInterface, already: number) {
254 for (let waited = 0; ; waited += POLL_MS) {
255 const texts = (await $.session.messages())
256 .filter((m) => m.role === 'user')
257 .map((m) => m.text.trim())
258 .filter((text) => text && !ENGINE_PREFIXES.some((prefix) => text.startsWith(prefix)))
259 if (texts.length > already) {
260 const text = texts.at(-1) as string
261 const replyTo = text.match(/<teammate-message teammate_id="([^"]+)"/)?.[1]
262 const prompt = text
263 .replace(/<teammate-message teammate_id="([^"]+)"[^>]*>/g, 'Message from $1:\n')
264 .replace(/<\/teammate-message>/g, '')
265 .trim()
266 return { prompt, count: texts.length, replyTo }
267 }
268 if (waited >= FOLLOW_UP_WAIT_MS) return { prompt: undefined, count: already, replyTo: undefined }
269 await $.clock.sleep(POLL_MS)
270 }
271}
272
273/**
274 * Takes a `pi-model: <pattern>` line out of a prompt: the way a caller picks
275 * the model pi runs, since the Agent tool's own `model` names Claude models
276 * only. The pattern is anything `pi --model` takes (`openrouter/<id>`, a
277 * fuzzy `*sonnet*`); only the first such line counts.
278 */
279export function modelOf(text: string): { prompt: string; model?: string } {
280 const line = text.match(/^[ \t]*pi-model:[ \t]*([A-Za-z0-9_.:\/*~@+-]+)[ \t]*$/im)
281 if (!line) return { prompt: text }
282 return { prompt: text.replace(line[0], '').trim(), model: line[1] }
283}
284
285/** Whether a spawn names this plugin's type, plain or plugin-qualified. */
286export function isPi(subagentType: string) {
287 return subagentType === 'pi' || subagentType.endsWith(':pi')
288}
289
290/**
291 * Yields a step's chunks and returns the result the engine expects of it.
292 *
293 * @param e the step being answered
294 * @param chunks what the person watches stream, in order
295 * @param answer the step's visible text
296 * @param toolUses the tool calls the chunks made
297 */
298export async function* respond(
299 e: { turnId: string; index: number },
300 chunks: TurnStepChunk[],
301 answer: string,
302 toolUses: TurnStepResult['toolUses'],
303): AsyncGenerator<TurnStepChunk, TurnStepResult> {
304 for (const chunk of chunks) yield chunk
305 const stop = chunks.at(-1)
306 const stopReason = stop?.kind === 'stop' ? stop.stopReason : 'end_turn'
307 const usage = stop?.kind === 'stop' ? stop.usage : null
308 return { turnId: e.turnId, index: e.index, answer, toolUses, stopReason, usage }
309}
310
311/**
312 * Rewrites pi's `--mode json` event stream as one small JSON line per thing a
313 * step shows: `{k:"msg"}` when an assistant message starts, `{k:"text"|
314 * "thinking", t}` per delta, `{k:"tool", t}` per tool pi starts, `{k:"usage",
315 * i, o, cr, cw}` per finished assistant message, `{k:"end"}`. `msg` carries
316 * the `provider/model` pi sent the request to.
317 *
318 * pi's own events carry whole tool outputs and transcripts, and
319 * `$.process.run` cuts a child's stdout at its output limit, so the step
320 * never reads them raw.
321 */
322const FILTER = String.raw`
323const rl = require('readline').createInterface({ input: process.stdin })
324const out = (o) => process.stdout.write(JSON.stringify(o) + '\n')
325rl.on('line', (line) => {
326 let e
327 try { e = JSON.parse(line) } catch { return }
328 const d = e.type === 'message_update' ? e.assistantMessageEvent : undefined
329 if (e.type === 'message_start' && e.message?.role === 'assistant') out({ k: 'msg', model: e.message.provider + '/' + e.message.model })
330 else if (e.type === 'message_end' && e.message?.role === 'assistant') {
331 const u = e.message.usage ?? {}
332 out({ k: 'usage', i: u.input ?? 0, o: u.output ?? 0, cr: u.cacheRead ?? 0, cw: u.cacheWrite ?? 0 })
333 }
334 else if (d?.type === 'text_delta') out({ k: 'text', t: d.delta })
335 else if (d?.type === 'thinking_delta') out({ k: 'thinking', t: d.delta })
336 else if (e.type === 'tool_execution_start') {
337 const a = JSON.stringify(e.args ?? {})
338 out({ k: 'tool', t: '\n▸ ' + e.toolName + ': ' + (a.length > 100 ? a.slice(0, 100) + '…' : a) + '\n' })
339 } else if (e.type === 'agent_end') out({ k: 'end' })
340})
341`
342
343/** One pi run, as it streams across the steps it spans. */
344export type Run = {
345 /** The detached pipeline's pid. */
346 pid: string
347 /** pi's session, `pi-<agentId>`, which also names its process for a kill. */
348 session: string
349 /** The file FILTER writes pi's events to. */
350 events: string
351 /** This run's number for the agent, which keeps tool call ids unique. */
352 id: number
353 /** How many event lines have been read. */
354 read: number
355 /** Every piece of text shown so far, across steps. */
356 shown: string
357 /** The text of pi's latest assistant message: its answer once it ends. */
358 report: string
359 /** Whether pi has ended. */
360 ended: boolean
361 /** Whether the answer must be handed back through HANDBACK_TOOL. */
362 handback: boolean
363 /** In a pi teammate, who sent the message pi is answering. */
364 replyTo?: string
365 /** How many steps have handed off to the next through PROGRESS_TOOL. */
366 steps: number
367 /** The `provider/model` pi reported for its latest message. */
368 model: string
369 /** Tokens pi used since the last step reported them. */
370 usage: { i: number; o: number; cr: number; cw: number }
371}
372
373/**
374 * Starts pi detached in `agentId`'s pi session, piped through FILTER into a
375 * file the steps read: `$.process.run` only returns once its child exits, so
376 * it cannot stream a child it waits on.
377 *
378 * @param prompt what to send pi
379 * @param id this run's number for the agent
380 * @param handback whether pi's answer must be handed back
381 * @param model a pi `--model` pattern; absent, pi's own default
382 */
383export async function start(
384 $: EngineInterface,
385 agentId: string,
386 prompt: string,
387 id: number,
388 handback: boolean,
389 model: string | undefined,
390): Promise<Run> {
391 const session = `pi-${agentId}`
392 const events = `/tmp/pi-agent-for-claude-${agentId}-${id}.jsonl`
393 const cwd = await $.session.cwd()
394 const started = await $.process.run([
395 'sh',
396 '-c',
397 'nohup sh -c \'pi -p --mode json --session-id "$1" $5 -- "$2" </dev/null 2>"$3.err" | node -e "$4" >"$3"\' sh "$1" "$2" "$3" "$4" "$5" >/dev/null 2>&1 & echo $!',
398 'sh',
399 session,
400 prompt,
401 events,
402 FILTER,
403 // Unquoted in the script so an absent model adds no argument; a pattern
404 // is one word (`openrouter/moonshotai/kimi-k2`), which modelOf enforces.
405 model ? `--model ${model}` : '',
406 ], { cwd })
407 return {
408 pid: started.stdout.trim(),
409 session,
410 events,
411 id,
412 read: 0,
413 shown: '',
414 report: '',
415 ended: false,
416 handback,
417 steps: 0,
418 model: '',
419 usage: { i: 0, o: 0, cr: 0, cw: 0 },
420 }
421}
422
423/**
424 * Streams a pi run's new output as this step's chunks, until pi ends or the
425 * step's `deadline` passes: pi's text as text, its thinking as thinking, and
426 * each tool it starts as a one-line `▸ tool: args` note in the text.
427 *
428 * Never throws: a throw here would hand the step to the model beneath, a
429 * Claude subagent with Claude's tools, so a failure ends the run as text. An
430 * aborted step kills pi, since a detached pi would outlive it.
431 *
432 * ponytail: no run time limit; pi runs until it ends or a step is aborted. A
433 * subagent stopped between two steps leaves its pi to finish on its own.
434 *
435 * @param deadline when this step must stop streaming, in `$.clock` ms
436 * @param signal the step's abort signal
437 * @returns the text this step showed, how many content blocks it used, and
438 * the kind of the last one
439 */
440export async function* pump(
441 $: EngineInterface,
442 run: Run,
443 deadline: number,
444 signal: AbortSignal,
445): AsyncGenerator<TurnStepChunk, { shown: string; blocks: number; lastKind?: 'text' | 'thinking' }> {
446 let shown = ''
447 let index = -1
448 let kind: 'text' | 'thinking' | undefined
449 let failure = ''
450 let finished = false
451
452 try {
453 while (!run.ended && !signal.aborted && (await $.clock.now()) < deadline) {
454 const { stdout } = await $.process.run([
455 'sh',
456 '-c',
457 // Liveness first: once pi is gone the tail after it has every event.
458 'kill -0 "$3" 2>/dev/null; alive=$?; tail -n +"$2" "$1" 2>/dev/null; [ $alive = 0 ] || printf "\\n{\\"k\\":\\"exited\\"}\\n"',
459 'sh',
460 run.events,
461 String(run.read + 1),
462 run.pid,
463 ])
464 const lines = stdout.split('\n')
465 lines.pop() // the line still being written, or '' after a complete one
466 for (const line of lines) {
467 run.read += 1
468 let event: { k: string; t?: string; model?: string; i?: number; o?: number; cr?: number; cw?: number }
469 try {
470 event = JSON.parse(line)
471 } catch {
472 continue
473 }
474 if (event.k === 'end' || event.k === 'exited') {
475 run.ended = true
476 continue
477 }
478 if (event.k === 'msg') {
479 run.report = ''
480 if (event.model) run.model = event.model
481 continue
482 }
483 if (event.k === 'usage') {
484 run.usage.i += event.i ?? 0
485 run.usage.o += event.o ?? 0
486 run.usage.cr += event.cr ?? 0
487 run.usage.cw += event.cw ?? 0
488 continue
489 }
490 const pieceKind = event.k === 'thinking' ? 'thinking' : 'text'
491 const text = event.t ?? ''
492 if (pieceKind !== kind) {
493 kind = pieceKind
494 index += 1
495 }
496 if (event.k === 'text') run.report += text
497 if (pieceKind === 'text') shown += text
498 yield { kind: pieceKind, index, text }
499 }
500 if (!run.ended) await $.clock.sleep(STREAM_POLL_MS)
501 }
502 finished = true
503 } catch (err) {
504 failure = `pi agent: the pi run failed: ${String(err)}`
505 finished = true
506 } finally {
507 // Closed by the engine, or aborted: pi must not outlive the subagent.
508 if (!run.ended && (!finished || signal.aborted)) {
509 run.ended = true
510 await $.process.run(['sh', '-c', 'pkill -f -- "--session-id $1 "; rm -f "$2" "$2.err"', 'sh', run.session, run.events])
511 .catch(() => undefined)
512 }
513 }
514 if (failure) run.ended = true
515
516 if (run.ended) {
517 if (!failure && !(run.shown + shown).trim()) {
518 const { stdout: err } = await $.process.run(['sh', '-c', 'cat "$1.err" 2>/dev/null', 'sh', run.events])
519 .catch(() => ({ stdout: '' }))
520 // pi colours its errors; the escapes would show as noise in a transcript.
521 const plain = err.replace(/\u001b\[[0-9;]*m/g, '').trim()
522 failure = `pi ended with no answer.\n${plain || '(no stderr)'}`
523 }
524 if (failure) {
525 shown += failure
526 run.report = failure
527 index += 1
528 yield { kind: 'text', index, text: failure }
529 }
530 await $.process.run(['rm', '-f', run.events, `${run.events}.err`]).catch(() => undefined)
531 }
532 run.shown += shown
533 return { shown, blocks: index + 1, lastKind: kind }
534}
535
536/**
537 * Decides what to send pi this step, and whether its answer must be handed
538 * back through HANDBACK_TOOL.
539 *
540 * The first step sends the spawn's prompt. A later one waits for a person's
541 * message past the first `already` in `agentId`'s own transcript and sends the
542 * newest. A follow-up (SendMessage, or typed into the agent's view) is only
543 * found in `subagents/agent-<id>.jsonl` - `$.session.messages()` reads the main
544 * loop's - and the engine can start the step before it writes the message
545 * there, so this polls briefly.
546 *
547 * ponytail: finds the file by name under ~/.claude/projects; the id is unique,
548 * but it leans on Claude Code's on-disk layout, not an API. Polls for up to
549 * FOLLOW_UP_WAIT_MS; raise it if a slow disk drops follow-ups.
550 *
551 * @param seed the prompt the spawn was given
552 * @param already how many of the person's messages pi has been sent
553 */
554export async function promptOf($: EngineInterface, agentId: string, seed: string, already: number) {
555 for (let waited = 0; ; waited += POLL_MS) {
556 const { texts, handback, bounced } = await transcriptOf($, agentId)
557 if (already === 0) return { prompt: seed, count: Math.max(texts.length, 1), handback, bounced: false }
558 if (texts.length > already) return { prompt: texts.at(-1), count: texts.length, handback, bounced: false }
559 if (bounced) return { prompt: undefined, count: already, handback: true, bounced }
560 if (waited >= FOLLOW_UP_WAIT_MS) return { prompt: undefined, count: already, handback, bounced }
561 await $.clock.sleep(POLL_MS)
562 }
563}
564
565/**
566 * Reads `agentId`'s transcript: the text of every message the person (or the
567 * spawn) sent it, oldest first; whether the engine asked it to hand its
568 * report back through HANDBACK_TOOL; and whether the newest message is the
569 * engine bouncing a report that was not.
570 */
571export async function transcriptOf($: EngineInterface, agentId: string) {
572 const { stdout } = await $.process.run([
573 'sh',
574 '-c',
575 'f=$(find "$HOME/.claude/projects" -name "agent-$1.jsonl" -print -quit) && [ -n "$f" ] && cat "$f"',
576 'sh',
577 agentId,
578 ])
579 const all = stdout
580 .split('\n')
581 .flatMap((line) => {
582 try {
583 return [JSON.parse(line)]
584 } catch {
585 return [] // the line Claude Code is still writing
586 }
587 })
588 .filter((entry) => entry.type === 'user')
589 .map((entry) => textOf(entry.message?.content))
590 .filter(Boolean)
591 return {
592 texts: all.filter((text) => !ENGINE_PREFIXES.some((prefix) => text.startsWith(prefix))),
593 handback: all.some((text) => text.includes(HANDBACK_TOOL)),
594 bounced: all.at(-1)?.startsWith('[handback-send-enforce]') ?? false,
595 }
596}
597
598/** Joins a message's text: a plain string, or the text blocks of a block list. */
599export function textOf(content: unknown): string {
600 if (typeof content === 'string') return content.trim()
601 if (!Array.isArray(content)) return ''
602 return content
603 .filter((block) => block?.type === 'text')
604 .map((block) => block.text)
605 .join('\n')
606 .trim()
607}
608