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
5 changes: 5 additions & 0 deletions .changeset/quiet-workers-exit.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@sveltejs/kit': patch
---

fix: exit build workers after completing their tasks while allowing synchronous exit handlers to run
59 changes: 59 additions & 0 deletions packages/kit/src/utils/fixtures/fork/index.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
import process from 'node:process';
import { setTimeout } from 'node:timers/promises';
import { forked } from '../../fork.js';

export const run = forked(
import.meta.url,
/**
* @param {{
* action: 'return' | 'undefined' | 'throw' | 'exit' | 'exit-code' | 'exit-error' | 'stdio';
* state?: Int32Array;
* code?: number;
* error?: unknown;
* }} options
*/
async ({ action, state, code = 0, error = new Error('exit handler failed') }) => {
if (state) {
setInterval(() => Atomics.add(state, 3, 1), 1);
process.on('exit', () => {
// Make resolving on the result message observably different from waiting for exit.
Atomics.wait(state, 2, 0, 50);
Atomics.store(state, 2, 1);
});
}

try {
await setTimeout(10);
if (state) Atomics.store(state, 0, 1);

switch (action) {
case 'undefined':
return undefined;
case 'throw':
throw new Error('callback failed');
case 'exit':
return process.exit(code);
case 'exit-code':
process.exitCode = code;
break;
case 'exit-error':
process.on('exit', () => {
throw error;
});
break;
case 'stdio':
for (let i = 0; i < 128; i += 1) {
process.stdout.write('out\n'.repeat(1024));
process.stderr.write('err\n'.repeat(1024));
}
}

return { answer: 42 };
} finally {
if (state) {
await setTimeout(10);
Atomics.store(state, 1, 1);
}
}
}
);
3 changes: 3 additions & 0 deletions packages/kit/src/utils/fixtures/fork/stdio.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
import { run } from './index.js';

await run({ action: 'stdio' });
51 changes: 41 additions & 10 deletions packages/kit/src/utils/fork.js
Original file line number Diff line number Diff line change
@@ -1,27 +1,45 @@
import { fileURLToPath } from 'node:url';
import { Worker, parentPort } from 'node:worker_threads';
import process from 'node:process';
import { setImmediate } from 'node:timers';

/**
* Runs a task in a subprocess so any dangling stuff gets killed upon completion.
* The subprocess needs to be the file `forked` is called in, and `forked` needs to be called eagerly at the top level.
* Runs a task in a worker that exits upon completion, stopping any dangling work.
* The worker needs to be the file `forked` is called in, and `forked` needs to be called eagerly at the top level.
* @template T
* @template U
* @param {string} module `import.meta.url` of the file
* @param {(opts: T) => Promise<U>} callback The function that is invoked in the subprocess
* @returns {(opts: T) => Promise<U>} A function that when called starts the subprocess
* @param {(opts: T) => Promise<U>} callback The function invoked in the worker. It must await all required work and asynchronous cleanup.
* @returns {(opts: T) => Promise<U>} A function that starts the worker and resolves after it has exited, including synchronous exit handlers
*/
export function forked(module, callback) {
if (process.env.SVELTEKIT_FORK && parentPort) {
parentPort.on(
'message',
/** @param {any} data */ async (data) => {
if (data?.type === 'args' && data.module === module) {
const payload = await callback(data.payload);

// Flush buffered output before exiting, rather than cutting off pending writes.
await Promise.all(
[process.stdout, process.stderr].map(
(stream) =>
new Promise((fulfil, reject) => {
stream.write('', (error) => (error ? reject(error) : fulfil(undefined)));
})
)
);

parentPort?.postMessage({
type: 'result',
module,
payload: await callback(data.payload)
payload
});

// Unlike worker.terminate(), this runs synchronous process.on('exit') handlers.
// Exit outside this async callback so errors in exit handlers are uncaught
// exceptions, rather than promise rejections after Node has started exiting.
setImmediate(() => process.exit());
}
}
);
Expand All @@ -43,6 +61,11 @@ export function forked(module, callback) {
}
});

/** @type {{ payload: U } | undefined} */
let result;
/** @type {{ error: unknown } | undefined} */
let failure;

worker.on(
'message',
/** @param {any} data */ (data) => {
Expand All @@ -55,19 +78,27 @@ export function forked(module, callback) {
}

if (data?.type === 'result' && data.module === module) {
worker.unref();
fulfil(data.payload);
result = data;
}
}
);

worker.once('error', reject);
worker.once('error', (error) => {
failure = { error };
});

worker.on('exit', (code) => {
if (code) {
worker.once('exit', (code) => {
if (failure) {
reject(failure.error);
} else if (code) {
const error = new Error(`Failed with code ${code}`);
error.stack = error.message;
reject(error);
} else if (!result) {
reject(new Error('Worker exited without returning a result'));
} else {
// All messages are delivered before 'exit'. Wait until cleanup has finished.
fulfil(result.payload);
}
});
});
Expand Down
58 changes: 58 additions & 0 deletions packages/kit/src/utils/fork.spec.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
import { execFile } from 'node:child_process';
import process from 'node:process';
import { setTimeout } from 'node:timers/promises';
import { fileURLToPath } from 'node:url';
import { promisify } from 'node:util';
import { expect, test } from 'vitest';
import { run } from './fixtures/fork/index.js';

test('waits for the callback, async finally and exit handlers, then stops background work', async () => {
const state = new Int32Array(new SharedArrayBuffer(4 * Int32Array.BYTES_PER_ELEMENT));

await expect(run({ action: 'return', state })).resolves.toEqual({ answer: 42 });
expect(Array.from(state.slice(0, 3))).toEqual([1, 1, 1]);

const ticks = Atomics.load(state, 3);
await setTimeout(20);
expect(Atomics.load(state, 3)).toBe(ticks);
});

test('accepts undefined as a result', async () => {
await expect(run({ action: 'undefined' })).resolves.toBeUndefined();
});

test('propagates callback errors after the worker has exited', async () => {
const state = new Int32Array(new SharedArrayBuffer(4 * Int32Array.BYTES_PER_ELEMENT));

await expect(run({ action: 'throw', state })).rejects.toThrow('callback failed');
expect(Array.from(state.slice(0, 3))).toEqual([1, 1, 1]);
});

test('rejects when the worker exits successfully without a result', async () => {
await expect(run({ action: 'exit' })).rejects.toThrow('Worker exited without returning a result');
});

test('rejects when the worker exits unsuccessfully without a result', async () => {
await expect(run({ action: 'exit', code: 2 })).rejects.toThrow('Failed with code 2');
});

test('rejects a nonzero exit code even after receiving a result', async () => {
await expect(run({ action: 'exit-code', code: 2 })).rejects.toThrow('Failed with code 2');
});

test('propagates exit handler errors even after receiving a result', async () => {
await expect(run({ action: 'exit-error' })).rejects.toThrow('exit handler failed');
});

test('preserves falsy values thrown by exit handlers', async () => {
await expect(run({ action: 'exit-error', error: null })).rejects.toBeNull();
});

test('delivers buffered stdout and stderr before exiting', async () => {
const { stdout, stderr } = await promisify(execFile)(process.execPath, [
fileURLToPath(new URL('./fixtures/fork/stdio.js', import.meta.url))
]);

expect(stdout).toBe('out\n'.repeat(128 * 1024));
expect(stderr).toBe('err\n'.repeat(128 * 1024));
});
Loading