feat(coding-agent/core): added isolated task execution with git worktree support

- Added isolated task execution feature with git worktree support.
- Implemented automatic patch generation and application for isolated tasks.
- Added worktree management with baseline capture and delta patching.
- Enhanced task result display to show isolation status and patch information.
This commit is contained in:
can1357
2026-01-21 23:57:33 +01:00
parent 60d51c99cd
commit da60f86f24
9 changed files with 404 additions and 18 deletions
+5
View File
@@ -1,6 +1,11 @@
# Changelog
## [Unreleased]
### Added
- Added `isolated` option to run tasks in isolated git worktrees
- Added automatic patch generation and application for isolated task execution
- Added worktree management for isolated task execution with baseline capture and delta patching
## [7.0.0] - 2026-01-21
### Added
@@ -41,6 +41,7 @@ import type {
/** Options for worker execution */
export interface ExecutorOptions {
cwd: string;
worktree?: string;
agent: AgentDefinition;
task: string;
description?: string;
@@ -205,6 +206,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
task,
index,
taskId,
worktree,
context,
modelOverride,
thinkingLevel,
@@ -689,6 +691,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
type: "start",
payload: {
cwd,
worktree,
task: fullTask,
systemPrompt: agent.systemPrompt,
model: resolvedModel,
@@ -13,11 +13,12 @@
* - Session artifacts for debugging
*/
import { mkdir, rm } from "node:fs/promises";
import { mkdir, rm, stat } from "node:fs/promises";
import { tmpdir } 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 { $ } from "bun";
import { nanoid } from "nanoid";
import type { Theme } from "../../../modes/interactive/theme/theme";
import taskDescriptionTemplate from "../../../prompts/tools/task.md" with { type: "text" };
@@ -39,6 +40,14 @@ import {
type TaskToolDetails,
taskSchema,
} from "./types";
import {
applyBaseline,
captureBaseline,
captureDeltaPatch,
cleanupWorktree,
ensureWorktree,
getRepoRoot,
} from "./worktree";
// Import review tools for side effects (registers subagent tool handlers)
import "../review";
@@ -153,7 +162,8 @@ export class TaskTool implements AgentTool<typeof taskSchema, TaskToolDetails, T
): Promise<AgentToolResult<TaskToolDetails>> {
const startTime = Date.now();
const { agents, projectAgentsDir } = await discoverAgents(this.session.cwd);
const { agent: agentName, context, model, output: outputSchema } = params;
const { agent: agentName, context, model, output: outputSchema, isolated } = params;
const isIsolated = isolated === true;
const isDefaultModelAlias = (value: string | undefined): boolean => {
if (!value) return true;
@@ -283,6 +293,30 @@ export class TaskTool implements AgentTool<typeof taskSchema, TaskToolDetails, T
};
}
let repoRoot: string | null = null;
let baseline = null as Awaited<ReturnType<typeof captureBaseline>> | 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,
},
};
}
}
// Derive artifacts directory
const sessionFile = this.session.getSessionFile();
const artifactsDir = sessionFile ? sessionFile.slice(0, -6) : null;
@@ -377,11 +411,8 @@ export class TaskTool implements AgentTool<typeof taskSchema, TaskToolDetails, T
}
emitProgress();
// Execute in parallel with concurrency limit
const { results: partialResults, aborted } = await mapWithConcurrencyLimit(
tasksWithContext,
MAX_CONCURRENCY,
async (task, index) => {
const runTask = async (task: (typeof tasksWithContext)[number], index: number) => {
if (!isIsolated) {
return runSubprocess({
cwd: this.session.cwd,
agent,
@@ -411,7 +442,83 @@ export class TaskTool implements AgentTool<typeof taskSchema, TaskToolDetails, T
settingsManager: this.session.settingsManager,
mcpManager: this.session.mcpManager,
});
},
}
const taskStart = Date.now();
let worktreeDir: string | undefined;
try {
if (!repoRoot || !baseline) {
throw new Error("Isolated task execution not initialized.");
}
worktreeDir = await ensureWorktree(repoRoot, task.taskId);
await applyBaseline(worktreeDir, baseline);
const result = await runSubprocess({
cwd: this.session.cwd,
worktree: worktreeDir,
agent,
task: task.task,
description: task.description,
index,
taskId: task.taskId,
context: undefined, // Already prepended above
modelOverride,
thinkingLevel: thinkingLevelOverride,
outputSchema: effectiveOutputSchema,
sessionFile,
persistArtifacts: !!artifactsDir,
artifactsDir: effectiveArtifactsDir,
enableLsp: false,
signal,
eventBus: undefined,
onProgress: (progress) => {
progressMap.set(index, {
...structuredClone(progress),
vars: tasksWithContext[index]?.vars,
});
emitProgress();
},
authStorage: this.session.authStorage,
modelRegistry: this.session.modelRegistry,
settingsManager: this.session.settingsManager,
mcpManager: this.session.mcpManager,
});
const patch = await captureDeltaPatch(worktreeDir, baseline);
const patchPath = path.join(effectiveArtifactsDir, `${task.taskId}.patch`);
await Bun.write(patchPath, patch);
return {
...result,
patchPath,
};
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
return {
index,
taskId: task.taskId,
agent: agent.name,
agentSource: agent.source,
task: task.task,
description: task.description,
exitCode: 1,
output: "",
stderr: message,
truncated: false,
durationMs: Date.now() - taskStart,
tokens: 0,
modelOverride,
error: message,
};
} finally {
if (worktreeDir) {
await cleanupWorktree(worktreeDir);
}
}
};
// Execute in parallel with concurrency limit
const { results: partialResults, aborted } = await mapWithConcurrencyLimit(
tasksWithContext,
MAX_CONCURRENCY,
runTask,
signal,
);
@@ -456,10 +563,75 @@ export class TaskTool implements AgentTool<typeof taskSchema, TaskToolDetails, T
// Collect output paths (artifacts already written by executor in real-time)
const outputPaths: string[] = [];
const patchPaths: string[] = [];
for (const result of results) {
if (result.outputPath) {
outputPaths.push(result.outputPath);
}
if (result.patchPath) {
patchPaths.push(result.patchPath);
}
}
let patchApplySummary = "";
let patchesApplied: boolean | null = null;
if (isIsolated) {
const patchesInOrder = results.map((result) => result.patchPath).filter(Boolean) as string[];
const missingPatch = results.some((result) => !result.patchPath);
if (!repoRoot || missingPatch) {
patchesApplied = false;
} else {
const patchStats = await Promise.all(
patchesInOrder.map(async (patchPath) => ({
patchPath,
size: (await stat(patchPath)).size,
})),
);
const nonEmptyPatches = patchStats.filter((patch) => patch.size > 0).map((patch) => patch.patchPath);
if (nonEmptyPatches.length === 0) {
patchesApplied = true;
} else {
const patchTexts = await Promise.all(
nonEmptyPatches.map(async (patchPath) => Bun.file(patchPath).text()),
);
const combinedPatch = patchTexts.map((text) => (text.endsWith("\n") ? text : `${text}\n`)).join("");
if (!combinedPatch.trim()) {
patchesApplied = true;
} else {
const combinedPatchPath = path.join(tmpdir(), `omp-task-combined-${nanoid()}.patch`);
try {
await Bun.write(combinedPatchPath, combinedPatch);
const checkResult = await $`git apply --check --binary ${combinedPatchPath}`
.cwd(repoRoot)
.quiet()
.nothrow();
if (checkResult.exitCode !== 0) {
patchesApplied = false;
} else {
const applyResult = await $`git apply --binary ${combinedPatchPath}`
.cwd(repoRoot)
.quiet()
.nothrow();
patchesApplied = applyResult.exitCode === 0;
}
} finally {
await rm(combinedPatchPath, { force: true });
}
}
}
}
if (patchesApplied) {
patchApplySummary = "\n\nApplied patches: yes";
} else {
const notification =
"<system-notification>Patches were not applied and must be handled manually.</system-notification>";
const patchList =
patchPaths.length > 0
? `\n\nPatch artifacts:\n${patchPaths.map((patch) => `- ${patch}`).join("\n")}`
: "";
patchApplySummary = `\n\n${notification}${patchList}`;
}
}
// Build final output - match plugin format
@@ -486,10 +658,12 @@ export class TaskTool implements AgentTool<typeof taskSchema, TaskToolDetails, T
const cancelledNote = aborted && cancelledCount > 0 ? ` (${cancelledCount} cancelled)` : "";
const summary = `${successCount}/${results.length} succeeded${cancelledNote} [${formatDuration(
totalDuration,
)}]\n\n${summaries.join("\n\n---\n\n")}${outputHint}${schemaNote}`;
)}]\n\n${summaries.join("\n\n---\n\n")}${outputHint}${schemaNote}${patchApplySummary}`;
// Cleanup temp directory if used
if (tempArtifactsDir) {
const shouldCleanupTempArtifacts =
tempArtifactsDir && (!isIsolated || patchesApplied === true || patchesApplied === null);
if (shouldCleanupTempArtifacts) {
await rm(tempArtifactsDir, { recursive: true, force: true });
}
@@ -380,6 +380,7 @@ export function renderCall(args: TaskParams, theme: Theme): Component {
const branch = theme.fg("dim", theme.tree.branch);
const last = theme.fg("dim", theme.tree.last);
const vertical = theme.fg("dim", theme.tree.vertical);
const showIsolated = args.isolated === true;
if (hasContext) {
lines.push(` ${branch} ${theme.fg("dim", "Context")}`);
@@ -387,11 +388,18 @@ export function renderCall(args: TaskParams, theme: Theme): Component {
const content = line ? theme.fg("muted", line) : "";
lines.push(` ${vertical} ${content}`);
}
lines.push(` ${last} ${theme.fg("dim", "Tasks")}: ${theme.fg("muted", `${args.tasks.length} agents`)}`);
const taskPrefix = showIsolated ? branch : last;
lines.push(` ${taskPrefix} ${theme.fg("dim", "Tasks")}: ${theme.fg("muted", `${args.tasks.length} agents`)}`);
if (showIsolated) {
lines.push(` ${last} ${theme.fg("dim", "Isolated")}: ${theme.fg("muted", "true")}`);
}
return new Text(lines.join("\n"), 0, 0);
}
lines.push(`${theme.fg("dim", "Tasks")}: ${theme.fg("muted", `${args.tasks.length} agents`)}`);
if (showIsolated) {
lines.push(`${theme.fg("dim", "Isolated")}: ${theme.fg("muted", "true")}`);
}
return new Text(lines.join("\n"), 0, 0);
}
@@ -722,6 +730,10 @@ function renderAgentResult(result: SingleResult, isLast: boolean, expanded: bool
lines.push(...renderOutputSection(result.output, continuePrefix, expanded, theme, 3, 12));
}
if (result.patchPath && !aborted && result.exitCode === 0) {
lines.push(`${continuePrefix}${theme.fg("dim", `Patch: ${result.patchPath}`)}`);
}
// Error message
if (result.error && !success) {
lines.push(`${continuePrefix}${theme.fg("error", truncate(result.error, 70, theme.format.ellipsis))}`);
@@ -739,6 +751,7 @@ export function renderResult(
theme: Theme,
): Component {
const { expanded, isPartial, spinnerFrame } = options;
const fallbackText = result.content.find((c) => c.type === "text")?.text ?? "";
const details = result.details;
if (!details) {
@@ -785,7 +798,22 @@ export function renderResult(
}
if (lines.length === 0) {
return new Text(theme.fg("dim", "No results"), 0, 0);
const text = fallbackText.trim() ? fallbackText : "No results";
return new Text(theme.fg("dim", truncate(text, 140, theme.format.ellipsis)), 0, 0);
}
if (fallbackText.trim()) {
const summaryLines = fallbackText.split("\n");
const markerIndex = summaryLines.findIndex(
(line) => line.includes("<system-notification>") || line.startsWith("Applied patches:"),
);
if (markerIndex >= 0) {
const extra = summaryLines.slice(markerIndex);
for (const line of extra) {
if (!line.trim()) continue;
lines.push(theme.fg("dim", line));
}
}
}
const indented = lines.map((line) => (line.trim() ? ` ${line}` : ""));
@@ -63,6 +63,11 @@ export const taskSchema = Type.Object({
description: "Model override for all tasks (fuzzy matching, e.g. 'sonnet', 'opus')",
}),
),
isolated: Type.Optional(
Type.Boolean({
description: "Run each task in an isolated git worktree",
}),
),
output: Type.Optional(
Type.Record(Type.String(), Type.Unknown(), {
description: "JTD schema for structured subagent output",
@@ -159,6 +164,8 @@ export interface SingleResult {
usage?: Usage;
/** Output path for the task result */
outputPath?: string;
/** Patch path for isolated worktree output */
patchPath?: string;
/** Data extracted by registered subprocess tool handlers (keyed by tool name) */
extractedToolData?: Record<string, unknown[]>;
/** Output metadata for Output tool integration */
@@ -84,6 +84,7 @@ export interface LspToolCallResponse {
export interface SubagentWorkerStartPayload {
cwd: string;
worktree?: string;
task: string;
systemPrompt: string;
model?: string;
@@ -538,9 +538,6 @@ async function runTask(runState: RunState, payload: SubagentWorkerStartPayload):
// Check for pre-start abort
checkAbort();
// Set working directory (CLI does this implicitly)
process.chdir(payload.cwd);
// Use serialized auth/models if provided, otherwise discover from disk
let authStorage: AuthStorage;
let modelRegistry: ModelRegistry;
@@ -574,7 +571,7 @@ async function runTask(runState: RunState, payload: SubagentWorkerStartPayload):
// Create session manager (equivalent to CLI's --session or --no-session)
const sessionManager = payload.sessionFile
? await SessionManager.open(payload.sessionFile)
: SessionManager.inMemory(payload.cwd);
: SessionManager.inMemory(payload.worktree ?? payload.cwd);
checkAbort();
// Use serialized settings if provided, otherwise use empty in-memory settings
@@ -585,9 +582,12 @@ async function runTask(runState: RunState, payload: SubagentWorkerStartPayload):
// Note: hasUI: false disables interactive features
const completionInstruction =
"When finished, call the complete tool exactly once. Do not end with a plain-text final answer.";
const worktreeNotice = payload.worktree
? `You will work under this working tree: ${payload.worktree}. CRITICAL: Do not touch the original repository; only make changes inside this worktree.`
: "";
const { session } = await createAgentSession({
cwd: payload.cwd,
cwd: payload.worktree ?? payload.cwd,
authStorage,
modelRegistry,
settingsManager,
@@ -597,7 +597,8 @@ async function runTask(runState: RunState, payload: SubagentWorkerStartPayload):
outputSchema: payload.outputSchema,
requireCompleteTool: true,
// Append system prompt (equivalent to CLI's --append-system-prompt)
systemPrompt: (defaultPrompt) => `${defaultPrompt}\n\n${payload.systemPrompt}\n\n${completionInstruction}`,
systemPrompt: (defaultPrompt) =>
`${defaultPrompt}\n\n${payload.systemPrompt}\n\n${worktreeNotice}\n\n${completionInstruction}`,
sessionManager,
hasUI: false,
// Pass spawn restrictions to nested tasks
@@ -0,0 +1,166 @@
import { randomUUID } from "node:crypto";
import { cp, mkdir, rm } from "node:fs/promises";
import { homedir, tmpdir } from "node:os";
import path from "node:path";
import { $ } from "bun";
export interface WorktreeBaseline {
repoRoot: string;
staged: string;
unstaged: string;
untracked: string[];
}
export function getEncodedProjectName(cwd: string): string {
return `--${cwd.replace(/^[/\\]/, "").replace(/[/\\:]/g, "-")}--`;
}
export async function getRepoRoot(cwd: string): Promise<string> {
const result = await $`git rev-parse --show-toplevel`.cwd(cwd).quiet().nothrow();
if (result.exitCode !== 0) {
throw new Error("Git repository not found for isolated task execution.");
}
const repoRoot = result.text().trim();
if (!repoRoot) {
throw new Error("Git repository root could not be resolved for isolated task execution.");
}
return repoRoot;
}
export async function ensureWorktree(baseCwd: string, taskId: string): Promise<string> {
const repoRoot = await getRepoRoot(baseCwd);
const encodedProject = getEncodedProjectName(repoRoot);
const worktreeDir = path.join(homedir(), ".omp", "wt", encodedProject, taskId);
await mkdir(path.dirname(worktreeDir), { recursive: true });
await $`git worktree remove -f ${worktreeDir}`.cwd(repoRoot).quiet().nothrow();
await rm(worktreeDir, { recursive: true, force: true });
await $`git worktree add --detach ${worktreeDir} HEAD`.cwd(repoRoot).quiet();
return worktreeDir;
}
export async function captureBaseline(repoRoot: string): Promise<WorktreeBaseline> {
const staged = await $`git diff --cached --binary`.cwd(repoRoot).quiet().text();
const unstaged = await $`git diff --binary`.cwd(repoRoot).quiet().text();
const untrackedRaw = await $`git ls-files --others --exclude-standard`.cwd(repoRoot).quiet().text();
const untracked = untrackedRaw
.split("\n")
.map((line) => line.trim())
.filter((line) => line.length > 0);
return {
repoRoot,
staged,
unstaged,
untracked,
};
}
async function writeTempPatchFile(patch: string): Promise<string> {
const tempPath = path.join(tmpdir(), `omp-task-patch-${randomUUID()}.patch`);
await Bun.write(tempPath, patch);
return tempPath;
}
async function applyPatch(
cwd: string,
patch: string,
options?: { cached?: boolean; env?: Record<string, string> },
): Promise<void> {
if (!patch.trim()) return;
const tempPath = await writeTempPatchFile(patch);
try {
const command = options?.cached ? $`git apply --cached --binary ${tempPath}` : $`git apply --binary ${tempPath}`;
let runner = command.cwd(cwd).quiet();
if (options?.env) {
runner = runner.env(options.env);
}
await runner;
} finally {
await rm(tempPath, { force: true });
}
}
export async function applyBaseline(worktreeDir: string, baseline: WorktreeBaseline): Promise<void> {
await applyPatch(worktreeDir, baseline.staged, { cached: true });
await applyPatch(worktreeDir, baseline.staged);
await applyPatch(worktreeDir, baseline.unstaged);
for (const entry of baseline.untracked) {
const source = path.join(baseline.repoRoot, entry);
const destination = path.join(worktreeDir, entry);
const exists = await Bun.file(source).exists();
if (!exists) continue;
await mkdir(path.dirname(destination), { recursive: true });
await cp(source, destination, { recursive: true });
}
}
async function applyPatchToIndex(cwd: string, patch: string, indexFile: string): Promise<void> {
if (!patch.trim()) return;
const tempPath = await writeTempPatchFile(patch);
try {
await $`git apply --cached --binary ${tempPath}`
.cwd(cwd)
.env({
GIT_INDEX_FILE: indexFile,
})
.quiet();
} finally {
await rm(tempPath, { force: true });
}
}
async function listUntracked(cwd: string): Promise<string[]> {
const raw = await $`git ls-files --others --exclude-standard`.cwd(cwd).quiet().text();
return raw
.split("\n")
.map((line) => line.trim())
.filter((line) => line.length > 0);
}
export async function captureDeltaPatch(worktreeDir: string, baseline: WorktreeBaseline): Promise<string> {
const tempIndex = path.join(tmpdir(), `omp-task-index-${randomUUID()}`);
try {
await $`git read-tree HEAD`.cwd(worktreeDir).env({
GIT_INDEX_FILE: tempIndex,
});
await applyPatchToIndex(worktreeDir, baseline.staged, tempIndex);
await applyPatchToIndex(worktreeDir, baseline.unstaged, tempIndex);
const diff = await $`git diff --binary`
.cwd(worktreeDir)
.env({
GIT_INDEX_FILE: tempIndex,
})
.quiet()
.text();
const currentUntracked = await listUntracked(worktreeDir);
const baselineUntracked = new Set(baseline.untracked);
const newUntracked = currentUntracked.filter((entry) => !baselineUntracked.has(entry));
if (newUntracked.length === 0) return diff;
const untrackedDiffs = await Promise.all(
newUntracked.map((entry) =>
$`git diff --binary --no-index /dev/null ${entry}`.cwd(worktreeDir).quiet().nothrow().text(),
),
);
return `${diff}${diff && !diff.endsWith("\n") ? "\n" : ""}${untrackedDiffs.join("\n")}`;
} finally {
await rm(tempIndex, { force: true });
}
}
export async function cleanupWorktree(dir: string): Promise<void> {
try {
const commonDirRaw = await $`git rev-parse --git-common-dir`.cwd(dir).quiet().nothrow().text();
const commonDir = commonDirRaw.trim();
if (commonDir) {
const resolvedCommon = path.resolve(dir, commonDir);
const repoRoot = path.dirname(resolvedCommon);
await $`git worktree remove -f ${dir}`.cwd(repoRoot).quiet().nothrow();
}
} finally {
await rm(dir, { recursive: true, force: true });
}
}
@@ -51,6 +51,7 @@ Agents with "Output: structured" have a fixed schema enforced via frontmatter; y
- `agent`: Agent type to use for all tasks
- `context`: Template with `{{placeholders}}` for multi-task. Each placeholder is filled from task vars.
- `model`: (optional) Model override for all tasks (fuzzy matching, e.g., "sonnet", "opus")
- `isolated`: (optional) Run each task in its own git worktree and return patches; patches are applied only if all apply cleanly.
- `tasks`: Array of `{id, description, vars}` - tasks to run in parallel (max {{MAX_PARALLEL_TASKS}}, {{MAX_CONCURRENCY}} concurrent)
- `id`: Short CamelCase identifier for display (max 20 chars, e.g., "SessionStore", "LspRefactor")
- `description`: Short human-readable description of what the task does