557c79f72c
- Added `parentToolCallId` and `index` fields to subagent session records and propagated them from task lifecycle and progress events. - Reworked subagent ordering to sort by parent-group order, spawn index, and stable creation order so out-of-order updates no longer reshuffled active HUD rows. - Updated task execution to pass spawn indices through sync and async paths, switched the `tool.task` icon to Octicons tasklist, and added a registry test for out-of-order progress ordering.
1448 lines
52 KiB
TypeScript
1448 lines
52 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 * as fs from "node:fs/promises";
|
|
import * as os from "node:os";
|
|
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, Snowflake } from "@oh-my-pi/pi-utils";
|
|
import type { ToolSession } from "..";
|
|
import { resolveAgentModelPatterns } from "../config/model-resolver";
|
|
import { MCPManager } from "../mcp/manager";
|
|
import type { Theme } from "../modes/theme/theme";
|
|
import planModeSubagentPrompt from "../prompts/system/plan-mode-subagent.md" with { type: "text" };
|
|
import subagentUserPromptTemplate from "../prompts/system/subagent-user-prompt.md" with { type: "text" };
|
|
import taskDescriptionTemplate from "../prompts/tools/task.md" with { type: "text" };
|
|
import taskSummaryTemplate from "../prompts/tools/task-summary.md" with { type: "text" };
|
|
import { truncateForPrompt } from "../tools/approval";
|
|
import { isIrcEnabled } from "../tools/irc";
|
|
import { formatBytes, formatDuration } from "../tools/render-utils";
|
|
import {
|
|
type AgentDefinition,
|
|
type AgentProgress,
|
|
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 type { LocalProtocolOptions } from "../internal-urls";
|
|
import { loadOverallPlanReference } from "../plan-mode/plan-handoff";
|
|
import { AgentRegistry } from "../registry/agent-registry";
|
|
import { generateCommitMessage } from "../utils/commit-message-generator";
|
|
import * as git from "../utils/git";
|
|
import { type DiscoveryResult, discoverAgents, getAgent } from "./discovery";
|
|
import { runSubprocess } from "./executor";
|
|
import { generateTaskName } from "./name-generator";
|
|
import { AgentOutputManager } from "./output-manager";
|
|
import { mapWithConcurrencyLimit, Semaphore } from "./parallel";
|
|
import { renderResult, renderCall as renderTaskCall } from "./render";
|
|
import { repairTaskParams } from "./repair-args";
|
|
import {
|
|
applyNestedPatches,
|
|
captureBaseline,
|
|
captureDeltaPatch,
|
|
cleanupIsolation,
|
|
cleanupTaskBranches,
|
|
commitToBranch,
|
|
ensureIsolation,
|
|
getRepoRoot,
|
|
type IsolationHandle,
|
|
mergeTaskBranches,
|
|
parseIsolationMode,
|
|
type WorktreeBaseline,
|
|
} from "./worktree";
|
|
|
|
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",
|
|
"search",
|
|
"find",
|
|
"web_search",
|
|
"ast_grep",
|
|
"yield",
|
|
"irc",
|
|
"ask",
|
|
"job",
|
|
"todo",
|
|
"recall",
|
|
"reflect",
|
|
"retain",
|
|
"memory_edit",
|
|
"render_mermaid",
|
|
"inspect_image",
|
|
"checkpoint",
|
|
"rewind",
|
|
"resolve",
|
|
"report_finding",
|
|
"search_tool_bm25",
|
|
]);
|
|
|
|
const PLAN_MODE_AGENT_TOOL_ALLOWLIST: ReadonlySet<string> = new Set(["ast_grep", "report_finding"]);
|
|
|
|
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[],
|
|
maxConcurrency: number,
|
|
isolationEnabled: boolean,
|
|
disabledAgents: string[],
|
|
batchEnabled: boolean,
|
|
asyncEnabled: boolean,
|
|
ircEnabled: boolean,
|
|
parentSpawns: string,
|
|
): string {
|
|
const spawningDisabled = parentSpawns === "";
|
|
let filteredAgents = disabledAgents.length > 0 ? agents.filter(a => !disabledAgents.includes(a.name)) : agents;
|
|
if (spawningDisabled) {
|
|
filteredAgents = [];
|
|
} else if (parentSpawns !== "*") {
|
|
const allowed = new Set(
|
|
parentSpawns
|
|
.split(",")
|
|
.map(s => s.trim())
|
|
.filter(Boolean),
|
|
);
|
|
filteredAgents = filteredAgents.filter(a => allowed.has(a.name));
|
|
}
|
|
const renderedAgents = filteredAgents.map(agent => ({
|
|
name: agent.name,
|
|
description: agent.description,
|
|
readOnly: isReadOnlyAgent(agent),
|
|
}));
|
|
return prompt.render(taskDescriptionTemplate, {
|
|
agents: renderedAgents,
|
|
spawningDisabled,
|
|
MAX_CONCURRENCY: maxConcurrency,
|
|
isolationEnabled,
|
|
batchEnabled,
|
|
asyncEnabled,
|
|
ircEnabled,
|
|
});
|
|
}
|
|
|
|
function createTaskModeError(text: string): AgentToolResult<TaskToolDetails> {
|
|
return {
|
|
content: [{ type: "text", text }],
|
|
details: { projectAgentsDir: null, results: [], totalDurationMs: 0 },
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Reject fields the current configuration does not accept. `schema` is never
|
|
* accepted (structured output comes from the agent definition's `output`
|
|
* frontmatter, the inherited session schema, or an eval-workflow
|
|
* `agent(..., schema)` call); `tasks`/`context` require `task.batch`.
|
|
*/
|
|
function validateShapeParams(batchEnabled: boolean, params: TaskParams): string | undefined {
|
|
if ((params as Record<string, unknown>).schema !== undefined) {
|
|
return "The task tool does not accept `schema`. Rely on the selected agent definition's `output` schema or the inherited session schema; workflows needing ad-hoc structured output use eval `agent(prompt, schema)`.";
|
|
}
|
|
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 \`assignment\`, or enable the task.batch setting.`;
|
|
}
|
|
}
|
|
return undefined;
|
|
}
|
|
|
|
/**
|
|
* Validate the spawn parameter contract against the wire shapes. `agent` is
|
|
* always required. With `task.batch` the model-facing shape is
|
|
* `{ agent, context, tasks[] }` — `tasks` non-empty with per-item assignments
|
|
* and unique ids, `context` non-empty, no top-level `assignment` alongside.
|
|
* The flat `{ agent, ...item }` form stays accepted at runtime under either
|
|
* setting (internal callers, stale transcripts). Returns a problem
|
|
* description, or undefined when valid.
|
|
*/
|
|
function validateSpawnParams(params: TaskParams, batchEnabled: boolean): string | undefined {
|
|
const agent = typeof params.agent === "string" ? params.agent.trim() : "";
|
|
if (!agent) {
|
|
return "Missing `agent`. Provide an agent type to spawn.";
|
|
}
|
|
const hasAssignment = typeof params.assignment === "string" && params.assignment.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 ({ id?, description?, assignment }).";
|
|
}
|
|
if (hasAssignment) {
|
|
return "Top-level `assignment` 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.assignment !== "string" || item.assignment.trim() === "") {
|
|
return `Task ${i + 1}${item?.id ? ` (\`${item.id}\`)` : ""} is missing \`assignment\`. Every task needs complete, self-contained instructions.`;
|
|
}
|
|
}
|
|
const seen = new Map<string, string>();
|
|
for (const item of tasks) {
|
|
const id = item.id?.trim();
|
|
if (!id) continue;
|
|
const key = id.toLowerCase();
|
|
const existing = seen.get(key);
|
|
if (existing !== undefined) {
|
|
return `Duplicate task id ${existing === id ? `\`${id}\`` : `\`${existing}\` / \`${id}\``}. Provided ids must be unique within a call (case-insensitive).`;
|
|
}
|
|
seen.set(key, id);
|
|
}
|
|
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 (!hasAssignment) {
|
|
return batchEnabled
|
|
? "Missing `tasks`. Provide a `tasks` array (one subagent per item) with a shared `context`."
|
|
: "Missing `assignment`. Provide complete, self-contained instructions for the agent.";
|
|
}
|
|
return undefined;
|
|
}
|
|
|
|
/**
|
|
* Normalize a validated call into its spawn list: the `tasks[]` batch when
|
|
* provided, otherwise the single top-level spawn.
|
|
*/
|
|
function resolveSpawnItems(params: TaskParams): TaskItem[] {
|
|
if (Array.isArray(params.tasks) && params.tasks.length > 0) {
|
|
return params.tasks;
|
|
}
|
|
return [{ id: params.id, description: params.description, assignment: params.assignment }];
|
|
}
|
|
|
|
/**
|
|
* Per-spawn params handed to the executor path: top-level call fields with the
|
|
* item's identity substituted in. `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): TaskParams {
|
|
const spawn: TaskParams = { agent: params.agent };
|
|
if (item.id !== undefined) spawn.id = item.id;
|
|
if (item.description !== undefined) spawn.description = item.description;
|
|
if (item.assignment !== undefined) spawn.assignment = item.assignment;
|
|
if (params.context !== undefined) spawn.context = params.context;
|
|
if (item.isolated !== undefined) {
|
|
spawn.isolated = item.isolated;
|
|
} else if ("isolated" in params) {
|
|
spawn.isolated = params.isolated;
|
|
}
|
|
return spawn;
|
|
}
|
|
|
|
/** 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;
|
|
}
|
|
|
|
// ═══════════════════════════════════════════════════════════════════════════
|
|
// 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.id === "string" && params.id.trim()) {
|
|
lines.push(`Task: ${truncateForPrompt(params.id)}`);
|
|
}
|
|
if (typeof params.assignment === "string") {
|
|
lines.push(`Assignment:\n${truncateForPrompt(params.assignment)}`);
|
|
}
|
|
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.id === "string" && firstTask.id.trim()) {
|
|
lines.push(`Task: ${truncateForPrompt(firstTask.id)}`);
|
|
}
|
|
if (typeof firstTask.assignment === "string") {
|
|
lines.push(`Assignment:\n${truncateForPrompt(firstTask.assignment)}`);
|
|
}
|
|
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 = true;
|
|
readonly loadMode = "discoverable";
|
|
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. Sized from
|
|
* `task.maxConcurrency` at first use; later setting changes do not resize it.
|
|
*/
|
|
#spawnSemaphore: Semaphore | undefined;
|
|
|
|
get parameters(): TaskToolSchemaInstance {
|
|
const isolationEnabled = this.session.settings.get("task.isolation.mode") !== "none";
|
|
return getTaskSchema({ isolationEnabled, batchEnabled: this.#isBatchEnabled() });
|
|
}
|
|
|
|
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 maxConcurrency = this.session.settings.get("task.maxConcurrency");
|
|
const isolationMode = this.session.settings.get("task.isolation.mode");
|
|
return renderDescription(
|
|
this.#discoveredAgents,
|
|
maxConcurrency,
|
|
isolationMode !== "none",
|
|
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 {
|
|
this.#spawnSemaphore ??= new Semaphore(this.session.settings.get("task.maxConcurrency"));
|
|
return this.#spawnSemaphore;
|
|
}
|
|
|
|
/**
|
|
* 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);
|
|
const batchEnabled = this.#isBatchEnabled();
|
|
const validationError = validateShapeParams(batchEnabled, params) ?? validateSpawnParams(params, batchEnabled);
|
|
if (validationError) {
|
|
return createTaskModeError(validationError);
|
|
}
|
|
|
|
const spawnItems = resolveSpawnItems(params);
|
|
const selectedAgent = this.#discoveredAgents.find(agent => agent.name === params.agent);
|
|
const asyncEnabled = this.session.settings.get("async.enabled");
|
|
const manager = asyncEnabled ? this.session.asyncJobManager : undefined;
|
|
if (!asyncEnabled || !manager || selectedAgent?.blocking === true) {
|
|
// Sync fallback: async execution disabled, orphaned host that never
|
|
// wired a job manager, or an agent definition that declares
|
|
// `blocking: true`. The session-scoped semaphore still bounds fan-out
|
|
// across parallel task calls.
|
|
if (asyncEnabled && !manager) {
|
|
logger.warn("task: no AsyncJobManager registered; falling back to sync execution");
|
|
}
|
|
return this.#executeSyncFanout(toolCallId, params, spawnItems, signal, onUpdate);
|
|
}
|
|
|
|
// Resolve agent ids up front so the immediate result can name them.
|
|
const outputManager =
|
|
this.session.agentOutputManager ?? new AgentOutputManager(this.session.getArtifactsDir ?? (() => null));
|
|
const agentLabel = params.agent ?? "task";
|
|
const agentSource = selectedAgent?.source ?? "bundled";
|
|
const spawns: Array<{ agentId: string; item: TaskItem; progress: AgentProgress }> = [];
|
|
for (let index = 0; index < spawnItems.length; index++) {
|
|
const item = spawnItems[index];
|
|
const agentId = await outputManager.allocate(item.id?.trim() || generateTaskName());
|
|
const assignment = (item.assignment ?? "").trim();
|
|
spawns.push({
|
|
agentId,
|
|
item,
|
|
progress: {
|
|
index,
|
|
id: agentId,
|
|
agent: agentLabel,
|
|
agentSource,
|
|
status: "pending",
|
|
task: renderSubagentUserPrompt(assignment),
|
|
assignment,
|
|
description: item.description,
|
|
recentTools: [],
|
|
recentOutput: [],
|
|
toolCount: 0,
|
|
requests: 0,
|
|
tokens: 0,
|
|
cost: 0,
|
|
durationMs: 0,
|
|
},
|
|
});
|
|
}
|
|
|
|
// Aggregate async state for the one tool call: every spawn's job reports
|
|
// into the shared progress snapshot; the call stays "running" until all
|
|
// jobs settle, then turns "failed" if any spawn failed. The single-spawn
|
|
// case passes the job's own suggestion through (pre-batch behavior).
|
|
const single = spawns.length === 1;
|
|
let settledCount = 0;
|
|
let failedCount = 0;
|
|
let primaryJobId = spawns[0].agentId;
|
|
const buildAsyncDetails = (state: "running" | "completed" | "failed", jobId: string): TaskToolDetails => ({
|
|
projectAgentsDir: null,
|
|
results: [],
|
|
totalDurationMs: 0,
|
|
progress: spawns.map(spawn => ({ ...spawn.progress })),
|
|
async: {
|
|
state: single ? state : settledCount < spawns.length ? "running" : failedCount > 0 ? "failed" : "completed",
|
|
jobId: single ? jobId : primaryJobId,
|
|
type: "task",
|
|
},
|
|
});
|
|
|
|
const ircEnabled = isIrcEnabled(this.session.settings, this.session.taskDepth ?? 0);
|
|
const started: Array<{ agentId: string; jobId: string; description?: string }> = [];
|
|
const failedSchedules: string[] = [];
|
|
for (const spawn of spawns) {
|
|
try {
|
|
const jobId = this.#registerSpawnJob({
|
|
manager,
|
|
toolCallId,
|
|
spawnParams: spawnParamsFor(params, spawn.item),
|
|
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, description: spawn.item.description });
|
|
} 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) {
|
|
return {
|
|
content: [
|
|
{
|
|
type: "text",
|
|
text: `Failed to start background task job${single ? "" : "s"}: ${failedSchedules.join("; ")}`,
|
|
},
|
|
],
|
|
details: { projectAgentsDir: null, results: [], totalDurationMs: 0 },
|
|
};
|
|
}
|
|
|
|
if (single) {
|
|
const { agentId, jobId, description } = started[0];
|
|
const coordinationHint = ircEnabled
|
|
? `DM \`${agentId}\` via \`irc\` to coordinate while it runs; use \`job\` only to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task.`
|
|
: `Use \`job\` to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task.`;
|
|
const descriptionSuffix = description ? ` — ${description}` : "";
|
|
onUpdate?.({
|
|
content: [{ type: "text", text: `Spawned agent \`${agentId}\`...` }],
|
|
details: buildAsyncDetails("running", jobId),
|
|
});
|
|
return {
|
|
content: [
|
|
{
|
|
type: "text",
|
|
text: `Spawned agent \`${agentId}\` (job \`${jobId}\`)${descriptionSuffix}. The result will be delivered when it yields. ${coordinationHint}`,
|
|
},
|
|
],
|
|
details: buildAsyncDetails("running", jobId),
|
|
};
|
|
}
|
|
|
|
const coordinationHint = ircEnabled
|
|
? `DM these ids via \`irc\` to coordinate while they run; use \`job\` only to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task.`
|
|
: `Use \`job\` to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task by id.`;
|
|
const scheduleFailureSummary =
|
|
failedSchedules.length > 0
|
|
? ` Failed to schedule ${failedSchedules.length} spawn${failedSchedules.length === 1 ? "" : "s"}: ${failedSchedules.join("; ")}.`
|
|
: "";
|
|
const startedListing = started
|
|
.map(({ agentId, jobId, description }) => {
|
|
const prefix = `- \`${agentId}\` (job \`${jobId}\`)`;
|
|
return description ? `${prefix} — ${description}` : prefix;
|
|
})
|
|
.join("\n");
|
|
onUpdate?.({
|
|
content: [{ type: "text", text: `Spawned ${started.length} agents...` }],
|
|
details: buildAsyncDetails("running", primaryJobId),
|
|
});
|
|
return {
|
|
content: [
|
|
{
|
|
type: "text",
|
|
text: `Spawned ${started.length} background agents using ${agentLabel}.${scheduleFailureSummary} Each result will be delivered when that agent yields.\n${startedListing}\n${coordinationHint}`,
|
|
},
|
|
],
|
|
details: buildAsyncDetails("running", primaryJobId),
|
|
};
|
|
}
|
|
|
|
/**
|
|
* 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: (state: "running" | "completed" | "failed", jobId: string) => TaskToolDetails;
|
|
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>;
|
|
onSettled?: (failed: boolean) => void;
|
|
}): string {
|
|
const { manager, toolCallId, spawnParams, agentId, progress, ircEnabled, buildDetails, onUpdate, onSettled } =
|
|
options;
|
|
const buildFollowUpHint = (aborted: boolean): string => {
|
|
if (aborted) {
|
|
return `\n\n${agentId} was aborted — transcript at history://${agentId}`;
|
|
}
|
|
const followUp = ircEnabled ? "message it via `irc` to follow up; " : "";
|
|
return `\n\n${agentId} is now idle — ${followUp}transcript at history://${agentId}`;
|
|
};
|
|
return manager.register(
|
|
"task",
|
|
agentId,
|
|
async ({ jobId: ownJobId, signal: runSignal, reportProgress, markRunning }) => {
|
|
const startedAt = Date.now();
|
|
const semaphore = this.#getSpawnSemaphore();
|
|
await semaphore.acquire();
|
|
if (runSignal.aborted) {
|
|
semaphore.release();
|
|
progress.status = "aborted";
|
|
onSettled?.(true);
|
|
throw new Error("Aborted before execution");
|
|
}
|
|
markRunning();
|
|
progress.status = "running";
|
|
await reportProgress(
|
|
`Running background task ${agentId}...`,
|
|
buildDetails("running", ownJobId) as unknown as Record<string, unknown>,
|
|
);
|
|
try {
|
|
const result = await this.#executeSync(
|
|
toolCallId,
|
|
spawnParams,
|
|
runSignal,
|
|
undefined,
|
|
agentId,
|
|
progress.index,
|
|
);
|
|
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;
|
|
onSettled?.(resultFailed);
|
|
const statusText = resultFailed
|
|
? `Background task ${agentId} failed.`
|
|
: `Background task ${agentId} complete.`;
|
|
await reportProgress(
|
|
statusText,
|
|
buildDetails(resultFailed ? "failed" : "completed", ownJobId) as unknown as Record<string, unknown>,
|
|
);
|
|
onUpdate?.({
|
|
content: [{ type: "text", text: statusText }],
|
|
details: buildDetails(resultFailed ? "failed" : "completed", ownJobId),
|
|
});
|
|
const deliveryText = `${finalText}${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("failed", ownJobId) as unknown as Record<string, unknown>);
|
|
onUpdate?.({
|
|
content: [{ type: "text", text: statusText }],
|
|
details: buildDetails("failed", ownJobId),
|
|
});
|
|
const message = error instanceof Error ? error.message : String(error);
|
|
const hint = AgentRegistry.global().get(agentId) ? buildFollowUpHint(false) : "";
|
|
throw new TaskJobError(`${message}${hint}`);
|
|
} finally {
|
|
semaphore.release();
|
|
}
|
|
},
|
|
{
|
|
id: agentId,
|
|
queued: true,
|
|
ownerId: this.session.getAgentId?.() ?? undefined,
|
|
onProgress: (text, details) => {
|
|
const progressDetails = (details as TaskToolDetails | undefined) ?? buildDetails("running", agentId);
|
|
onUpdate?.({ content: [{ type: "text", text }], details: progressDetails });
|
|
},
|
|
},
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Sync fallback fan-out (no job manager, or a `blocking: true` agent): 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,
|
|
spawnItems: TaskItem[],
|
|
signal?: AbortSignal,
|
|
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>,
|
|
): Promise<AgentToolResult<TaskToolDetails>> {
|
|
const semaphore = this.#getSpawnSemaphore();
|
|
if (spawnItems.length === 1) {
|
|
await semaphore.acquire();
|
|
try {
|
|
return await this.#executeSync(
|
|
toolCallId,
|
|
spawnParamsFor(params, spawnItems[0]),
|
|
signal,
|
|
onUpdate,
|
|
undefined,
|
|
0,
|
|
);
|
|
} finally {
|
|
semaphore.release();
|
|
}
|
|
}
|
|
|
|
const startTime = Date.now();
|
|
const latestProgress = new Map<number, AgentProgress>();
|
|
const emitCombined = () => {
|
|
onUpdate?.({
|
|
content: [{ type: "text", text: `Running ${spawnItems.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 { results: payloads } = await mapWithConcurrencyLimit(
|
|
spawnItems,
|
|
spawnItems.length,
|
|
async (item, index, workerSignal) => {
|
|
await semaphore.acquire();
|
|
try {
|
|
const itemOnUpdate: AgentToolUpdateCallback<TaskToolDetails> | undefined = onUpdate
|
|
? update => {
|
|
const progress = update.details?.progress?.[0];
|
|
if (progress) {
|
|
latestProgress.set(index, { ...progress, index });
|
|
emitCombined();
|
|
}
|
|
}
|
|
: undefined;
|
|
return await this.#executeSync(
|
|
toolCallId,
|
|
spawnParamsFor(params, item),
|
|
workerSignal,
|
|
itemOnUpdate,
|
|
undefined,
|
|
index,
|
|
);
|
|
} finally {
|
|
semaphore.release();
|
|
}
|
|
},
|
|
signal,
|
|
);
|
|
|
|
const results: SingleResult[] = [];
|
|
const contentParts: string[] = [];
|
|
const outputPaths: string[] = [];
|
|
const usageTotals = createUsageTotals();
|
|
let hasUsage = false;
|
|
let projectAgentsDir: string | null = null;
|
|
for (let index = 0; index < spawnItems.length; index++) {
|
|
const payload = payloads[index];
|
|
if (!payload) {
|
|
contentParts.push(`Task ${spawnItems[index].id?.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 {
|
|
content: [{ type: "text", text: contentParts.join("\n\n") }],
|
|
details: {
|
|
projectAgentsDir,
|
|
results,
|
|
totalDurationMs: Date.now() - startTime,
|
|
usage: hasUsage ? usageTotals : undefined,
|
|
outputPaths: outputPaths.length > 0 ? outputPaths : undefined,
|
|
},
|
|
};
|
|
}
|
|
|
|
/**
|
|
* 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,
|
|
): Promise<AgentToolResult<TaskToolDetails>> {
|
|
return this.#runSpawn(toolCallId, params, signal, onUpdate, preAllocatedId, spawnIndex);
|
|
}
|
|
|
|
/** Spawn a fresh subagent and run it to completion. */
|
|
async #runSpawn(
|
|
toolCallId: string,
|
|
params: TaskParams,
|
|
signal?: AbortSignal,
|
|
onUpdate?: AgentToolUpdateCallback<TaskToolDetails>,
|
|
preAllocatedId?: string,
|
|
spawnIndex = 0,
|
|
): Promise<AgentToolResult<TaskToolDetails>> {
|
|
const startTime = Date.now();
|
|
const { agents, projectAgentsDir } = await discoverAgents(this.session.cwd);
|
|
const agentName = params.agent ?? "";
|
|
const sharedContext = this.#isBatchEnabled() ? params.context?.trim() || undefined : undefined;
|
|
const assignment = (params.assignment ?? "").trim();
|
|
const isolationMode = this.session.settings.get("task.isolation.mode");
|
|
const isolationRequested = "isolated" in params ? params.isolated === true : false;
|
|
const isIsolated = isolationMode !== "none" && isolationRequested;
|
|
const mergeMode = this.session.settings.get("task.isolation.merge");
|
|
const commitStyle = this.session.settings.get("task.isolation.commits");
|
|
const taskDepth = this.session.taskDepth ?? 0;
|
|
const subagentLspEnabled = (this.session.enableLsp ?? true) && this.session.settings.get("task.enableLsp");
|
|
|
|
if (isolationMode === "none" && "isolated" in params) {
|
|
return {
|
|
content: [{ type: "text", text: "Task isolation is disabled." }],
|
|
details: { projectAgentsDir, results: [], totalDurationMs: 0 },
|
|
};
|
|
}
|
|
|
|
// Validate agent exists
|
|
const agent = getAgent(agents, agentName);
|
|
if (!agent) {
|
|
const available = agents.map(a => a.name).join(", ") || "none";
|
|
return {
|
|
content: [{ type: "text", text: `Unknown agent "${agentName}". Available: ${available}` }],
|
|
details: { projectAgentsDir, results: [], totalDurationMs: 0 },
|
|
};
|
|
}
|
|
|
|
// Check if agent is disabled in settings
|
|
const disabledAgents = this.session.settings.get("task.disabledAgents") as string[];
|
|
if (disabledAgents.length > 0 && disabledAgents.includes(agentName)) {
|
|
const enabled = agents.filter(a => !disabledAgents.includes(a.name)).map(a => a.name);
|
|
return {
|
|
content: [
|
|
{
|
|
type: "text",
|
|
text: `Agent "${agentName}" is disabled in settings. Enable it via /agents, or use a different agent type.${enabled.length > 0 ? ` Available: ${enabled.join(", ")}` : ""}`,
|
|
},
|
|
],
|
|
details: { projectAgentsDir, results: [], totalDurationMs: 0 },
|
|
};
|
|
}
|
|
|
|
const planModeState = this.session.getPlanModeState?.();
|
|
const planModeBaseTools = ["read", "search", "find", "lsp", "web_search"];
|
|
const planModeTools = [
|
|
...planModeBaseTools,
|
|
...(agent.tools ?? []).filter(
|
|
tool => PLAN_MODE_AGENT_TOOL_ALLOWLIST.has(tool) && !planModeBaseTools.includes(tool),
|
|
),
|
|
];
|
|
const effectiveAgent: typeof agent = planModeState?.enabled
|
|
? {
|
|
...agent,
|
|
systemPrompt: `${planModeSubagentPrompt}\n\n${agent.systemPrompt}`,
|
|
tools: planModeTools,
|
|
spawns: undefined,
|
|
}
|
|
: agent;
|
|
|
|
// Apply per-agent model override from settings (highest priority)
|
|
const agentModelOverrides = this.session.settings.get("task.agentModelOverrides");
|
|
const settingsModelOverride = agentModelOverrides[agentName];
|
|
const parentActiveModelPattern = this.session.getActiveModelString?.();
|
|
const modelOverride = resolveAgentModelPatterns({
|
|
settingsOverride: settingsModelOverride,
|
|
agentModel: effectiveAgent.model,
|
|
settings: this.session.settings,
|
|
activeModelPattern: parentActiveModelPattern,
|
|
fallbackModelPattern: this.session.getModelString?.(),
|
|
});
|
|
const thinkingLevelOverride = effectiveAgent.thinkingLevel;
|
|
|
|
// Output schema priority: agent frontmatter > inherited parent session.
|
|
// The task call itself never carries a schema; workflows needing ad-hoc
|
|
// structured output go through eval agent(prompt, schema).
|
|
const effectiveOutputSchema = effectiveAgent.output ?? this.session.outputSchema;
|
|
|
|
let repoRoot: string | null = null;
|
|
let baseline: WorktreeBaseline | null = null;
|
|
if (isIsolated) {
|
|
try {
|
|
repoRoot = await getRepoRoot(this.session.cwd);
|
|
baseline = await captureBaseline(repoRoot);
|
|
} catch (err) {
|
|
const message = err instanceof Error ? err.message : String(err);
|
|
return {
|
|
content: [{ type: "text", text: `Isolated task execution requires a git repository. ${message}` }],
|
|
details: { projectAgentsDir, results: [], totalDurationMs: Date.now() - startTime },
|
|
};
|
|
}
|
|
}
|
|
|
|
const preferredIsolationBackend = parseIsolationMode(isolationMode);
|
|
|
|
// Derive artifacts directory
|
|
const sessionFile = this.session.getSessionFile();
|
|
const artifactsDir = sessionFile ? sessionFile.slice(0, -6) : null;
|
|
const tempArtifactsDir = artifactsDir ? null : path.join(os.tmpdir(), `omp-task-${Snowflake.next()}`);
|
|
const effectiveArtifactsDir = artifactsDir || tempArtifactsDir!;
|
|
|
|
const localProtocolOptions: LocalProtocolOptions = this.session.localProtocolOptions ?? {
|
|
getArtifactsDir: this.session.getArtifactsDir ?? (() => null),
|
|
getSessionId: this.session.getSessionId ?? (() => null),
|
|
};
|
|
|
|
// Subagents adopt the parent's ArtifactManager so artifact IDs are unique
|
|
// across the whole tree and outputs land flat in the parent's dir.
|
|
const parentArtifactManager = this.session.getArtifactManager?.() ?? undefined;
|
|
|
|
// When the session is executing an approved plan, hand the overall plan to
|
|
// every subagent so they share the main agent's plan context. Skipped in
|
|
// plan mode (read-only exploration uses planModeSubagentPrompt instead) and
|
|
// when no plan file exists at the session's reference path.
|
|
const planReference = planModeState?.enabled
|
|
? undefined
|
|
: await loadOverallPlanReference(
|
|
this.session.getPlanReferencePath?.() ?? "local://PLAN.md",
|
|
localProtocolOptions,
|
|
);
|
|
|
|
try {
|
|
// Check self-recursion prevention
|
|
if (this.#blockedAgent && agentName === this.#blockedAgent) {
|
|
return {
|
|
content: [
|
|
{
|
|
type: "text",
|
|
text: `Cannot spawn ${this.#blockedAgent} agent from within itself (recursion prevention). Use a different agent type.`,
|
|
},
|
|
],
|
|
details: { projectAgentsDir, results: [], totalDurationMs: Date.now() - startTime },
|
|
};
|
|
}
|
|
|
|
// Check spawn restrictions from parent
|
|
const parentSpawns = this.session.getSessionSpawns() ?? "*";
|
|
const allowedSpawns = parentSpawns.split(",").map(s => s.trim());
|
|
const isSpawnAllowed = (): boolean => {
|
|
if (parentSpawns === "") return false; // Empty = deny all
|
|
if (parentSpawns === "*") return true; // Wildcard = allow all
|
|
return allowedSpawns.includes(agentName);
|
|
};
|
|
|
|
if (!isSpawnAllowed()) {
|
|
const allowed = parentSpawns === "" ? "none (spawns disabled for this agent)" : parentSpawns;
|
|
return {
|
|
content: [{ type: "text", text: `Cannot spawn '${agentName}'. Allowed: ${allowed}` }],
|
|
details: { projectAgentsDir, results: [], totalDurationMs: Date.now() - startTime },
|
|
};
|
|
}
|
|
|
|
await fs.mkdir(effectiveArtifactsDir, { recursive: true });
|
|
|
|
// Allocate a unique ID across the session to prevent artifact collisions
|
|
let agentId: string;
|
|
if (preAllocatedId) {
|
|
agentId = preAllocatedId;
|
|
} else {
|
|
const outputManager =
|
|
this.session.agentOutputManager ?? new AgentOutputManager(this.session.getArtifactsDir ?? (() => null));
|
|
agentId = await outputManager.allocate(params.id?.trim() || generateTaskName());
|
|
}
|
|
|
|
const availableSkills = [...(this.session.skills ?? [])];
|
|
// Resolve autoload skills from agent definition against available skills
|
|
const resolvedAutoloadSkills =
|
|
agent.autoloadSkills?.length && availableSkills.length > 0
|
|
? agent.autoloadSkills
|
|
.map(name => availableSkills.find(s => s.name === name))
|
|
.filter((s): s is NonNullable<typeof s> => s !== undefined)
|
|
: [];
|
|
const contextFiles = this.session.contextFiles?.filter(
|
|
file => path.basename(file.path).toLowerCase() !== "agents.md",
|
|
);
|
|
const promptTemplates = this.session.promptTemplates;
|
|
const parentEvalSessionId = this.session.getEvalSessionId?.() ?? undefined;
|
|
const mcpManager = this.session.mcpManager ?? MCPManager.instance();
|
|
|
|
// Progress tracking for the single agent
|
|
let latestProgress: AgentProgress = {
|
|
index: spawnIndex,
|
|
id: agentId,
|
|
agent: agentName,
|
|
agentSource: agent.source,
|
|
status: "pending",
|
|
task: renderSubagentUserPrompt(assignment),
|
|
assignment,
|
|
recentTools: [],
|
|
recentOutput: [],
|
|
toolCount: 0,
|
|
requests: 0,
|
|
tokens: 0,
|
|
cost: 0,
|
|
durationMs: 0,
|
|
modelOverride,
|
|
description: params.description,
|
|
};
|
|
const emitProgress = () => {
|
|
onUpdate?.({
|
|
content: [{ type: "text", text: `Running agent ${agentId}...` }],
|
|
details: {
|
|
projectAgentsDir,
|
|
results: [],
|
|
totalDurationMs: Date.now() - startTime,
|
|
progress: [latestProgress],
|
|
},
|
|
});
|
|
};
|
|
emitProgress();
|
|
|
|
const buildCommitMessageFn = () =>
|
|
commitStyle === "ai" && this.session.modelRegistry
|
|
? async (diff: string) => {
|
|
return generateCommitMessage(
|
|
diff,
|
|
this.session.modelRegistry!,
|
|
this.session.settings,
|
|
this.session.getSessionId?.() ?? undefined,
|
|
);
|
|
}
|
|
: undefined;
|
|
|
|
const sharedRunOptions = {
|
|
cwd: this.session.cwd,
|
|
agent: effectiveAgent,
|
|
task: renderSubagentUserPrompt(assignment),
|
|
assignment,
|
|
context: sharedContext,
|
|
planReference,
|
|
description: params.description,
|
|
index: spawnIndex,
|
|
parentToolCallId: toolCallId,
|
|
id: agentId,
|
|
taskDepth,
|
|
modelOverride,
|
|
parentActiveModelPattern,
|
|
thinkingLevel: thinkingLevelOverride,
|
|
outputSchema: effectiveOutputSchema,
|
|
sessionFile,
|
|
persistArtifacts: !!artifactsDir,
|
|
artifactsDir: effectiveArtifactsDir,
|
|
enableLsp: subagentLspEnabled,
|
|
signal,
|
|
eventBus: this.session.eventBus,
|
|
onProgress: (progress: AgentProgress) => {
|
|
// Shallow snapshot; recentTools is mutated in place by the
|
|
// executor, the rest is reassigned or immutable. A deep clone
|
|
// here cost O(extractedToolData) per progress event.
|
|
latestProgress = { ...progress, recentTools: progress.recentTools.slice() };
|
|
emitProgress();
|
|
},
|
|
authStorage: this.session.authStorage,
|
|
modelRegistry: this.session.modelRegistry,
|
|
settings: this.session.settings,
|
|
mcpManager,
|
|
contextFiles,
|
|
skills: availableSkills,
|
|
autoloadSkills: resolvedAutoloadSkills,
|
|
workspaceTree: this.session.workspaceTree,
|
|
promptTemplates,
|
|
rules: this.session.rules,
|
|
preloadedExtensionPaths: this.session.extensionPaths,
|
|
preloadedCustomToolPaths: this.session.customToolPaths,
|
|
localProtocolOptions,
|
|
parentArtifactManager,
|
|
parentHindsightSessionState: this.session.getHindsightSessionState?.(),
|
|
parentMnemopiSessionState: this.session.getMnemopiSessionState?.(),
|
|
parentTelemetry: this.session.getTelemetry?.(),
|
|
parentEvalSessionId,
|
|
};
|
|
|
|
const runTask = async (): Promise<SingleResult> => {
|
|
if (!isIsolated) {
|
|
return runSubprocess(sharedRunOptions);
|
|
}
|
|
|
|
const taskStart = Date.now();
|
|
let isolationHandle: IsolationHandle | undefined;
|
|
try {
|
|
if (!repoRoot || !baseline) {
|
|
throw new Error("Isolated task execution not initialized.");
|
|
}
|
|
const taskBaseline = structuredClone(baseline);
|
|
|
|
isolationHandle = await ensureIsolation(repoRoot, agentId, preferredIsolationBackend);
|
|
const isolationDir = isolationHandle.mergedDir;
|
|
|
|
// Isolated runs re-discover extensions/custom tools inside the
|
|
// worktree instead of reusing the parent's source paths.
|
|
const result = await runSubprocess({
|
|
...sharedRunOptions,
|
|
worktree: isolationDir,
|
|
preloadedExtensionPaths: undefined,
|
|
preloadedCustomToolPaths: undefined,
|
|
});
|
|
if (mergeMode === "branch" && result.exitCode === 0) {
|
|
try {
|
|
const commitResult = await commitToBranch(
|
|
isolationDir,
|
|
taskBaseline,
|
|
agentId,
|
|
params.description,
|
|
buildCommitMessageFn(),
|
|
);
|
|
return {
|
|
...result,
|
|
branchName: commitResult?.branchName,
|
|
nestedPatches: commitResult?.nestedPatches,
|
|
};
|
|
} catch (mergeErr) {
|
|
// Agent succeeded but branch commit failed — clean up stale branch
|
|
const branchName = `omp/task/${agentId}`;
|
|
await git.branch.tryDelete(repoRoot, branchName);
|
|
const msg = mergeErr instanceof Error ? mergeErr.message : String(mergeErr);
|
|
return { ...result, error: `Merge failed: ${msg}` };
|
|
}
|
|
}
|
|
if (result.exitCode === 0) {
|
|
try {
|
|
const delta = await captureDeltaPatch(isolationDir, taskBaseline);
|
|
const patchPath = path.join(effectiveArtifactsDir, `${agentId}.patch`);
|
|
await Bun.write(patchPath, delta.rootPatch);
|
|
return {
|
|
...result,
|
|
patchPath,
|
|
nestedPatches: delta.nestedPatches,
|
|
};
|
|
} catch (patchErr) {
|
|
const msg = patchErr instanceof Error ? patchErr.message : String(patchErr);
|
|
return { ...result, error: `Patch capture failed: ${msg}` };
|
|
}
|
|
}
|
|
return result;
|
|
} catch (err) {
|
|
const message = err instanceof Error ? err.message : String(err);
|
|
return {
|
|
index: spawnIndex,
|
|
id: agentId,
|
|
agent: agent.name,
|
|
agentSource: agent.source,
|
|
task: renderSubagentUserPrompt(assignment),
|
|
assignment,
|
|
description: params.description,
|
|
exitCode: 1,
|
|
output: "",
|
|
stderr: message,
|
|
truncated: false,
|
|
durationMs: Date.now() - taskStart,
|
|
tokens: 0,
|
|
requests: 0,
|
|
modelOverride,
|
|
error: message,
|
|
};
|
|
} finally {
|
|
if (isolationHandle) {
|
|
await cleanupIsolation(isolationHandle);
|
|
}
|
|
}
|
|
};
|
|
|
|
const result = await runTask();
|
|
|
|
let mergeSummary = "";
|
|
let changesApplied: boolean | null = null;
|
|
let hadAnyChanges = false;
|
|
let mergedBranchForNestedPatches = false;
|
|
if (isIsolated && repoRoot) {
|
|
try {
|
|
if (mergeMode === "branch") {
|
|
if (!result.branchName || result.exitCode !== 0 || result.aborted) {
|
|
changesApplied = true;
|
|
mergeSummary = "\n\nNo changes to apply.";
|
|
} else {
|
|
const mergeResult = await mergeTaskBranches(repoRoot, [
|
|
{ branchName: result.branchName, taskId: result.id, description: result.description },
|
|
]);
|
|
mergedBranchForNestedPatches = mergeResult.merged.includes(result.branchName);
|
|
changesApplied = mergeResult.failed.length === 0;
|
|
hadAnyChanges = changesApplied && mergeResult.merged.length > 0;
|
|
|
|
if (changesApplied) {
|
|
mergeSummary = hadAnyChanges
|
|
? `\n\nMerged branch: ${result.branchName}`
|
|
: "\n\nNo changes to apply.";
|
|
} else {
|
|
const conflictPart = mergeResult.conflict ? `\nConflict: ${mergeResult.conflict}` : "";
|
|
mergeSummary = `\n\n<system-notification>Branch merge failed: ${result.branchName}.${conflictPart}\nThe unmerged branch remains for manual resolution.</system-notification>`;
|
|
}
|
|
if (mergeResult.stashConflict) {
|
|
mergeSummary += `\n\n<system-notification>${mergeResult.stashConflict}</system-notification>`;
|
|
}
|
|
|
|
// Clean up the merged branch (keep failed ones for manual resolution)
|
|
if (changesApplied) {
|
|
await cleanupTaskBranches(repoRoot, [result.branchName]);
|
|
}
|
|
}
|
|
} else {
|
|
// Patch mode: apply the patch from a successful run. A failed or
|
|
// aborted run has nothing to apply and must not block the result.
|
|
const succeeded = result.exitCode === 0 && !result.error && !result.aborted;
|
|
if (!succeeded) {
|
|
changesApplied = true;
|
|
hadAnyChanges = false;
|
|
} else if (!result.patchPath) {
|
|
changesApplied = false;
|
|
hadAnyChanges = false;
|
|
} else {
|
|
const patchText = await Bun.file(result.patchPath).text();
|
|
if (!patchText.trim()) {
|
|
changesApplied = true;
|
|
hadAnyChanges = false;
|
|
} else {
|
|
const normalized = patchText.endsWith("\n") ? patchText : `${patchText}\n`;
|
|
changesApplied = await git.patch.canApplyText(repoRoot, normalized);
|
|
if (changesApplied) {
|
|
try {
|
|
await git.patch.applyText(repoRoot, normalized);
|
|
hadAnyChanges = true;
|
|
} catch {
|
|
changesApplied = false;
|
|
hadAnyChanges = false;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
if (changesApplied) {
|
|
mergeSummary = hadAnyChanges ? "\n\nApplied patches: yes" : "\n\nNo changes to apply.";
|
|
} else {
|
|
const notification =
|
|
"<system-notification>Patches were not applied and must be handled manually.</system-notification>";
|
|
const patchList = result.patchPath ? `\n\nPatch artifact:\n- ${result.patchPath}` : "";
|
|
mergeSummary = `\n\n${notification}${patchList}`;
|
|
}
|
|
}
|
|
} catch (mergeErr) {
|
|
const msg = mergeErr instanceof Error ? mergeErr.message : String(mergeErr);
|
|
changesApplied = false;
|
|
hadAnyChanges = false;
|
|
mergeSummary = `\n\n<system-notification>Merge phase failed: ${msg}\nTask outputs are preserved but changes were not applied.</system-notification>`;
|
|
}
|
|
}
|
|
|
|
// Apply nested repo patches (separate from parent git)
|
|
if (isIsolated && repoRoot && (mergeMode === "branch" || changesApplied !== false)) {
|
|
const nestedPatches = result.nestedPatches ?? [];
|
|
const eligible =
|
|
nestedPatches.length > 0 &&
|
|
result.exitCode === 0 &&
|
|
!result.aborted &&
|
|
(mergeMode !== "branch" || mergedBranchForNestedPatches);
|
|
if (eligible) {
|
|
try {
|
|
await applyNestedPatches(repoRoot, nestedPatches, buildCommitMessageFn());
|
|
} catch {
|
|
// Nested patch failures are non-fatal to the parent merge
|
|
mergeSummary +=
|
|
"\n\n<system-notification>Some nested repository patches failed to apply.</system-notification>";
|
|
}
|
|
}
|
|
}
|
|
|
|
// Cleanup temp directory if used
|
|
const shouldCleanupTempArtifacts =
|
|
tempArtifactsDir && (!isIsolated || changesApplied === true || changesApplied === null);
|
|
if (shouldCleanupTempArtifacts) {
|
|
await fs.rm(tempArtifactsDir, { recursive: true, force: true });
|
|
}
|
|
|
|
return this.#buildResultPayload(result, projectAgentsDir, Date.now() - startTime, mergeSummary);
|
|
} catch (err) {
|
|
return {
|
|
content: [{ type: "text", text: `Task execution failed: ${err}` }],
|
|
details: { projectAgentsDir, results: [], totalDurationMs: Date.now() - startTime },
|
|
};
|
|
}
|
|
}
|
|
|
|
/** 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;
|
|
}
|
|
const summary = prompt.render(taskSummaryTemplate, {
|
|
agentName: result.agent,
|
|
id: result.id,
|
|
status,
|
|
duration: formatDuration(totalDurationMs),
|
|
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,
|
|
},
|
|
};
|
|
}
|
|
}
|