// @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; resolve(): void } { let resolve!: () => void; const promise = new Promise((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(command: PiRpcCommand, options?: PiRpcRequestOptions) { this.requests.push(command); if (this.pendingType === command.type) { this.pendingType = null; return await new Promise>((_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 { 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 { 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(); 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(); 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(); 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 => { 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(); 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(); 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(); 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(); 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(); 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('refreshes idle worker resources immediately and defers an active run until settlement', async () => { const workers = new Map(); 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-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); await pool.refreshResources(); expect(workers.get('conversation-idle')).toHaveLength(2); expect(workers.get('conversation-idle')![0]!.stopped).toBe(true); expect(workers.get('conversation-running')).toHaveLength(1); workers.get('conversation-running')![0]!.emit({ type: 'agent_settled' }); await expect.poll(() => workers.get('conversation-running')).toHaveLength(2); }); it('rebuilds stale idle workers before prompt and lets running workers settle first', async () => { const workers = new Map(); 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(); 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(); 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(); 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'); }); });