Files
makelore/electron/coding-runtime/pi/subagent.ts

447 lines
15 KiB
TypeScript

import { randomUUID } from 'node:crypto';
import type { PublicUsage, SubagentDetailsV1 } from '../contracts';
import { PiProcessBudget, type PiProcessLease } from './worker-pool';
export type PiSubagentMode = SubagentDetailsV1['mode'];
export type PiSubagentToolProfile = SubagentDetailsV1['tasks'][number]['toolProfile'];
export interface PiSubagentTaskRequest {
agentId: string;
task: string;
toolProfile: PiSubagentToolProfile;
}
export interface PiSubagentDispatchRequest {
mode: PiSubagentMode;
tasks: PiSubagentTaskRequest[];
}
export interface PiSubagentParentIdentity {
conversationId: string;
workerGeneration: number;
runId: string;
projectId: string;
}
export interface PiSubagentChildOpenInput extends PiSubagentParentIdentity {
dispatchId: string;
taskId: string;
agentId: string;
toolProfile: PiSubagentToolProfile;
}
export interface PiSubagentChildResult {
summary: string;
usage?: PublicUsage;
}
export interface PiSubagentChild {
readonly id: string;
run(
prompt: string,
signal: AbortSignal,
onProgress?: (summary: string) => void,
): Promise<PiSubagentChildResult>;
stop(): Promise<void>;
}
export interface PiSubagentDispatchResult {
details: SubagentDetailsV1;
}
export interface PiSubagentDispatchOptions {
signal?: AbortSignal;
onUpdate?: (details: SubagentDetailsV1) => void;
}
export interface PiSubagentSchedulerOptions {
openChild(input: PiSubagentChildOpenInput): Promise<PiSubagentChild>;
processBudget: PiProcessBudget;
reclaimProcessCapacity?(signal?: AbortSignal): Promise<boolean>;
createId?: (kind: 'dispatch' | 'task') => string;
}
interface DispatchRecord {
identity: PiSubagentParentIdentity;
controller: AbortController;
children: Set<PiSubagentChild>;
flight: Promise<PiSubagentDispatchResult>;
}
interface SemaphoreWaiter {
signal: AbortSignal;
resolve(release: () => void): void;
reject(error: Error): void;
abort(): void;
}
const MAX_TASKS_PER_DISPATCH = 8;
const MAX_AGENT_ID_LENGTH = 128;
const MAX_TASK_LENGTH = 6_000;
const MAX_SUMMARY_LENGTH = 4_000;
export class PiSubagentChildError extends Error {
constructor(readonly code: string) {
super(code);
this.name = 'PiSubagentChildError';
}
}
class FifoSemaphore {
private readonly waiters: SemaphoreWaiter[] = [];
private active = 0;
constructor(readonly maximum: number) {
if (!Number.isSafeInteger(maximum) || maximum <= 0) {
throw new Error('Subagent concurrency must be a positive safe integer');
}
}
acquire(signal: AbortSignal): Promise<() => void> {
if (signal.aborted) return Promise.reject(new PiSubagentChildError('SUBAGENT_ABORTED'));
if (this.active < this.maximum) return Promise.resolve(this.issuePermit());
return new Promise<() => void>((resolve, reject) => {
const waiter: SemaphoreWaiter = {
signal,
resolve,
reject,
abort: () => {
const index = this.waiters.indexOf(waiter);
if (index >= 0) this.waiters.splice(index, 1);
reject(new PiSubagentChildError('SUBAGENT_ABORTED'));
},
};
signal.addEventListener('abort', waiter.abort, { once: true });
this.waiters.push(waiter);
});
}
private issuePermit(): () => void {
this.active += 1;
let released = false;
return () => {
if (released) return;
released = true;
this.active -= 1;
this.advance();
};
}
private advance(): void {
while (this.active < this.maximum && this.waiters.length > 0) {
const waiter = this.waiters.shift() as SemaphoreWaiter;
waiter.signal.removeEventListener('abort', waiter.abort);
if (waiter.signal.aborted) {
waiter.reject(new PiSubagentChildError('SUBAGENT_ABORTED'));
continue;
}
waiter.resolve(this.issuePermit());
}
}
}
function asRecord(value: unknown): Record<string, unknown> | null {
return value !== null && typeof value === 'object' && !Array.isArray(value)
? value as Record<string, unknown>
: null;
}
function boundedText(value: unknown, label: string, maximum: number): string {
if (typeof value !== 'string') throw new Error(`${label} must be a string`);
const normalized = value.trim();
if (!normalized) throw new Error(`${label} must not be empty`);
if (normalized.length > maximum) throw new Error(`${label} is too long`);
return normalized;
}
export function parsePiSubagentDispatchRequest(value: unknown): PiSubagentDispatchRequest {
const record = asRecord(value);
if (!record || !['single', 'parallel', 'chain'].includes(String(record.mode))) {
throw new Error('Subagent dispatch mode is invalid');
}
if (!Array.isArray(record.tasks) || record.tasks.length === 0) {
throw new Error('Subagent dispatch requires at least one task');
}
if (record.tasks.length > MAX_TASKS_PER_DISPATCH) {
throw new Error('Subagent dispatch accepts at most 8 tasks');
}
if (record.mode === 'single' && record.tasks.length !== 1) {
throw new Error('Single subagent dispatch requires exactly one task');
}
const tasks = record.tasks.map((candidate, index) => {
const task = asRecord(candidate);
if (!task) throw new Error(`Subagent task ${index + 1} is invalid`);
if (task.toolProfile !== 'read-only' && task.toolProfile !== 'coding') {
throw new Error(`Subagent task ${index + 1} tool profile is invalid`);
}
return {
agentId: boundedText(task.agentId, `Subagent task ${index + 1} agent id`, MAX_AGENT_ID_LENGTH),
task: boundedText(task.task, `Subagent task ${index + 1} prompt`, MAX_TASK_LENGTH),
toolProfile: task.toolProfile,
};
});
return { mode: record.mode as PiSubagentMode, tasks };
}
function safeSummary(value: string): string {
return value.length <= MAX_SUMMARY_LENGTH ? value : `${value.slice(0, MAX_SUMMARY_LENGTH - 1)}`;
}
function safeUsage(value: PublicUsage | undefined): PublicUsage | undefined {
if (!value) return undefined;
const fields = [
value.inputTokens,
value.outputTokens,
value.cacheReadTokens,
value.cacheWriteTokens,
].filter((field) => field !== undefined);
if (fields.some((field) => !Number.isFinite(field) || field < 0)) return undefined;
return structuredClone(value);
}
function publicErrorCode(error: unknown, aborted: boolean): string {
if (aborted) return 'SUBAGENT_ABORTED';
if (error instanceof PiSubagentChildError && /^[A-Z][A-Z0-9_]{0,63}$/.test(error.code)) {
return error.code;
}
return 'SUBAGENT_CHILD_FAILED';
}
function parentKey(identity: PiSubagentParentIdentity): string {
return `${identity.conversationId}:${identity.workerGeneration}:${identity.runId}`;
}
export class PiSubagentScheduler {
private readonly openChild: PiSubagentSchedulerOptions['openChild'];
private readonly processBudget: PiProcessBudget;
private readonly childPermits: FifoSemaphore;
private readonly reclaimProcessCapacity: ((signal?: AbortSignal) => Promise<boolean>) | undefined;
private readonly createId: NonNullable<PiSubagentSchedulerOptions['createId']>;
private readonly dispatches = new Map<string, DispatchRecord>();
private readonly parentDispatches = new Map<string, Set<string>>();
private closing = false;
constructor(options: PiSubagentSchedulerOptions) {
this.openChild = options.openChild;
this.processBudget = options.processBudget;
this.childPermits = new FifoSemaphore(4);
this.reclaimProcessCapacity = options.reclaimProcessCapacity;
this.createId = options.createId ?? ((kind) => `${kind}-${randomUUID()}`);
}
dispatch(
input: PiSubagentParentIdentity & { request: unknown },
options: PiSubagentDispatchOptions = {},
): Promise<PiSubagentDispatchResult> {
if (this.closing) return Promise.reject(new Error('Subagent scheduler is shutting down'));
const request = parsePiSubagentDispatchRequest(input.request);
const dispatchId = this.createId('dispatch');
if (this.dispatches.has(dispatchId)) {
return Promise.reject(new Error(`Duplicate subagent dispatch id: ${dispatchId}`));
}
const identity: PiSubagentParentIdentity = {
conversationId: input.conversationId,
workerGeneration: input.workerGeneration,
runId: input.runId,
projectId: input.projectId,
};
const controller = new AbortController();
const externalAbort = () => controller.abort();
options.signal?.addEventListener('abort', externalAbort, { once: true });
if (options.signal?.aborted) controller.abort();
const tasks: SubagentDetailsV1['tasks'] = request.tasks.map((task) => ({
taskId: this.createId('task'),
agentId: task.agentId,
toolProfile: task.toolProfile,
status: 'queued',
}));
const details: SubagentDetailsV1 = {
schema: 'subagent.v1',
dispatchId,
mode: request.mode,
tasks,
};
const record = {
identity,
controller,
children: new Set<PiSubagentChild>(),
} as DispatchRecord;
const flight = this.runDispatch(record, request, details, options.onUpdate)
.finally(() => {
options.signal?.removeEventListener('abort', externalAbort);
this.dispatches.delete(dispatchId);
const key = parentKey(identity);
const ids = this.parentDispatches.get(key);
ids?.delete(dispatchId);
if (ids?.size === 0) this.parentDispatches.delete(key);
});
record.flight = flight;
this.dispatches.set(dispatchId, record);
const key = parentKey(identity);
const ids = this.parentDispatches.get(key) ?? new Set<string>();
ids.add(dispatchId);
this.parentDispatches.set(key, ids);
this.emit(details, options.onUpdate);
return flight;
}
abortParent(identity: PiSubagentParentIdentity): void {
for (const dispatchId of this.parentDispatches.get(parentKey(identity)) ?? []) {
this.dispatches.get(dispatchId)?.controller.abort();
}
}
async close(): Promise<void> {
if (this.closing) {
await Promise.allSettled([...this.dispatches.values()].map((record) => record.flight));
return;
}
this.closing = true;
for (const record of this.dispatches.values()) record.controller.abort();
await Promise.allSettled([...this.dispatches.values()].map((record) => record.flight));
}
private async runDispatch(
record: DispatchRecord,
request: PiSubagentDispatchRequest,
details: SubagentDetailsV1,
onUpdate: PiSubagentDispatchOptions['onUpdate'],
): Promise<PiSubagentDispatchResult> {
if (request.mode === 'chain') {
let previous = '';
for (let index = 0; index < request.tasks.length; index += 1) {
const task = request.tasks[index] as PiSubagentTaskRequest;
if (record.controller.signal.aborted) {
this.markRemaining(details, index, 'aborted', onUpdate);
break;
}
const prompt = task.task.split('{previous}').join(previous);
await this.runTask(record, task, details, index, prompt, onUpdate);
const projected = details.tasks[index];
if (projected?.status !== 'complete') {
this.markRemaining(
details,
index + 1,
record.controller.signal.aborted ? 'aborted' : 'skipped',
onUpdate,
);
break;
}
previous = projected.summary ?? '';
}
} else {
await Promise.all(request.tasks.map((task, index) => (
this.runTask(record, task, details, index, task.task, onUpdate)
)));
}
return { details: structuredClone(details) };
}
private async runTask(
record: DispatchRecord,
task: PiSubagentTaskRequest,
details: SubagentDetailsV1,
index: number,
prompt: string,
onUpdate: PiSubagentDispatchOptions['onUpdate'],
): Promise<void> {
const projected = details.tasks[index];
if (!projected) return;
let releaseChild: (() => void) | undefined;
let processLease: PiProcessLease | undefined;
let child: PiSubagentChild | undefined;
try {
releaseChild = await this.childPermits.acquire(record.controller.signal);
processLease = await this.acquireProcessLease(record.controller.signal);
if (record.controller.signal.aborted) throw new PiSubagentChildError('SUBAGENT_ABORTED');
projected.status = 'running';
this.emit(details, onUpdate);
child = await this.openChild({
...record.identity,
dispatchId: details.dispatchId,
taskId: projected.taskId,
agentId: task.agentId,
toolProfile: task.toolProfile,
});
record.children.add(child);
const result = await child.run(prompt, record.controller.signal, (summary) => {
projected.summary = safeSummary(summary);
this.emit(details, onUpdate);
});
if (record.controller.signal.aborted) throw new PiSubagentChildError('SUBAGENT_ABORTED');
projected.status = 'complete';
projected.summary = safeSummary(result.summary);
const usage = safeUsage(result.usage);
if (usage) projected.usage = usage;
} catch (error) {
const aborted = record.controller.signal.aborted
|| (error instanceof PiSubagentChildError && error.code === 'SUBAGENT_ABORTED');
projected.status = aborted ? 'aborted' : 'error';
projected.errorCode = publicErrorCode(error, aborted);
} finally {
if (child) {
record.children.delete(child);
await child.stop().catch(() => undefined);
}
processLease?.release();
releaseChild?.();
this.emit(details, onUpdate);
}
}
private async acquireProcessLease(signal: AbortSignal): Promise<PiProcessLease> {
const reservation = new AbortController();
const abortReservation = () => reservation.abort();
signal.addEventListener('abort', abortReservation, { once: true });
if (signal.aborted) reservation.abort();
const budgetWasFull = this.processBudget.activeCount >= this.processBudget.maxProcesses;
const processLeaseFlight = this.processBudget.acquire(reservation.signal, 'child');
try {
if (budgetWasFull && this.reclaimProcessCapacity) {
const outcome = await Promise.race([
processLeaseFlight.then((lease) => ({ lease })),
this.reclaimProcessCapacity(reservation.signal).then(() => null),
]);
if (outcome) return outcome.lease;
}
return await processLeaseFlight;
} catch (error) {
reservation.abort();
const orphanedLease = await processLeaseFlight.catch(() => undefined);
orphanedLease?.release();
throw error;
} finally {
signal.removeEventListener('abort', abortReservation);
reservation.abort();
}
}
private markRemaining(
details: SubagentDetailsV1,
from: number,
status: 'aborted' | 'skipped',
onUpdate: PiSubagentDispatchOptions['onUpdate'],
): void {
for (let index = from; index < details.tasks.length; index += 1) {
const task = details.tasks[index];
if (!task || task.status !== 'queued') continue;
task.status = status;
if (status === 'aborted') task.errorCode = 'SUBAGENT_ABORTED';
}
this.emit(details, onUpdate);
}
private emit(
details: SubagentDetailsV1,
onUpdate: PiSubagentDispatchOptions['onUpdate'],
): void {
if (!onUpdate) return;
try {
onUpdate(structuredClone(details));
} catch {
// UI projection observers must not affect child lifecycle.
}
}
}