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…

https://github.com/user-attachments/assets/0bd4f5e4-c2d5-4e49-a882-30bf01357248
Run 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
Experimental: built on Claude Code's early-access function hooks and some undocumented engine behaviour, so Claude Code updates may break it.
Add to ~/.claude/settings.json (shell variables alone don't reach teammate panes):
{
"env": {
"CLAUDE_CODE_ENABLE_FUNCTION_HOOKS": "1",
"CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS": "1"
},
"teammateMode": "tmux"
}
Install from a Claude Code session, then restart (inside tmux for teammate panes):
/plugin marketplace add FazalAAli/pi-agent-for-claude
/plugin install pi-agent-for-claude@pi-agent-for-claude
To develop it, load a clone from disk instead: claude --plugin-dir /path/to/pi-agent-for-claude.
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; otherwise 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.
Claude Code starts a real subagent; the plugin replaces only its model requests (turn.step) with a detached pi --mode json run, streaming pi's text, thinking and tool calls. Hook calls get 10 s, so long runs span 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.
— answered by pi line, but check yourself. For teammates, the variable must be in settings env.~/.claude/projects/, exact engine message text, and ps to detect teammates.pi_progress calls appear in the transcript about every 7 s on long runs.claude --debug and search the log for pi-agent-for-claude and hook failed../smoke.sh runs an end-to-end check against real pi (costs a few cents). claude plugin validate . checks the manifest and hooks module.
MIT. See LICENSE.
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