fix: close Pi subagent review gaps
This commit is contained in:
@@ -353,8 +353,12 @@ export class PiSubagentScheduler {
|
||||
let child: PiSubagentChild | undefined;
|
||||
try {
|
||||
releaseChild = await this.childPermits.acquire(record.controller.signal);
|
||||
await this.reclaimIdleCapacity(record.controller.signal);
|
||||
processLease = await this.processBudget.acquire(record.controller.signal);
|
||||
const budgetWasFull = this.processBudget.activeCount >= this.processBudget.maxProcesses;
|
||||
const processLeaseFlight = this.processBudget.acquire(record.controller.signal, 'child');
|
||||
if (budgetWasFull && this.reclaimProcessCapacity) {
|
||||
await this.reclaimProcessCapacity();
|
||||
}
|
||||
processLease = await processLeaseFlight;
|
||||
if (record.controller.signal.aborted) throw new PiSubagentChildError('SUBAGENT_ABORTED');
|
||||
projected.status = 'running';
|
||||
this.emit(details, onUpdate);
|
||||
@@ -391,14 +395,6 @@ export class PiSubagentScheduler {
|
||||
}
|
||||
}
|
||||
|
||||
private async reclaimIdleCapacity(signal: AbortSignal): Promise<void> {
|
||||
while (!signal.aborted
|
||||
&& this.processBudget.activeCount >= this.processBudget.maxProcesses
|
||||
&& this.reclaimProcessCapacity) {
|
||||
if (!await this.reclaimProcessCapacity()) break;
|
||||
}
|
||||
}
|
||||
|
||||
private markRemaining(
|
||||
details: SubagentDetailsV1,
|
||||
from: number,
|
||||
|
||||
@@ -116,6 +116,7 @@ export interface PiProcessLease {
|
||||
|
||||
interface PiProcessBudgetWaiter {
|
||||
signal?: AbortSignal;
|
||||
priority: 'normal' | 'child';
|
||||
resolve(lease: PiProcessLease): void;
|
||||
reject(error: Error): void;
|
||||
onAbort?: () => void;
|
||||
@@ -134,11 +135,19 @@ export class PiProcessBudget {
|
||||
get activeCount(): number { return this.active; }
|
||||
get waitingCount(): number { return this.waiters.length; }
|
||||
|
||||
acquire(signal?: AbortSignal): Promise<PiProcessLease> {
|
||||
acquire(
|
||||
signal?: AbortSignal,
|
||||
priority: 'normal' | 'child' = 'normal',
|
||||
): Promise<PiProcessLease> {
|
||||
if (signal?.aborted) return Promise.reject(new Error('Pi process budget acquisition cancelled'));
|
||||
if (this.active < this.maxProcesses) return Promise.resolve(this.issueLease());
|
||||
return new Promise<PiProcessLease>((resolve, reject) => {
|
||||
const waiter: PiProcessBudgetWaiter = { resolve, reject, ...(signal ? { signal } : {}) };
|
||||
const waiter: PiProcessBudgetWaiter = {
|
||||
resolve,
|
||||
reject,
|
||||
priority,
|
||||
...(signal ? { signal } : {}),
|
||||
};
|
||||
if (signal) {
|
||||
waiter.onAbort = () => {
|
||||
const index = this.waiters.indexOf(waiter);
|
||||
@@ -166,7 +175,9 @@ export class PiProcessBudget {
|
||||
|
||||
private advance(): void {
|
||||
while (this.active < this.maxProcesses && this.waiters.length > 0) {
|
||||
const waiter = this.waiters.shift() as PiProcessBudgetWaiter;
|
||||
const childIndex = this.waiters.findIndex(({ priority }) => priority === 'child');
|
||||
const [waiter] = this.waiters.splice(childIndex >= 0 ? childIndex : 0, 1);
|
||||
if (!waiter) return;
|
||||
if (waiter.onAbort && waiter.signal) {
|
||||
waiter.signal.removeEventListener('abort', waiter.onAbort);
|
||||
}
|
||||
|
||||
@@ -65,7 +65,9 @@ export type PiWorkerProcessOptions = {
|
||||
export function buildPiRpcArgs(
|
||||
sessionDir: string,
|
||||
additionalArgs: readonly string[] = [],
|
||||
tools: readonly string[] = ['read', 'bash', 'edit', 'write', 'grep', 'find', 'ls', 'ask_user'],
|
||||
tools: readonly string[] = [
|
||||
'read', 'bash', 'edit', 'write', 'grep', 'find', 'ls', 'ask_user', 'subagent',
|
||||
],
|
||||
): string[] {
|
||||
return [
|
||||
'--mode', 'rpc',
|
||||
|
||||
Reference in New Issue
Block a user