Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/evict-dead-warmed-processes.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@livekit/agents': patch
---

Remove warmed processes from the pool queue when they exit before receiving a job.
133 changes: 115 additions & 18 deletions agents/src/ipc/proc_pool.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,26 @@ async function flushMicrotasks(ticks = 10): Promise<void> {
}
}

function createControlledJoinExecutor() {
let resolveJoin: () => void = () => {};
const joinPromise = new Promise<void>((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<void, Error>
Expand Down Expand Up @@ -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<void>((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);

Expand All @@ -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<void, Error>
> => {
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<void, Error>
> => {
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<void, Error>
> => {
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();
}
Expand Down
20 changes: 15 additions & 5 deletions agents/src/ipc/proc_pool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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();
Expand All @@ -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 {
Expand Down Expand Up @@ -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();
});
Expand Down