Files
oh-my-pi/packages/coding-agent/src/task/index.ts
T
pr-eval 424458e99d chore(secrets): dropped drive-by changes unrelated to secret placeholders
Reverted branch-side edits to spawn-policy prompts/tests, settings tab
groups, mermaid cache typing, prewalk todo gating, and packages/ai test
churn back to merge-base content; trimmed their changelog entries. These
repaired stale CI against an older main and are stale or conflicting
against current main.
2026-07-23 17:56:29 +02:00

1535 lines
58 KiB
TypeScript

/**
* Task tool - Delegate tasks to specialized agents.
*
* Discovers agent definitions from:
* - Bundled agents (shipped with omp-coding-agent)
* - ~/.omp/agent/agents/*.md (user-level)
* - .omp/agents/*.md (project-level)
*
* Supports:
* - Single agent spawn per call (parallelism = parallel task calls)
* - Batch spawning + shared context per call when `task.batch` is enabled
* - Background execution through AsyncJobManager when `async.enabled` is enabled
* - Progress tracking via JSON events
* - Session artifacts for debugging
*/
import path from "node:path";
import type { AgentTool, AgentToolResult, AgentToolUpdateCallback } from "@oh-my-pi/pi-agent-core";
import type { Usage } from "@oh-my-pi/pi-ai";
import { $env, logger, prompt } from "@oh-my-pi/pi-utils";
import type { ToolSession } from "..";
import type { Theme } from "../modes/theme/theme";
import subagentUserPromptTemplate from "../prompts/system/subagent-user-prompt.md" with { type: "text" };
import taskDescriptionTemplate from "../prompts/tools/task.md" with { type: "text" };
import taskAsyncContractTemplate from "../prompts/tools/task-async-contract.md" with { type: "text" };
import taskSummaryTemplate from "../prompts/tools/task-summary.md" with { type: "text" };
import { truncateForPrompt } from "../tools/approval";
import { isIrcEnabled } from "../tools/hub";
import { formatBytes, formatDuration } from "../tools/render-utils";
import { resolveSpawnPolicy } from "./spawn-policy";
import {
type AgentDefinition,
type AgentProgress,
canSpawnAtDepth,
getTaskSchema,
type SingleResult,
type TaskItem,
type TaskParams,
type TaskToolDetails,
type TaskToolSchemaInstance,
} from "./types";
// Import review tools for side effects (registers subagent tool handlers)
import "../tools/review";
import type { AsyncJobManager } from "../async";
import { hasResolvableTranscript } from "../internal-urls/registry-helpers";
import { AgentRegistry } from "../registry/agent-registry";
import { type DiscoveryResult, discoverAgents } from "./discovery";
import { generateTaskName } from "./name-generator";
import { AgentOutputManager } from "./output-manager";
import { mapWithConcurrencyLimitAllSettled, Semaphore } from "./parallel";
import { renderResult, renderCall as renderTaskCall } from "./render";
import { repairTaskParams } from "./repair-args";
import { resolveEffectiveSubagentPolicy, runStructuredSubagent, StructuredSubagentError } from "./structured-subagent";
function renderSubagentUserPrompt(assignment: string): string {
return prompt.render(subagentUserPromptTemplate, {
assignment: assignment.trim(),
});
}
function createUsageTotals(): Usage {
return {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
}
function addUsageTotals(target: Usage, usage: Partial<Usage>): void {
const input = usage.input ?? 0;
const output = usage.output ?? 0;
const cacheRead = usage.cacheRead ?? 0;
const cacheWrite = usage.cacheWrite ?? 0;
const totalTokens = usage.totalTokens ?? input + output + cacheRead + cacheWrite;
const cost =
usage.cost ??
({
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
total: 0,
} satisfies Usage["cost"]);
target.input += input;
target.output += output;
target.cacheRead += cacheRead;
target.cacheWrite += cacheWrite;
target.totalTokens += totalTokens;
target.cost.input += cost.input;
target.cost.output += cost.output;
target.cost.cacheRead += cost.cacheRead;
target.cost.cacheWrite += cost.cacheWrite;
target.cost.total += cost.total;
}
// Re-export types and utilities
export { loadBundledAgents as BUNDLED_AGENTS } from "./agents";
export { discoverCommands, expandCommand, getCommand } from "./commands";
export { discoverAgents, getAgent } from "./discovery";
export { AgentOutputManager } from "./output-manager";
export type {
AgentDefinition,
AgentProgress,
SingleResult,
SubagentEventPayload,
SubagentLifecyclePayload,
SubagentProgressPayload,
TaskParams,
TaskToolDetails,
} from "./types";
export {
TASK_SUBAGENT_EVENT_CHANNEL,
TASK_SUBAGENT_LIFECYCLE_CHANNEL,
TASK_SUBAGENT_PROGRESS_CHANNEL,
taskSchema,
} from "./types";
// Built-in tools whose approval tier is "read" (see tool classes' `approval`).
// An agent is read-only iff its declared tools are a non-empty subset of this set.
// Fail-safe: any unknown tool makes the agent not read-only.
export const READ_ONLY_TOOL_NAMES: ReadonlySet<string> = new Set([
"read",
"grep",
"glob",
"web_search",
"ast_grep",
"yield",
"hub",
"ask",
"todo",
"recall",
"reflect",
"retain",
"memory_edit",
"inspect_image",
"checkpoint",
"rewind",
]);
export function isReadOnlyAgent(agent: AgentDefinition): boolean {
return !!agent.tools?.length && agent.tools.every(tool => READ_ONLY_TOOL_NAMES.has(tool));
}
/**
* Preview text for a child result. Falls back to "(no output)" — annotated
* with the request count when the child actually did work, so the parent can
* tell a no-op child from one that burned requests before being cancelled.
*/
export function formatResultOutputFallback(result: Pick<SingleResult, "output" | "stderr" | "requests">): string {
const base = result.output.trim() || result.stderr.trim();
if (base) return base;
return result.requests > 0 ? `(no output) after ${result.requests} req` : "(no output)";
}
/**
* Render the tool description from a cached agent list and current settings.
*/
function renderDescription(
agents: AgentDefinition[],
isolationEnabled: boolean,
applyIsolatedChanges: boolean,
disabledAgents: string[],
batchEnabled: boolean,
asyncEnabled: boolean,
ircEnabled: boolean,
parentSpawns: string,
): string {
const spawnPolicy = resolveSpawnPolicy(parentSpawns);
const spawningDisabled = !spawnPolicy.enabled;
let filteredAgents = disabledAgents.length > 0 ? agents.filter(a => !disabledAgents.includes(a.name)) : agents;
if (spawningDisabled) {
filteredAgents = [];
} else if (spawnPolicy.allowedAgents !== null) {
const allowed = new Set(spawnPolicy.allowedAgents);
filteredAgents = filteredAgents.filter(a => allowed.has(a.name));
}
const renderedAgents = filteredAgents.map(agent => ({
name: agent.name,
description: agent.description,
readOnly: isReadOnlyAgent(agent),
blocking: agent.blocking === true,
}));
return prompt.render(taskDescriptionTemplate, {
agents: renderedAgents,
spawningDisabled,
defaultAgent: spawnPolicy.defaultAgent,
isolationEnabled,
applyIsolatedChanges,
batchEnabled,
asyncEnabled,
hasBlockingAgents: renderedAgents.some(agent => agent.blocking),
ircEnabled,
});
}
function createTaskModeError(text: string): AgentToolResult<TaskToolDetails> {
return {
content: [{ type: "text", text }],
details: { projectAgentsDir: null, results: [], totalDurationMs: 0 },
};
}
/**
* Reject legacy fields and shape/configuration combinations the current tool
* cannot accept. `outputSchema` is a first-class per-spawn field; stale
* `schema` remains an eval-only alias and is rejected.
*/
function validateShapeParams(batchEnabled: boolean, params: TaskParams): string | undefined {
if (Object.hasOwn(params, "schema")) {
return "The task tool uses `outputSchema`; rename the stale `schema` field.";
}
if (!batchEnabled) {
const disallowed = (["tasks", "context"] as const).filter(field => params[field] !== undefined);
if (disallowed.length > 0) {
return `task.batch is disabled, so the task tool does not accept ${disallowed.map(f => `\`${f}\``).join(" or ")}. Spawn one agent per call with \`task\`, or enable the task.batch setting.`;
}
}
return undefined;
}
/**
* Validate the spawn parameter contract against the wire shapes. With
* `task.batch` the model-facing shape is `{ context, tasks[] }` — `tasks`
* non-empty with per-item `task` instructions and unique names, `context`
* non-empty, no top-level `task` alongside. The flat `{ agent?, ...item }`
* form stays accepted at runtime under either setting (internal callers, stale
* transcripts). Missing `agent` values resolve against the session spawn
* policy later, in `spawnParamsFor`. Returns a problem description, or
* undefined when valid.
*/
function hasInvalidModelSelector(model: unknown): boolean {
if (model === undefined) return false;
const selectors = typeof model === "string" ? [model] : Array.isArray(model) ? model : undefined;
const materializedSelectors = selectors ? Array.from(selectors) : [];
return (
!selectors ||
materializedSelectors.length === 0 ||
materializedSelectors.some(
selector => typeof selector !== "string" || !selector.split(",").some(pattern => pattern.trim()),
)
);
}
function validateSpawnParams(params: TaskParams, batchEnabled: boolean): string | undefined {
const hasTask = typeof params.task === "string" && params.task.trim() !== "";
const tasks = params.tasks;
if (batchEnabled && tasks !== undefined) {
if (!Array.isArray(tasks) || tasks.length === 0) {
return "Missing `tasks`. Provide at least one task item ({ name?, agent?, task }).";
}
if (hasTask) {
return "Top-level `task` is not part of the batch shape. Put the work in `tasks[]` items.";
}
for (let i = 0; i < tasks.length; i++) {
const item = tasks[i];
if (!item || typeof item.task !== "string" || item.task.trim() === "") {
return `Task ${i + 1}${item?.name ? ` (\`${item.name}\`)` : ""} is missing \`task\`. Every task needs complete, self-contained instructions.`;
}
if (hasInvalidModelSelector(item.model)) {
return `Task ${i + 1}${item.name ? ` (\`${item.name}\`)` : ""} has an invalid \`model\`. Provide a non-empty selector or a non-empty array of non-empty selectors.`;
}
}
const seen = new Map<string, string>();
for (const item of tasks) {
const name = item.name?.trim();
if (!name) continue;
const key = name.toLowerCase();
const existing = seen.get(key);
if (existing !== undefined) {
return `Duplicate task name ${existing === name ? `\`${name}\`` : `\`${existing}\` / \`${name}\``}. Provided names must be unique within a call (case-insensitive).`;
}
seen.set(key, name);
}
if (typeof params.context !== "string" || params.context.trim() === "") {
return "Missing `context`. Provide the shared background for this batch — goal, constraints, and any contract the tasks share.";
}
return undefined;
}
if (!hasTask) {
return batchEnabled
? "Missing `tasks`. Provide a `tasks` array (one subagent per item) with a shared `context`."
: "Missing `task`. Provide complete, self-contained instructions for the agent.";
}
if (hasInvalidModelSelector(params.model)) {
return "Invalid `model`. Provide a non-empty selector or a non-empty array of non-empty selectors.";
}
return undefined;
}
/**
* Normalize a validated call into its spawn list: the `tasks[]` batch when
* provided, otherwise the single top-level spawn. The flat form's `isolated`
* flag is only materialized when the caller sent one — `#runSpawn`
* distinguishes an absent key from an explicit value.
*/
function resolveSpawnItems(params: TaskParams): TaskItem[] {
if (Array.isArray(params.tasks) && params.tasks.length > 0) {
return params.tasks;
}
const item: TaskItem = { name: params.name, agent: params.agent, task: params.task, model: params.model };
if ("outputSchema" in params) item.outputSchema = params.outputSchema;
if ("schemaMode" in params) item.schemaMode = params.schemaMode;
if ("isolated" in params) item.isolated = params.isolated;
return [item];
}
/**
* Per-spawn params handed to the executor path: top-level call fields with the
* item's identity substituted in. Each spawn's `agent` resolves here —
* the item's own value, else `defaultAgent` from the session spawn policy.
* `tasks` never leaks into a spawn; the shared `context` rides along
* unchanged. Keys are only materialized when present — `#runSpawn`
* distinguishes an absent `isolated` from an explicit one. The item's
* `isolated` (batch form) wins over the top-level flag (flat form).
*/
function spawnParamsFor(params: TaskParams, item: TaskItem, defaultAgent: string): TaskParams {
const spawn: TaskParams = { agent: item.agent?.trim() || defaultAgent };
if (item.name !== undefined) spawn.name = item.name;
if (item.task !== undefined) spawn.task = item.task;
if (item.model !== undefined) spawn.model = item.model;
if (params.context !== undefined) spawn.context = params.context;
if ("outputSchema" in item) spawn.outputSchema = item.outputSchema;
if ("schemaMode" in item) spawn.schemaMode = item.schemaMode;
if (item.isolated !== undefined) {
spawn.isolated = item.isolated;
} else if ("isolated" in params) {
spawn.isolated = params.isolated;
}
return spawn;
}
/** One sync-executed spawn: its item, position in the original call, and (for mixed calls) a pre-claimed agent id. */
interface SyncSpawnRef {
item: TaskItem;
index: number;
preAllocatedId?: string;
}
/** Merged view of a sync spawn set's payloads: joined text plus flattened results/usage/paths. */
interface MergedSyncPayloads {
contentParts: string[];
results: SingleResult[];
usage?: Usage;
outputPaths?: string[];
projectAgentsDir: string | null;
}
/**
* Merge per-spawn sync payloads into one result view. `index` is each spawn's
* position in the original call so batch rows keep stable ordering; a missing
* payload (cancelled before start) becomes an explanatory content line.
*/
function mergeSyncPayloads(
spawns: SyncSpawnRef[],
payloads: (AgentToolResult<TaskToolDetails> | undefined)[],
): MergedSyncPayloads {
const results: SingleResult[] = [];
const contentParts: string[] = [];
const outputPaths: string[] = [];
const usageTotals = createUsageTotals();
let hasUsage = false;
let projectAgentsDir: string | null = null;
for (let position = 0; position < spawns.length; position++) {
const payload = payloads[position];
const { item, index } = spawns[position];
if (!payload) {
contentParts.push(`Task ${item.name?.trim() || `#${index + 1}`}: cancelled before start.`);
continue;
}
projectAgentsDir ??= payload.details?.projectAgentsDir ?? null;
const text = payload.content.find(part => part.type === "text")?.text;
if (text) contentParts.push(text);
for (const result of payload.details?.results ?? []) {
results.push({ ...result, index });
if (result.usage) {
addUsageTotals(usageTotals, result.usage);
hasUsage = true;
}
if (result.outputPath) outputPaths.push(result.outputPath);
}
}
return {
contentParts,
results,
usage: hasUsage ? usageTotals : undefined,
outputPaths: outputPaths.length > 0 ? outputPaths : undefined,
projectAgentsDir,
};
}
/** Generic worker agent types; several in one call usually means a more specific type exists. */
const GENERIC_SPAWN_AGENTS: ReadonlySet<string> = new Set(["task", "sonic"]);
/**
* Advisory — never a rejection — nudging the spawner toward tailored
* specific agent types when one call resolves ≥2 items to a generic
* `task`/`sonic` worker and the spawner still holds spawn capacity
* (DepthCapacity: it currently has the `task` tool). `agentNames` are the
* per-item resolved agent types. Returns undefined when no nudge applies.
*/
export function buildSpecializationAdvisory(agentNames: string[], depthCapacity: boolean): string | undefined {
if (!depthCapacity) return undefined;
const generics = agentNames.filter(name => GENERIC_SPAWN_AGENTS.has(name));
if (generics.length < 2) return undefined;
return (
`Tip: this call spawned ${generics.length} generic \`${generics[0]}\` workers. ` +
`Check the agent list for a closer specialist type — e.g. read-only research belongs on ` +
`\`agent: "scout"\`, which runs on a faster model.`
);
}
/**
* Suggestion — never a rejection — nudging the spawner to coordinate via the
* hub when one call creates ≥2 live siblings and it still holds spawn
* capacity. Returns undefined when there is nothing to coordinate or peer
* messaging is unavailable.
*/
export function buildCoordinationAdvisory(
items: TaskItem[],
depthCapacity: boolean,
ircEnabled: boolean,
): string | undefined {
if (!depthCapacity || !ircEnabled || items.length < 2) return undefined;
return (
`Coordinate: ${items.length} siblings are running together. If their work overlaps, have them ` +
`message each other via \`hub\` (by id, or "all" to broadcast) before editing shared files — ` +
`live coordination beats a serial handoff. Check \`hub\` op:"list" to see who is doing what.`
);
}
/**
* Compose the non-blocking advisory appended to a `task` result: the
* specialization nudge (from the per-item resolved agent types), plus — only
* when some spawns keep running after this call (`willRunAsync`) — the
* coordination suggestion over those still-live spawns (`items`). Coordination
* is gated on async because a sync spawn has already finished by the time the
* call returns, so a "coordinate while they run" hint would misfire. Returns
* undefined when neither applies.
*/
export function composeSpawnAdvisory(args: {
agents: string[];
items: TaskItem[];
depthCapacity: boolean;
ircEnabled: boolean;
willRunAsync: boolean;
}): string | undefined {
return (
[
buildSpecializationAdvisory(args.agents, args.depthCapacity),
args.willRunAsync ? buildCoordinationAdvisory(args.items, args.depthCapacity, args.ircEnabled) : undefined,
]
.filter(Boolean)
.join("\n\n") || undefined
);
}
/** Sentinel for async jobs whose subagent finished with a failing result; progress is already updated. */
class TaskJobError extends Error {}
/**
* Process-level memo for create-time agent discovery, keyed by resolved cwd.
*
* `TaskTool.create` runs for every (sub)agent session in this process and the
* walk-up + plugin-registry scan in `discoverAgents` is identical for a given
* cwd, so repeat creations reuse the first scan. Execution-time discovery
* (`#runSpawn`) intentionally stays fresh. The memo also tracks the live
* `discoverAgents` binding: test spies swap that binding, which invalidates
* the memo automatically.
*/
const discoveryMemo = new Map<string, Promise<DiscoveryResult>>();
let discoveryMemoFn: typeof discoverAgents | undefined;
function discoverAgentsForCreate(cwd: string): Promise<DiscoveryResult> {
const fn = discoverAgents;
if (discoveryMemoFn !== fn) {
discoveryMemoFn = fn;
discoveryMemo.clear();
}
const key = path.resolve(cwd);
let pending = discoveryMemo.get(key);
if (!pending) {
pending = fn(cwd);
discoveryMemo.set(key, pending);
pending.catch(() => {
if (discoveryMemo.get(key) === pending) discoveryMemo.delete(key);
});
}
return pending;
}
function formatModelForApproval(model: unknown): string | undefined {
const selectors = typeof model === "string" ? [model] : Array.isArray(model) ? model : [];
const normalized = selectors.filter(
(selector): selector is string => typeof selector === "string" && !!selector.trim(),
);
return normalized.length > 0 ? truncateForPrompt(normalized.join(" → ")) : undefined;
}
// ═══════════════════════════════════════════════════════════════════════════
// Tool Class
// ═══════════════════════════════════════════════════════════════════════════
/**
* Task tool - Delegate tasks to specialized agents.
*
* Each call spawns one subagent — or, with `task.batch`, one per `tasks[]`
* item. When `async.enabled` is on, spawns run as AsyncJobManager jobs; when
* disabled, the tool blocks until every spawn finishes.
*/
export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetails, Theme> {
readonly name = "task";
readonly approval = "exec" as const;
readonly formatApprovalDetails = (args: unknown): string[] => {
const params = args as Partial<TaskParams>;
const lines: string[] = [];
if (typeof params.agent === "string") {
lines.push(`Agent: ${truncateForPrompt(params.agent)}`);
}
if (typeof params.name === "string" && params.name.trim()) {
lines.push(`Name: ${truncateForPrompt(params.name)}`);
}
const model = formatModelForApproval(params.model);
if (model) lines.push(`Model: ${model}`);
if (typeof params.task === "string") {
lines.push(`Task:\n${truncateForPrompt(params.task)}`);
}
if (typeof params.context === "string" && params.context.trim()) {
lines.push(`Context:\n${truncateForPrompt(params.context)}`);
}
const tasks = Array.isArray(params.tasks) ? params.tasks : [];
const firstTask = tasks[0];
if (firstTask) {
if (typeof firstTask.name === "string" && firstTask.name.trim()) {
lines.push(`Name: ${truncateForPrompt(firstTask.name)}`);
}
if (typeof firstTask.agent === "string" && firstTask.agent.trim()) {
lines.push(`Agent: ${truncateForPrompt(firstTask.agent)}`);
}
const itemModel = formatModelForApproval(firstTask.model);
if (itemModel) lines.push(`Model: ${itemModel}`);
if (typeof firstTask.task === "string") {
lines.push(`Task:\n${truncateForPrompt(firstTask.task)}`);
}
if (tasks.length > 1) {
lines.push(`+${tasks.length - 1} more task${tasks.length === 2 ? "" : "s"}`);
}
}
return lines;
};
readonly label = "Task";
readonly summary = "Spawn subagents to complete delegated tasks";
readonly strict = false;
readonly loadMode = "essential";
// Arktype validates model calls against the active wire schema, but the flat
// single-spawn schema carries `"+": "delete"`: a batch `{ context, tasks[] }`
// payload has those keys stripped, then fails on the now-missing `task` with
// the misleading `task must be a string (was missing)`. That preempts the
// tool's own actionable shape checks (`validateShapeParams` /
// `validateSpawnParams`), which never run. Lenient validation forwards the
// raw args to `execute()` on any arktype failure so those checks surface the
// real reason ("enable task.batch, or use the flat `task` shape"). Valid
// calls still normalize through arktype; `execute()` resolves `agent`
// defaults independently, so the success path is unchanged.
readonly lenientArgValidation = true;
readonly renderResult = renderResult;
// Suppress the streaming call preview once a (partial or final) result exists
// so the task renders as ONE block that transitions in place — not a pending
// call frame stacked above the result frame. Mirrors `taskToolRenderer`.
readonly mergeCallAndResult = true;
readonly #discoveredAgents: AgentDefinition[];
readonly #blockedAgent: string | undefined;
/**
* One semaphore per TaskTool instance (i.e. per session): bounds concurrent
* subagents across parallel `task` calls within the session. Resized in
* place from `task.maxConcurrency` before every acquire/release so a
* mid-session settings change (UI toggle, `/settings`) applies to both new
* spawns and work already parked in the semaphore queue.
*/
#spawnSemaphore: Semaphore | undefined;
get parameters(): TaskToolSchemaInstance {
const planMode = this.session.getPlanModeState?.()?.enabled === true;
const isolationEnabled = !planMode && this.session.settings.get("task.isolation.mode") !== "none";
const defaultAgent = resolveSpawnPolicy(this.session.getSessionSpawns()).defaultAgent;
return getTaskSchema({ isolationEnabled, batchEnabled: this.#isBatchEnabled(), defaultAgent });
}
renderCall(args: unknown, options: Parameters<typeof renderTaskCall>[1], theme: Theme) {
return renderTaskCall(repairTaskParams(args as TaskParams), options, theme);
}
/** Dynamic description that reflects current disabled-agent settings */
get description(): string {
const disabledAgents = this.session.settings.get("task.disabledAgents") as string[];
const planMode = this.session.getPlanModeState?.()?.enabled === true;
const isolationMode = this.session.settings.get("task.isolation.mode");
return renderDescription(
this.#discoveredAgents,
!planMode && isolationMode !== "none",
this.session.settings.get("task.isolation.apply"),
disabledAgents,
this.#isBatchEnabled(),
this.session.settings.get("async.enabled"),
isIrcEnabled(this.session.settings, this.session.taskDepth ?? 0),
this.session.getSessionSpawns() ?? "*",
);
}
private constructor(
private readonly session: ToolSession,
discoveredAgents: AgentDefinition[],
) {
this.#blockedAgent = $env.PI_BLOCKED_AGENT;
this.#discoveredAgents = discoveredAgents;
}
#isBatchEnabled(): boolean {
return this.session.settings.get("task.batch");
}
#getSpawnSemaphore(): Semaphore {
const max = this.session.settings.get("task.maxConcurrency");
if (this.#spawnSemaphore) {
this.#spawnSemaphore.resize(max);
} else {
this.#spawnSemaphore = new Semaphore(max);
}
return this.#spawnSemaphore;
}
#releaseSpawnSemaphore(): void {
this.#getSpawnSemaphore().release();
}
/**
* Resolve the shared policy before detached work exists. The resulting
* policy intentionally stays local: executor dispatch resolves again from
* normalized task params rather than smuggling internal policy over the
* task wire contract.
*/
#resolveSpawnPreflight(params: TaskParams) {
return resolveEffectiveSubagentPolicy({
session: this.session,
invocationKind: "task",
assignment: (params.task ?? "").trim(),
context: this.#isBatchEnabled() ? params.context?.trim() || undefined : undefined,
agent: params.agent,
model: params.model,
...(Object.hasOwn(params, "outputSchema") ? { outputSchema: params.outputSchema } : {}),
...(Object.hasOwn(params, "schemaMode") ? { schemaMode: params.schemaMode } : {}),
...("isolated" in params ? { isolation: { requested: params.isolated } } : {}),
blockedAgent: this.#blockedAgent,
enableLsp: (this.session.enableLsp ?? true) && this.session.settings.get("task.enableLsp"),
enableIrc: isIrcEnabled(this.session.settings, this.session.taskDepth ?? 0),
maxRuntimeMs: this.session.settings.get("task.maxRuntimeMs"),
});
}
/**
* Create a TaskTool instance with async agent discovery.
*/
static async create(session: ToolSession): Promise<TaskTool> {
const { agents } = await discoverAgentsForCreate(session.cwd);
return new TaskTool(session, agents);
}
async execute(
toolCallId: string,
rawParams: unknown,
signal?: AbortSignal,
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>,
): Promise<AgentToolResult<TaskToolDetails>> {
const params = repairTaskParams(rawParams as TaskParams);
// Schema defaults fill `agent` for model calls, but internal callers
// and stale transcripts can bypass arktype. `spawnParamsFor` resolves each
// item's agent type against the session's actual default agent.
const defaultAgent = resolveSpawnPolicy(this.session.getSessionSpawns()).defaultAgent;
const batchEnabled = this.#isBatchEnabled();
const validationError = validateShapeParams(batchEnabled, params) ?? validateSpawnParams(params, batchEnabled);
if (validationError) {
return createTaskModeError(validationError);
}
const spawnItems = resolveSpawnItems(params);
const normalizedSpawnParams = spawnItems.map(item => spawnParamsFor(params, item, defaultAgent));
const resolvedAgents = normalizedSpawnParams.map(spawn => spawn.agent ?? defaultAgent);
// Execution mode is per item: an item whose agent type declares
// `blocking: true` runs inline on this turn (the parent waits on its
// result); every other item becomes a background job when async
// execution is available.
const provisionalBlocking = resolvedAgents.map(
name => this.#discoveredAgents.find(agent => agent.name === name)?.blocking === true,
);
const asyncEnabled = this.session.settings.get("async.enabled");
const manager = asyncEnabled ? this.session.asyncJobManager : undefined;
const provisionalAsyncItems = manager ? spawnItems.filter((_, index) => !provisionalBlocking[index]) : [];
const depthCapacity = canSpawnAtDepth(
this.session.settings.get("task.maxRecursionDepth") ?? 2,
this.session.taskDepth ?? 0,
);
const ircEnabled = isIrcEnabled(this.session.settings, this.session.taskDepth ?? 0);
if (!manager || provisionalAsyncItems.length === 0) {
// Sync fallback: async execution disabled, orphaned host that never
// wired a job manager, or every item's agent type declares
// `blocking: true`. `runStructuredSubagent` performs its own shared
// preflight before reserving an id in these inline paths.
if (asyncEnabled && !this.session.asyncJobManager) {
logger.warn("task: no AsyncJobManager registered; falling back to sync execution");
}
const advisory = this.session.suppressSpawnAdvisory
? undefined
: composeSpawnAdvisory({
agents: resolvedAgents,
items: provisionalAsyncItems,
depthCapacity,
ircEnabled,
willRunAsync: false,
});
const result = await this.#executeSyncFanout(
toolCallId,
params,
spawnItems.map((item, index) => ({ item, index })),
defaultAgent,
signal,
onUpdate,
);
if (!advisory) return result;
let appended = false;
const content = result.content.map(part => {
if (!appended && part.type === "text" && typeof part.text === "string") {
appended = true;
return { ...part, text: `${part.text}\n\n${advisory}` };
}
return part;
});
if (!appended) content.push({ type: "text", text: advisory });
return { ...result, content };
}
// Async jobs are otherwise registered before their body can reach
// `runStructuredSubagent`. Resolve the shared policy first so policy
// failures remain synchronous and cannot leave a queued invalid job.
const preflights = await Promise.all(
normalizedSpawnParams.map(async spawn => {
try {
return { policy: await this.#resolveSpawnPreflight(spawn) };
} catch (error) {
return { error: error instanceof StructuredSubagentError ? error.message : String(error) };
}
}),
);
const preflightFailures = preflights
.map((preflight, index) => ("error" in preflight ? { index, error: preflight.error } : undefined))
.filter((failure): failure is { index: number; error: string } => failure !== undefined);
const renderPreflightFailures = () =>
preflightFailures
.map(({ index, error }) => {
const item = spawnItems[index]!;
return `Task ${item.name?.trim() || `#${index + 1}`} failed preflight: ${error}`;
})
.join("\n");
if (preflightFailures.length === spawnItems.length) {
return createTaskModeError(renderPreflightFailures());
}
const validIndices = preflights.flatMap((preflight, index) => (preflight.policy ? [index] : []));
const validSpawns = validIndices.map(index => ({ item: spawnItems[index]!, index }));
const itemBlocking = preflights.map(preflight => preflight.policy?.effectiveAgent.blocking === true);
const asyncItems = validIndices.filter(index => !itemBlocking[index]).map(index => spawnItems[index]!);
// Coordination only makes sense for spawns that keep running after this
// call returns (the async subset). Blocking items have already completed
// by then, so a "coordinate while they run" hint would misfire.
const advisory = this.session.suppressSpawnAdvisory
? undefined
: composeSpawnAdvisory({
agents: validIndices.map(index => resolvedAgents[index]!),
items: asyncItems,
depthCapacity,
ircEnabled,
willRunAsync: asyncItems.length > 0,
});
// Returns a fresh result (copied content array, copied text part) rather
// than mutating the caller's — task results are short-lived here, but an
// in-place edit on a shared/cached AgentToolResult would be a hidden trap.
const withAdvisory = (result: AgentToolResult<TaskToolDetails>): AgentToolResult<TaskToolDetails> => {
if (!advisory) return result;
let appended = false;
const content = result.content.map(part => {
if (!appended && part.type === "text" && typeof part.text === "string") {
appended = true;
return { ...part, text: `${part.text}\n\n${advisory}` };
}
return part;
});
if (!appended) content.push({ type: "text", text: advisory });
return { ...result, content };
};
const withPreflightFailures = (result: AgentToolResult<TaskToolDetails>): AgentToolResult<TaskToolDetails> => {
if (preflightFailures.length === 0) return result;
const failures = renderPreflightFailures();
let prepended = false;
const content = result.content.map(part => {
if (!prepended && part.type === "text" && typeof part.text === "string") {
prepended = true;
return { ...part, text: `${failures}\n\n${part.text}` };
}
return part;
});
if (!prepended) content.unshift({ type: "text", text: failures });
return { ...result, content };
};
if (asyncItems.length === 0) {
return withPreflightFailures(
withAdvisory(
await this.#executeSyncFanout(toolCallId, params, validSpawns, defaultAgent, signal, onUpdate),
),
);
}
// Async IDs are claimed before job registration, so retain the fallback
// manager on the session rather than recreating it for every call.
let outputManager = this.session.agentOutputManager;
if (!outputManager) {
outputManager = new AgentOutputManager(this.session.getArtifactsDir ?? (() => null));
this.session.agentOutputManager = outputManager;
}
const callStartedAt = Date.now();
const spawns: Array<{
agentId: string;
item: TaskItem;
index: number;
blocking: boolean;
progress: AgentProgress;
}> = [];
for (const index of validIndices) {
const item = spawnItems[index]!;
const agentType = resolvedAgents[index]!;
const preflight = preflights[index]!;
const policy = preflight.policy;
if (!policy) continue;
const agentSource = policy.agent.source;
const agentId = await outputManager.allocate(item.name?.trim() || generateTaskName());
const assignment = (item.task ?? "").trim();
spawns.push({
agentId,
item,
index,
blocking: itemBlocking[index],
progress: {
index,
id: agentId,
agent: agentType,
agentSource,
status: "pending",
task: renderSubagentUserPrompt(assignment),
assignment,
recentTools: [],
recentOutput: [],
toolCount: 0,
requests: 0,
tokens: 0,
cost: 0,
durationMs: 0,
},
});
}
const asyncSpawns = spawns.filter(spawn => !spawn.blocking);
const syncSpawns = spawns.filter(spawn => spawn.blocking);
const agentLabel = [...new Set(asyncSpawns.map(spawn => spawn.progress.agent))].join(", ");
// Aggregate state for the one tool call. Async spawns report into the
// shared progress snapshot through their jobs: the async half stays
// "running" until every job settles, then turns "failed" if any spawn
// failed. Blocking spawns run inline below and land in `results` before
// the call returns, so post-return job updates never drop them.
let settledCount = 0;
let failedCount = preflightFailures.length;
let primaryJobId = asyncSpawns[0].agentId;
const syncResults: SingleResult[] = [];
let syncUsage: Usage | undefined;
let syncOutputPaths: string[] | undefined;
let syncProjectAgentsDir: string | null = null;
const buildAsyncDetails = (): TaskToolDetails => ({
projectAgentsDir: syncProjectAgentsDir,
results: [...syncResults],
totalDurationMs: Date.now() - callStartedAt,
usage: syncUsage,
outputPaths: syncOutputPaths,
progress: spawns.map(spawn => ({ ...spawn.progress })),
async: {
state: settledCount < asyncSpawns.length ? "running" : failedCount > 0 ? "failed" : "completed",
jobId: primaryJobId,
type: "task",
},
});
const started: Array<{ agentId: string; jobId: string }> = [];
const failedSchedules: string[] = [];
for (const spawn of asyncSpawns) {
try {
const jobId = this.#registerSpawnJob({
manager,
toolCallId,
spawnParams: spawnParamsFor(params, spawn.item, defaultAgent),
agentId: spawn.agentId,
progress: spawn.progress,
ircEnabled,
buildDetails: buildAsyncDetails,
onUpdate,
onSettled: failed => {
settledCount += 1;
if (failed) failedCount += 1;
},
});
if (started.length === 0) primaryJobId = jobId;
started.push({ agentId: spawn.agentId, jobId });
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
failedSchedules.push(`${spawn.agentId}: ${message}`);
spawn.progress.status = "failed";
settledCount += 1;
failedCount += 1;
}
}
if (started.length === 0 && syncSpawns.length === 0) {
return withPreflightFailures({
content: [
{
type: "text",
text: `Failed to start background task job${failedSchedules.length === 1 ? "" : "s"}: ${failedSchedules.join("; ")}`,
},
],
details: { projectAgentsDir: null, results: [], totalDurationMs: 0 },
});
}
const scheduleFailureSummary =
failedSchedules.length > 0
? ` Failed to schedule ${failedSchedules.length} spawn${failedSchedules.length === 1 ? "" : "s"}: ${failedSchedules.join("; ")}.`
: "";
const coordinationHint = [
started.length === 1
? ircEnabled
? `DM \`${started[0].agentId}\` via \`hub\` send to coordinate while it runs; use \`hub\` only to inspect (\`jobs\`), wait, or cancel a stuck task.`
: `Use \`hub\` to inspect (\`jobs\`), wait, or cancel a stuck task.`
: ircEnabled
? `DM these ids via \`hub\` send to coordinate while they run; use \`hub\` only to inspect (\`jobs\`), wait, or cancel a stuck task.`
: `Use \`hub\` to inspect (\`jobs\`), wait, or cancel a stuck task by id.`,
taskAsyncContractTemplate.trim(),
].join("\n");
if (syncSpawns.length === 0) {
if (spawns.length === 1) {
const { agentId, jobId } = started[0];
onUpdate?.({
content: [{ type: "text", text: `Spawned agent \`${agentId}\`...` }],
details: buildAsyncDetails(),
});
return withPreflightFailures(
withAdvisory({
content: [
{
type: "text",
text: `Spawned agent \`${agentId}\` (job \`${jobId}\`). Its result auto-delivers on yield unless a settled \`hub jobs\`/\`wait\` snapshot consumes it first. ${coordinationHint}`,
},
],
details: buildAsyncDetails(),
}),
);
}
const startedListing = started.map(({ agentId, jobId }) => `- \`${agentId}\` (job \`${jobId}\`)`).join("\n");
onUpdate?.({
content: [{ type: "text", text: `Spawned ${started.length} agents...` }],
details: buildAsyncDetails(),
});
return withPreflightFailures(
withAdvisory({
content: [
{
type: "text",
text: `Spawned ${started.length} background agents using ${agentLabel}.${scheduleFailureSummary} Each result auto-delivers on yield unless a settled \`hub jobs\`/\`wait\` snapshot consumes it first.\n${startedListing}\n${coordinationHint}`,
},
],
details: buildAsyncDetails(),
}),
);
}
// Mixed call: the async jobs above already run detached; the blocking
// subset runs inline and gates the call's return — exactly what each
// agent type declares (`blocking: true` = the parent waits on it).
const syncLabel = syncSpawns.map(spawn => `\`${spawn.agentId}\``).join(", ");
onUpdate?.({
content: [
{
type: "text",
text: `Running ${syncLabel} inline; ${started.length} background agent${started.length === 1 ? "" : "s"} spawned...`,
},
],
details: buildAsyncDetails(),
});
const payloads = await this.#runSyncSpawns({
toolCallId,
params,
defaultAgent,
signal,
spawns: syncSpawns.map(spawn => ({ item: spawn.item, index: spawn.index, preAllocatedId: spawn.agentId })),
onItemProgress: onUpdate
? (index, progress) => {
const spawn = spawns.find(candidate => candidate.index === index);
if (spawn) spawn.progress = { ...progress, index };
onUpdate({
content: [{ type: "text", text: `Running ${syncLabel} inline...` }],
details: buildAsyncDetails(),
});
}
: undefined,
});
const merged = mergeSyncPayloads(
syncSpawns.map(spawn => ({ item: spawn.item, index: spawn.index })),
payloads,
);
syncResults.push(...merged.results);
syncUsage = merged.usage;
syncOutputPaths = merged.outputPaths;
syncProjectAgentsDir = merged.projectAgentsDir;
// Settle the inline spawns' progress rows from their merged results so
// post-return job updates carry final statuses, not the last snapshot.
for (let position = 0; position < syncSpawns.length; position++) {
const spawn = syncSpawns[position];
const result = merged.results.find(r => r.id === spawn.agentId);
if (result) {
spawn.progress.status = result.aborted
? "aborted"
: result.exitCode === 0 && !result.error
? "completed"
: "failed";
spawn.progress.durationMs = result.durationMs;
} else {
spawn.progress.status = payloads[position] ? "failed" : "aborted";
}
}
const spawnedSummary =
started.length > 0
? `Spawned ${started.length} background agent${started.length === 1 ? "" : "s"}.${scheduleFailureSummary} Each result auto-delivers on yield unless a settled \`hub jobs\`/\`wait\` snapshot consumes it first.\n${started.map(({ agentId, jobId }) => `- \`${agentId}\` (job \`${jobId}\`)`).join("\n")}\n${coordinationHint}`
: scheduleFailureSummary.trim();
const text = [merged.contentParts.join("\n\n"), spawnedSummary]
.filter(section => section.trim().length > 0)
.join("\n\n");
return withPreflightFailures(
withAdvisory({
content: [{ type: "text", text: text.length > 0 ? text : "No results." }],
details: buildAsyncDetails(),
}),
);
}
/**
* Register one background job that runs a single spawn to completion and
* delivers its yield text. The job body mirrors the sync path; `buildDetails`
* supplies the (possibly batch-shared) progress snapshot and `onSettled`
* feeds the caller's aggregate counters.
*/
#registerSpawnJob(options: {
manager: AsyncJobManager;
toolCallId: string;
spawnParams: TaskParams;
agentId: string;
progress: AgentProgress;
ircEnabled: boolean;
buildDetails: () => TaskToolDetails;
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>;
onSettled?: (failed: boolean) => void;
}): string {
const { manager, toolCallId, spawnParams, agentId, progress, ircEnabled, buildDetails, onUpdate, onSettled } =
options;
const buildFollowUpHint = async (aborted: boolean): Promise<string> => {
if (aborted) {
const ref = AgentRegistry.global().get(agentId);
const transcript = (await hasResolvableTranscript(agentId))
? `transcript at history://${agentId}`
: "transcript unavailable";
if (ref?.status === "idle" || ref?.status === "parked") {
const followUp = ircEnabled ? "message it via `hub` to resume; " : "";
return `\n\n${agentId} was stopped but is still resumable — ${followUp}${transcript}`;
}
return `\n\n${agentId} was aborted — ${transcript}`;
}
const followUp = ircEnabled ? "message it via `hub` to follow up; " : "";
return `\n\n${agentId} is now idle — ${followUp}transcript at history://${agentId}`;
};
return manager.register(
"task",
agentId,
async ({ signal: runSignal, reportProgress, markRunning }) => {
const startedAt = Date.now();
const semaphore = this.#getSpawnSemaphore();
let semaphoreHeld = false;
// Every release funnels through here: the flag flips before the
// release so no path — acquire-time abort, executor failure, or a
// future refactor that reorders the branches — can return a permit
// twice. Releasing a permit this job never acquired would steal one
// from a running job and let a later spawn start past
// task.maxConcurrency.
const releasePermit = () => {
if (!semaphoreHeld) return;
semaphoreHeld = false;
this.#releaseSpawnSemaphore();
};
try {
await semaphore.acquire(runSignal);
semaphoreHeld = true;
} catch {
// Fall through so an acquire-time abort goes through the same
// path as the post-acquire race below: progress + onSettled
// have to fire even when the spawn never reached the executor,
// otherwise the batch aggregate state stays "running" forever.
}
const acquiredAt = Date.now();
if (!semaphoreHeld || runSignal.aborted) {
releasePermit();
progress.status = "aborted";
onSettled?.(true);
throw new Error("Aborted before execution");
}
try {
markRunning();
progress.status = "running";
await reportProgress(
`Running background task ${agentId}...`,
buildDetails() as unknown as Record<string, unknown>,
);
const forwardSyncProgress: AgentToolUpdateCallback<TaskToolDetails> = async update => {
const nextProgress = update.details?.progress?.[0];
if (nextProgress) {
// The job body owns status and identity (id/index/agent);
// copy only the live metrics the subagent streams so the
// polling row reflects the resolved model, reasoning level,
// and running counters without reverting the "running"
// status back to the subagent's initial "pending" snapshot.
progress.resolvedModel = nextProgress.resolvedModel;
progress.resolvedModelIsFallback = nextProgress.resolvedModelIsFallback;
progress.tokens = nextProgress.tokens;
progress.requests = nextProgress.requests;
progress.contextTokens = nextProgress.contextTokens;
progress.contextWindow = nextProgress.contextWindow;
progress.cost = nextProgress.cost;
progress.toolCount = nextProgress.toolCount;
progress.currentTool = nextProgress.currentTool;
progress.lastIntent = nextProgress.lastIntent;
progress.recentTools = nextProgress.recentTools.slice();
progress.recentOutput = nextProgress.recentOutput.slice();
progress.retryState = nextProgress.retryState;
progress.retryFailure = nextProgress.retryFailure;
}
const updateText =
update.content.find(part => part.type === "text")?.text ?? `Running background task ${agentId}...`;
await reportProgress(updateText, buildDetails() as unknown as Record<string, unknown>);
};
const result = await this.#executeSync(
toolCallId,
spawnParams,
runSignal,
forwardSyncProgress,
agentId,
progress.index,
true,
{ invokedAt: startedAt, acquiredAt },
);
const finalText = result.content.find(part => part.type === "text")?.text ?? "(no output)";
const singleResult = result.details?.results[0];
// A missing result means the sync path failed at the tool level
// (results: []) — treat it as a failure, not success.
const resultFailed = !singleResult || (singleResult.aborted ?? false) || singleResult.exitCode !== 0;
progress.status = singleResult?.aborted ? "aborted" : resultFailed ? "failed" : "completed";
progress.durationMs = singleResult?.durationMs ?? Math.max(0, Date.now() - startedAt);
progress.tokens = singleResult?.tokens ?? 0;
progress.requests = singleResult?.requests ?? 0;
progress.contextTokens = singleResult?.contextTokens;
progress.contextWindow = singleResult?.contextWindow;
progress.cost = singleResult?.usage?.cost.total ?? 0;
progress.extractedToolData = singleResult?.extractedToolData;
progress.retryFailure = singleResult?.retryFailure;
progress.retryState = undefined;
if (singleResult?.resolvedModel) {
progress.resolvedModel = singleResult.resolvedModel;
progress.resolvedModelIsFallback = singleResult.resolvedModelIsFallback;
} else {
delete progress.resolvedModel;
delete progress.resolvedModelIsFallback;
}
onSettled?.(resultFailed);
const statusText = resultFailed
? `Background task ${agentId} failed.`
: `Background task ${agentId} complete.`;
await reportProgress(statusText, buildDetails() as unknown as Record<string, unknown>);
const deliveryText = `${finalText}${await buildFollowUpHint(singleResult?.aborted === true)}`;
if (resultFailed) {
// Mark the job itself failed; the failed agent stays interrogable.
throw new TaskJobError(deliveryText);
}
return deliveryText;
} catch (error) {
if (error instanceof TaskJobError) {
throw error;
}
progress.status = "failed";
progress.durationMs = Math.max(0, Date.now() - startedAt);
onSettled?.(true);
const statusText = `Background task ${agentId} failed.`;
await reportProgress(statusText, buildDetails() as unknown as Record<string, unknown>);
const message = error instanceof Error ? error.message : String(error);
const hint = AgentRegistry.global().get(agentId) ? await buildFollowUpHint(false) : "";
throw new TaskJobError(`${message}${hint}`);
} finally {
releasePermit();
}
},
{
id: agentId,
agentId,
queued: true,
ownerId: this.session.getAgentId?.() ?? undefined,
onProgress: text => {
onUpdate?.({ content: [{ type: "text", text }], details: buildDetails() });
},
},
);
}
/**
* Sync fan-out (async unavailable, or every item's agent type is
* `blocking: true`): run every spawn to completion inline and merge the
* per-spawn payloads into a single tool result. The session-scoped
* semaphore still bounds concurrency across parallel task calls.
*/
async #executeSyncFanout(
toolCallId: string,
params: TaskParams,
spawns: SyncSpawnRef[],
defaultAgent: string,
signal?: AbortSignal,
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>,
): Promise<AgentToolResult<TaskToolDetails>> {
if (spawns.length === 1) {
const spawn = spawns[0]!;
const semaphore = this.#getSpawnSemaphore();
const invokedAt = Date.now();
await semaphore.acquire(signal);
const acquiredAt = Date.now();
try {
return await this.#executeSync(
toolCallId,
spawnParamsFor(params, spawn.item, defaultAgent),
signal,
onUpdate,
spawn.preAllocatedId,
spawn.index,
false,
{ invokedAt, acquiredAt },
);
} finally {
this.#releaseSpawnSemaphore();
}
}
const startTime = Date.now();
const latestProgress = new Map<number, AgentProgress>();
const emitCombined = () => {
onUpdate?.({
content: [{ type: "text", text: `Running ${spawns.length} agents...` }],
details: {
projectAgentsDir: null,
results: [],
totalDurationMs: Date.now() - startTime,
progress: Array.from(latestProgress.entries())
.sort((a, b) => a[0] - b[0])
.map(([, progress]) => progress),
},
});
};
const payloads = await this.#runSyncSpawns({
toolCallId,
params,
defaultAgent,
signal,
spawns,
onItemProgress: onUpdate
? (index, progress) => {
latestProgress.set(index, { ...progress, index });
emitCombined();
}
: undefined,
});
const merged = mergeSyncPayloads(spawns, payloads);
return {
content: [{ type: "text", text: merged.contentParts.join("\n\n") }],
details: {
projectAgentsDir: merged.projectAgentsDir,
results: merged.results,
totalDurationMs: Date.now() - startTime,
usage: merged.usage,
outputPaths: merged.outputPaths,
},
};
}
/**
* Run a set of spawns to completion inline, bounded by the session spawn
* semaphore. `preAllocatedId` reuses an id claimed up front (mixed calls);
* `index` is each item's position in the original call so progress rows and
* merged results keep stable ordering. Per-item progress snapshots flow
* through `onItemProgress`. Returns per-spawn payloads in input order;
* `undefined` marks a spawn cancelled before it started.
*/
async #runSyncSpawns(args: {
toolCallId: string;
params: TaskParams;
defaultAgent: string;
spawns: SyncSpawnRef[];
signal?: AbortSignal;
onItemProgress?: (index: number, progress: AgentProgress) => void;
}): Promise<(AgentToolResult<TaskToolDetails> | undefined)[]> {
const { toolCallId, params, defaultAgent, spawns, signal, onItemProgress } = args;
const semaphore = this.#getSpawnSemaphore();
const { results } = await mapWithConcurrencyLimitAllSettled(
spawns,
spawns.length,
async (spawn, _position, workerSignal) => {
const invokedAt = Date.now();
let semaphoreHeld = false;
try {
await semaphore.acquire(workerSignal);
semaphoreHeld = true;
} catch (error) {
if (workerSignal.aborted) return undefined;
throw error;
}
const acquiredAt = Date.now();
try {
const itemOnUpdate: AgentToolUpdateCallback<TaskToolDetails> | undefined = onItemProgress
? update => {
const progress = update.details?.progress?.[0];
if (progress) onItemProgress(spawn.index, progress);
}
: undefined;
return await this.#executeSync(
toolCallId,
spawnParamsFor(params, spawn.item, defaultAgent),
workerSignal,
itemOnUpdate,
spawn.preAllocatedId,
spawn.index,
false,
{ invokedAt, acquiredAt },
);
} finally {
if (semaphoreHeld) this.#releaseSpawnSemaphore();
}
},
signal,
);
return results.map((settled, position) => {
if (!settled) return undefined;
if (settled.status === "fulfilled") return settled.value;
const message = settled.reason instanceof Error ? settled.reason.message : String(settled.reason);
const item = spawns[position].item;
return {
content: [
{
type: "text",
text: `Task ${item.name?.trim() || `#${spawns[position].index + 1}`} failed: ${message}`,
},
],
details: { projectAgentsDir: null, results: [], totalDurationMs: 0 },
};
});
}
/**
* Synchronous execution of one spawn. Used as the body of every
* async job and directly by the sync fallback (no job manager / blocking
* agent) and by in-process callers that need the result inline (e.g. the
* commit flow's analyze_files tool).
*/
async #executeSync(
toolCallId: string,
params: TaskParams,
signal?: AbortSignal,
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>,
preAllocatedId?: string,
spawnIndex = 0,
detached = false,
launchTiming?: { invokedAt: number; acquiredAt: number },
): Promise<AgentToolResult<TaskToolDetails>> {
return this.#runSpawn(toolCallId, params, signal, onUpdate, preAllocatedId, spawnIndex, detached, launchTiming);
}
/** Spawn a fresh subagent and run it to completion. */
async #runSpawn(
toolCallId: string,
params: TaskParams,
signal?: AbortSignal,
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>,
preAllocatedId?: string,
spawnIndex = 0,
detached = false,
launchTiming?: { invokedAt: number; acquiredAt: number },
): Promise<AgentToolResult<TaskToolDetails>> {
const startTime = Date.now();
const assignment = (params.task ?? "").trim();
const context = this.#isBatchEnabled() ? params.context?.trim() || undefined : undefined;
let latestProgress: AgentProgress | undefined;
try {
const execution = await runStructuredSubagent({
session: this.session,
invocationKind: "task",
assignment,
context,
agent: params.agent,
model: params.model,
...(Object.hasOwn(params, "outputSchema") ? { outputSchema: params.outputSchema } : {}),
...(Object.hasOwn(params, "schemaMode") ? { schemaMode: params.schemaMode } : {}),
identity: { id: preAllocatedId, label: params.name },
index: spawnIndex,
parentToolCallId: toolCallId,
detached,
invokedAt: launchTiming?.invokedAt,
acquiredAt: launchTiming?.acquiredAt,
...("isolated" in params ? { isolation: { requested: params.isolated } } : {}),
blockedAgent: this.#blockedAgent,
enableLsp: (this.session.enableLsp ?? true) && this.session.settings.get("task.enableLsp"),
enableIrc: isIrcEnabled(this.session.settings, this.session.taskDepth ?? 0),
maxRuntimeMs: this.session.settings.get("task.maxRuntimeMs"),
signal,
onProgress: progress => {
latestProgress = { ...progress, recentTools: progress.recentTools.slice() };
onUpdate?.({
content: [{ type: "text", text: `Running agent ${progress.id}...` }],
details: {
projectAgentsDir: null,
results: [],
totalDurationMs: Date.now() - startTime,
progress: [latestProgress],
},
});
},
});
return this.#buildResultPayload(
execution.result,
execution.policy.discovery.projectAgentsDir,
Date.now() - startTime,
execution.mergeSummary,
);
} catch (error) {
const message = error instanceof StructuredSubagentError ? error.message : String(error);
return {
content: [{ type: "text", text: `Task execution failed: ${message}` }],
details: {
projectAgentsDir: null,
results: [],
totalDurationMs: Date.now() - startTime,
...(latestProgress ? { progress: [latestProgress] } : {}),
},
};
}
}
/** Build the tool result (summary text + details) for a settled run. */
#buildResultPayload(
result: SingleResult,
projectAgentsDir: string | null,
totalDurationMs: number,
mergeSummary: string,
): AgentToolResult<TaskToolDetails> {
const status = result.aborted
? "cancelled"
: result.exitCode === 0 && result.error
? "merge failed"
: result.exitCode === 0
? "completed"
: `failed (exit ${result.exitCode})`;
const output = formatResultOutputFallback(result);
const outputCharCount = result.outputMeta?.charCount ?? output.length;
const fullOutputThreshold = 5000;
let preview = output;
let truncated = false;
if (outputCharCount > fullOutputThreshold) {
const slice = output.slice(0, fullOutputThreshold);
const lastNewline = slice.lastIndexOf("\n");
preview = lastNewline >= 0 ? slice.slice(0, lastNewline) : slice;
truncated = true;
}
// A stopped-but-adopted agent (soft-budget stop) stays messageable; tell
// the parent so it can resume via irc instead of redoing the work.
const refStatus = AgentRegistry.global().get(result.id)?.status;
const resumable = result.aborted && (refStatus === "idle" || refStatus === "parked");
const summary = prompt.render(taskSummaryTemplate, {
agentName: result.agent,
id: result.id,
status,
duration: formatDuration(totalDurationMs),
abortReason: result.aborted ? result.abortReason : undefined,
resumable,
preview,
truncated,
meta: result.outputMeta
? {
lineCount: result.outputMeta.lineCount,
charSize: formatBytes(result.outputMeta.charCount),
}
: undefined,
mergeSummary,
});
return {
content: [{ type: "text", text: summary }],
details: {
projectAgentsDir,
results: [result],
totalDurationMs,
usage: result.usage,
outputPaths: result.outputPath ? [result.outputPath] : undefined,
},
};
}
}