From 87f4d75123c76cd18b4fb6ca636787c8d9ade68f Mon Sep 17 00:00:00 2001 From: Pranjal Date: Mon, 31 Aug 2026 16:07:58 +0530 Subject: [PATCH] fix(agents): evict dead warmed processes --- .changeset/evict-dead-warmed-processes.md | 5 + agents/src/ipc/proc_pool.test.ts | 133 +++++++++++++++++++--- agents/src/ipc/proc_pool.ts | 20 +++- 3 files changed, 135 insertions(+), 23 deletions(-) create mode 100644 .changeset/evict-dead-warmed-processes.md diff --git a/.changeset/evict-dead-warmed-processes.md b/.changeset/evict-dead-warmed-processes.md new file mode 100644 index 0000000000..64af66de09 --- /dev/null +++ b/.changeset/evict-dead-warmed-processes.md @@ -0,0 +1,5 @@ +--- +'@livekit/agents': patch +--- + +Remove warmed processes from the pool queue when they exit before receiving a job. diff --git a/agents/src/ipc/proc_pool.test.ts b/agents/src/ipc/proc_pool.test.ts index 491335b44a..fa55e50352 100644 --- a/agents/src/ipc/proc_pool.test.ts +++ b/agents/src/ipc/proc_pool.test.ts @@ -30,6 +30,26 @@ async function flushMicrotasks(ticks = 10): Promise { } } +function createControlledJoinExecutor() { + let resolveJoin: () => void = () => {}; + const joinPromise = new Promise((resolve) => { + resolveJoin = resolve; + }); + const executor: JobExecutor = { + ...createMockExecutor(), + join: vi.fn(() => joinPromise), + }; + return { executor, resolveJoin }; +} + +function mockJobProcExecutor(executor: JobExecutor) { + return vi + .spyOn(jobProcExecutorModule, 'JobProcExecutor') + .mockImplementation(function MockJobProcExecutor(this: unknown) { + return executor as unknown as jobProcExecutorModule.JobProcExecutor; + } as unknown as typeof jobProcExecutorModule.JobProcExecutor); +} + describe('ProcPool warmed process lock handling', () => { it('releases lock token from the dequeued warmed process entry', async (): Promise< Throws @@ -87,21 +107,8 @@ describe('ProcPool warmed process lock handling', () => { const pool = new ProcPool('agent', 1, 1000, 1000, undefined, 0, 0); const initUnlock = vi.fn(); const procUnlock = vi.fn(); - - let joinResolve: () => void = () => {}; - const joinPromise = new Promise((resolve) => { - joinResolve = resolve; - }); - const mockProc: JobExecutor = { - ...createMockExecutor(), - join: vi.fn(() => joinPromise), - }; - - const jobProcExecutorSpy = vi - .spyOn(jobProcExecutorModule, 'JobProcExecutor') - .mockImplementation(function MockJobProcExecutor(this: unknown) { - return mockProc as unknown as jobProcExecutorModule.JobProcExecutor; - } as unknown as typeof jobProcExecutorModule.JobProcExecutor); + const { executor: mockProc, resolveJoin } = createControlledJoinExecutor(); + const jobProcExecutorSpy = mockJobProcExecutor(mockProc); pool.initMutex.lock = vi.fn(async () => initUnlock); @@ -114,12 +121,102 @@ describe('ProcPool warmed process lock handling', () => { expect(pool.warmedProcQueue.items.length).toBe(1); expect(mockProc.join).toHaveBeenCalledTimes(1); - joinResolve(); + const warmedProcEntry = await pool.warmedProcQueue.get(); + warmedProcEntry.unlock(); + resolveJoin(); await watchPromise; - // finally block must not double-release. + // The watcher must not release a lock token already claimed by a consumer. expect(initUnlock).toHaveBeenCalledTimes(1); - expect(procUnlock).not.toHaveBeenCalled(); + expect(procUnlock).toHaveBeenCalledTimes(1); + } finally { + jobProcExecutorSpy.mockRestore(); + } + }); + + it('evicts a warmed process that exits before it is dequeued', async (): Promise< + Throws + > => { + const pool = new ProcPool('agent', 1, 1000, 1000, undefined, 0, 0); + const initUnlock = vi.fn(); + const procUnlock = vi.fn(); + const { executor: mockProc, resolveJoin } = createControlledJoinExecutor(); + const jobProcExecutorSpy = mockJobProcExecutor(mockProc); + + pool.initMutex.lock = vi.fn(async () => initUnlock); + + try { + const watchPromise = pool.procWatchTask(procUnlock); + await flushMicrotasks(); + + expect(pool.warmedProcQueue.items).toHaveLength(1); + + resolveJoin(); + await watchPromise; + + expect(pool.warmedProcQueue.items).toHaveLength(0); + expect(procUnlock).toHaveBeenCalledTimes(1); + } finally { + jobProcExecutorSpy.mockRestore(); + } + }); + + it('keeps other warmed processes queued when one exits', async (): Promise< + Throws + > => { + const pool = new ProcPool('agent', 2, 1000, 1000, undefined, 0, 0); + const initUnlock = vi.fn(); + const healthyProcUnlock = vi.fn(); + const exitedProcUnlock = vi.fn(); + const healthyProc = createMockExecutor(); + const { executor: exitedProc, resolveJoin } = createControlledJoinExecutor(); + const jobProcExecutorSpy = mockJobProcExecutor(exitedProc); + + pool.initMutex.lock = vi.fn(async () => initUnlock); + await pool.warmedProcQueue.put({ proc: healthyProc, unlock: healthyProcUnlock }); + + try { + const watchPromise = pool.procWatchTask(exitedProcUnlock); + await flushMicrotasks(); + + expect(pool.warmedProcQueue.items).toHaveLength(2); + + resolveJoin(); + await watchPromise; + + expect(pool.warmedProcQueue.items).toEqual([ + { proc: healthyProc, unlock: healthyProcUnlock }, + ]); + expect(exitedProcUnlock).toHaveBeenCalledTimes(1); + expect(healthyProcUnlock).not.toHaveBeenCalled(); + } finally { + jobProcExecutorSpy.mockRestore(); + } + }); + + it('does not double-release a queued process lock during close', async (): Promise< + Throws + > => { + const pool = new ProcPool('agent', 1, 1000, 1000, undefined, 0, 0); + const initUnlock = vi.fn(); + const procUnlock = vi.fn(); + const { executor: mockProc, resolveJoin } = createControlledJoinExecutor(); + const jobProcExecutorSpy = mockJobProcExecutor(mockProc); + + pool.initMutex.lock = vi.fn(async () => initUnlock); + pool.started = true; + + try { + const watchPromise = pool.procWatchTask(procUnlock); + await flushMicrotasks(); + + const closePromise = pool.close(); + expect(pool.warmedProcQueue.items).toHaveLength(0); + + resolveJoin(); + await Promise.all([closePromise, watchPromise]); + + expect(procUnlock).toHaveBeenCalledTimes(1); } finally { jobProcExecutorSpy.mockRestore(); } diff --git a/agents/src/ipc/proc_pool.ts b/agents/src/ipc/proc_pool.ts index 3dffacbe5a..b2315005b4 100644 --- a/agents/src/ipc/proc_pool.ts +++ b/agents/src/ipc/proc_pool.ts @@ -98,7 +98,7 @@ export class ProcPool { const unlock = await this.initMutex.lock(); let initReleased = false; - let procUnlockTransferred = false; + let warmedProcEntry: { proc: JobExecutor; unlock: () => void } | undefined; try { if (this.closed) { return; @@ -107,8 +107,9 @@ export class ProcPool { await proc.start(); try { await proc.initialize(); - await this.warmedProcQueue.put({ proc, unlock: procUnlock }); - procUnlockTransferred = true; + const entry = { proc, unlock: procUnlock }; + await this.warmedProcQueue.put(entry); + warmedProcEntry = entry; // Release initMutex after enqueue — holding it through join() serialises // the pool to concurrency 1 since child procs are one-shot. unlock(); @@ -122,8 +123,15 @@ export class ProcPool { if (!initReleased) { unlock(); } - if (!procUnlockTransferred) { + if (!warmedProcEntry) { procUnlock(); + } else { + // A missing entry has already been claimed by launchJob() or close(). + const entryIndex = this.warmedProcQueue.items.indexOf(warmedProcEntry); + if (entryIndex !== -1) { + this.warmedProcQueue.items.splice(entryIndex, 1); + warmedProcEntry.unlock(); + } } } } finally { @@ -169,7 +177,9 @@ export class ProcPool { } this.closed = true; this.controller.abort(); - this.warmedProcQueue.items.forEach((e) => { + // Claim queued entries before closing them so their watchers cannot release the same slots. + const warmedProcs = this.warmedProcQueue.items.splice(0); + warmedProcs.forEach((e) => { e.unlock(); e.proc.close(); });