import { randomUUID } from 'node:crypto'; import type { PublicUsage, SubagentDetailsV1 } from '../contracts'; import { PiProcessBudget, type PiProcessLease } from './worker-pool'; import type { PiWorkerStopReason } from './worker-process'; 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; stop(reason: PiWorkerStopReason): Promise; } export interface PiSubagentDispatchResult { details: SubagentDetailsV1; } export interface PiSubagentDispatchOptions { signal?: AbortSignal; onUpdate?: (details: SubagentDetailsV1) => void; } export interface PiSubagentSchedulerOptions { openChild(input: PiSubagentChildOpenInput): Promise; processBudget: PiProcessBudget; reclaimProcessCapacity?(signal?: AbortSignal): Promise; createId?: (kind: 'dispatch' | 'task') => string; } interface DispatchRecord { identity: PiSubagentParentIdentity; controller: AbortController; children: Set; flight: Promise; } 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'); } } get activeCount(): number { return this.active; } get waitingCount(): number { return this.waiters.length; } 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 | null { return value !== null && typeof value === 'object' && !Array.isArray(value) ? value as Record : 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) | undefined; private readonly createId: NonNullable; private readonly dispatches = new Map(); private readonly parentDispatches = new Map>(); 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 { 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(), } 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(); 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(); } } getDiagnostics(): { activeChildPermits: number; waitingChildPermits: number; activeDispatches: number; activeParents: number; } { return { activeChildPermits: this.childPermits.activeCount, waitingChildPermits: this.childPermits.waitingCount, activeDispatches: this.dispatches.size, activeParents: this.parentDispatches.size, }; } async close(): Promise { 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 { 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 { const projected = details.tasks[index]; if (!projected) return; let releaseChild: (() => void) | undefined; let processLease: PiProcessLease | undefined; let child: PiSubagentChild | undefined; let stopReason: PiWorkerStopReason = 'subagent_complete'; 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'); if (aborted) stopReason = 'subagent_abort'; projected.status = aborted ? 'aborted' : 'error'; projected.errorCode = publicErrorCode(error, aborted); } finally { if (child) { record.children.delete(child); await child.stop(stopReason).catch(() => undefined); } processLease?.release(); releaseChild?.(); this.emit(details, onUpdate); } } private async acquireProcessLease(signal: AbortSignal): Promise { 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. } } }