import { randomUUID } from "node:crypto"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import type { AgentToolResult } from "@earendil-works/pi-agent-core"; import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent"; import { canInvokeAgent, effectiveAgentInvocation, resolveAgentName, type AgentConfig, type AgentInvocationOrigin, type AgentScope } from "../../agents/agents.ts"; import { getArtifactsDir, getProjectChainRunsDir } from "../../shared/artifacts.ts"; import { ChainClarifyComponent, type ChainClarifyResult } from "./chain-clarify.ts"; import { resolveEffectiveThinking, toModelInfo, type ModelInfo } from "../../shared/model-info.ts"; import { executeChain } from "./chain-execution.ts"; import { beginForegroundChild, finishForegroundChild, forgetForegroundSteeringCleanup, foregroundSchedulingSettled, foregroundSteeringCleanupKey, foregroundSteeringPaths, retainForegroundSchedulingOwner, settleForegroundSchedulingOwner, trackForegroundSteeringCleanup, updateForegroundChild, } from "./foreground-control.ts"; import { resolveExecutionAgentScope } from "../../agents/agent-scope.ts"; import { handleManagementAction } from "../../agents/agent-management.ts"; import { buildDoctorReport } from "../../extension/doctor.ts"; import { clearPendingForegroundControlNotices } from "../../extension/control-notices.ts"; import { runSync } from "./execution.ts"; import { handleWatchdogToolAction, WATCHDOG_TOOL_ACTIONS } from "../../watchdog/tool-actions.ts"; import type { MainWatchdogRuntime } from "../../watchdog/runtime.ts"; import { buildModelCandidates, normalizeParentModel, resolveEffectiveSubagentModel, resolveModelCandidate, type ParentModel } from "../shared/model-fallback.ts"; import type { ModelScopeConfig } from "../shared/model-scope.ts"; import { aggregateParallelOutputs } from "../shared/parallel-utils.ts"; import { recordRun } from "../shared/run-history.ts"; import { buildChainInstructions, writeInitialProgressFile, getStepAgents, isParallelStep, isDynamicParallelStep, resolveStepBehavior, suppressProgressForReadOnlyTask, taskDisallowsFileUpdates, type ChainStep, type ResolvedStepBehavior, type SequentialStep, type StepOverrides, } from "../../shared/settings.ts"; import { discoverAvailableSkills, normalizeSkillInput } from "../../agents/skills.ts"; import { buildAsyncRunnerSteps, executeAsyncChain, executeAsyncSingle, formatAsyncStartedMessage, isAsyncAvailable } from "../background/async-execution.ts"; import { collectRequestedAgentNames, duplicateNames, escapeRegExp, firstChainAgent, firstRawChainTask, getRequestedModeLabel, isAsyncRunNotFound, isExactResumeError, isResumeAmbiguity, nestedRunAgent, nestedRunSessionFile, pathWithin, resolveAsyncEventGoal, resolveRequestedCwd, resumeTargetExact, canonicalizeAgentName, } from "./executor-helpers.ts"; import type { AgentDefaultContextPolicy } from "./executor-validation.ts"; import { applySingleAgentLaunchDefaults, buildRequestedModeError, collectChainSessionFiles, collectChainThinkingOverrides, preflightForkSessionsForStaticTasks, resolveAgentDefaultContextPolicy, resolveEffectiveToolBudget, resolveRunTimeout, shouldForkAgent, toExecutionErrorResult, withForkThinkingNotes, withResolvedContext, wrapChainTasksForFork, } from "./executor-validation.ts"; import type { ExecutionContextData, ExecutorDeps, SubagentParamsLike, } from "./executor-types.ts"; import * as control from "./executor-control.ts"; const { removeForegroundControlIfIdle, getForegroundControl, formatForegroundActivity, nestedResolutionScopeForExecutor, trustedSessionRootsForStatus, spawnBudgetErrorResult, withSpawnBudgetStatus, hasActiveSubagentChildren, countRequestedSubagentSpawns, foregroundStatusResult, trimRememberedForegroundRuns, foregroundChildActivityFromProgress, rememberForegroundRun, applyControlEventToRememberedForegroundRun, updateRememberedForegroundChild, resolveForegroundResumeTarget, resolveResumeTarget, getAsyncInterruptTarget, emitControlNotification, interruptAsyncRun, appendStepToAsyncChain, validateNestedSessionFile, resolveNestedResumeTarget, waitForNestedControlResult, sendNestedControlRequest, directNestedAsyncInterrupt, directNestedAsyncSteer, interruptNestedRun, resumeLiveNestedRun, steerNestedRun, resumeAsyncRun, resultSummaryForIntercom, formatFailedSingleRunOutput, createForegroundControlNotifier, emitForegroundResultIntercom, maybeBuildForegroundIntercomReceipt, rememberParentModel, } = control; import type { ScheduledRunAction } from "../background/scheduled-runs.ts"; import { enqueueChainAppendRequest, readPendingChainAppendRequests, runnerStepOutputNames } from "../background/chain-append.ts"; import { ChainOutputValidationError, validateChainOutputBindingsWithContext } from "../shared/chain-outputs.ts"; import { validateExecutionAcceptance } from "../shared/acceptance.ts"; import { createForkContextResolver, forkedChildRequiresThinkingOff } from "../../shared/fork-context.ts"; import { resolveCurrentSessionId } from "../../shared/session-identity.ts"; import { applyIntercomBridgeToAgent, INTERCOM_BRIDGE_MARKER, resolveIntercomBridge, resolveIntercomSessionTarget, resolveSubagentIntercomTarget, type IntercomBridgeState } from "../../intercom/intercom-bridge.ts"; import { formatControlIntercomMessage, formatControlNoticeMessage, resolveControlConfig, shouldNotifyControlEvent } from "../shared/subagent-control.ts"; import { resolveTurnBudgetConfig } from "../shared/turn-budget.ts"; import { formatSpawnBudget, getSpawnBudgetSnapshot, grantSpawnBudget, preflightSpawnBudget, preflightSpawnBudgetGrant, reserveSpawnBudget } from "../shared/spawn-budget.ts"; import { validateToolBudgetConfig } from "../shared/tool-budget.ts"; import { usageBudgetExceededMessage, usageBudgetState, validateUsageBudgetConfig } from "../shared/usage-budget.ts"; import { intersectSubagentCapabilityCeilings, resolveCurrentSubagentCapabilityCeiling, type ResolvedSubagentCapabilityCeiling } from "../shared/capability-ceiling.ts"; import { isAgentContractV1 } from "../shared/agent-contract.ts"; import { finalizeSingleOutput, injectSingleOutputInstruction, normalizeSingleOutputOverride, resolveSingleOutputPath, validateFileOnlyOutputMode } from "../shared/single-output.ts"; import { cleanupStructuredOutputRuntime, createStructuredOutputRuntime } from "../shared/structured-output.ts"; import { compactForegroundDetails, getSingleResultOutput, mapConcurrent, readStatus, resolveChildCwd, sumResultsCost, sumResultsUsage } from "../../shared/utils.ts"; import { DEFAULT_GLOBAL_CONCURRENCY_LIMIT, Semaphore } from "../shared/parallel-utils.ts"; import { formatParallelHandoffError, formatParallelHandoffReference, parallelHandoffPath, writeParallelHandoffGroup } from "../shared/parallel-handoff.ts"; import { summarizeContextModes, type ContextMode, type ContextSummary } from "../shared/context-mode.ts"; import { attachNestedChildrenToResultChildren, buildSubagentResultIntercomPayload, deliverSubagentResultIntercomEvent, formatSubagentResultReceipt, resolveSubagentResultStatus, stripDetailsOutputsForIntercomReceipt, } from "../../intercom/result-intercom.ts"; import { applySteeringRecoveryAgentConfig, buildRevivedAsyncTask, resolveAsyncResumeTarget, resolveAsyncRunLocation } from "../background/async-resume.ts"; import { deliverCheckpointDecisionRequest, deliverInterruptRequest, requestAsyncSteer } from "../background/control-channel.ts"; import { waitForSteeringAction } from "../background/steering.ts"; import { steerAsyncRun } from "./async-steering-action.ts"; import { stopAsyncRun } from "./async-stop-action.ts"; import { reconcileAsyncRun } from "../background/stale-run-reconciler.ts"; import { resolveAsyncRootResultPath } from "../background/chain-root-attachment.ts"; import { attachRootChildrenToSteps, createNestedRoute, findNestedControlResult, resolveInheritedNestedRouteFromEnv, resolveNestedAsyncDir, resolveNestedParentAddressFromEnv, snapshotNestedEventFiles, updateForegroundNestedProjection, writeNestedControlRequest, writeNestedEvent, type NestedRunResolutionScope } from "../shared/nested-events.ts"; import { resolveSubagentRunId, type ResolvedSubagentRunId } from "../background/run-id-resolver.ts"; import { formatNestedRunStatusLines } from "../shared/nested-render.ts"; import { inspectSubagentStatus } from "../background/run-status.ts"; import { applyForceTopLevelAsyncOverride } from "../background/top-level-async.ts"; import { cleanupWorktrees, createWorktrees, diffWorktrees, findWorktreeTaskCwdConflict, formatWorktreeDiffSummary, formatWorktreeTaskCwdConflict, type WorktreeSetup, } from "../shared/worktree.ts"; import { type AgentProgress, type AsyncStatus, type AcceptanceInput, type AgentContract, type ArtifactConfig, type ArtifactPaths, type ControlConfig, type ControlEvent, type Details, type ExtensionConfig, type ForegroundRunControl, type IntercomEventBus, type JsonSchemaObject, type MaxOutputConfig, type NestedRouteInfo, type NestedRunSummary, type ResolvedControlConfig, type ResolvedTurnBudget, type ResolvedToolBudget, type SingleResult, type SubagentForegroundCompleteEvent, type ToolBudgetConfig, type TurnBudgetConfig, type UsageBudgetConfig, type SubagentRunMode, type SubagentState, ASYNC_DIR, DEFAULT_ARTIFACT_CONFIG, RESULTS_DIR, SUBAGENT_ACTIONS, SUBAGENT_CONTROL_EVENT, SUBAGENT_CONTROL_INTERCOM_EVENT, SUBAGENT_FOREGROUND_COMPLETE_EVENT, checkSubagentDepth, resolveTopLevelParallelConcurrency, resolveTopLevelParallelMaxTasks, resolveChildMaxSubagentDepth, resolveCurrentMaxSubagentDepth, wrapForkTask, } from "../../shared/types.ts"; const MUTATING_MANAGEMENT_ACTIONS = new Set(["create", "update", "delete", "eject", "disable", "enable", "reset", "grant-spawn-budget", "watchdog.configure"]); interface TaskParam { agent: string; task: string; cwd?: string; count?: number; output?: string | boolean; outputMode?: "inline" | "file-only"; reads?: string[] | boolean; progress?: boolean; model?: string; skill?: string | string[] | boolean; outputSchema?: JsonSchemaObject; acceptance?: AcceptanceInput; agentContract?: AgentContract; toolBudget?: ToolBudgetConfig; } export async function runChainPath(data: ExecutionContextData, deps: ExecutorDeps): Promise> { const { params, effectiveCwd, agents, ctx, signal, runId, shareEnabled, sessionDirForIndex, sessionFileForIndex, sessionFileForTask, thinkingOverrideForTask, artifactsDir, artifactConfig, onUpdate, sessionRoot, controlConfig, contextPolicy, } = data; const onControlEvent = createForegroundControlNotifier(data, deps); const childIntercomTarget = data.intercomBridge.active ? resolveSubagentIntercomTarget : undefined; const foregroundControl = deps.state.foregroundControls.get(runId); const normalized = normalizeSkillInput(params.skill); const chainSkills = normalized === false ? [] : (normalized ?? []); const chain = wrapChainTasksForFork(params.chain as ChainStep[], contextPolicy); const currentMaxSubagentDepth = resolveCurrentMaxSubagentDepth(deps.config.maxSubagentDepth); const chainCtx = normalizeParentModel(ctx.model) || !data.parentModel ? ctx : { ...ctx, model: data.parentModel }; const chainResult = await executeChain({ chain, task: params.task, agents, ctx: chainCtx, modelScope: data.modelScope, intercomEvents: deps.pi.events, signal, runId, cwd: effectiveCwd, shareEnabled, sessionDirForIndex, sessionFileForIndex, sessionFileForTask, thinkingOverrideForTask, contextForAgent: contextPolicy.contextForAgent, artifactsDir, artifactConfig, includeProgress: params.includeProgress, clarify: params.clarify, onUpdate, onControlEvent, controlConfig, agentContract: params.agentContract, childIntercomTarget: childIntercomTarget ? (agent, index) => childIntercomTarget(runId, agent, index) : undefined, orchestratorIntercomTarget: data.intercomBridge.active ? data.intercomBridge.orchestratorTarget : undefined, foregroundControl, nestedRoute: foregroundControl?.nestedRoute, chainSkills, chainDir: params.chainDir ?? getProjectChainRunsDir(effectiveCwd), dynamicFanoutMaxItems: deps.config.chain?.dynamicFanout?.maxItems, maxSubagentDepth: currentMaxSubagentDepth, worktreeSetupHook: deps.config.worktreeSetupHook, worktreeSetupHookTimeoutMs: deps.config.worktreeSetupHookTimeoutMs, worktreeBaseDir: deps.config.worktreeBaseDir, timeoutMs: data.timeoutMs, deadlineAt: data.deadlineAt, turnBudget: data.turnBudget, onDetachedExit: (index, result) => { try { updateRememberedForegroundChild(deps.state, { runId, mode: "chain", cwd: effectiveCwd, sessionId: data.parentSessionId, index, result, events: deps.pi.events }); } finally { removeForegroundControlIfIdle(deps.state, runId); } }, onForegroundChildSettled: () => { removeForegroundControlIfIdle(deps.state, runId); }, toolBudget: data.toolBudget, usageBudget: data.usageBudget, configToolBudget: data.configToolBudget, globalConcurrencyLimit: deps.config.globalConcurrencyLimit, capabilityCeiling: data.capabilityCeiling, }); if (chainResult.requestedAsync) { if (!isAsyncAvailable()) { return { content: [{ type: "text", text: "Background mode requires upstream jiti for TypeScript execution but it could not be found. Ensure the pi-agents-flow package dependencies are installed." }], isError: true, details: { mode: "chain" as const, results: [] }, }; } const id = randomUUID(); const parentModel = data.parentModel; const asyncCtx = { pi: deps.pi, cwd: ctx.cwd, currentSessionId: deps.state.currentSessionId!, parentSessionId: ctx.sessionManager.getSessionId() ?? undefined, currentModelProvider: parentModel?.provider, currentModel: parentModel, modelScope: data.modelScope, interactive: ctx.hasUI, }; const rawAsyncChain = chainResult.requestedAsync.chain; const asyncChain = wrapChainTasksForFork(rawAsyncChain, contextPolicy); const firstAgent = firstChainAgent(rawAsyncChain); return executeAsyncChain(id, { chain: asyncChain, task: params.task, goal: resolveAsyncEventGoal(params.task, rawAsyncChain, firstAgent ? shouldForkAgent(contextPolicy, firstAgent) : false), agents, ctx: asyncCtx, availableModels: ctx.modelRegistry.getAvailable().map(toModelInfo), cwd: effectiveCwd, maxOutput: params.maxOutput, artifactsDir: artifactConfig.enabled ? artifactsDir : undefined, artifactConfig, shareEnabled, sessionRoot, chainSkills: chainResult.requestedAsync.chainSkills, sessionFilesByFlatIndex: collectChainSessionFiles(asyncChain, sessionFileForTask, deps.config.chain?.dynamicFanout?.maxItems), thinkingOverridesByFlatIndex: collectChainThinkingOverrides(asyncChain, thinkingOverrideForTask, deps.config.chain?.dynamicFanout?.maxItems), contextForAgent: contextPolicy.contextForAgent, dynamicFanoutMaxItems: deps.config.chain?.dynamicFanout?.maxItems, maxSubagentDepth: currentMaxSubagentDepth, waitToolEnabled: deps.waitToolEnabled, worktreeSetupHook: deps.config.worktreeSetupHook, worktreeSetupHookTimeoutMs: deps.config.worktreeSetupHookTimeoutMs, worktreeBaseDir: deps.config.worktreeBaseDir, controlConfig, agentContract: params.agentContract, controlIntercomTarget: data.intercomBridge.active ? data.intercomBridge.orchestratorTarget : undefined, childIntercomTarget: data.intercomBridge.active ? (agent, index) => resolveSubagentIntercomTarget(id, agent, index) : undefined, nestedRoute: data.nestedRoute, timeoutMs: data.timeoutMs, turnBudget: data.turnBudget, toolBudget: data.toolBudget, usageBudget: data.usageBudget, configToolBudget: data.configToolBudget, capabilityCeiling: data.capabilityCeiling, globalConcurrencyLimit: deps.config.globalConcurrencyLimit, }); } const rawChainDetails = chainResult.details ? { ...chainResult.details, runId, timeoutMs: data.timeoutMs } : undefined; if (foregroundControl && rawChainDetails) { updateForegroundNestedProjection(foregroundControl); attachRootChildrenToSteps(runId, rawChainDetails.results, foregroundControl.nestedChildren); rawChainDetails.totalCost = sumResultsCost(rawChainDetails.results); rawChainDetails.usageBudget = usageBudgetState(data.usageBudget, rawChainDetails.totalCost); } const chainDetails = rawChainDetails ? compactForegroundDetails(rawChainDetails) : undefined; if (chainDetails) rememberForegroundRun(deps.state, { runId, mode: "chain", cwd: effectiveCwd, sessionId: data.parentSessionId, results: chainDetails.results, checkpoint: chainDetails.checkpoint }); const intercomReceipt = chainDetails && !chainDetails.results.some((result) => result.interrupted || result.detached) ? await maybeBuildForegroundIntercomReceipt({ pi: deps.pi, intercomBridge: data.intercomBridge, runId, mode: "chain", details: chainDetails, ...(foregroundControl?.nestedChildren?.length ? { nestedChildren: foregroundControl.nestedChildren } : {}), }) : null; if (intercomReceipt) { return { ...chainResult, content: [{ type: "text", text: intercomReceipt.text }], details: intercomReceipt.details, }; } return chainDetails ? { ...chainResult, details: chainDetails } : chainResult; } interface ForegroundParallelRunInput { tasks: TaskParam[]; taskTexts: string[]; taskDescriptions: string[]; agents: AgentConfig[]; ctx: ExtensionContext; state: SubagentState; intercomEvents: IntercomEventBus; parentSessionId: string | null; capabilityCeiling?: ResolvedSubagentCapabilityCeiling; signal: AbortSignal; runId: string; sessionDirForIndex: (idx?: number) => string | undefined; sessionFileForIndex: (idx?: number) => string | undefined; sessionFileForTask: ForkSessionFileForTask; thinkingOverrideForTask: ForkThinkingOverrideForTask; shareEnabled: boolean; artifactConfig: ArtifactConfig; artifactsDir: string; outputBaseDir: string; maxOutput?: MaxOutputConfig; paramsCwd: string; progressDir: string; maxSubagentDepths: number[]; waitToolEnabled?: boolean; availableModels: ModelInfo[]; modelScope?: ModelScopeConfig; parentModel?: ParentModel; modelOverrides: (string | undefined)[]; behaviors: Array>; firstProgressIndex: number; controlConfig: ResolvedControlConfig; contextPolicy: AgentDefaultContextPolicy; onControlEvent?: (event: ControlEvent) => void; childIntercomTarget?: (agent: string, index: number) => string | undefined; orchestratorIntercomTarget?: string; foregroundControl?: SubagentState["foregroundControls"] extends Map ? T : never; concurrencyLimit: number; globalSemaphore?: Semaphore; liveResults: (SingleResult | undefined)[]; liveProgress: (AgentProgress | undefined)[]; onUpdate?: (r: AgentToolResult
) => void; worktreeSetup?: WorktreeSetup; timeoutMs?: number; deadlineAt?: number; turnBudget?: ResolvedTurnBudget; usageBudget?: UsageBudgetConfig; toolBudgets: (ResolvedToolBudget | undefined)[]; agentContract?: AgentContract; }