/** Common-segment pipe request handler (prompt / interrupt / reply) over the subagent port and notice buffer. */ import { formatSettlementNotice } from './vocab.ts'; import type { PipeRequest, PipeResponse } from './pipe-channel.ts'; import type { SubagentPortBox } from './subagent-core.ts'; export interface MachineRequest { id: string; from: string | null; push: boolean; sinceTs: number; } export interface PipeHandlerSession { paneId: string; port: SubagentPortBox; claimSettleNotice: (key: string) => boolean; deliverNotice: (content: string, paneId?: string) => Promise; sendUserMessageIn: (content: string) => Promise; sendUserMessageAs: (content: string, mode: 'steer' | 'followUp') => Promise; abort: () => void; setPendingMachineRequest: (req: MachineRequest | null) => void; /** P0 (RFC §4.4): apply a role switch requested over the pipe (master → worker). */ applyRoleSwitch: (role: string, switchedBy: string) => Promise<{ ok: boolean; message: string }>; } export async function handlePipeRequest( req: PipeRequest, s: PipeHandlerSession, ): Promise { switch (req.type) { case 'ping': return { type: 'ok', id: req.id, detail: s.paneId }; case 'prompt': case 'follow_up': { s.setPendingMachineRequest({ id: req.id, from: req.from ?? null, push: req.push === true, sinceTs: Date.now(), }); // follow_up + steer → deliver between tool calls (B3). Initial prompt // stays followUp so an idle worker starts a new run (D96). if (req.steer === true) await s.sendUserMessageAs(req.text, 'steer'); else await s.sendUserMessageIn(req.text); return { type: 'ok', id: req.id }; } case 'interrupt': s.abort(); s.setPendingMachineRequest(null); return { type: 'ok', id: req.id }; case 'role': { const res = await s.applyRoleSwitch(req.role, s.paneId || 'pipe'); return res.ok ? { type: 'ok', id: req.id, detail: res.message } : { type: 'error', id: req.id, message: res.message }; } case 'reply': { s.port.current?.applyReplySession(req.paneId, req.sessionFile); // Consume before the finish notice: otherwise listRunningSubs still reports the pane and the // next master settle injects "still running" for an agent just announced finished. s.port.current?.consumeReply(req.paneId, req.text); // B8: the claim key `${paneId}:${requestId}` must match the poll loop's (subagent-poller.ts): // a child echoes the id of the pipe request it answered (`id: req.id` in index.ts's settle // push) and the parent hands that same id to startPoller (`prompt-` / `fu-`). // A mismatch surfaces as two "finished" notices for one settlement. if (s.claimSettleNotice(`${req.paneId}:${req.id}`)) { const notes = s.port.current?.reconcileOnReply(req.paneId) ?? []; const statLine = await s.port.current?.settleStatLine(req.paneId) ?? null; const head = formatSettlementNotice(`${req.paneId}`, req.text); const body = (req.sessionFile ? `${head}\nSession: ${req.sessionFile}` : head) + (statLine ? `\n${statLine}` : '') + (notes.length ? `\n${notes.join('\n')}` : ''); await s.deliverNotice(body, req.paneId); } return { type: 'ok', id: req.id }; } default: { // PipeRequest is closed; parsePipeLine may still forward an unknown type. const unknownReq: { type: string; id: string } = req; return { type: 'error', id: unknownReq.id, message: `unknown type ${unknownReq.type}` }; } } }