fix: close Pi worker rebuild races
This commit is contained in:
@@ -354,6 +354,84 @@ describe('Pi worker pool', () => {
|
||||
]));
|
||||
});
|
||||
|
||||
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({
|
||||
|
||||
Reference in New Issue
Block a user