TODO Next: subagent parallelism (task tasks batch)

task now accepts { description, prompt, kind } or
{ tasks: TaskSpec[], kind? } (up to 8) and fans out with
Promise.all. Each subagent keeps its own context window and
the progress panel receives start/step/result/end per id, so
independent searches overlap in wall time instead of queueing.

- src/subagent.ts: runOne extracted, createTaskTool uses union
  schema (single | batch), batch validates worker channel once
  then Promise.all, results joined as headings.
- docs/agents.md: delegation section notes tasks batch.
- 800 pass (added subagent-parallel.test.ts).
This commit is contained in:
asepharyana
2026-09-09 10:51:01 +07:00
parent 2587e03beb
commit 7edaeafca0
5 changed files with 262 additions and 96 deletions
+1 -1
View File
@@ -14,7 +14,7 @@ TODO.md status 2026-09-08:
1. **Now flip + tests** — patch TODO.md [x], test: hot-reload skill mid-session callable, pruned decision recoverable. Verifikasi: bun test. 1. **Now flip + tests** — patch TODO.md [x], test: hot-reload skill mid-session callable, pruned decision recoverable. Verifikasi: bun test.
2. **Derive tool-name lists** — tandai mutating di definisi tool (tools.ts/tools-extra.ts/tools-git.ts/tools-net.ts), TOOL_SETS & MUTATING_TOOLS derive + test coverage. Risiko: silently ungated write jika lupa. 2. **Derive tool-name lists** — tandai mutating di definisi tool (tools.ts/tools-extra.ts/tools-git.ts/tools-net.ts), TOOL_SETS & MUTATING_TOOLS derive + test coverage. Risiko: silently ungated write jika lupa.
3. **MCP tanpa schema tax** — DONE: phi 3 meta-tools mcp_list/mcp_inspect/mcp_call (mcpExpose=phi/direct/auto), prompt hanya nama server, permission+guard lewat built-in, keep direct registration untuk server 2-tool. Test: server 20-tool 0 schema sampai mcp_call (mcp.test.ts). 3. **MCP tanpa schema tax** — DONE: phi 3 meta-tools mcp_list/mcp_inspect/mcp_call (mcpExpose=phi/direct/auto), prompt hanya nama server, permission+guard lewat built-in, keep direct registration untuk server 2-tool. Test: server 20-tool 0 schema sampai mcp_call (mcp.test.ts).
4. **Subagent parallelism** — task terima beberapa investigations, run Promise.all dengan panel fan-out, test overlap waktu. 4. **Subagent parallelism** — DONE: task `tasks` union + Promise.all fan-out (8 cap), panel per-subagent, test overlap waktu (subagent-parallel.test.ts).
5. **Undo a turn** — snapshot file-tool edits pre-prompt (cap 100), /undo restores files+conversation atau keduanya, /redo, bash tidak ter-snapshot (docs), test edit reverted + record hilang. 5. **Undo a turn** — snapshot file-tool edits pre-prompt (cap 100), /undo restores files+conversation atau keduanya, /redo, bash tidak ter-snapshot (docs), test edit reverted + record hilang.
6. **Maintenance polish** — pricing source+date, estimateTokens label everywhere /cost, listPaths notice staleness / refresh, MUTATING derive dari #2. 6. **Maintenance polish** — pricing source+date, estimateTokens label everywhere /cost, listPaths notice staleness / refresh, MUTATING derive dari #2.
+2 -2
View File
@@ -61,8 +61,8 @@ have.
Two independent searches run sequentially. The panel already renders several agents; the loop Two independent searches run sequentially. The panel already renders several agents; the loop
does not fan out. does not fan out.
- [ ] `task` accepts several investigations and runs them together - [x] `task` accepts several investigations and runs them together (`src/subagent.ts`: `tasks: TaskSpec[]` union, `runOne` + `Promise.all` fan-out up to 8, panel emits start/step/result/end per subagent)
- [ ] Test: two delegated searches overlap in time rather than queueing - [x] Test: two delegated searches overlap in time rather than queueing (`test/subagent-parallel.test.ts`: delayed greps overlap < 2*delay, both headings in one tool result)
### Undo a turn ### Undo a turn
+1 -1
View File
@@ -175,7 +175,7 @@ around it. The tool descriptions say the same thing, so the rule survives compac
accident. The read-only kinds work everywhere. No subagent holds `web_fetch`; network access accident. The read-only kinds work everywhere. No subagent holds `web_fetch`; network access
stays with the main agent, where the approval prompt says what it is for. stays with the main agent, where the approval prompt says what it is for.
When not to delegate: a single grep, or anything you must supervise step by step — keep that in For independent pieces of work, pass `tasks: [{description, prompt}, …]` to run them in parallel instead of calling `task` several times — each subagent keeps its own window and the panel fans out. When not to delegate: a single grep, or anything you must supervise step by step — keep that in
your own turn, where every call is on screen. A worker wins when the intermediate steps are your own turn, where every call is on screen. A worker wins when the intermediate steps are
noise: a mechanical rename across twenty files, a test scaffold written to match an existing noise: a mechanical rename across twenty files, a test scaffold written to match an existing
suite, a cleanup whose shape you already know. suite, a cleanup whose shape you already know.
+157 -92
View File
@@ -139,6 +139,96 @@ let counter = 0;
* asked for a worker and got an explorer would be told the task failed for the * asked for a worker and got an explorer would be told the task failed for the
* wrong reason. * wrong reason.
*/ */
export type TaskSpec = { description: string; prompt: string; kind?: SubagentKind };
async function runOne(
spec: TaskSpec,
id: string,
opts: {
model: LanguageModel;
subagentModel?: LanguageModel;
cwd?: string;
maxSteps?: number;
report?: SubagentReporter;
approve?: SubagentApproval;
onUsage?: (usage: { kind: SubagentKind; inputTokens: number; outputTokens: number }) => void;
},
abortSignal?: AbortSignal,
): Promise<string> {
const flavour: SubagentKind = spec.kind ?? 'explore';
if (flavour === 'worker' && !opts.approve) {
throw new Error('The worker kind needs an approval channel, which this session has not provided.');
}
const report = opts.report;
report?.({ type: 'start', id, kind: flavour, description: spec.description });
let steps = 0;
let text = '';
let usedTokens: { inputTokens: number; outputTokens: number } | undefined;
try {
const model = flavour === 'explore' ? (opts.subagentModel ?? opts.model) : opts.model;
const result = streamText({
model,
system: PROMPTS[flavour](opts.cwd ?? process.cwd()),
messages: [{ role: 'user', content: spec.prompt }],
tools: TOOLS[flavour],
stopWhen: isStepCount(opts.maxSteps ?? 20),
...(opts.approve
? {
toolApproval: async ({ toolCall }: { toolCall: { toolName: string; input: unknown } }) => {
const approved = await opts.approve!(toolCall);
return approved
? undefined
: { type: 'denied' as const, reason: 'The user denied this call. Stop and report it.' };
},
}
: {}),
...(abortSignal ? { abortSignal } : {}),
});
const sink = () => {};
void result.responseMessages.then(undefined, sink);
void result.usage.then(undefined, sink);
void result.steps.then(undefined, sink);
void result.finalStep.then(undefined, sink);
void result.finishReason.then(undefined, sink);
for await (const part of result.stream) {
if (part.type === 'tool-call') {
steps++;
report?.({ type: 'step', id, tool: part.toolName, summary: summarize(part.input) });
} else if (part.type === 'tool-result') {
report?.({ type: 'result', id, tool: part.toolName, summary: outcome(part.output), ok: true });
} else if (part.type === 'tool-error') {
const message = part.error instanceof Error ? part.error.message : String(part.error);
report?.({ type: 'result', id, tool: part.toolName, summary: outcome(message), ok: false });
} else if (part.type === 'text-delta') {
text += part.text;
} else if (part.type === 'error') {
const message = part.error instanceof Error ? part.error.message : String(part.error);
throw part.error instanceof Error ? part.error : new Error(message);
}
}
try {
const usage = await result.usage;
usedTokens = { inputTokens: usage.inputTokens ?? 0, outputTokens: usage.outputTokens ?? 0 };
} catch {
// A run that errored before producing usage has nothing to account for.
}
} catch (e) {
const message = e instanceof Error ? e.message : String(e);
report?.({ type: 'error', id, message });
throw e;
}
const trimmed = text.trim();
report?.({ type: 'end', id, ok: trimmed.length > 0, steps });
if (usedTokens) opts.onUsage?.({ kind: flavour, ...usedTokens });
return trimmed || 'Subagent returned no findings.';
}
export function createTaskTool(opts: { export function createTaskTool(opts: {
model: LanguageModel; model: LanguageModel;
/** Cheaper model for `explore`, which is search rather than reasoning. Defaults to `model`. */ /** Cheaper model for `explore`, which is search rather than reasoning. Defaults to `model`. */
@@ -154,6 +244,34 @@ export function createTaskTool(opts: {
onUsage?: (usage: { kind: SubagentKind; inputTokens: number; outputTokens: number }) => void; onUsage?: (usage: { kind: SubagentKind; inputTokens: number; outputTokens: number }) => void;
}) { }) {
const canWrite = opts.approve !== undefined; const canWrite = opts.approve !== undefined;
const kindEnum = canWrite ? (['explore', 'review', 'worker'] as const) : (['explore', 'review'] as const);
const singleSchema = z.object({
description: z.string().describe('Short label shown to the user, 3-6 words'),
prompt: z.string().describe('Self-contained instructions: what to do, where, and what to report'),
kind: z.enum(kindEnum as unknown as [string, ...string[]]).optional().describe(
canWrite
? 'explore: read-only research. review: read-only critique. worker: makes changes. Default explore.'
: 'explore: find and report. review: critique code for defects. Default explore.',
),
});
const batchSchema = z.object({
tasks: z
.array(
z.object({
description: z.string().describe('Short label shown to the user, 3-6 words'),
prompt: z.string().describe('Self-contained instructions: what to do, where, and what to report'),
kind: z.enum(kindEnum as unknown as [string, ...string[]]).optional().describe(
canWrite
? 'explore: read-only research. review: read-only critique. worker: makes changes. Default explore.'
: 'explore: find and report. review: critique code for defects. Default explore.',
),
}),
)
.min(1)
.max(8)
.describe('Several independent investigations to run in parallel. Use this instead of calling task multiple times.'),
kind: z.enum(kindEnum as unknown as [string, ...string[]]).optional().describe('Default kind for tasks that omit it'),
});
return tool({ return tool({
description: description:
@@ -165,101 +283,48 @@ export function createTaskTool(opts: {
'the user exactly as yours are. Use it for a self-contained task whose intermediate steps you do not ' + 'the user exactly as yours are. Use it for a self-contained task whose intermediate steps you do not ' +
'need to see; keep work you must supervise step by step in your own turn.' 'need to see; keep work you must supervise step by step in your own turn.'
: '') + : '') +
'\nDo not delegate something you can answer with a single grep.', '\nDo not delegate something you can answer with a single grep. For independent pieces of work, pass `tasks` to run them in parallel instead of calling task several times sequentially.',
inputSchema: z.object({ inputSchema: z.union([singleSchema, batchSchema]),
description: z.string().describe('Short label shown to the user, 3-6 words'), execute: async (input, { abortSignal }) => {
prompt: z.string().describe('Self-contained instructions: what to do, where, and what to report'), const asBatch = input as { tasks?: TaskSpec[]; description?: string; prompt?: string; kind?: SubagentKind };
kind: z if (asBatch.tasks && Array.isArray(asBatch.tasks)) {
.enum(canWrite ? ['explore', 'review', 'worker'] : ['explore', 'review']) const specs: TaskSpec[] = asBatch.tasks.map((t) => ({
.optional() description: t.description,
.describe( prompt: t.prompt,
canWrite kind: (t.kind ?? asBatch.kind ?? 'explore') as SubagentKind,
? 'explore: read-only research. review: read-only critique. worker: makes changes. Default explore.' }));
: 'explore: find and report. review: critique code for defects. Default explore.', // Validate worker channel before fanning out, so the error is immediate.
), for (const s of specs) {
}), if (s.kind === 'worker' && !opts.approve) {
execute: async ({ description, prompt, kind }, { abortSignal }) => { throw new Error('The worker kind needs an approval channel, which this session has not provided.');
const flavour: SubagentKind = kind ?? 'explore';
if (flavour === 'worker' && !opts.approve) {
throw new Error('The worker kind needs an approval channel, which this session has not provided.');
}
const id = `sub${++counter}`;
const report = opts.report;
report?.({ type: 'start', id, kind: flavour, description });
let steps = 0;
let text = '';
let usedTokens: { inputTokens: number; outputTokens: number } | undefined;
try {
// `explore` is search, not reasoning, so it runs on the cheaper model when
// one is configured. `review` and `worker` keep the parent's: they judge
// and they change, both of which want the full model.
const model = flavour === 'explore' ? (opts.subagentModel ?? opts.model) : opts.model;
const result = streamText({
model,
system: PROMPTS[flavour](opts.cwd ?? process.cwd()),
messages: [{ role: 'user', content: prompt }],
tools: TOOLS[flavour],
stopWhen: isStepCount(opts.maxSteps ?? 20),
...(opts.approve
? {
toolApproval: async ({ toolCall }: { toolCall: { toolName: string; input: unknown } }) => {
const approved = await opts.approve!(toolCall);
return approved
? undefined
: { type: 'denied' as const, reason: 'The user denied this call. Stop and report it.' };
},
}
: {}),
...(abortSignal ? { abortSignal } : {}),
});
const sink = () => {};
void result.responseMessages.then(undefined, sink);
void result.usage.then(undefined, sink);
void result.steps.then(undefined, sink);
void result.finalStep.then(undefined, sink);
void result.finishReason.then(undefined, sink);
for await (const part of result.stream) {
if (part.type === 'tool-call') {
steps++;
report?.({ type: 'step', id, tool: part.toolName, summary: summarize(part.input) });
} else if (part.type === 'tool-result') {
report?.({ type: 'result', id, tool: part.toolName, summary: outcome(part.output), ok: true });
} else if (part.type === 'tool-error') {
const message = part.error instanceof Error ? part.error.message : String(part.error);
report?.({ type: 'result', id, tool: part.toolName, summary: outcome(message), ok: false });
} else if (part.type === 'text-delta') {
text += part.text;
} else if (part.type === 'error') {
// A provider failure arrives as a stream part, not a throw, so it has to
// be rethrown here or the subagent silently returns nothing.
const message = part.error instanceof Error ? part.error.message : String(part.error);
throw part.error instanceof Error ? part.error : new Error(message);
} }
} }
const ids = specs.map(() => `sub${++counter}`);
try { const innerOpts = {
const usage = await result.usage; model: opts.model,
usedTokens = { inputTokens: usage.inputTokens ?? 0, outputTokens: usage.outputTokens ?? 0 }; subagentModel: opts.subagentModel,
} catch { cwd: opts.cwd,
// A run that errored before producing usage has nothing to account for. maxSteps: opts.maxSteps,
} report: opts.report,
} catch (e) { approve: opts.approve,
const message = e instanceof Error ? e.message : String(e); onUsage: opts.onUsage,
report?.({ type: 'error', id, message }); };
throw e; const results = await Promise.all(
specs.map((spec, i) =>
runOne(spec, ids[i]!, innerOpts, abortSignal).catch((e) => {
const msg = e instanceof Error ? e.message : String(e);
return `Subagent "${spec.description}" failed: ${msg}`;
}),
),
);
return results.map((r, i) => `## ${specs[i]!.description}\n${r}`).join('\n\n');
} }
const spec: TaskSpec = {
const trimmed = text.trim(); description: (input as { description: string }).description,
report?.({ type: 'end', id, ok: trimmed.length > 0, steps }); prompt: (input as { prompt: string }).prompt,
// Settled after the stream closes; a failed run reports nothing rather than kind: (input as { kind?: SubagentKind }).kind,
// a half count. The parent prices these against the subagent's own model id. };
if (usedTokens) opts.onUsage?.({ kind: flavour, ...usedTokens }); return runOne(spec, `sub${++counter}`, opts as Parameters<typeof runOne>[2], abortSignal);
return trimmed || 'Subagent returned no findings.';
}, },
}); });
} }
+101
View File
@@ -0,0 +1,101 @@
import { expect, test } from 'bun:test';
import { MockLanguageModelV4, simulateReadableStream } from 'ai/test';
import type { LanguageModelV4StreamPart } from '@ai-sdk/provider';
import { mkdtempSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { Session } from '../src/session';
import { createTaskTool, type SubagentEvent } from '../src/subagent';
const usage = { inputTokens: 10, outputTokens: 10 } as any;
const parts = (id: string, toolName: string, input: unknown): LanguageModelV4StreamPart[] => [
{ type: 'tool-input-start', id, toolName },
{ type: 'tool-input-end', id },
{ type: 'tool-call', toolCallId: id, toolName, input: JSON.stringify(input) },
{ type: 'finish', finishReason: { unified: 'tool-calls', raw: 'tool_use' }, usage },
];
const txt = (body: string): LanguageModelV4StreamPart[] => [
{ type: 'text-start', id: '0' },
{ type: 'text-delta', id: '0', delta: body },
{ type: 'text-end', id: '0' },
{ type: 'finish', finishReason: { unified: 'stop', raw: 'stop' }, usage },
];
function inTempDir<T>(fn: (dir: string) => Promise<T>): Promise<T> {
const orig = process.cwd();
const dir = mkdtempSync(join(tmpdir(), 'shiro-par-'));
process.chdir(dir);
return fn(dir).finally(() => {
process.chdir(orig);
rmSync(dir, { recursive: true, force: true });
});
}
test('two tasks dispatched together overlap in time rather than queueing', () =>
inTempDir(async () => {
await Bun.write('a.ts', 'export const a = 1;\n');
await Bun.write('b.ts', 'export const b = 2;\n');
const events: SubagentEvent[] = [];
let seen = 0;
const delay = 50;
const model = new MockLanguageModelV4({
doStream: async () => {
const n = seen++;
if (n === 0) {
return { stream: simulateReadableStream({ chunks: parts('c1', 'task', { tasks: [{ description: 'find a', prompt: 'Find a' }, { description: 'find b', prompt: 'Find b' }] }), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
}
// Each subagent does: one grep (delayed by tool), then text
// We simulate delay at model level to observe overlap
if (n === 1 || n === 2) {
await new Promise((r) => setTimeout(r, delay));
return { stream: simulateReadableStream({ chunks: parts(`s${n}`, 'grep', { pattern: n === 1 ? 'a' : 'b', include: '**/*.ts' }), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
}
if (n === 3 || n === 4) return { stream: simulateReadableStream({ chunks: txt(n === 3 ? 'found a' : 'found b'), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
return { stream: simulateReadableStream({ chunks: txt('both found'), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
},
});
const session = new Session({
model,
askApproval: async () => 'deny' as const,
extraTools: { task: createTaskTool({ model, report: (e) => events.push(e) }) },
autoApprove: ['task'],
});
const t0 = Date.now();
for await (const _ of session.send('go')) void _;
const dt = Date.now() - t0;
// Sequential would be ~delay + delay serially; parallel keeps it near one delay.
// Allow generous headroom for CI jitter but require overlap.
expect(dt).toBeLessThan(delay * 2 + 80);
// Both findings reach the parent as one tool result with headings
const toolMsg = session.messages.find((m) => m.role === 'tool');
const body = JSON.stringify(toolMsg);
expect(body).toContain('find a');
expect(body).toContain('find b');
// Panel saw two starts before two ends (fan-out)
const starts = events.filter((e) => e.type === 'start').map((e) => (e as Extract<SubagentEvent, { type: 'start' }>).description);
expect(starts).toContain('find a');
expect(starts).toContain('find b');
}));
test('a single task still works through the same path', () =>
inTempDir(async () => {
await Bun.write('a.ts', 'export const a = 1;\n');
let seen = 0;
const model = new MockLanguageModelV4({
doStream: async (opts) => {
const n = seen++;
if (n === 0) return { stream: simulateReadableStream({ chunks: parts('c1', 'task', { description: 'find a', prompt: 'Find a' }), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
if (n === 1) return { stream: simulateReadableStream({ chunks: parts('s1', 'grep', { pattern: 'a', include: '**/*.ts' }), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
if (n === 2) return { stream: simulateReadableStream({ chunks: txt('found a at a.ts'), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
return { stream: simulateReadableStream({ chunks: txt('done'), chunkDelayInMs: null, initialDelayInMs: null }) } as any;
},
});
const session = new Session({ model, askApproval: async () => 'deny' as const, extraTools: { task: createTaskTool({ model }) }, autoApprove: ['task'] });
for await (const _ of session.send('go')) void _;
expect(JSON.stringify(session.messages)).toContain('found a at a.ts');
}));