/** * Workflow subagent runner. * * Each `agent()` call in a workflow script becomes one isolated in-process * AgentSession created here: in-memory session, normal trust-aware resources * and extensions, recursive orchestration/user-prompt tools denied, and an * optional one-shot `structured_output` tool when a schema is supplied. * * `runAgent()` never throws: every failure mode (session creation, provider * errors, aborts, missing structured output) settles into an `AgentOutcome`. */ import { type AgentSession, type AgentSessionEvent, type AgentSessionEventListener, createAgentSession, DefaultResourceLoader, type ExtensionAPI, type ExtensionContext, SessionManager, SettingsManager, type ToolDefinition, } from "@earendil-works/pi-coding-agent"; import { AgentToolRenderLedger } from "../shared/agent-tool-renderer.ts"; import { bindChildSessionExtensions, childToolPolicy, createChildResources, shutdownAndDisposeChildSession, } from "../shared/child-session.ts"; import { createToolCallTimeoutGuard } from "../shared/tool-call-timeout.ts"; import { type AgentUsage, emptyUsage, type TranscriptEntry } from "./model.ts"; import { childToolsWithStructuredOutput, createStructuredOutputTool, STRUCTURED_OUTPUT_SYSTEM_INSTRUCTION, } from "../shared/structured-output.ts"; import { buildWorkflowAgentPrompt } from "./prompt.ts"; import { AgentProgressProjection, type ProgressAssistantMessage, transcriptFromMessages, } from "./progress-projection.ts"; import { createReplayFilesystemBoundary, type ReplayFilesystemBoundaryOptions, } from "./replay-safety.ts"; import { truncateUtf8 } from "./serialization.ts"; import { bindWorkflowToolRenderer } from "./tool-renderer.ts"; const AGENT_OUTPUT_MAX_BYTES = 64 * 1024; export type WorkflowModel = NonNullable; export type ThinkingLevel = ReturnType; type AgentMessage = AgentSession["messages"][number]; type ToolTimingEvent = Extract< AgentSessionEvent, { type: "tool_execution_start" | "tool_execution_end" } >; export interface ToolExecutionTiming { startedAt?: number; finishedAt?: number; durationMs?: number; } export interface AgentOutcome { ok: boolean; /** Final assistant text (may be empty when only structured output was produced). */ output: string; /** Captured structured_output payload when a schema was supplied. */ structured?: unknown; error?: string; aborted: boolean; usage: AgentUsage; model?: string; contextWindow?: number; transcript: TranscriptEntry[]; } export interface AgentProgress { preview: string; usage: AgentUsage; model?: string; contextWindow?: number; transcript: TranscriptEntry[]; } export type WorkflowAgentSessionFactory = ( options: Parameters[0], ) => Promise<{ session: AgentSession }>; export interface RunAgentOptions { prompt: string; schema?: unknown; model?: WorkflowModel; thinkingLevel?: ThinkingLevel; cwd: string; loader: DefaultResourceLoader; settingsManager: SettingsManager; /** Optional per-run manager reused by one logical workflow operator. */ sessionManager?: SessionManager; modelRegistry: ExtensionContext["modelRegistry"]; /** Agent Type allowlist; childToolPolicy can only narrow capabilities. */ tools?: readonly string[]; signal?: AbortSignal; onProgress?: (progress: AgentProgress) => void; /** Canonical repository boundary required before this call can be journaled. */ replayFilesystemBoundary?: ReplayFilesystemBoundaryOptions; /** Test-only override for the per-tool execution timeout. */ toolCallTimeoutMs?: number; /** Test-only override for the end-to-end abort/shutdown deadline. */ shutdownTimeoutMs?: number; /** Test seam for lifecycle races; production always uses createAgentSession. */ sessionFactory?: WorkflowAgentSessionFactory; } /** Build a fresh extension runtime for each concurrent workflow child. */ export function createWorkflowResources( cwd: string, variant: "plain" | "structured", projectTrusted: boolean, agentTypePrompt?: string, ) { const appendSystemPrompt = [ ...(agentTypePrompt ? [agentTypePrompt] : []), ...(variant === "structured" ? [STRUCTURED_OUTPUT_SYSTEM_INSTRUCTION] : []), ]; return createChildResources({ cwd, projectTrusted, ...(appendSystemPrompt.length > 0 ? { appendSystemPrompt } : {}), }); } export function workflowChildTools( tools: readonly string[] | undefined, structured: boolean, ) { return childToolsWithStructuredOutput(tools, structured); } interface WorkflowToolSession { getAllTools(): Array<{ name: string; sourceInfo?: { path: string; source: string; origin: string }; }>; getToolDefinition(name: string): ToolDefinition | undefined; subscribe(listener: AgentSessionEventListener): () => void; } /** Guard current tools and tools registered by extensions at later agent starts. */ export function guardWorkflowChildTools( session: WorkflowToolSession, timeoutMs?: number, replayFilesystemBoundary?: ReplayFilesystemBoundaryOptions, ) { const boundary = replayFilesystemBoundary ? createReplayFilesystemBoundary(replayFilesystemBoundary) : undefined; const timeout = createToolCallTimeoutGuard(timeoutMs); const apply = () => { // Keep the replay boundary outermost: if the timeout rejects before a // cooperative tool finishes aborting, path revalidation still completes // before the child can use or journal the timeout result. timeout.apply(session); boundary?.apply(session); }; apply(); return session.subscribe((event) => { if (event.type === "agent_start") apply(); }); } type AssistantMessage = ProgressAssistantMessage; export { transcriptFromMessages }; export interface AssistantSettlement { stopReason: AssistantMessage["stopReason"]; errorMessage?: string; } /** * Observe assistant message_end events instead of rescanning mutable Pi state: * overflow recovery temporarily removes the failed assistant from that state. */ export function observeAssistantSettlement( previous: AssistantSettlement | undefined, message: AgentMessage | undefined, ) { if (message?.role !== "assistant") return previous; return { stopReason: message.stopReason, errorMessage: message.errorMessage, }; } export function agentFailureMessage( settlement: AssistantSettlement | undefined, promptErrorMessage?: string, ) { if ( settlement?.stopReason !== "error" && settlement?.errorMessage === undefined && promptErrorMessage === undefined ) { return undefined; } return settlement?.errorMessage ?? promptErrorMessage ?? "Agent failed"; } /** Record lifecycle timings without inferring completion from message timestamps. */ export function recordToolExecutionTiming( timings: Map, event: ToolTimingEvent, observedAt = Date.now(), ) { const previous = timings.get(event.toolCallId); if (event.type === "tool_execution_start") { if (previous?.startedAt !== undefined) return; timings.set(event.toolCallId, { ...previous, startedAt: observedAt }); return; } if (previous?.finishedAt !== undefined) return; const durationMs = previous?.startedAt === undefined ? undefined : Math.max(0, observedAt - previous.startedAt); timings.set(event.toolCallId, { ...previous, finishedAt: observedAt, ...(durationMs === undefined ? {} : { durationMs }), }); } function errorText(error: unknown): string { return (error instanceof Error ? error.message : String(error)).slice( 0, 16 * 1024, ); } export async function runAgent( options: RunAgentOptions, ): Promise { let structured: unknown; let settled = false; let customTools: ToolDefinition[] | undefined; let session: AgentSession | undefined; let unsubscribeToolGuards: (() => void) | undefined; let aborted = false; let abortOperation: Promise | undefined; let rejectForAbort: ((error: Error) => void) | undefined; let rejectForProjectionFailure: ((error: Error) => void) | undefined; const abortRace = new Promise((_resolve, reject) => { rejectForAbort = reject; }); const projectionFailureRace = new Promise((_resolve, reject) => { rejectForProjectionFailure = reject; }); // The same promise covers startup and prompting. Keep it observed even when // a pre-aborted run returns before either race is installed. void abortRace.catch(() => {}); void projectionFailureRace.catch(() => {}); const abortError = () => options.signal?.reason instanceof Error ? options.signal.reason : new Error("Agent was aborted"); const onAbort = () => { if (aborted) return; aborted = true; if (session) { try { abortOperation ??= session.abort(); void abortOperation.catch(() => {}); } catch (error) { abortOperation = Promise.reject(error); void abortOperation.catch(() => {}); } } rejectForAbort?.(abortError()); }; if (options.signal) { if (options.signal.aborted) onAbort(); else options.signal.addEventListener("abort", onAbort, { once: true }); } try { customTools = options.schema !== undefined ? [ createStructuredOutputTool(options.schema, (value) => { if (!settled) structured = value; }), ] : undefined; const childTools = workflowChildTools( options.tools, customTools !== undefined, ); if (aborted) throw abortError(); const sessionCreation = (options.sessionFactory ?? createAgentSession)({ cwd: options.cwd, ...(options.model ? { model: options.model } : {}), ...(options.thinkingLevel ? { thinkingLevel: options.thinkingLevel } : {}), resourceLoader: options.loader, settingsManager: options.settingsManager, sessionManager: options.sessionManager ?? SessionManager.inMemory(options.cwd), ...(customTools ? { customTools } : {}), ...childToolPolicy(childTools), }); // Promise.race cannot cancel the factory. If cancellation wins, retain // ownership of a late session and dispose it without re-entering startup. void sessionCreation .then( ({ session: lateSession }) => { if (!aborted) return; return shutdownAndDisposeChildSession(lateSession, { abort: true, timeoutMs: options.shutdownTimeoutMs, }); }, () => undefined, ) .catch(() => {}); ({ session } = await Promise.race([sessionCreation, abortRace])); if (aborted) throw abortError(); await Promise.race([ bindChildSessionExtensions(session, childTools), abortRace, ]); if (aborted) throw abortError(); unsubscribeToolGuards = guardWorkflowChildTools( session, options.toolCallTimeoutMs, options.replayFilesystemBoundary, ); } catch (error) { settled = true; unsubscribeToolGuards?.(); options.signal?.removeEventListener("abort", onAbort); const cleanup = session ? await shutdownAndDisposeChildSession(session, { abort: true, abortOperation, timeoutMs: options.shutdownTimeoutMs, }) : undefined; const cleanupError = cleanup?.errors.join("; "); if (aborted) { return { ok: false, output: "", error: cleanupError ? `Agent was aborted; Cleanup failed: ${cleanupError}` : "Agent was aborted", aborted: true, usage: emptyUsage(), model: options.model?.id, contextWindow: options.model?.contextWindow, transcript: [], }; } return { ok: false, output: "", error: `Failed to create agent session: ${errorText(error)}${cleanupError ? `; cleanup failed: ${cleanupError}` : ""}`, aborted: false, usage: emptyUsage(), model: options.model?.id, contextWindow: options.model?.contextWindow, transcript: [], }; } const childSession = session; const toolRenderer = new AgentToolRenderLedger(); let usage = emptyUsage(); let modelId = childSession.model?.id ?? options.model?.id; let contextWindow = childSession.model?.contextWindow; let assistantSettlement: AssistantSettlement | undefined; let promptErrorMessage: string | undefined; const toolTimings = new Map(); const captureToolRenderData = (messages: readonly AgentMessage[]) => { for (const message of messages) { if (message.role === "assistant") { for (const part of message.content) { if (part.type !== "toolCall") continue; toolRenderer.start( part.id, part.name, part.arguments, childSession.getToolDefinition(part.name), ); } } else if (message.role === "toolResult") { toolRenderer.end( message.toolCallId, message.toolName, message, message.isError, ); } } }; const refreshModel = (latestAssistant?: AssistantMessage) => { const sessionModel = childSession.model; modelId = sessionModel?.id ?? modelId; contextWindow = sessionModel?.contextWindow ?? contextWindow; if (!latestAssistant) return; const responseMatchesSession = !sessionModel || (latestAssistant.provider === sessionModel.provider && latestAssistant.model === sessionModel.id); const reportedId = latestAssistant.responseModel ?? latestAssistant.model; const reportedModel = responseMatchesSession ? options.modelRegistry.find(latestAssistant.provider, reportedId) : undefined; if (reportedModel) { modelId = reportedModel.id; contextWindow = reportedModel.contextWindow; } }; const sampleContextTokens = () => { const context = childSession.getContextUsage(); if ( typeof context?.contextWindow === "number" && Number.isFinite(context.contextWindow) && context.contextWindow > 0 ) { contextWindow = context.contextWindow; } if ( typeof context?.tokens === "number" && Number.isFinite(context.tokens) && context.tokens >= 0 ) { return context.tokens; } return context?.tokens === null ? null : undefined; }; const projection = new AgentProgressProjection(); let lastProjection = projection.snapshot(toolTimings); const snapshotProjection = () => { const snapshot = projection.snapshot(toolTimings); usage = snapshot.usage; refreshModel(snapshot.latestAssistant); lastProjection = snapshot; return snapshot; }; const reconcileProjection = () => { projection.replace(childSession.messages, sampleContextTokens()); return snapshotProjection(); }; const emitProgress = () => { const snapshot = snapshotProjection(); options.onProgress?.({ preview: snapshot.preview, usage, model: modelId, contextWindow, transcript: bindWorkflowToolRenderer(snapshot.transcript, toolRenderer), }); }; let compactionReconcileQueued = false; const queueCompactionReconcile = () => { if (compactionReconcileQueued) return; compactionReconcileQueued = true; queueMicrotask(() => { compactionReconcileQueued = false; if (settled) return; // Pi may synchronously remove a restored overflow assistant immediately // after compaction_end. Read canonical state after that mutation. try { reconcileProjection(); emitProgress(); } catch (error) { promptErrorMessage ??= `Failed to reconcile agent progress after compaction: ${errorText(error)}`; rejectForProjectionFailure?.(new Error(promptErrorMessage)); try { abortOperation ??= childSession.abort(); void abortOperation.catch(() => {}); } catch { // The projection failure remains authoritative; bounded cleanup gets // another chance to abort and dispose the child below. } } }); }; const unsubscribe = childSession.subscribe((event) => { if (settled) return; if (event.type === "tool_execution_start") { toolRenderer.start( event.toolCallId, event.toolName, event.args, childSession.getToolDefinition(event.toolName), ); } else if (event.type === "tool_execution_update") { toolRenderer.update( event.toolCallId, event.toolName, event.args, event.partialResult, ); } else if (event.type === "tool_execution_end") { toolRenderer.end( event.toolCallId, event.toolName, event.result, event.isError, ); } if (event.type === "message_end") { assistantSettlement = observeAssistantSettlement( assistantSettlement, event.message, ); // Pi publishes the finalized message before tool lifecycle delivery. // Hydrate native rendering from this one message without reading history. captureToolRenderData([event.message]); projection.append(event.message); } if ( event.type === "tool_execution_start" || event.type === "tool_execution_end" ) { recordToolExecutionTiming(toolTimings, event); } else if (event.type === "compaction_end") { queueCompactionReconcile(); return; } else if (event.type !== "message_end") { return; } emitProgress(); }); let output = ""; let transcript: TranscriptEntry[] = []; let cleanupErrors: string[] = []; try { // Operator reuse may restore messages before this activation subscribes. // Hydrate both bounded projections once; ordinary progress never rescans it. projection.replace(childSession.messages, sampleContextTokens()); captureToolRenderData(childSession.messages); snapshotProjection(); if (!aborted) { // Pi owns transport liveness and retries. Quiet model output is not // evidence of a stalled request (thinking and retry backoff can be silent). await Promise.race([ childSession.prompt(buildWorkflowAgentPrompt(options.prompt)), abortRace, projectionFailureRace, ]); } } catch (error) { promptErrorMessage ??= errorText(error); } finally { options.signal?.removeEventListener("abort", onAbort); settled = true; unsubscribe(); unsubscribeToolGuards?.(); try { const finalProjection = reconcileProjection(); output = truncateUtf8(finalProjection.preview, AGENT_OUTPUT_MAX_BYTES); transcript = bindWorkflowToolRenderer( finalProjection.transcript, toolRenderer, ); } catch (error) { promptErrorMessage ??= `Failed to reconcile final agent progress: ${errorText(error)}`; // A broken canonical accessor must not prevent owned-session cleanup. // Preserve the last projection that was fully observed instead. try { output = truncateUtf8(lastProjection.preview, AGENT_OUTPUT_MAX_BYTES); transcript = bindWorkflowToolRenderer( lastProjection.transcript, toolRenderer, ); } catch (fallbackError) { promptErrorMessage += `; failed to render last progress: ${errorText(fallbackError)}`; } } const cleanup = await shutdownAndDisposeChildSession(childSession, { abort: aborted || promptErrorMessage !== undefined, abortOperation, timeoutMs: options.shutdownTimeoutMs, }); cleanupErrors = cleanup.errors; } const cleanupError = cleanupErrors.length > 0 ? `Cleanup failed: ${cleanupErrors.join("; ")}` : undefined; if (aborted || assistantSettlement?.stopReason === "aborted") { return { ok: false, output, structured, error: cleanupError ? `Agent was aborted; ${cleanupError}` : "Agent was aborted", aborted: true, usage, model: modelId, contextWindow, transcript, }; } const failureMessage = agentFailureMessage(assistantSettlement, promptErrorMessage) ?? cleanupError; if (failureMessage !== undefined) { return { ok: false, output, structured, error: failureMessage, aborted: false, usage, model: modelId, contextWindow, transcript, }; } if (options.schema !== undefined && structured === undefined) { return { ok: false, output, error: "Agent finished without calling structured_output; no structured result matching the schema was produced.", aborted: false, usage, model: modelId, contextWindow, transcript, }; } if (options.schema === undefined && assistantSettlement === undefined) { return { ok: false, output, structured, error: "Agent finished without an assistant response.", aborted: false, usage, model: modelId, contextWindow, transcript, }; } return { ok: true, output, structured, aborted: false, usage, model: modelId, contextWindow, transcript, }; }