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 { buildChainWorktreeTaskCwdError, buildParallelModeError, buildParallelWorktreeTaskCwdError, resolveSingleRunOutputBaseDir } from "./executor-path-parallel-helpers.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 function runAsyncPath(data: ExecutionContextData, deps: ExecutorDeps): AgentToolResult
| null { const { params, effectiveCwd, agents, ctx, shareEnabled, sessionRoot, sessionFileForIndex, sessionFileForTask, thinkingOverrideForTask, artifactConfig, artifactsDir, effectiveAsync, controlConfig, intercomBridge, nestedRoute, contextPolicy, } = data; const hasChain = (params.chain?.length ?? 0) > 0; const hasTasks = (params.tasks?.length ?? 0) > 0; const hasSingle = !hasChain && !hasTasks && Boolean(params.agent); if (!effectiveAsync) return null; if (hasChain && params.chain) { const chainWorktreeTaskCwdError = buildChainWorktreeTaskCwdError(params.chain as ChainStep[], effectiveCwd); if (chainWorktreeTaskCwdError) { return { content: [{ type: "text", text: chainWorktreeTaskCwdError }], isError: true, details: { mode: "chain" as const, results: [] }, }; } } if (hasTasks && params.tasks) { const maxParallelTasks = resolveTopLevelParallelMaxTasks(deps.config.parallel?.maxTasks); if (params.tasks.length > maxParallelTasks) { return buildParallelModeError(`Max ${maxParallelTasks} tasks`); } if (params.worktree) { const worktreeTaskCwdError = buildParallelWorktreeTaskCwdError(params.tasks, effectiveCwd); if (worktreeTaskCwdError) return buildParallelModeError(worktreeTaskCwdError); } } if (!isAsyncAvailable()) { return { content: [{ type: "text", text: "Async 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: "single" 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 availableModels: ModelInfo[] = ctx.modelRegistry.getAvailable().map(toModelInfo); const currentMaxSubagentDepth = resolveCurrentMaxSubagentDepth(deps.config.maxSubagentDepth); const currentProvider = parentModel?.provider; const controlIntercomTarget = intercomBridge.active ? intercomBridge.orchestratorTarget : undefined; const childIntercomTarget = intercomBridge.active ? (agent: string, index: number) => resolveSubagentIntercomTarget(id, agent, index) : undefined; if (hasTasks && params.tasks) { const skillOverrides = params.tasks.map((task) => normalizeSkillInput(task.skill)); const parallelTasks = params.tasks.map((task, index) => ({ agent: task.agent, task: shouldForkAgent(contextPolicy, task.agent) ? wrapForkTask(task.task) : task.task, cwd: task.cwd, ...(task.model !== undefined ? { model: task.model } : {}), ...(skillOverrides[index] !== undefined ? { skill: skillOverrides[index] } : {}), ...(task.output !== undefined && task.output !== true ? { output: task.output } : {}), ...(task.outputMode !== undefined ? { outputMode: task.outputMode } : {}), ...(task.reads !== undefined && task.reads !== true ? { reads: task.reads } : {}), ...(task.progress !== undefined ? { progress: task.progress } : {}), ...(task.toolBudget !== undefined ? { toolBudget: task.toolBudget } : {}), ...(task.outputSchema !== undefined ? { outputSchema: task.outputSchema } : {}), ...(task.agentContract !== undefined ? { agentContract: task.agentContract } : {}), ...(task.acceptance !== undefined ? { acceptance: task.acceptance } : {}), })); return executeAsyncChain(id, { chain: [{ parallel: parallelTasks, concurrency: resolveTopLevelParallelConcurrency(params.concurrency, deps.config.parallel?.concurrency), worktree: params.worktree, }], resultMode: "parallel", goal: params.tasks[0]?.task ?? "", agents, ctx: asyncCtx, availableModels, cwd: effectiveCwd, maxOutput: params.maxOutput, artifactsDir: artifactConfig.enabled ? artifactsDir : undefined, artifactConfig, shareEnabled, sessionRoot, chainSkills: [], sessionFilesByFlatIndex: params.tasks.map((task, index) => sessionFileForTask(task.agent, index, task.model)), thinkingOverridesByFlatIndex: params.tasks.map((task, index) => thinkingOverrideForTask(task.agent, index, task.model)), contextForAgent: contextPolicy.contextForAgent, maxSubagentDepth: currentMaxSubagentDepth, waitToolEnabled: deps.waitToolEnabled, worktreeSetupHook: deps.config.worktreeSetupHook, worktreeSetupHookTimeoutMs: deps.config.worktreeSetupHookTimeoutMs, worktreeBaseDir: deps.config.worktreeBaseDir, controlConfig, agentContract: params.agentContract, controlIntercomTarget, childIntercomTarget, nestedRoute, timeoutMs: data.timeoutMs, turnBudget: data.turnBudget, toolBudget: data.toolBudget, usageBudget: data.usageBudget, configToolBudget: data.configToolBudget, capabilityCeiling: data.capabilityCeiling, globalConcurrencyLimit: deps.config.globalConcurrencyLimit, }); } if (hasChain && params.chain) { const normalized = normalizeSkillInput(params.skill); const chainSkills = normalized === false ? [] : (normalized ?? []); const rawChain = params.chain as ChainStep[]; const chain = wrapChainTasksForFork(rawChain, contextPolicy); return executeAsyncChain(id, { chain, task: params.task, goal: resolveAsyncEventGoal(params.task, rawChain), agents, ctx: asyncCtx, availableModels, cwd: effectiveCwd, maxOutput: params.maxOutput, artifactsDir: artifactConfig.enabled ? artifactsDir : undefined, artifactConfig, shareEnabled, sessionRoot, chainSkills, sessionFilesByFlatIndex: collectChainSessionFiles(chain, sessionFileForTask, deps.config.chain?.dynamicFanout?.maxItems), thinkingOverridesByFlatIndex: collectChainThinkingOverrides(chain, 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, childIntercomTarget, nestedRoute, timeoutMs: data.timeoutMs, turnBudget: data.turnBudget, toolBudget: data.toolBudget, usageBudget: data.usageBudget, configToolBudget: data.configToolBudget, capabilityCeiling: data.capabilityCeiling, globalConcurrencyLimit: deps.config.globalConcurrencyLimit, }); } if (hasSingle) { const a = agents.find((x) => x.name === params.agent); if (!a) { return { content: [{ type: "text", text: `Unknown agent: ${params.agent}` }], isError: true, details: { mode: "single" as const, results: [] }, }; } const rawOutput = params.output !== undefined ? params.output : a.output; const effectiveOutput = normalizeSingleOutputOverride(rawOutput, a.output); const effectiveOutputMode = params.outputMode ?? "inline"; const normalizedSkills = normalizeSkillInput(params.skill); const skills = normalizedSkills === false ? [] : normalizedSkills; const maxSubagentDepth = resolveChildMaxSubagentDepth(currentMaxSubagentDepth, a.maxSubagentDepth); const modelOverride = resolveEffectiveSubagentModel(params.model as string | undefined, a.model, parentModel, availableModels, currentProvider, { scope: data.modelScope }); return executeAsyncSingle(id, { agent: params.agent!, task: shouldForkAgent(contextPolicy, params.agent!) ? wrapForkTask(params.task ?? "") : (params.task ?? ""), goal: params.task ?? "", agentConfig: a, ctx: asyncCtx, availableModels, cwd: effectiveCwd, maxOutput: params.maxOutput, artifactsDir: artifactConfig.enabled ? artifactsDir : undefined, artifactConfig, shareEnabled, sessionRoot, sessionFile: sessionFileForTask(params.agent!, 0, modelOverride), context: contextPolicy.contextForAgent(params.agent!), skills, output: effectiveOutput, outputMode: effectiveOutputMode, outputBaseDir: resolveSingleRunOutputBaseDir(deps, artifactsDir, id), modelOverride, thinkingOverride: thinkingOverrideForTask(params.agent!, 0, modelOverride), maxSubagentDepth, waitToolEnabled: deps.waitToolEnabled, worktreeSetupHook: deps.config.worktreeSetupHook, worktreeSetupHookTimeoutMs: deps.config.worktreeSetupHookTimeoutMs, worktreeBaseDir: deps.config.worktreeBaseDir, controlConfig, controlIntercomTarget, childIntercomTarget: childIntercomTarget ? (agent, index) => childIntercomTarget(agent, index) : undefined, nestedRoute, agentContract: params.agentContract, structuredOutputSchema: params.outputSchema, acceptance: params.acceptance, timeoutMs: data.timeoutMs, turnBudget: data.turnBudget, toolBudget: data.toolBudget, usageBudget: data.usageBudget, configToolBudget: data.configToolBudget, capabilityCeiling: data.capabilityCeiling, }); } return null; }