141 lines
5.7 KiB
TypeScript
141 lines
5.7 KiB
TypeScript
/**
|
|
* Parallel delegation tool — splits a large task into sub-tasks that run
|
|
* concurrently across multiple subagents. Mirrors `tools/parallel_delegate.rs`.
|
|
*/
|
|
import type { JsonValue } from "@zesdex/domain";
|
|
import { type Tool, type ToolCtx } from "./mod.ts";
|
|
import { argStr, optBool, optInt } from "./util.ts";
|
|
|
|
/** Delegate a large task to multiple subagents running in parallel. */
|
|
export class ParallelDelegate implements Tool {
|
|
name = "parallel_delegate";
|
|
description = "Delegate a large task to multiple subagents running in parallel for 3x faster completion";
|
|
|
|
parameters = {
|
|
type: "object",
|
|
properties: {
|
|
task: { type: "string", description: "The task to be split and delegated to parallel agents" },
|
|
directives: {
|
|
type: "array",
|
|
items: {
|
|
type: "object",
|
|
properties: {
|
|
directive: { type: "string" },
|
|
access: { type: "string", enum: ["read", "write", "full"] },
|
|
},
|
|
required: ["directive"],
|
|
},
|
|
},
|
|
synthesize: { type: "boolean", default: true },
|
|
max_parallel: { type: "integer", default: 3, description: "Maximum number of parallel agents (default: 3, max: 8)" },
|
|
},
|
|
required: ["task"],
|
|
};
|
|
|
|
async run(ctx: ToolCtx, args: JsonValue): Promise<string> {
|
|
const task = argStr(args, "task");
|
|
const synthesize = optBool(args, "synthesize", true);
|
|
const maxParallel = Math.min(optInt(args, "max_parallel", 3), 8);
|
|
|
|
let directives: Array<{ directive: string; access: string }>;
|
|
const explicit = args && typeof args === "object" && !Array.isArray(args)
|
|
? (args as Record<string, JsonValue>).directives
|
|
: undefined;
|
|
if (Array.isArray(explicit)) {
|
|
directives = explicit
|
|
.filter((d): d is Record<string, JsonValue> => typeof d === "object" && d !== null)
|
|
.map((d) => ({
|
|
directive: typeof d.directive === "string" ? d.directive : "",
|
|
access: typeof d.access === "string" ? d.access : "write",
|
|
}))
|
|
.filter((d) => d.directive !== "");
|
|
} else {
|
|
// Let the LLM decide the decomposition instead of guessing keywords.
|
|
directives = await aiDecomposeTask(task, maxParallel);
|
|
}
|
|
|
|
if (directives.length === 0) {
|
|
throw new Error("no directives could be derived for the task");
|
|
}
|
|
|
|
// Delegate to subagent engine if available, else report the breakdown.
|
|
try {
|
|
const { runParallelDelegation } = await import("../../../subagent/infrastructure/delegate.ts");
|
|
return await runParallelDelegation(task, directives, synthesize, ctx);
|
|
} catch (e) {
|
|
const msg = (e as Error).message;
|
|
if (msg.includes("not yet") || msg.includes("Cannot find")) {
|
|
let out = `## Parallel Delegation Complete\n\n**Task:** ${task}\n**Parallel agents:** ${directives.length}\n\n`;
|
|
directives.forEach((d, i) => { out += `---\n### Agent ${i}: [${d.access}]\n\n${d.directive}\n`; });
|
|
return out;
|
|
}
|
|
throw e;
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Have the LLM decompose a task into focused, non-overlapping parallel
|
|
* sub-directives. This is the AI deciding how to split the work — no keyword
|
|
* guessing. On LLM failure (no provider / error), degrades to a single honest
|
|
* full-task directive rather than fabricated "part N" stubs.
|
|
*/
|
|
async function aiDecomposeTask(
|
|
task: string,
|
|
maxParallel: number,
|
|
): Promise<Array<{ directive: string; access: string }>> {
|
|
try {
|
|
const { resolveConfig } = await import("../../../subagent/infrastructure/config_resolver.ts");
|
|
const { buildProviderService } = await import("../../../subagent/infrastructure/http_provider.ts");
|
|
const { baseUrl, apiKey, model } = await resolveConfig();
|
|
const svc = buildProviderService(baseUrl, apiKey, model);
|
|
|
|
const systemPrompt =
|
|
"You decompose a large task into independent, non-overlapping subtasks that can be worked on " +
|
|
"in parallel. Return a JSON array of up to the requested count of objects, each with a " +
|
|
"'directive' (a concrete, self-contained action) and an 'access' tier of exactly one of " +
|
|
"read | write | full. read = inspect/search only; write = edit files; full = edit + run shell. " +
|
|
"Prefer fewer, genuinely parallel directives over many that touch the same files. " +
|
|
"Return ONLY the JSON array, no prose.";
|
|
|
|
const { message } = await svc.chat(
|
|
[
|
|
{ role: "system", content: systemPrompt },
|
|
{ role: "user", content: `Task: ${task}\nMax subtasks: ${maxParallel}\n` },
|
|
],
|
|
undefined,
|
|
1024,
|
|
0.3,
|
|
);
|
|
|
|
const parsed = parseDirectivesJson(message.content ?? "");
|
|
if (parsed.length > 0) return parsed.slice(0, maxParallel);
|
|
} catch {
|
|
/* fall through to the single-directive fallback */
|
|
}
|
|
// Honest fallback: run the whole task as one agent rather than guessing a split.
|
|
return [{ directive: task, access: "full" }];
|
|
}
|
|
|
|
/** Parse an LLM JSON array of directives, tolerant of stray prose/markdown fences. */
|
|
export function parseDirectivesJson(raw: string): Array<{ directive: string; access: string }> {
|
|
const fenced = raw.match(/```(?:json)?\s*([\s\S]*?)```/i);
|
|
const body = fenced ? fenced[1]! : raw;
|
|
const start = body.indexOf("[");
|
|
const end = body.lastIndexOf("]");
|
|
if (start < 0 || end <= start) return [];
|
|
try {
|
|
const arr = JSON.parse(body.slice(start, end + 1)) as unknown;
|
|
if (!Array.isArray(arr)) return [];
|
|
return arr
|
|
.filter((d): d is Record<string, unknown> => typeof d === "object" && d !== null)
|
|
.map((d) => ({
|
|
directive: typeof d.directive === "string" ? d.directive.trim() : "",
|
|
access: ["read", "write", "full"].includes(String(d.access)) ? String(d.access) : "write",
|
|
}))
|
|
.filter((d) => d.directive !== "");
|
|
} catch {
|
|
return [];
|
|
}
|
|
}
|