SLOPSHOPPER

pi-agent-for-claude

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…

newguardtoolprocessagents
★ 23v0.1.0MITupdated 2026-09-18FazalAAli/pi-agent-for-claude
A shopper browsing a rack in a slop shop
README

https://github.com/user-attachments/assets/0bd4f5e4-c2d5-4e49-a882-30bf01357248

pi agent for claude

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.

Requirements

  • Claude Code 2.1.275+
  • pi, logged in to at least one provider
  • Node.js
  • macOS or Linux
  • tmux, for teammate panes

Setup

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.

Usage

Ask Claude for pi in plain language; it picks the pi-agent-for-claude:pi agent type.

  • Spawn: "Have pi review the error handling in src/api."
  • Pick a model: "Have pi run with model kimi k2 and …". Claude adds a pi-model: <pattern> line, passed to pi --model; otherwise pi uses its default. pi --list-models <search> lists models.
  • Follow up: messages to a running pi agent continue its pi session and model.
  • Teammate pane: name it. "Spawn a pi agent named piper to …"

Answers end with — answered by pi, <provider/model>, read from pi's own events. Trust it over what the model says about itself.

Security

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.

How it works

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.

Limitations

  • No hooks, no pi. With function hooks off, a "pi" agent runs as Claude Haiku. Claude is told to flag answers missing the — answered by pi line, but check yourself. For teammates, the variable must be in settings env.
  • Relies on undocumented engine behaviour: the 10 s hook budget, subagent transcript files under ~/.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.
  • Cost figures for pi agents are priced as Claude tokens; check your provider's billing.
  • Teammate shutdown requests go to pi as plain messages; close the pane yourself.

Troubleshooting

  • "pi ended with no answer" plus a pi error: usually an unresolvable model pattern, or a provider pi isn't logged in to.
  • Anything else: run claude --debug and search the log for pi-agent-for-claude and hook failed.

Testing

./smoke.sh runs an end-to-end check against real pi (costs a few cents). claude plugin validate . checks the manifest and hooks module.

License

MIT. See LICENSE.

Source 1 files
hooks/register.ts 608 lines
1import 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