diff --git a/packages/coding-agent/src/task/commands.ts b/packages/coding-agent/src/task/commands.ts index 3d61ece2c..9c393a239 100644 --- a/packages/coding-agent/src/task/commands.ts +++ b/packages/coding-agent/src/task/commands.ts @@ -120,7 +120,8 @@ export function getCommand(commands: WorkflowCommand[], name: string): WorkflowC * Replaces $@ with the provided input. */ export function expandCommand(command: WorkflowCommand, input: string): string { - return command.instructions.replace(/\$@/g, input); + // Function replacement so `$`-patterns in user input ($$, $&, ...) stay literal. + return command.instructions.replace(/\$@/g, () => input); } /** diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index 0337e5579..5fcc075ce 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -1285,59 +1285,67 @@ export async function runSubprocess(options: ExecutorOptions): Promise { - const subagentPrompt = prompt.render(subagentSystemPromptTemplate, { - agent: agent.systemPrompt, - context: options.context?.trim() ?? "", - planReference: options.planReference?.content ?? "", - planReferencePath: options.planReference?.path ?? "", - worktree: worktree ?? "", - outputSchema: normalizedOutputSchema, - contextFile: contextFileForPrompt, - ircPeers: ircEnabled ? renderIrcPeerRoster(id) : "", - ircSelfId: ircEnabled ? id : "", - }); - return defaultPrompt.length === 0 - ? [subagentPrompt] - : [...defaultPrompt.slice(0, -1), subagentPrompt, defaultPrompt[defaultPrompt.length - 1]]; - }, - sessionManager, - hasUI: false, - spawns: spawnsEnv, - taskDepth: childDepth, - parentHindsightSessionState: options.parentHindsightSessionState, - parentMnemopiSessionState: options.parentMnemopiSessionState, - parentTaskPrefix: id, - agentId: id, - agentDisplayName: agent.name, - enableLsp: lspEnabled, - skipPythonPreflight, - enableMCP, - mcpManager: options.mcpManager, - customTools: mcpProxyTools.length > 0 ? mcpProxyTools : undefined, - localProtocolOptions: options.localProtocolOptions, - telemetry: subagentTelemetry, - parentEvalSessionId: options.parentEvalSessionId, - }), - ); + const sessionPromise = createAgentSession({ + cwd: worktree ?? cwd, + authStorage, + modelRegistry, + settings: subagentSettings, + model, + thinkingLevel: effectiveThinkingLevel, + toolNames, + outputSchema, + requireYieldTool: true, + contextFiles: options.contextFiles, + skills: options.skills, + promptTemplates: options.promptTemplates, + workspaceTree: options.workspaceTree, + rules: options.rules, + preloadedExtensionPaths: options.preloadedExtensionPaths, + preloadedCustomToolPaths: options.preloadedCustomToolPaths, + systemPrompt: defaultPrompt => { + const subagentPrompt = prompt.render(subagentSystemPromptTemplate, { + agent: agent.systemPrompt, + context: options.context?.trim() ?? "", + planReference: options.planReference?.content ?? "", + planReferencePath: options.planReference?.path ?? "", + worktree: worktree ?? "", + outputSchema: normalizedOutputSchema, + contextFile: contextFileForPrompt, + ircPeers: ircEnabled ? renderIrcPeerRoster(id) : "", + ircSelfId: ircEnabled ? id : "", + }); + return defaultPrompt.length === 0 + ? [subagentPrompt] + : [...defaultPrompt.slice(0, -1), subagentPrompt, defaultPrompt[defaultPrompt.length - 1]]; + }, + sessionManager, + hasUI: false, + spawns: spawnsEnv, + taskDepth: childDepth, + parentHindsightSessionState: options.parentHindsightSessionState, + parentMnemopiSessionState: options.parentMnemopiSessionState, + parentTaskPrefix: id, + agentId: id, + agentDisplayName: agent.name, + enableLsp: lspEnabled, + skipPythonPreflight, + enableMCP, + mcpManager: options.mcpManager, + customTools: mcpProxyTools.length > 0 ? mcpProxyTools : undefined, + localProtocolOptions: options.localProtocolOptions, + telemetry: subagentTelemetry, + parentEvalSessionId: options.parentEvalSessionId, + }); + let session: AgentSession; + try { + ({ session } = await awaitAbortable(sessionPromise)); + } catch (err) { + // Abort raced session startup. The session may still resolve later + // holding live LSP/MCP child processes — dispose it when it does so + // a cancelled subagent cannot leak them. + void sessionPromise.then(created => created.session.dispose()).catch(() => {}); + throw err; + } activeSession = session; diff --git a/packages/coding-agent/src/task/index.ts b/packages/coding-agent/src/task/index.ts index 26f2d116c..a88b8fe49 100644 --- a/packages/coding-agent/src/task/index.ts +++ b/packages/coding-agent/src/task/index.ts @@ -242,6 +242,57 @@ function validateTaskModeParams(simpleMode: TaskSimpleMode, params: TaskParams): return "task.simple is set to independent, so the task tool does not accept `context` or `schema`. Put all required background and output expectations inside each task assignment or the selected agent definition."; } +/** Sentinel for async jobs whose subagent finished with a failing result; batch counters are already updated. */ +class TaskJobError extends Error {} + +/** + * Validate task ids: every task needs a non-empty id and ids must be unique + * (case-insensitive). Returns a problem description, or undefined when valid. + */ +function validateTaskIds(tasks: TaskParams["tasks"]): string | undefined { + const missingTaskIndexes: number[] = []; + const idIndexes = new Map(); + + for (let i = 0; i < tasks.length; i++) { + const id = tasks[i]?.id; + if (typeof id !== "string" || id.trim() === "") { + missingTaskIndexes.push(i); + continue; + } + const normalizedId = id.toLowerCase(); + const indexes = idIndexes.get(normalizedId); + if (indexes) { + indexes.push(i); + } else { + idIndexes.set(normalizedId, [i]); + } + } + + const duplicateIds: Array<{ id: string; indexes: number[] }> = []; + for (const [normalizedId, indexes] of idIndexes.entries()) { + if (indexes.length > 1) { + duplicateIds.push({ + id: tasks[indexes[0]]?.id ?? normalizedId, + indexes, + }); + } + } + + if (missingTaskIndexes.length === 0 && duplicateIds.length === 0) { + return undefined; + } + + const problems: string[] = []; + if (missingTaskIndexes.length > 0) { + problems.push(`Missing task ids at indexes: ${missingTaskIndexes.join(", ")}`); + } + if (duplicateIds.length > 0) { + const details = duplicateIds.map(entry => `${entry.id} (indexes ${entry.indexes.join(", ")})`).join("; "); + problems.push(`Duplicate task ids detected (case-insensitive): ${details}`); + } + return `Invalid tasks: ${problems.join(". ")}`; +} + // ═══════════════════════════════════════════════════════════════════════════ // Tool Class // ═══════════════════════════════════════════════════════════════════════════ @@ -363,6 +414,11 @@ export class TaskTool implements AgentTool null)); const uniqueIds = await outputManager.allocateBatch(taskItems.map(t => t.id)); @@ -396,9 +452,13 @@ export class TaskTool implements AgentTool { + // Shallow copies: top-level fields are reassigned (never mutated in + // place) and the large nested payloads (extractedToolData) are + // immutable once attached — structuredClone here cost O(batch × payload) + // per progress event. return Array.from(progressByTaskId.values()) .sort((a, b) => a.index - b.index) - .map(progress => structuredClone(progress)); + .map(progress => ({ ...progress })); }; const buildAsyncDetails = (state: "running" | "completed" | "failed", jobId: string): TaskToolDetails => ({ @@ -424,6 +484,7 @@ export class TaskTool implements AgentTool { + async ({ signal: runSignal, reportProgress, markRunning }) => { const startedAt = Date.now(); const progress = progressByTaskId.get(taskItem.id); await semaphore.acquire(); @@ -447,8 +508,11 @@ export class TaskTool implements AgentTool part.type === "text")?.text ?? "(no output)"; const singleResult = result.details?.results[0]; + // A missing per-task result means #executeSync failed at the + // tool level (results: []) — treat it as a failure, not success. + const resultFailed = + !singleResult || (singleResult.aborted ?? false) || singleResult.exitCode !== 0; if (progress) { - progress.status = singleResult?.aborted - ? "aborted" - : (singleResult?.exitCode ?? 0) === 0 - ? "completed" - : "failed"; + progress.status = singleResult?.aborted ? "aborted" : resultFailed ? "failed" : "completed"; progress.durationMs = singleResult?.durationMs ?? Math.max(0, Date.now() - startedAt); progress.tokens = singleResult?.tokens ?? 0; progress.contextTokens = singleResult?.contextTokens; @@ -478,7 +542,7 @@ export class TaskTool implements AgentTool { const progressDetails = @@ -543,6 +615,7 @@ export class TaskTool implements AgentTool(); - - for (let i = 0; i < tasks.length; i++) { - const id = tasks[i]?.id; - if (typeof id !== "string" || id.trim() === "") { - missingTaskIndexes.push(i); - continue; - } - const normalizedId = id.toLowerCase(); - const indexes = idIndexes.get(normalizedId); - if (indexes) { - indexes.push(i); - } else { - idIndexes.set(normalizedId, [i]); - } - } - - const duplicateIds: Array<{ id: string; indexes: number[] }> = []; - for (const [normalizedId, indexes] of idIndexes.entries()) { - if (indexes.length > 1) { - duplicateIds.push({ - id: tasks[indexes[0]]?.id ?? normalizedId, - indexes, - }); - } - } - - if (missingTaskIndexes.length > 0 || duplicateIds.length > 0) { - const problems: string[] = []; - if (missingTaskIndexes.length > 0) { - problems.push(`Missing task ids at indexes: ${missingTaskIndexes.join(", ")}`); - } - if (duplicateIds.length > 0) { - const details = duplicateIds.map(entry => `${entry.id} (indexes ${entry.indexes.join(", ")})`).join("; "); - problems.push(`Duplicate task ids detected (case-insensitive): ${details}`); - } + const taskIdProblem = validateTaskIds(tasks); + if (taskIdProblem) { return { - content: [{ type: "text", text: `Invalid tasks: ${problems.join(". ")}` }], + content: [{ type: "text", text: taskIdProblem }], details: { projectAgentsDir, results: [], @@ -951,7 +989,11 @@ export class TaskTool implements AgentTool { + const runTask = async ( + task: (typeof tasksWithUniqueIds)[number], + index: number, + workerSignal?: AbortSignal, + ) => { if (!isIsolated) { return runSubprocess({ cwd: this.session.cwd, @@ -973,12 +1015,13 @@ export class TaskTool implements AgentTool { - progressMap.set(index, { - ...structuredClone(progress), - }); + // 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. + progressMap.set(index, { ...progress, recentTools: progress.recentTools.slice() }); emitProgress(); }, authStorage: this.session.authStorage, @@ -1034,12 +1077,10 @@ export class TaskTool implements AgentTool { - progressMap.set(index, { - ...structuredClone(progress), - }); + progressMap.set(index, { ...progress, recentTools: progress.recentTools.slice() }); emitProgress(); }, authStorage: this.session.authStorage, @@ -1226,6 +1267,9 @@ export class TaskTool implements AgentToolBranch merge failed. ${mergedPart}${failedPart}${conflictPart}\nUnmerged branches remain for manual resolution.`; } + if (mergeResult.stashConflict) { + mergeSummary += `\n\n${mergeResult.stashConflict}`; + } } // Clean up merged branches (keep failed ones for manual resolution) @@ -1234,9 +1278,11 @@ export class TaskTool implements AgentTool result.patchPath).filter(Boolean) as string[]; - const missingPatch = results.some(result => !result.patchPath); + // Patch mode: apply patches from successful tasks. Failed or + // aborted siblings must not block completed work from landing. + const successfulResults = results.filter(r => r.exitCode === 0 && !r.error && !r.aborted); + const patchesInOrder = successfulResults.map(result => result.patchPath).filter(Boolean) as string[]; + const missingPatch = successfulResults.some(result => !result.patchPath); if (missingPatch) { changesApplied = false; hadAnyChanges = false; diff --git a/packages/coding-agent/src/task/parallel.ts b/packages/coding-agent/src/task/parallel.ts index 1569f9fbf..4a061e88f 100644 --- a/packages/coding-agent/src/task/parallel.ts +++ b/packages/coding-agent/src/task/parallel.ts @@ -20,13 +20,13 @@ export interface ParallelResult { * * @param items - Items to process * @param concurrency - Maximum concurrent operations - * @param fn - Async function to execute for each item + * @param fn - Async function to execute for each item; receives a worker signal that fires on abort or fail-fast so in-flight siblings can cancel * @param signal - Optional abort signal to stop scheduling new work */ export async function mapWithConcurrencyLimit( items: T[], concurrency: number, - fn: (item: T, index: number) => Promise, + fn: (item: T, index: number, signal: AbortSignal) => Promise, signal?: AbortSignal, ): Promise> { const normalizedConcurrency = Number.isFinite(concurrency) ? Math.floor(concurrency) : items.length; @@ -52,7 +52,7 @@ export async function mapWithConcurrencyLimit( const index = nextIndex++; if (index >= items.length) return; try { - results[index] = await fn(items[index], index); + results[index] = await fn(items[index], index, workerSignal); } catch (error) { // On abort, the fn itself handles it and returns a result // Only propagate non-abort errors diff --git a/packages/coding-agent/src/task/worktree.ts b/packages/coding-agent/src/task/worktree.ts index 7bca9c163..a10220e41 100644 --- a/packages/coding-agent/src/task/worktree.ts +++ b/packages/coding-agent/src/task/worktree.ts @@ -5,6 +5,7 @@ import * as path from "node:path"; import * as natives from "@oh-my-pi/pi-natives"; import { getWorktreeDir, hashPath, logger, Snowflake } from "@oh-my-pi/pi-utils"; import * as git from "../utils/git"; +import { mapWithConcurrencyLimit } from "./parallel"; const { IsoBackendKind } = natives; type IsoBackendKind = natives.IsoBackendKind; @@ -82,16 +83,16 @@ async function discoverNestedRepos(repoRoot: string): Promise { async function captureUntrackedPatch(repoRoot: string, untracked: readonly string[]): Promise { if (untracked.length === 0) return ""; const nullPath = getGitNoIndexNullPath(); - const untrackedDiffs = await Promise.all( - untracked.map(entry => - git.diff(repoRoot, { - allowFailure: true, - binary: true, - noIndex: { left: nullPath, right: entry }, - }), - ), + // Bound concurrent git spawns; large untracked sets would otherwise fork one + // process per file at once. + const { results: untrackedDiffs } = await mapWithConcurrencyLimit([...untracked], 8, entry => + git.diff(repoRoot, { + allowFailure: true, + binary: true, + noIndex: { left: nullPath, right: entry }, + }), ); - return untrackedDiffs.filter(diff => diff.trim()).join("\n"); + return untrackedDiffs.filter((diff): diff is string => !!diff?.trim()).join("\n"); } async function captureRepoBaseline(repoRoot: string): Promise { @@ -427,6 +428,8 @@ export interface MergeBranchResult { merged: string[]; failed: string[]; conflict?: string; + /** Set when cherry-picks landed on HEAD but restoring the stashed working tree failed. */ + stashConflict?: string; } /** @@ -438,64 +441,69 @@ export async function mergeTaskBranches( repoRoot: string, branches: Array<{ branchName: string; taskId: string; description?: string }>, ): Promise { - const merged: string[] = []; - const failed: string[] = []; + // Serialize against other in-process git mutations on this repo: concurrent + // background merges interleaving stash push/pop + cherry-pick would corrupt + // the working tree (lost uncommitted changes, mixed-up stash entries). + return git.withRepoLock(repoRoot, async () => { + const merged: string[] = []; + const failed: string[] = []; - // Stash dirty working tree so cherry-pick can operate on a clean HEAD. - // Without this, cherry-pick refuses to run when uncommitted changes exist. - const didStash = await git.stash.push(repoRoot, "omp-task-merge"); + // Stash dirty working tree so cherry-pick can operate on a clean HEAD. + // Without this, cherry-pick refuses to run when uncommitted changes exist. + const didStash = await git.stash.push(repoRoot, "omp-task-merge"); - let conflictResult: MergeBranchResult | undefined; + let conflictResult: MergeBranchResult | undefined; - try { - for (const { branchName } of branches) { - try { - await git.cherryPick(repoRoot, branchName); - } catch (err) { + try { + for (const { branchName } of branches) { try { - await git.cherryPick.abort(repoRoot); - } catch { - /* no state to abort */ - } - const stderr = - err instanceof git.GitCommandError - ? err.result.stderr.trim() - : err instanceof Error - ? err.message - : String(err); - failed.push(branchName); - conflictResult = { - merged, - failed: [...failed, ...branches.slice(merged.length + failed.length).map(b => b.branchName)], - conflict: `${branchName}: ${stderr}`, - }; - break; - } - - merged.push(branchName); - } - } finally { - if (didStash) { - try { - await git.stash.pop(repoRoot, { index: true }); - } catch { - // Stash-pop conflicts mean the replayed changes clash with the user's - // uncommitted edits. Treat this as a merge failure so the caller preserves - // recovery branches instead of reporting success and deleting them. - logger.warn("Failed to restore stashed changes after task merge; stash entry preserved"); - if (!conflictResult) { + await git.cherryPick(repoRoot, branchName); + } catch (err) { + try { + await git.cherryPick.abort(repoRoot); + } catch { + /* no state to abort */ + } + const stderr = + err instanceof git.GitCommandError + ? err.result.stderr.trim() + : err instanceof Error + ? err.message + : String(err); + failed.push(branchName); conflictResult = { merged, - failed: merged, - conflict: - "stash pop: cherry-picked changes conflict with uncommitted edits. Run `git stash pop` and resolve manually.", + failed: [...failed, ...branches.slice(merged.length + failed.length).map(b => b.branchName)], + conflict: `${branchName}: ${stderr}`, }; + break; + } + + merged.push(branchName); + } + } finally { + if (didStash) { + try { + await git.stash.pop(repoRoot, { index: true }); + } catch { + // Stash-pop conflicts mean the replayed changes clash with the user's + // uncommitted edits. The cherry-picked commits are already on HEAD, so + // the merged branches DID land — report them as merged and surface the + // stash conflict separately instead of claiming they are unmerged. + logger.warn("Failed to restore stashed changes after task merge; stash entry preserved"); + const stashConflict = + "stash pop: cherry-picked changes conflict with uncommitted edits. The merged commits are on HEAD; run `git stash pop` and resolve manually."; + if (conflictResult) { + conflictResult.stashConflict = stashConflict; + } else { + conflictResult = { merged, failed: [], stashConflict }; + } } } } - } - return conflictResult ?? { merged, failed }; + return conflictResult ?? { merged, failed }; + }); } /** Clean up temporary task branches. */ diff --git a/packages/coding-agent/test/task/commands.test.ts b/packages/coding-agent/test/task/commands.test.ts new file mode 100644 index 000000000..b85206ffb --- /dev/null +++ b/packages/coding-agent/test/task/commands.test.ts @@ -0,0 +1,18 @@ +import { describe, expect, it } from "bun:test"; +import { expandCommand, type WorkflowCommand } from "@oh-my-pi/pi-coding-agent/task/commands"; + +function makeCommand(instructions: string): WorkflowCommand { + return { name: "test", description: "test", instructions, source: "project", filePath: "test.md" }; +} + +describe("expandCommand", () => { + it("substitutes $@ with the input", () => { + expect(expandCommand(makeCommand("Do: $@ and again $@"), "fix the bug")).toBe( + "Do: fix the bug and again fix the bug", + ); + }); + + it("keeps $-patterns in user input literal", () => { + expect(expandCommand(makeCommand("Run $@"), "echo $$ $& $' $` $@")).toBe("Run echo $$ $& $' $` $@"); + }); +});