Files
makelore/tests/unit/pi-worker-pool.test.ts
2026-09-01 12:40:08 +08:00

1224 lines
46 KiB
TypeScript

// @vitest-environment node
import { describe, expect, it } from 'vitest';
import type { PrepareConversationInput } from '../../electron/coding-runtime/contracts';
import {
PiProcessBudget,
PiWorkerPool,
type PiConversationWorker,
type PiWorkerOpenResult,
type PiWorkerPoolEvent,
} from '../../electron/coding-runtime/pi/worker-pool';
import type { PiProcessError } from '../../electron/coding-runtime/pi/process-errors';
import { PiProcessError as PiProcessFailure } from '../../electron/coding-runtime/pi/process-errors';
import type {
PiRpcCommand,
PiRpcEvent,
PiRpcRequestOptions,
PiRpcResponse,
} from '../../electron/coding-runtime/pi/rpc-client';
import type { PiRuntimeTelemetryEvent } from '../../electron/coding-runtime/pi/telemetry';
import { PiSubagentScheduler } from '../../electron/coding-runtime/pi/subagent';
import type {
PiWorkerProofFailure,
PiWorkerStopReason,
} from '../../electron/coding-runtime/pi/worker-process';
const MODEL = {
model: {
accountId: 'account-a',
modelId: 'model-a',
thinkingLevel: 'medium' as const,
},
modelResolution: 'resolved' as const,
};
function conversation(conversationId: string): PrepareConversationInput {
return {
conversationId,
projectId: 'project-a',
agentId: 'agent-a',
title: conversationId,
model: MODEL,
};
}
function deferred(): { promise: Promise<void>; resolve(): void } {
let resolve!: () => void;
const promise = new Promise<void>((done) => { resolve = done; });
return { promise, resolve };
}
class FakeWorker implements PiConversationWorker {
readonly generation = 1;
readonly requests: PiRpcCommand[] = [];
stopped = false;
readonly stopReasons: PiWorkerStopReason[] = [];
private readonly eventListeners = new Set<(event: PiRpcEvent) => void>();
private readonly invalidationListeners = new Set<(error: PiProcessError) => void>();
private timeoutType: string | null = null;
private pendingType: string | null = null;
private stateData: unknown = {
isStreaming: true,
isCompacting: false,
pendingMessageCount: 0,
retryAttempt: 0,
};
private lateResult: ((result: {
response?: PiRpcResponse;
error?: PiProcessError;
}) => void) | undefined;
constructor(readonly id: string) {}
async request<T = unknown>(command: PiRpcCommand, options?: PiRpcRequestOptions) {
this.requests.push(command);
if (this.pendingType === command.type) {
this.pendingType = null;
return await new Promise<PiRpcResponse<T>>((_resolve, reject) => {
options?.signal?.addEventListener('abort', () => reject(new PiProcessFailure(
'PI_RPC_ABORTED',
`fake ${command.type} was aborted after authoritative settlement`,
)), { once: true });
});
}
if (this.timeoutType === command.type) {
this.timeoutType = null;
this.lateResult = (options as PiRpcRequestOptions & {
onLateResult?(result: {
response?: PiRpcResponse;
error?: PiProcessError;
}): void;
} | undefined)?.onLateResult;
throw new PiProcessFailure('PI_RPC_TIMEOUT', `fake ${command.type} confirmation timeout`);
}
return {
type: 'response' as const,
id: 'fake',
success: true,
...(command.type === 'get_state'
? { data: structuredClone(this.stateData) as T }
: {}),
};
}
setState(state: unknown): void {
this.stateData = structuredClone(state);
}
timeoutNext(type: string): void {
this.timeoutType = type;
}
pendNext(type: string): void {
this.pendingType = type;
}
completeLateSuccess(): void {
this.lateResult?.({
response: { type: 'response', id: 'fake-late', success: true },
});
this.lateResult = undefined;
}
completeLateFailure(): void {
this.lateResult?.({
error: new PiProcessFailure('PI_RPC_RESPONSE_ERROR', 'fake late rejection'),
});
this.lateResult = undefined;
}
async send(command: PiRpcCommand): Promise<void> {
this.requests.push(command);
}
subscribe(listener: (event: PiRpcEvent) => void): () => void {
this.eventListeners.add(listener);
return () => this.eventListeners.delete(listener);
}
subscribeInvalidation(listener: (error: PiProcessError) => void): () => void {
this.invalidationListeners.add(listener);
return () => this.invalidationListeners.delete(listener);
}
emit(event: PiRpcEvent): void {
for (const listener of this.eventListeners) listener(event);
}
invalidate(error = new PiProcessFailure('PI_RPC_EXITED', 'fake worker crashed')): void {
for (const listener of this.invalidationListeners) listener(error);
}
async injectFailureForProof(failure: PiWorkerProofFailure): Promise<void> {
this.invalidate(new PiProcessFailure(
failure === 'protocol_invalidation' ? 'PI_RPC_PROTOCOL_ERROR' : 'PI_RPC_EXITED',
'injected proof failure',
));
}
async stop(reason: PiWorkerStopReason) {
this.stopped = true;
this.stopReasons.push(reason);
return { mode: 'stdin-close' as const, code: 0, signal: null };
}
}
describe('Pi worker pool', () => {
it('does not charge logical threads as OS processes in shared-server mode', async () => {
const processBudget = new PiProcessBudget(1);
const workers: FakeWorker[] = [];
const pool = new PiWorkerPool({
processMode: 'shared',
processBudget,
maxIdle: 4,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`thread-${input.conversationId}`);
workers.push(worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all([
pool.prepare(conversation('conversation-a')),
pool.prepare(conversation('conversation-b')),
pool.prepare(conversation('conversation-c')),
]);
expect(processBudget.activeCount).toBe(0);
expect(workers).toHaveLength(3);
await pool.shutdown();
expect(workers.every((worker) => worker.stopped)).toBe(true);
});
it('retains top-level ownership after a mutation confirmation timeout', async () => {
const workers = new Map<string, FakeWorker>();
const pool = new PiWorkerPool({
maxRunning: 2,
maxIdle: 4,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.set(input.conversationId, worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all([
pool.prepare(conversation('conversation-target')),
pool.prepare(conversation('conversation-sibling')),
]);
workers.get('conversation-target')!.timeoutNext('prompt');
const target = pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-target',
command: { type: 'prompt', message: 'slow preflight' },
});
await expect(target.accepted).rejects.toMatchObject({ code: 'PI_RPC_TIMEOUT' });
expect(pool.getActiveRun('conversation-target')).toMatchObject({ runId: 'run-target' });
expect(pool.getResilienceProofDiagnostics().runs).toEqual({ active: 1, waiting: 0 });
expect(() => pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-overlap',
command: { type: 'compact' },
})).toThrow('Conversation already has a top-level run');
const sibling = pool.startTopLevel({
conversationId: 'conversation-sibling',
runId: 'run-sibling',
command: { type: 'prompt', message: 'independent sibling' },
});
await expect(sibling.accepted).resolves.toMatchObject({ success: true });
workers.get('conversation-sibling')!.emit({ type: 'agent_settled' });
workers.get('conversation-target')!.completeLateSuccess();
expect(pool.getActiveRun('conversation-target')).toMatchObject({ runId: 'run-target' });
workers.get('conversation-target')!.emit({ type: 'agent_settled' });
expect(pool.getActiveRun('conversation-target')).toBeNull();
expect(pool.getResilienceProofDiagnostics().runs).toEqual({ active: 0, waiting: 0 });
});
it('releases uncertain top-level ownership after a late explicit failure', async () => {
const worker = new FakeWorker('worker-target');
const pool = new PiWorkerPool({
maxIdle: 2,
openWorker: async () => ({
worker,
session: { piSessionId: 'session-target', sessionKey: 'key-target' },
}),
});
await pool.prepare(conversation('conversation-target'));
worker.timeoutNext('compact');
const ticket = pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-compact',
command: { type: 'compact' },
});
await expect(ticket.accepted).rejects.toMatchObject({ code: 'PI_RPC_TIMEOUT' });
expect(pool.getActiveRun('conversation-target')).toMatchObject({ runId: 'run-compact' });
worker.completeLateFailure();
await expect.poll(() => pool.getActiveRun('conversation-target')).toBeNull();
expect(pool.getResilienceProofDiagnostics().runs).toEqual({ active: 0, waiting: 0 });
});
it('settles compact from its authoritative RPC success without agent_settled', async () => {
const worker = new FakeWorker('worker-target');
const events: PiWorkerPoolEvent[] = [];
const pool = new PiWorkerPool({
maxIdle: 2,
openWorker: async () => ({
worker,
session: { piSessionId: 'session-target', sessionKey: 'key-target' },
}),
});
pool.subscribe((event) => events.push(event));
await pool.prepare(conversation('conversation-target'));
const ticket = pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-compact-success',
command: { type: 'compact' },
});
await expect(ticket.accepted).resolves.toMatchObject({ success: true });
expect(pool.getActiveRun('conversation-target')).toBeNull();
expect(pool.getState('conversation-target')?.state).toBe('idle');
expect(events).toContainEqual(expect.objectContaining({
type: 'top-level.settled',
runId: 'run-compact-success',
}));
});
it('treats agent settlement as authoritative when the RPC confirmation is still pending', async () => {
const worker = new FakeWorker('worker-target');
const pool = new PiWorkerPool({
maxIdle: 2,
openWorker: async () => ({
worker,
session: { piSessionId: 'session-target', sessionKey: 'key-target' },
}),
});
await pool.prepare(conversation('conversation-target'));
worker.pendNext('prompt');
const ticket = pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-settled-first',
command: { type: 'prompt', message: 'settle before response' },
});
await expect.poll(() => worker.requests.some(({ type }) => type === 'prompt')).toBe(true);
worker.emit({ type: 'agent_settled' });
await expect(ticket.accepted).resolves.toMatchObject({ success: true });
expect(pool.getActiveRun('conversation-target')).toBeNull();
expect(pool.getResilienceProofDiagnostics().runs).toEqual({ active: 0, waiting: 0 });
});
it('settles from authoritative idle state when agent_settled is missing and ignores a late duplicate', async () => {
const worker = new FakeWorker('worker-target');
const events: PiWorkerPoolEvent[] = [];
const pool = new PiWorkerPool({
maxIdle: 2,
settlementProbeIntervalMs: 5,
settlementProbeTimeoutMs: 20,
terminalSettlementTimeoutMs: 100,
openWorker: async () => ({
worker,
session: { piSessionId: 'session-target', sessionKey: 'key-target' },
}),
});
pool.subscribe((event) => events.push(event));
await pool.prepare(conversation('conversation-target'));
const ticket = pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-missing-settled',
command: { type: 'prompt', message: 'finish without settlement event' },
});
await expect(ticket.accepted).resolves.toMatchObject({ success: true });
worker.setState({
isStreaming: false,
isCompacting: false,
pendingMessageCount: 0,
retryAttempt: 0,
});
await expect.poll(() => pool.getActiveRun('conversation-target')).toBeNull();
expect(pool.getState('conversation-target')?.state).toBe('idle');
expect(events.filter((event) => event.type === 'top-level.settled')).toEqual([
expect.objectContaining({ source: 'state_probe' }),
]);
worker.emit({ type: 'agent_settled' });
expect(events.filter((event) => event.type === 'top-level.settled')).toHaveLength(1);
await pool.shutdown();
});
it('fails only the target thread when an accepted prompt rejects during settlement', async () => {
const workers = new Map<string, FakeWorker>();
const events: PiWorkerPoolEvent[] = [];
const pool = new PiWorkerPool({
maxIdle: 2,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.set(input.conversationId, worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
pool.subscribe((event) => events.push(event));
await Promise.all([
pool.prepare(conversation('conversation-target')),
pool.prepare(conversation('conversation-sibling')),
]);
const ticket = pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-post-accept-failure',
command: { type: 'prompt', message: 'accepted then failed' },
});
await expect(ticket.accepted).resolves.toMatchObject({ success: true });
workers.get('conversation-target')!.emit({
type: 'makelore_thread_error',
code: 'PROMPT_FAILED_AFTER_ACCEPTANCE',
});
expect(pool.getActiveRun('conversation-target')).toBeNull();
expect(pool.getState('conversation-target')).toMatchObject({
state: 'crashed',
failureCode: 'PI_RPC_PROTOCOL_ERROR',
});
expect(pool.getState('conversation-sibling')).toMatchObject({ state: 'ready' });
expect(events).toContainEqual(expect.objectContaining({
type: 'worker.crashed',
conversationId: 'conversation-target',
}));
expect(workers.get('conversation-target')!.requests.filter(({ type }) => type === 'prompt'))
.toHaveLength(1);
await pool.shutdown();
});
it('bounds a contradictory terminal phase without replaying the accepted prompt', async () => {
const worker = new FakeWorker('worker-target');
const pool = new PiWorkerPool({
maxIdle: 2,
settlementProbeIntervalMs: 5,
settlementProbeTimeoutMs: 20,
terminalSettlementTimeoutMs: 25,
openWorker: async () => ({
worker,
session: { piSessionId: 'session-target', sessionKey: 'key-target' },
}),
});
await pool.prepare(conversation('conversation-target'));
const ticket = pool.startTopLevel({
conversationId: 'conversation-target',
runId: 'run-terminal-stall',
command: { type: 'prompt', message: 'terminal stall' },
});
await expect(ticket.accepted).resolves.toMatchObject({ success: true });
worker.emit({
type: 'message_end',
message: { role: 'assistant', content: [], stopReason: 'stop' },
});
await expect.poll(() => pool.getState('conversation-target')?.state).toBe('crashed');
expect(pool.getActiveRun('conversation-target')).toBeNull();
expect(worker.requests.filter(({ type }) => type === 'prompt')).toHaveLength(1);
await pool.shutdown();
});
it('injects a proof failure only into the requested current generation', async () => {
const workers = new Map<string, FakeWorker>();
const pool = new PiWorkerPool({
maxIdle: 4,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.set(input.conversationId, worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all([
pool.prepare(conversation('conversation-target')),
pool.prepare(conversation('conversation-other')),
]);
await expect(pool.injectFailureForProof('conversation-target', 'unexpected_exit'))
.resolves.toEqual({ generation: 1 });
expect(pool.getState('conversation-target')).toMatchObject({ state: 'crashed', generation: 1 });
expect(pool.getState('conversation-other')).toMatchObject({ state: 'ready', generation: 1 });
expect(workers.get('conversation-other')?.stopped).toBe(false);
});
it('single-flights prepare per Conversation and never shares its worker with another Conversation', async () => {
const gate = deferred();
const opened: string[] = [];
const pool = new PiWorkerPool({
openWorker: async ({ conversation: input }): Promise<PiWorkerOpenResult> => {
opened.push(input.conversationId);
await gate.promise;
return {
worker: new FakeWorker(`worker-${input.conversationId}`),
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
maxIdle: 6,
});
const first = pool.prepare(conversation('conversation-a'));
const duplicate = pool.prepare(conversation('conversation-a'));
const other = pool.prepare(conversation('conversation-b'));
await expect.poll(() => opened).toEqual(['conversation-a', 'conversation-b']);
gate.resolve();
const [firstState, duplicateState, otherState] = await Promise.all([first, duplicate, other]);
expect(firstState).toEqual(duplicateState);
expect(firstState).toMatchObject({
conversationId: 'conversation-a',
workerId: 'worker-conversation-a',
state: 'ready',
generation: 1,
});
expect(otherState).toMatchObject({
conversationId: 'conversation-b',
workerId: 'worker-conversation-b',
state: 'ready',
generation: 1,
});
});
it('starts only four top-level runs and advances the remaining queue fairly on agent_settled', async () => {
const workers = new Map<string, FakeWorker>();
const pool = new PiWorkerPool({
maxIdle: 6,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.set(input.conversationId, worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
const ids = ['a', 'b', 'c', 'd', 'e', 'f'].map((id) => `conversation-${id}`);
await Promise.all(ids.map((id) => pool.prepare(conversation(id))));
const runs = ids.map((conversationId, index) => pool.startTopLevel({
conversationId,
runId: `run-${index + 1}`,
command: { type: 'prompt', message: conversationId },
}));
expect(runs.map((run) => run.queuePosition)).toEqual([undefined, undefined, undefined, undefined, 1, 2]);
await expect.poll(() => ids.map((id) => workers.get(id)!.requests.length))
.toEqual([1, 1, 1, 1, 0, 0]);
workers.get('conversation-a')!.emit({ type: 'agent_end' });
expect(workers.get('conversation-e')!.requests).toHaveLength(0);
workers.get('conversation-a')!.emit({ type: 'agent_settled' });
await expect.poll(() => workers.get('conversation-e')!.requests.length).toBe(1);
expect(workers.get('conversation-f')!.requests).toHaveLength(0);
workers.get('conversation-b')!.emit({ type: 'agent_settled' });
await expect.poll(() => workers.get('conversation-f')!.requests.length).toBe(1);
});
it('evicts only the least-recent idle worker and reopens it on demand', async () => {
const workers = new Map<string, FakeWorker[]>();
const pool = new PiWorkerPool({
maxIdle: 2,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}-${(workers.get(input.conversationId)?.length ?? 0) + 1}`);
workers.set(input.conversationId, [...(workers.get(input.conversationId) ?? []), worker]);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await pool.prepare(conversation('conversation-a'));
await pool.prepare(conversation('conversation-b'));
await pool.prepare(conversation('conversation-c'));
await expect.poll(() => workers.get('conversation-a')![0]!.stopped).toBe(true);
expect(pool.getState('conversation-a')).toBeNull();
expect(pool.getState('conversation-b')?.state).toBe('ready');
expect(pool.getState('conversation-c')?.state).toBe('ready');
const reopened = await pool.prepare(conversation('conversation-a'));
expect(reopened).toMatchObject({
workerId: 'worker-conversation-a-2',
generation: 2,
state: 'ready',
});
});
it('shares a fair total-process budget and starts the next worker only after a lease is released', async () => {
const processBudget = new PiProcessBudget(2);
const opened: string[] = [];
const pool = new PiWorkerPool({
maxIdle: 3,
processBudget,
openWorker: async ({ conversation: input }) => {
opened.push(input.conversationId);
return {
worker: new FakeWorker(`worker-${input.conversationId}`),
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all([
pool.prepare(conversation('conversation-a')),
pool.prepare(conversation('conversation-b')),
]);
const third = pool.prepare(conversation('conversation-c'));
await expect.poll(() => processBudget.waitingCount).toBe(1);
expect(opened).toEqual(['conversation-a', 'conversation-b']);
await pool.dispose('conversation-a', 'test_injection');
await expect(third).resolves.toMatchObject({ conversationId: 'conversation-c', state: 'ready' });
expect(opened).toEqual(['conversation-a', 'conversation-b', 'conversation-c']);
expect(processBudget.activeCount).toBe(2);
await pool.shutdown();
expect(processBudget.activeCount).toBe(0);
});
it('releases the process lease even when stopping an idle worker fails', async () => {
const processBudget = new PiProcessBudget(1);
const pool = new PiWorkerPool({
processBudget,
maxIdle: 1,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
worker.stop = async () => {
worker.stopped = true;
throw new Error('stop failed');
};
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await pool.prepare(conversation('conversation-stop-failure'));
await expect(pool.dispose('conversation-stop-failure', 'test_injection')).rejects.toThrow('stop failed');
expect(processBudget.activeCount).toBe(0);
expect(processBudget.waitingCount).toBe(0);
});
it('suspends and later resumes the oldest queued parent to make child capacity', async () => {
const processBudget = new PiProcessBudget(8);
const workers = new Map<string, FakeWorker[]>();
const queuedStopGate = deferred();
let stoppingQueuedWorkers = 0;
const pool = new PiWorkerPool({
processBudget,
maxRunning: 4,
maxIdle: 8,
openWorker: async ({ conversation: input, generation, existingSession }) => {
const worker = new FakeWorker(`worker-${input.conversationId}-${generation}`);
if (generation === 1 && Number(input.conversationId.split('-').at(-1)) > 4) {
worker.stop = async () => {
worker.stopped = true;
stoppingQueuedWorkers += 1;
await queuedStopGate.promise;
return { mode: 'stdin-close' as const, code: 0, signal: null };
};
}
workers.set(input.conversationId, [...(workers.get(input.conversationId) ?? []), worker]);
return {
worker,
session: existingSession ?? {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
const ids = Array.from({ length: 8 }, (_, index) => `conversation-${index + 1}`);
for (const id of ids) await pool.prepare(conversation(id));
const tickets = ids.map((conversationId, index) => pool.startTopLevel({
conversationId,
runId: `run-${index + 1}`,
command: { type: 'prompt', message: conversationId },
}));
for (const ticket of tickets) void ticket.accepted.catch(() => undefined);
await expect.poll(() => ids.map((id) => workers.get(id)?.[0]?.requests.length ?? 0))
.toEqual([1, 1, 1, 1, 0, 0, 0, 0]);
expect(ids.slice(4).map((id) => pool.getState(id)?.state))
.toEqual(['queued', 'queued', 'queued', 'queued']);
expect(processBudget.activeCount).toBe(8);
const scheduler = new PiSubagentScheduler({
processBudget,
reclaimProcessCapacity: (signal) => pool.reclaimIdleWorker(signal),
openChild: async (input) => ({
id: input.taskId,
async run() { return { summary: 'child complete' }; },
async stop() {},
}),
});
const children = scheduler.dispatch({
conversationId: ids[0]!,
workerGeneration: 1,
runId: 'run-child',
projectId: 'project-a',
request: {
mode: 'parallel',
tasks: Array.from({ length: 4 }, (_, index) => ({
agentId: `agent-${index + 1}`,
task: `inspect ${index + 1}`,
toolProfile: 'read-only',
})),
},
});
await expect.poll(() => stoppingQueuedWorkers).toBe(4);
expect(processBudget.activeCount).toBe(8);
expect(processBudget.waitingCount).toBe(4);
const firstQueuedId = ids[4]!;
let firstQueuedAccepted = false;
void tickets[4]!.accepted.then(
() => { firstQueuedAccepted = true; },
() => undefined,
);
workers.get(ids[0]!)?.[0]?.emit({ type: 'agent_settled' });
await new Promise((resolvePromise) => setImmediate(resolvePromise));
expect(firstQueuedAccepted).toBe(false);
expect(workers.get(ids[4]!)?.length).toBe(1);
expect(workers.get(ids[4]!)?.[0]?.requests).toHaveLength(0);
expect(pool.getState(firstQueuedId)).toMatchObject({
state: 'running',
generation: 1,
session: {
piSessionId: `session-${firstQueuedId}`,
sessionKey: `key-${firstQueuedId}`,
},
});
queuedStopGate.resolve();
await expect(children).resolves.toMatchObject({
details: {
tasks: Array.from({ length: 4 }, () => ({ status: 'complete' })),
},
});
expect(ids.slice(4).map((id) => workers.get(id)?.[0]?.stopped))
.toEqual([true, true, true, true]);
await expect.poll(() => workers.get(firstQueuedId)?.length).toBe(2);
await expect.poll(() => workers.get(firstQueuedId)?.[1]?.requests.length).toBe(1);
await expect(tickets[4]!.accepted).resolves.toMatchObject({ success: true });
expect(processBudget.activeCount).toBe(5);
expect(processBudget.waitingCount).toBe(0);
for (let index = 1; index < 4; index += 1) {
const queuedId = ids[index + 4]!;
workers.get(ids[index]!)?.[0]?.emit({ type: 'agent_settled' });
await expect.poll(() => workers.get(queuedId)?.length).toBe(2);
await expect.poll(() => workers.get(queuedId)?.[1]?.requests.length).toBe(1);
await expect(tickets[index + 4]!.accepted).resolves.toMatchObject({ success: true });
}
expect(pool.getState(firstQueuedId)).toMatchObject({
state: 'running',
generation: 2,
session: {
piSessionId: `session-${firstQueuedId}`,
sessionKey: `key-${firstQueuedId}`,
},
});
expect(ids.slice(4).map((id) => workers.get(id)?.length)).toEqual([2, 2, 2, 2]);
expect(processBudget.activeCount).toBe(8);
await scheduler.close();
await pool.shutdown();
expect(processBudget.activeCount).toBe(0);
expect(processBudget.waitingCount).toBe(0);
});
it('never evicts a running worker when the warm-idle LRU exceeds its cap', async () => {
const workers = new Map<string, FakeWorker>();
const pool = new PiWorkerPool({
maxIdle: 1,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.set(input.conversationId, worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await pool.prepare(conversation('conversation-running'));
const running = pool.startTopLevel({
conversationId: 'conversation-running',
runId: 'run-running',
command: { type: 'prompt', message: 'keep alive' },
});
await running.accepted;
await pool.prepare(conversation('conversation-old-idle'));
await pool.prepare(conversation('conversation-new-idle'));
expect(workers.get('conversation-running')!.stopped).toBe(false);
expect(pool.getState('conversation-running')).toMatchObject({ state: 'running' });
expect(workers.get('conversation-old-idle')!.stopped).toBe(true);
expect(pool.getState('conversation-new-idle')).toMatchObject({ state: 'ready' });
});
it('cleans only the crashed generation and releases its permit for the next Conversation', async () => {
const workers = new Map<string, FakeWorker>();
const processBudget = new PiProcessBudget(8);
const pool = new PiWorkerPool({
maxRunning: 2,
maxIdle: 3,
processBudget,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.set(input.conversationId, worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all(['a', 'b', 'c'].map((id) => pool.prepare(conversation(`conversation-${id}`))));
pool.startTopLevel({ conversationId: 'conversation-a', runId: 'run-a', command: { type: 'prompt', message: 'a' } });
pool.startTopLevel({ conversationId: 'conversation-b', runId: 'run-b', command: { type: 'prompt', message: 'b' } });
pool.startTopLevel({ conversationId: 'conversation-c', runId: 'run-c', command: { type: 'prompt', message: 'c' } });
await expect.poll(() => workers.get('conversation-a')!.requests.length).toBe(1);
const cancelled: string[] = [];
for (const kind of ['command', 'interaction', 'child'] as const) {
pool.trackGenerationResource({
conversationId: 'conversation-a',
kind,
id: `${kind}-a`,
cancel: () => { cancelled.push(kind); },
});
}
pool.trackGenerationResource({
conversationId: 'conversation-b',
kind: 'child',
id: 'child-b',
cancel: () => { cancelled.push('other'); },
});
workers.get('conversation-a')!.invalidate();
expect(cancelled.sort()).toEqual(['child', 'command', 'interaction']);
expect(pool.getState('conversation-a')).toMatchObject({ state: 'crashed', generation: 1 });
expect(pool.getState('conversation-b')).toMatchObject({ state: 'running', generation: 1 });
await expect.poll(() => workers.get('conversation-c')!.requests.length).toBe(1);
await expect.poll(() => processBudget.activeCount).toBe(2);
expect(cancelled).not.toContain('other');
});
it('rebuilds stale idle workers before prompt and lets running workers settle first', async () => {
const workers = new Map<string, FakeWorker[]>();
const revisions: Array<{ conversationId: string; provider: number; resources: number }> = [];
const telemetry: PiRuntimeTelemetryEvent[] = [];
const replacementReasons: string[] = [];
const pool = new PiWorkerPool({
maxIdle: 4,
onTelemetry: (event) => telemetry.push(event),
openWorker: async ({ conversation: input, revision }) => {
const worker = new FakeWorker(`worker-${input.conversationId}-${(workers.get(input.conversationId)?.length ?? 0) + 1}`);
workers.set(input.conversationId, [...(workers.get(input.conversationId) ?? []), worker]);
revisions.push({ conversationId: input.conversationId, ...revision });
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
pool.subscribe((event) => {
if (event.type === 'worker.replaced') replacementReasons.push(event.reason);
});
await Promise.all([
pool.prepare(conversation('conversation-running')),
pool.prepare(conversation('conversation-idle')),
]);
pool.startTopLevel({
conversationId: 'conversation-running',
runId: 'run-running',
command: { type: 'prompt', message: 'running' },
});
await expect.poll(() => workers.get('conversation-running')![0]!.requests.length).toBe(1);
pool.markProviderStale();
const idleTicket = pool.startTopLevel({
conversationId: 'conversation-idle',
runId: 'run-idle',
command: { type: 'prompt', message: 'idle' },
});
await expect.poll(() => workers.get('conversation-idle')?.length).toBe(2);
expect(workers.get('conversation-idle')![0]!.stopped).toBe(true);
await expect.poll(() => workers.get('conversation-idle')![1]!.requests.length).toBe(1);
await idleTicket.accepted;
expect(telemetry.filter(({ runRef }) => runRef === 'runidle')).toEqual([
expect.objectContaining({ milestone: 'worker.queue_wait', workerGeneration: 2, cold: false }),
expect.objectContaining({ milestone: 'prompt.accepted', workerGeneration: 2, cold: false }),
]);
expect(workers.get('conversation-running')).toHaveLength(1);
expect(workers.get('conversation-running')![0]!.stopped).toBe(false);
workers.get('conversation-running')![0]!.emit({ type: 'agent_settled' });
await expect.poll(() => workers.get('conversation-running')?.length).toBe(2);
expect(workers.get('conversation-running')![0]!.stopped).toBe(true);
expect(revisions).toEqual(expect.arrayContaining([
{ conversationId: 'conversation-idle', provider: 2, resources: 1 },
{ conversationId: 'conversation-running', provider: 2, resources: 1 },
]));
expect(replacementReasons).toEqual([
'stale_resource_rebuild',
'stale_resource_rebuild',
]);
expect(workers.get('conversation-idle')![0]!.stopReasons).toEqual([
'stale_resource_rebuild',
]);
expect(workers.get('conversation-running')![0]!.stopReasons).toEqual([
'stale_resource_rebuild',
]);
});
it('records an explicit recover replacement reason without changing the session binding', async () => {
const workers: FakeWorker[] = [];
const replacements: Array<{ reason: string; generation: number }> = [];
const pool = new PiWorkerPool({
openWorker: async ({ conversation: input, generation, existingSession }) => {
const worker = new FakeWorker(`worker-${generation}`);
workers.push(worker);
return {
worker,
session: existingSession ?? {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
pool.subscribe((event) => {
if (event.type === 'worker.replaced') {
replacements.push({ reason: event.reason, generation: event.generation });
}
});
const before = await pool.prepare(conversation('conversation-recover'));
const recovered = await pool.recover('conversation-recover');
expect(recovered.session).toEqual(before.session);
expect(replacements).toEqual([{ reason: 'recover', generation: 2 }]);
expect(workers[0]!.stopReasons).toEqual(['recover']);
});
it('re-applies the idle LRU after a running stale worker rebuilds on settle', async () => {
const workers = new Map<string, FakeWorker[]>();
const pool = new PiWorkerPool({
maxIdle: 1,
openWorker: async ({ conversation: input, generation, existingSession }) => {
const worker = new FakeWorker(`worker-${input.conversationId}-${generation}`);
workers.set(input.conversationId, [...(workers.get(input.conversationId) ?? []), worker]);
return {
worker,
session: existingSession ?? {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await pool.prepare(conversation('conversation-running'));
const running = pool.startTopLevel({
conversationId: 'conversation-running',
runId: 'run-running',
command: { type: 'prompt', message: 'running' },
});
await running.accepted;
await pool.prepare(conversation('conversation-idle'));
pool.markProviderStale();
workers.get('conversation-running')![0]!.emit({ type: 'agent_settled' });
await expect.poll(() => workers.get('conversation-running')?.length).toBe(2);
await expect.poll(() => workers.get('conversation-idle')![0]!.stopped).toBe(true);
expect(pool.getState('conversation-idle')).toBeNull();
expect(pool.getState('conversation-running')).toMatchObject({ state: 'ready', generation: 2 });
});
it('reuses an in-flight rebuild for recover and leaves no unowned replacement worker', async () => {
const generationTwoGate = deferred();
const workers: FakeWorker[] = [];
const openedGenerations: number[] = [];
const pool = new PiWorkerPool({
maxIdle: 2,
openWorker: async ({ conversation: input, generation, existingSession }) => {
openedGenerations.push(generation);
const worker = new FakeWorker(`worker-${input.conversationId}-${generation}`);
workers.push(worker);
if (generation === 2) await generationTwoGate.promise;
return {
worker,
session: existingSession ?? {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await pool.prepare(conversation('conversation-a'));
pool.markProviderStale();
const ticket = pool.startTopLevel({
conversationId: 'conversation-a',
runId: 'run-a',
command: { type: 'prompt', message: 'do not replay' },
});
const ticketOutcome = ticket.accepted.then(
() => 'resolved',
(error: unknown) => error instanceof Error ? error.message : String(error),
);
await expect.poll(() => openedGenerations).toEqual([1, 2]);
const recovered = pool.recover('conversation-a');
await new Promise((resolve) => setTimeout(resolve, 10));
expect(openedGenerations).toEqual([1, 2]);
generationTwoGate.resolve();
await expect(recovered).resolves.toMatchObject({ generation: 2, state: 'ready' });
expect(await ticketOutcome).toMatch(/recovering|cancelled/);
await pool.shutdown();
expect(workers).toHaveLength(2);
expect(workers.every((worker) => worker.stopped)).toBe(true);
});
it('rejects queued work and stops every parent worker during app shutdown', async () => {
const workers: FakeWorker[] = [];
const pool = new PiWorkerPool({
maxRunning: 1,
maxIdle: 3,
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.push(worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all(['a', 'b', 'c'].map((id) => pool.prepare(conversation(`conversation-${id}`))));
pool.startTopLevel({ conversationId: 'conversation-a', runId: 'run-a', command: { type: 'prompt', message: 'a' } });
const queued = pool.startTopLevel({
conversationId: 'conversation-b',
runId: 'run-b',
command: { type: 'prompt', message: 'b' },
});
let childCancelled = false;
pool.trackGenerationResource({
conversationId: 'conversation-a',
kind: 'child',
id: 'child-a',
cancel: () => { childCancelled = true; },
});
await pool.shutdown();
await expect(queued.accepted).rejects.toThrow('shutting down');
expect(childCancelled).toBe(true);
expect(workers.every((worker) => worker.stopped)).toBe(true);
expect(pool.getState('conversation-a')).toBeNull();
await expect(pool.prepare(conversation('conversation-after-quit')))
.rejects.toThrow('shutting down');
});
it('waits for an in-flight fork open and stops that process before shutdown completes', async () => {
const forkGate = deferred();
const workers: FakeWorker[] = [];
let forkOpenStarted = false;
const pool = new PiWorkerPool({
maxIdle: 3,
openWorker: async ({ conversation: input, existingSession }) => {
if (input.conversationId === 'conversation-fork') {
forkOpenStarted = true;
await forkGate.promise;
}
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.push(worker);
return {
worker,
session: existingSession ?? {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await pool.prepare(conversation('conversation-source'));
const fork = pool.fork('conversation-source', conversation('conversation-fork'));
await expect.poll(() => forkOpenStarted).toBe(true);
const forkOutcome = fork.then(
() => 'resolved',
(error: unknown) => error instanceof Error ? error.message : String(error),
);
let shutdownCompleted = false;
const shutdown = pool.shutdown().then(() => { shutdownCompleted = true; });
await Promise.resolve();
expect(shutdownCompleted).toBe(false);
forkGate.resolve();
await shutdown;
expect(await forkOutcome).toContain('shutting down');
expect(workers).toHaveLength(2);
expect(workers.every((worker) => worker.stopped)).toBe(true);
});
it('recovers the target session with a new generation and disposes no sibling worker', async () => {
const workers = new Map<string, FakeWorker[]>();
const pool = new PiWorkerPool({
maxIdle: 4,
openWorker: async ({ conversation: input, existingSession }) => {
const worker = new FakeWorker(`worker-${input.conversationId}-${(workers.get(input.conversationId)?.length ?? 0) + 1}`);
workers.set(input.conversationId, [...(workers.get(input.conversationId) ?? []), worker]);
return {
worker,
session: existingSession ?? {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all([
pool.prepare(conversation('conversation-a')),
pool.prepare(conversation('conversation-b')),
]);
workers.get('conversation-a')![0]!.invalidate();
const recovered = await pool.recover('conversation-a');
expect(recovered).toMatchObject({
workerId: 'worker-conversation-a-2',
generation: 2,
state: 'ready',
session: { piSessionId: 'session-conversation-a', sessionKey: 'key-conversation-a' },
});
expect(workers.get('conversation-a')![0]!.stopped).toBe(true);
expect(workers.get('conversation-b')![0]!.stopped).toBe(false);
await pool.dispose('conversation-a', 'test_injection');
expect(workers.get('conversation-a')![1]!.stopped).toBe(true);
expect(pool.getState('conversation-a')).toBeNull();
expect(pool.getState('conversation-b')).toMatchObject({ state: 'ready', generation: 1 });
});
it('records privacy-safe queue, acceptance, and authoritative settled spans', async () => {
let now = 0;
const telemetry: PiRuntimeTelemetryEvent[] = [];
const workers = new Map<string, FakeWorker>();
const pool = new PiWorkerPool({
maxRunning: 1,
maxIdle: 2,
now: () => now,
onTelemetry: (event) => telemetry.push(event),
openWorker: async ({ conversation: input }) => {
const worker = new FakeWorker(`worker-${input.conversationId}`);
workers.set(input.conversationId, worker);
return {
worker,
session: {
piSessionId: `session-${input.conversationId}`,
sessionKey: `key-${input.conversationId}`,
},
};
},
});
await Promise.all([
pool.prepare(conversation('f47ac10b-58cc-4372-a567-0e02b2c3d479')),
pool.prepare(conversation('8b1a9953-c461-4d88-9c3e-7e1f8f3f2c11')),
]);
const first = pool.startTopLevel({
conversationId: 'f47ac10b-58cc-4372-a567-0e02b2c3d479',
runId: 'run-first-1234567890',
command: { type: 'prompt', message: 'private first prompt' },
});
now = 5;
const second = pool.startTopLevel({
conversationId: '8b1a9953-c461-4d88-9c3e-7e1f8f3f2c11',
runId: 'run-second-1234567890',
command: { type: 'prompt', message: 'private second prompt' },
});
await first.accepted;
now = 25;
workers.get('f47ac10b-58cc-4372-a567-0e02b2c3d479')!.emit({ type: 'agent_settled' });
await second.accepted;
expect(telemetry.map(({ milestone }) => milestone)).toEqual([
'worker.queue_wait',
'prompt.accepted',
'agent.settled',
'worker.queue_wait',
'prompt.accepted',
]);
expect(telemetry[0]).toMatchObject({ durationMs: 0, workerGeneration: 1, cold: true });
expect(telemetry[2]).toMatchObject({
milestone: 'agent.settled', durationMs: 20, workerGeneration: 1, cold: true,
});
expect(telemetry[3]).toMatchObject({ durationMs: 20, workerGeneration: 1, cold: true });
const serialized = JSON.stringify(telemetry);
expect(serialized).not.toContain('private first prompt');
expect(serialized).not.toContain('private second prompt');
expect(serialized).not.toContain('f47ac10b-58cc-4372-a567-0e02b2c3d479');
expect(serialized).not.toContain('8b1a9953-c461-4d88-9c3e-7e1f8f3f2c11');
});
});