import { type ChildProcess, spawn } from "node:child_process"; import { accessSync, constants, existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync, } from "node:fs"; import { homedir, tmpdir } from "node:os"; import { delimiter, extname, isAbsolute, join } from "node:path"; import { fileURLToPath } from "node:url"; import { applyCavemanRpc } from "./caveman-rpc.js"; import { parseHerdrOutputFile } from "./herdr-output.js"; import { createOutputFilePath, writeInitialEntry } from "./output-file.js"; import type { SelectedAgentLaunchConfig } from "./selected-agent-launch-config.js"; import type { AgentContextMode } from "./types.js"; export interface HerdrRunArgs { agentId: string; cwd: string; parentSessionId: string; prompt: string; outputFile?: string; selectedAgentType?: string; selectedAgentLaunchConfig?: SelectedAgentLaunchConfig; context?: AgentContextMode; inheritContext?: boolean; maxTurns?: number; signal?: AbortSignal; piCommand?: string; resultWaitTimeoutMs?: number; approve?: boolean; cavemanEvents?: CavemanEventBus; onRunTag?: (tag: string) => void; onWarning?: (warnings: string[]) => void; } interface HerdrDependencies { hasCommand: (command: string) => boolean; spawnProcess: typeof spawn; exists: (path: string) => boolean; readFile: (path: string) => string; } interface CavemanEventBus { on(event: string, handler: (data: unknown) => void): () => void; emit(event: string, data: unknown): void; } function canExecute(path: string): boolean { try { accessSync(path, constants.X_OK); return true; } catch { return false; } } function hasPathSeparator(command: string): boolean { return command.includes("/") || command.includes("\\"); } function hasCommand(command: string): boolean { if (!command || command.includes("\0")) { return false; } if (isAbsolute(command) || hasPathSeparator(command)) { return canExecute(command); } return (process.env.PATH ?? "") .split(delimiter) .filter(Boolean) .some((dir) => canExecute(join(dir, command))); } const DEFAULT_DEPS: HerdrDependencies = { hasCommand, spawnProcess: spawn, exists: existsSync, readFile: (path) => readFileSync(path, "utf-8"), }; const PANE_ID_TOKEN_RE = /[A-Za-z0-9_.:-]+/g; const PANE_ID_RE = /^[-A-Za-z0-9_.:]+$/; const WHITESPACE_RE = /\s+/; let deps: HerdrDependencies = DEFAULT_DEPS; export function setHerdrRunnerDependenciesForTest( next: Partial ): void { deps = { ...DEFAULT_DEPS, ...next }; } export function resetHerdrRunnerDependenciesForTest(): void { deps = DEFAULT_DEPS; } function requireCommand(command: string, label: string): void { if (!deps.hasCommand(command)) { throw new Error( `Cannot use runnerBackend "herdr": required ${label} command not found: ${command}.` ); } } function getPiPackageIntercomPackagePath(): string { return join( homedir(), ".pi", "agent", "npm", "node_modules", "pi-intercom", "package.json" ); } function requirePiIntercom(cwd: string): void { if ( deps.exists("/tmp/pi-intercom/index.ts") || deps.exists(join(cwd, "node_modules/pi-intercom/index.ts")) || deps.exists(join(cwd, "node_modules/pi-intercom/package.json")) || deps.exists(getPiPackageIntercomPackagePath()) ) { return; } throw new Error( 'Cannot use runnerBackend "herdr": pi-intercom is required but was not found.' ); } function getHerdrCommand(): string { const command = process.env.PI_SUBAGENTS_HERDR_COMMAND ?? "herdr"; requireCommand(command, "Herdr"); return command; } function getPiCommand(explicit?: string): string { const command = explicit ?? process.env.PI_SUBAGENTS_PI_COMMAND ?? "pi"; requireCommand(command, "Pi CLI"); return command; } function getParentPaneId(): string { const paneId = process.env.HERDR_PANE_ID; if (!paneId) { throw new Error( 'Cannot use runnerBackend "herdr": HERDR_PANE_ID is required to split next to the parent pane.' ); } if (!PANE_ID_RE.test(paneId)) { throw new Error( `Cannot use runnerBackend "herdr": invalid HERDR_PANE_ID: ${paneId}` ); } return paneId; } export function resolveDefaultHerdrChildExtensionPath( runtimeModuleUrl: URL = new URL(import.meta.url) ): string { const runtimeExtension = extname(fileURLToPath(runtimeModuleUrl)); const childExtension = runtimeExtension === ".ts" ? ".ts" : ".js"; return fileURLToPath( new URL(`./herdr-child-extension${childExtension}`, runtimeModuleUrl) ); } function getChildExtensionPath(): string { const overridePath = process.env.PI_SUBAGENTS_HERDR_CHILD_EXTENSION; if (overridePath) { if (deps.exists(overridePath)) { return overridePath; } throw new Error( `Cannot use runnerBackend "herdr": PI_SUBAGENTS_HERDR_CHILD_EXTENSION does not exist: ${overridePath}` ); } return resolveDefaultHerdrChildExtensionPath(); } const RESULT_POLL_INTERVAL_MS = 25; const RESULT_WAIT_TIMEOUT_MS = 30 * 60 * 1000; const ABORT_KILL_TIMEOUT_MS = 1000; const NO_ASSISTANT_RESULT_PREFIX = "Herdr child did not produce a final assistant result in output file"; const DIAGNOSTIC_TAIL_CHARS = 2000; const UNSUPPORTED_HERDR_BRIDGE_TOOL_NAMES = new Set([ "message_parent", "ask_parent", ]); function formatTail(label: string, text: string): string { if (!text.trim()) { return ""; } const tail = text.slice(-DIAGNOSTIC_TAIL_CHARS).trimEnd(); return `\n${label} tail:\n${tail}`; } function getOutputFileTail(path: string): string { try { return deps.readFile(path); } catch { return ""; } } function getNoAssistantResultMessage( path: string, beforeTimeout = false, outputTail = "" ): string { return `${NO_ASSISTANT_RESULT_PREFIX}${beforeTimeout ? " before timeout" : ""}: ${path}${formatTail("output file", outputTail)}`; } function getResultFromOutputFile(path: string): string { const metadata = parseHerdrOutputFile(path); if (metadata.finalEntryFound) { return metadata.responseText ?? ""; } if (metadata.firstMalformedError) { throw metadata.firstMalformedError; } throw new Error(getNoAssistantResultMessage(path)); } function isNoAssistantResultError(error: unknown, path: string): boolean { return ( error instanceof Error && (error.message === getNoAssistantResultMessage(path) || error.message === getNoAssistantResultMessage(path, true)) ); } function waitForAbortableDelay( ms: number, signal?: AbortSignal ): Promise { return new Promise((resolve, reject) => { if (signal?.aborted) { reject(new Error("Herdr child aborted.")); return; } const timeout = setTimeout(() => { signal?.removeEventListener("abort", onAbort); resolve(); }, ms); const onAbort = () => { clearTimeout(timeout); reject(new Error("Herdr child aborted.")); }; signal?.addEventListener("abort", onAbort, { once: true }); }); } async function waitForAssistantResult( outputFile: string, timeoutMs = RESULT_WAIT_TIMEOUT_MS, signal?: AbortSignal ): Promise { const deadline = Date.now() + timeoutMs; while (true) { if (signal?.aborted) { throw new Error("Herdr child aborted."); } try { return getResultFromOutputFile(outputFile); } catch (error) { if (!isNoAssistantResultError(error, outputFile)) { throw error; } } if (Date.now() >= deadline) { throw new Error( getNoAssistantResultMessage( outputFile, true, getOutputFileTail(outputFile) ) ); } await waitForAbortableDelay(RESULT_POLL_INTERVAL_MS, signal); } } function waitForExit( child: ChildProcess, signal?: AbortSignal, label = "Herdr child" ): Promise { return new Promise((resolve, reject) => { let abortKillTimeout: NodeJS.Timeout | undefined; const cleanup = () => { signal?.removeEventListener("abort", onAbort); if (abortKillTimeout) { clearTimeout(abortKillTimeout); } }; const rejectAborted = () => { reject(new Error(`${label} aborted.`)); }; const onAbort = () => { child.kill("SIGTERM"); abortKillTimeout = setTimeout(() => { child.kill("SIGKILL"); }, ABORT_KILL_TIMEOUT_MS); abortKillTimeout.unref?.(); rejectAborted(); }; if (signal?.aborted) { child.kill("SIGTERM"); rejectAborted(); return; } signal?.addEventListener("abort", onAbort, { once: true }); child.once("error", (err) => { cleanup(); reject(err); }); child.once("exit", (code, sig) => { cleanup(); if (signal?.aborted) { reject(new Error(`${label} aborted.`)); return; } if (code === 0) { resolve(); return; } reject(new Error(`${label} exited with ${sig ?? code}.`)); }); }); } async function spawnAndCaptureStdout( command: string, args: string[], options: Parameters[2], signal?: AbortSignal, label = command ): Promise { if (signal?.aborted) { throw new Error(`${label} aborted.`); } const child = deps.spawnProcess(command, args, options); const chunks: Buffer[] = []; const stderrChunks: Buffer[] = []; const stdoutDone = new Promise((resolve, reject) => { if (!child.stdout) { resolve(); return; } child.stdout.once("end", resolve); child.stdout.once("close", resolve); child.stdout.once("error", reject); }); child.stdout?.on("data", (chunk) => { chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(String(chunk))); }); child.stderr?.on("data", (chunk) => { stderrChunks.push( Buffer.isBuffer(chunk) ? chunk : Buffer.from(String(chunk)) ); }); child.stderr?.resume(); try { await waitForExit(child, signal, label); } catch (error) { const stdout = Buffer.concat(chunks).toString("utf-8"); const stderr = Buffer.concat(stderrChunks).toString("utf-8"); if (error instanceof Error) { throw new Error( `${error.message}${formatTail("stdout", stdout)}${formatTail("stderr", stderr)}` ); } throw error; } await stdoutDone; return Buffer.concat(chunks).toString("utf-8"); } async function closeHerdrPane( herdrCommand: string, paneId: string, cwd: string, env: NodeJS.ProcessEnv ): Promise { await spawnAndCaptureStdout( herdrCommand, ["pane", "close", paneId], { cwd, env, stdio: "ignore" }, undefined, "Herdr pane close" ); } function getJsonPaneId(value: unknown): string | undefined { if (!value || typeof value !== "object") { return undefined; } const record = value as Record; if (typeof record.pane_id === "string") { return record.pane_id; } return getJsonPaneId(record.result) ?? getJsonPaneId(record.pane); } function parsePaneId(stdout: string): string { const trimmed = stdout.trim(); try { const paneId = getJsonPaneId(JSON.parse(trimmed) as unknown); if (paneId) { return paneId; } } catch { // Fall back to legacy plain-text parsing below. } const candidates = [ trimmed, ...trimmed.split(WHITESPACE_RE), ...[...trimmed.matchAll(PANE_ID_TOKEN_RE)].map((match) => match[0]), ]; const paneId = candidates .reverse() .find((candidate) => PANE_ID_RE.test(candidate)); if (!paneId) { throw new Error( `Cannot use runnerBackend "herdr": failed to parse pane id from Herdr split output: ${JSON.stringify(stdout)}` ); } return paneId; } function shellQuote(value: string): string { return `'${value.replaceAll("'", `'\\''`)}'`; } function createPromptFile(name: string, prompt: string): string { const dir = mkdtempSync(join(tmpdir(), "pi-subagents-herdr-prompt-")); const path = join(dir, name); writeFileSync(path, prompt, "utf-8"); return path; } function createSystemPromptFile(systemPrompt: string): string { return createPromptFile("system-prompt.md", systemPrompt); } function createUserPromptFile(prompt: string): string { return createPromptFile("user-prompt.md", prompt); } function cleanupTempFile(path: string | undefined): void { if (!path) { return; } rmSync(join(path, ".."), { force: true, recursive: true }); } function getModelInput(config: SelectedAgentLaunchConfig): string | undefined { const resolvedModelInput = config.model ? `${config.model.provider}/${config.model.id}` : undefined; if (config.modelInput && config.modelInput === resolvedModelInput) { return config.modelInput; } return resolvedModelInput; } function appendSelectedAgentArgs( piArgs: string[], config: SelectedAgentLaunchConfig, systemPromptPath: string, context: AgentContextMode, parentSessionId: string ): void { if (config.explicitMaxTurns) { throw new Error( 'Cannot use runnerBackend "herdr" with max_turns/maxTurns: Pi CLI has no max-turns support and child-side enforcement is not implemented.' ); } piArgs.push("--system-prompt", systemPromptPath); const disallowed = new Set(config.disallowedTools ?? []); const tools = config.builtinToolNames.filter((name) => !disallowed.has(name)); const unsupportedBridgeTools = tools.filter((name) => UNSUPPORTED_HERDR_BRIDGE_TOOL_NAMES.has(name) ); if (unsupportedBridgeTools.length > 0) { throw new Error( `Cannot use runnerBackend "herdr" with selected-agent parent bridge tools (${unsupportedBridgeTools.join(", ")}): Herdr child parent bridge tools are not implemented.` ); } if (tools.length > 0) { piArgs.push("--tools", tools.join(",")); } else { piArgs.push("--no-tools"); } if (Array.isArray(config.extensions)) { throw new Error( 'Cannot use runnerBackend "herdr" with selected-agent extension allowlists: Pi CLI extension allowlist support is not implemented.' ); } if (config.extensions === false) { piArgs.push("--no-extensions"); } if (config.noSkills) { piArgs.push("--no-skills"); } const modelInput = getModelInput(config); if (modelInput) { piArgs.push("--model", modelInput); } if (config.thinkingLevel) { piArgs.push("--thinking", config.thinkingLevel); } piArgs.push("--no-context-files"); if (context === "fork") { piArgs.push("--fork", parentSessionId); } } export async function runHerdrAgent(args: HerdrRunArgs): Promise<{ responseText: string; outputFile: string; warnings: string[]; }> { if (args.selectedAgentType && !args.selectedAgentLaunchConfig) { throw new Error( `Cannot use runnerBackend "herdr" with selected agent type "${args.selectedAgentType}": Pi CLI selected-agent parity is not implemented, and launching would silently widen or drop prompt/tool/skill/extension/context semantics.` ); } if (args.maxTurns != null) { throw new Error( 'Cannot use runnerBackend "herdr" with max_turns/maxTurns: Pi CLI has no max-turns support and child-side enforcement is not implemented.' ); } requirePiIntercom(args.cwd); const herdrCommand = getHerdrCommand(); const parentPaneId = getParentPaneId(); const piCommand = getPiCommand(args.piCommand); const childExtensionPath = getChildExtensionPath(); const outputFile = args.outputFile ?? createOutputFilePath(args.cwd, args.agentId, args.parentSessionId); writeInitialEntry(outputFile, args.agentId, args.prompt, args.cwd); const warnings: string[] = []; const addWarning = (warning: string) => { warnings.push(warning); args.onWarning?.([...warnings]); }; let selectedAgentSystemPrompt = args.selectedAgentLaunchConfig?.systemPrompt; if (args.selectedAgentLaunchConfig) { for (const warning of args.selectedAgentLaunchConfig.agentConfig .frontmatterWarnings ?? []) { addWarning(warning); } const caveman = args.selectedAgentLaunchConfig.agentConfig.caveman; if (caveman !== undefined && selectedAgentSystemPrompt !== undefined) { const result = await applyCavemanRpc( args.cavemanEvents, selectedAgentSystemPrompt, caveman ); selectedAgentSystemPrompt = result.systemPrompt; args.onRunTag?.(result.tag); if (result.warning) { addWarning(result.warning); } } } const systemPromptPath = selectedAgentSystemPrompt ? createSystemPromptFile(selectedAgentSystemPrompt) : undefined; const userPromptPath = createUserPromptFile(args.prompt); const env = { ...process.env, PI_SUBAGENT_ID: args.agentId, PI_SUBAGENT_OUTPUT_FILE: outputFile, PI_SUBAGENT_PARENT_SESSION_ID: args.parentSessionId, PI_INTERCOM_REQUIRED: "1", }; let paneId: string | undefined; try { const context = args.context ?? (args.inheritContext ? "fork" : "fresh"); const piArgs = [piCommand, args.approve ? "--approve" : "--no-approve"]; piArgs.push("--name", `pi-subagent-${args.agentId}`); if (args.selectedAgentLaunchConfig && systemPromptPath) { appendSelectedAgentArgs( piArgs, args.selectedAgentLaunchConfig, systemPromptPath, context, args.parentSessionId ); } piArgs.push("--extension", childExtensionPath, `@${userPromptPath}`); const childCommand = piArgs.map(shellQuote).join(" "); const splitStdout = await spawnAndCaptureStdout( herdrCommand, [ "pane", "split", parentPaneId, "--direction", "right", "--cwd", args.cwd, "--env", `PI_SUBAGENT_ID=${args.agentId}`, "--env", `PI_SUBAGENT_OUTPUT_FILE=${outputFile}`, "--env", `PI_SUBAGENT_PARENT_SESSION_ID=${args.parentSessionId}`, "--env", "PI_INTERCOM_REQUIRED=1", "--no-focus", ], { cwd: args.cwd, env, stdio: ["ignore", "pipe", "pipe"] }, args.signal, "Herdr pane split" ); paneId = parsePaneId(splitStdout); await spawnAndCaptureStdout( herdrCommand, ["pane", "run", paneId, childCommand], { cwd: args.cwd, env, stdio: "ignore" }, args.signal, "Herdr pane run" ); const responseText = await waitForAssistantResult( outputFile, args.resultWaitTimeoutMs, args.signal ); return { responseText, outputFile, warnings }; } catch (error) { if (paneId && args.signal?.aborted) { try { await closeHerdrPane(herdrCommand, paneId, args.cwd, env); } catch { // Preserve the original abort error. } } throw error; } finally { cleanupTempFile(systemPromptPath); cleanupTempFile(userPromptPath); } }