Skip to content
Merged
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
52 changes: 52 additions & 0 deletions src/__tests__/server.workerPool.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,58 @@ describe('buildPersistentPool', () => {
await expect(t2).resolves.toBe('result-2');
await expect(t3).resolves.toBe('result-3');
});

it('should reset slot.active and reuse slot after a worker crash during task execution', async () => {
jest.useFakeTimers();
const instances: any[] = [];

MockWorker.mockImplementation((): any => {
const listeners: Record<string, any> = {};
const workerInstance = {
on: jest.fn((event: string, cb: any): any => {
listeners[event] = cb;

return workerInstance;
}),
off: jest.fn((event: string, _cb: any): any => {
delete listeners[event];

return workerInstance;
}),
removeAllListeners: jest.fn(),
terminate: jest.fn().mockResolvedValue(0),
postMessage: jest.fn(),
emit: (event: string, value: any) => {
if (listeners[event]) {
listeners[event](value);
}
}
};

instances.push(workerInstance);

return workerInstance as any;
});
const pool = buildPersistentPool(1);
// Run task 1 which will crash
const task1 = pool.runTask({ moduleSpecifier: 'fail-spec', args: {} });

instances[0].emit('exit', 1);
await expect(task1).rejects.toThrow('Persistent worker exited unexpectedly with code 1');

// Fast-forward past the backoff timer to trigger worker respawn
jest.advanceTimersByTime(100);

// Run task 2, verify it can execute in the recovered slot
const task2 = pool.runTask({ moduleSpecifier: 'next-spec', args: {} });
const secondWorker = instances[1];

expect(secondWorker).toBeDefined();
secondWorker.emit('message', { success: true, payload: 'recovered-success' });
await expect(task2).resolves.toBe('recovered-success');

jest.useRealTimers();
});
});

