import { runRepoSearchSubagent } from "../subagent/repo-search/runner.js"; import type { RepoSearchRunConfig } from "../subagent/repo-search/types.js"; import { runSubagent } from "../shared/subagent/run.js"; import { buildWorkerTask, MULTI_TASK_WORKER_PROMPT } from "./prompt.js"; import type { MultiTaskBatch, MultiTaskWorker } from "./types.js"; const MAX_VISIBLE_TOOL_CALLS = 8; function researchTask(worker: MultiTaskWorker): string { return [ worker.task, "", "Repository search scopes:", ...worker.paths.map((path) => `- ${path}`), ].join("\n"); } async function runImplementation( batch: MultiTaskBatch, worker: MultiTaskWorker, extensionPaths: string[], emitProgress: () => void, ): Promise { const result = await runSubagent({ cwd: batch.cwd, title: `Multi Task · ${worker.id}`, model: worker.model, thinkingLevel: worker.thinkingLevel, task: buildWorkerTask(worker.task, worker.paths), systemPrompt: MULTI_TASK_WORKER_PROMPT, capability: "all", availableTools: batch.implementationTools, extensionPaths, loadDefaultResources: true, // Batch progress lives in one aggregated tool card, so workers always // run as managed RPC children regardless of the Repo Search presentation. presentation: "manual", parentSessionId: batch.parentSessionId, keepOpen: batch.keepOpen, signal: worker.controller.signal, env: { PI_MULTI_TASK_ALLOWED_PATHS: JSON.stringify(worker.paths), }, onUpdate: (update) => { worker.progress = update.status; worker.toolCalls = update.toolCalls.slice(-MAX_VISIBLE_TOOL_CALLS); worker.subagentId = update.subagentId; worker.reusable = update.reusable; worker.turn = update.turn; emitProgress(); }, }); worker.output = result.output; worker.runDir = result.runDir; worker.subagentId = result.subagentId; worker.reusable = result.reusable; worker.turn = result.turn; } async function runResearch( batch: MultiTaskBatch, worker: MultiTaskWorker, config: RepoSearchRunConfig, emitProgress: () => void, ): Promise { const result = await runRepoSearchSubagent({ cwd: batch.cwd, task: researchTask(worker), config: { ...config, presentation: "manual" }, keepOpen: batch.keepOpen, parentSessionId: batch.parentSessionId, signal: worker.controller.signal, onUpdate: (details) => { worker.progress = "searching"; worker.toolCalls = details.toolCalls.slice(-MAX_VISIBLE_TOOL_CALLS); worker.subagentId = details.subagentId; worker.reusable = details.reusable; worker.turn = details.turn; emitProgress(); }, }); worker.output = result.details.output; worker.runDir = result.details.runDir; worker.subagentId = result.details.subagentId; worker.reusable = result.details.reusable; worker.turn = result.details.turn; } export async function executeWorker(options: { batch: MultiTaskBatch; worker: MultiTaskWorker; extensionPaths: string[]; researchConfig?: RepoSearchRunConfig; emitProgress: () => void; }): Promise { const { batch, worker, extensionPaths, researchConfig, emitProgress } = options; if (batch.cancelRequested) { worker.status = "cancelled"; worker.progress = "cancelled"; emitProgress(); return; } worker.status = "running"; worker.startedAt = new Date().toISOString(); emitProgress(); try { if (worker.kind === "research") { if (!researchConfig) throw new Error("research worker 缺少 Repo Search 配置"); await runResearch(batch, worker, researchConfig, emitProgress); } else { await runImplementation(batch, worker, extensionPaths, emitProgress); } worker.status = "completed"; worker.progress = "completed"; } catch (error) { worker.status = batch.cancelRequested ? "cancelled" : "failed"; worker.progress = worker.status; worker.error = error instanceof Error ? error.message : String(error); } finally { worker.completedAt = new Date().toISOString(); emitProgress(); } }