describe('buildTransientPool', () => {
Expand Down
121 changes: 97 additions & 24 deletions src/server.workerPool.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { Worker } from 'node:worker_threads';
import { fileURLToPath } from 'node:url';
import { availableParallelism } from 'node:os';
import { formatUnknownError } from './logger';
import { formatUnknownError, log } from './logger';

/**
* Payload for a task execution, including module details, arguments, and configuration options.
Expand Down Expand Up @@ -49,6 +49,25 @@ interface QueuedTask {
reject: (reason: unknown) => void;
}

/**
* Worker instance. Props and methods for worker instance status and respawn behavior.
*
* @interface Workers
*
* @property worker Optional worker instance.
* @property active Is worker active.
* @property reject Optional promise reject for worker errors.
* @property consecutiveCrashes Count consecutive crashes.
* @property respawnTimer Optional timer for worker respawn.
*/
interface Workers {
worker?: Worker;
active: boolean;
reject?: (reason: unknown) => void;
consecutiveCrashes: number;
respawnTimer?: NodeJS.Timeout;
}

/**
* Throttled worker thread pool for parallel execution.
*
Expand Down Expand Up @@ -79,6 +98,21 @@ const poolRegistry = new Map<PoolKind, WorkerPoolInstance>();
*/
const MAX_QUEUE_CAP = 50;

/**
* Maximum number of consecutive worker crashes before shutdown.
*/
const MAX_CONSECUTIVE_CRASHES = 5;

/**
* Minimum backoff time for worker respawn.
*/
const MIN_BACKOFF_MS = 50;

/**
* Maximum backoff time for worker respawn.
*/
const MAX_BACKOFF_MS = 2000;

/**
* Resolves the location of the worker entry script safely across bundling and testing frameworks.
*/
Expand Down Expand Up @@ -271,41 +305,73 @@ const buildTransientPool = (maxWorkers = Math.max(1, availableParallelism() - 1)
*/
const buildPersistentPool = (maxWorkers = Math.max(1, availableParallelism() - 1)): WorkerPoolInstance => {
const queue: QueuedTask[] = [];
const workers: { worker: Worker; active: boolean; reject?: (reason: unknown) => void }[] = [];
const workers: Workers[] = [];
const workerScript = getWorkerScriptPath();
const poolAbort = createPoolAbort();

const spawnWorker = (index: number) => {
const slot = workers[index];

if (!slot || poolAbort.isAborted()) {
return;
}

const worker = new Worker(workerScript);

slot.worker = worker;
worker.on('error', () => handleWorkerCrash(index));
worker.on('exit', (code: number) => handleWorkerCrash(index, code));
next();
};

const handleWorkerCrash = (index: number, exitCode?: number) => {
const slot = workers[index];

if (!slot) {
if (!slot || poolAbort.isAborted()) {
return;
}

slot.worker?.removeAllListeners();

if (slot.reject) {
delete slot.worker;

const reject = slot.reject;

delete slot.reject;
slot.active = false;

if (reject) {
const message = exitCode !== undefined
? `Persistent worker exited unexpectedly with code ${exitCode}`
: 'Persistent worker thread crashed';

slot.reject(new Error(message));
} else {
slot.active = false;
reject(new Error(message));
}

slot.worker = new Worker(workerScript);
slot.worker.on('error', () => handleWorkerCrash(index));
slot.worker.on('exit', (code: number) => handleWorkerCrash(index, code));
next();
slot.consecutiveCrashes += 1;
if (slot.consecutiveCrashes > MAX_CONSECUTIVE_CRASHES) {
log.warn(`Persistent worker slot ${index} reached max crash limit (${MAX_CONSECUTIVE_CRASHES}); halting respawn.`);

return;
}

const backoffMs = Math.min(
MAX_BACKOFF_MS,
MIN_BACKOFF_MS * 2 ** (slot.consecutiveCrashes - 1) + Math.random() * 20
);

slot.respawnTimer = setTimeout(() => {
delete slot.respawnTimer;
spawnWorker(index);
}, backoffMs);
};

const next = (): void => {
if (poolAbort.isAborted() || queue.length === 0) {
return;
}

const idleWorkerSlot = workers.find(worker => !worker.active);
const idleWorkerSlot = workers.find(worker => !worker.active && worker.worker);

if (!idleWorkerSlot) {
return;
Expand All @@ -325,6 +391,8 @@ const buildPersistentPool = (maxWorkers = Math.max(1, availableParallelism() - 1

const onMessage = (message: WorkerIpcMessage) => {
cleanup();
idleWorkerSlot.consecutiveCrashes = 0;

if (message && message.success) {
resolve(message.payload);
} else {
Expand All @@ -338,28 +406,25 @@ const buildPersistentPool = (maxWorkers = Math.max(1, availableParallelism() - 1
};

const cleanup = () => {
worker.off('message', onMessage);
worker.off('error', onError);
worker?.off('message', onMessage);
worker?.off('error', onError);
idleWorkerSlot.active = false;
delete idleWorkerSlot.reject;

next();
};

worker.on('message', onMessage);
worker.on('error', onError);
worker?.on('message', onMessage);
worker?.on('error', onError);

// Send the task data via postMessage channel to be captured by persistent listeners
worker.postMessage(payload);
worker?.postMessage(payload);
};

// Pre-spawn and warm up the permanent thread containers
for (let i = 0; i < maxWorkers; i++) {
const worker = new Worker(workerScript); // Instantiated WITHOUT initial workerData

workers.push({ worker, active: false });
worker.on('error', () => handleWorkerCrash(i));
worker.on('exit', (code: number) => handleWorkerCrash(i, code));
workers.push({ active: false, consecutiveCrashes: 0 });
spawnWorker(i);
}

return {
Expand All @@ -385,8 +450,16 @@ const buildPersistentPool = (maxWorkers = Math.max(1, availableParallelism() - 1
queue.length = 0;

for (const slot of workers) {
slot.worker.removeAllListeners();
await slot.worker.terminate();
if (slot.respawnTimer) {
clearTimeout(slot.respawnTimer);
delete slot.respawnTimer;
}

slot.worker?.removeAllListeners();
await slot.worker?.terminate();

delete slot.worker;

slot.active = false;
delete slot.reject;
}
Expand Down
Loading