From 08faa59567d87698ce9f20d57555fb6df6519b34 Mon Sep 17 00:00:00 2001 From: Basit Chonka Date: Fri, 18 Sep 2026 15:36:01 +0200 Subject: [PATCH 1/6] terminate worker if execution exceeds 2mins --- packages/shell-bson-parser/.eslintrc.cjs | 2 +- packages/shell-bson-parser/src/index.spec.ts | 37 ++++++++++++++++++- .../shell-bson-parser/src/worker-client.ts | 30 ++++++++++++++- .../test/fixtures/slow-worker.mjs | 9 +++++ packages/shell-bson-parser/tsconfig-lint.json | 2 +- 5 files changed, 75 insertions(+), 5 deletions(-) create mode 100644 packages/shell-bson-parser/test/fixtures/slow-worker.mjs diff --git a/packages/shell-bson-parser/.eslintrc.cjs b/packages/shell-bson-parser/.eslintrc.cjs index 4b264cafc..5ead6c7f1 100644 --- a/packages/shell-bson-parser/.eslintrc.cjs +++ b/packages/shell-bson-parser/.eslintrc.cjs @@ -5,5 +5,5 @@ module.exports = { tsconfigRootDir: __dirname, project: ['./tsconfig-lint.json'], }, - ignorePatterns: ['webpack.worker.config.cjs'], + ignorePatterns: ['webpack.worker.config.cjs', 'test/fixtures/**'], }; diff --git a/packages/shell-bson-parser/src/index.spec.ts b/packages/shell-bson-parser/src/index.spec.ts index 365f89bdf..0de9493b6 100644 --- a/packages/shell-bson-parser/src/index.spec.ts +++ b/packages/shell-bson-parser/src/index.spec.ts @@ -7,7 +7,7 @@ import { fileURLToPath } from 'url'; import * as WebWorkerModule from 'web-worker'; import * as api from './index.js'; -import { terminateWorker } from './worker-client.js'; +import { terminateWorker, callWorker } from './worker-client.js'; import { restrictGlobalScope, restrictObjectPrototype, @@ -206,4 +206,39 @@ describe('shell-bson-parser with webworker processing', function () { // It should not modify the default object proto expect(Object.prototype).to.have.property('__proto__'); }); + + describe('execution timeout', function () { + const initialWorkerScriptUrl = process.env.TEST_WORKER_SCRIPT_URL; + + beforeEach(function () { + terminateWorker(); + process.env.TEST_WORKER_SCRIPT_URL = '../test/fixtures/slow-worker.mjs'; + }); + + afterEach(function () { + terminateWorker(); + if (initialWorkerScriptUrl) { + process.env.TEST_WORKER_SCRIPT_URL = initialWorkerScriptUrl; + } else { + delete process.env.TEST_WORKER_SCRIPT_URL + } + }); + + it('rejects a request whose worker thread is wedged past the timeout', async function () { + try { + await callWorker([1000]); + expect.fail('Expected callWorker to throw an error due to timeout'); + } catch (err) { + expect((err as Error)?.message).to.equal( + 'Worker execution timed out after 500ms', + ); + } + }); + + it('spins up a fresh worker for the next call after a timeout kill', async function () { + await callWorker([1000]).catch(() => {}); // timeouts out + const result = await callWorker([0]); + expect(result).to.equal('done'); + }); + }); }); diff --git a/packages/shell-bson-parser/src/worker-client.ts b/packages/shell-bson-parser/src/worker-client.ts index 8a6ba2a57..f88b54a77 100644 --- a/packages/shell-bson-parser/src/worker-client.ts +++ b/packages/shell-bson-parser/src/worker-client.ts @@ -4,12 +4,26 @@ const WebWorker = (WebWorkerModule as unknown as { default: typeof Worker }) import { trackBSON, untrackBSON } from './structured-clone-bson.js'; import type { WorkerResponse } from './worker-types.js'; +/** Default execution timeout for worker requests */ +const DEFAULT_EXECUTION_TIMEOUT_MS = 120_000; + +function getExecutionTimeoutMs(): number { + return process.env.TEST_EXECUTION_TIMEOUT_MS + ? Number(process.env.TEST_EXECUTION_TIMEOUT_MS) + : DEFAULT_EXECUTION_TIMEOUT_MS; +} + + let worker: Worker | null = null; let blobUrl: string | null = null; let nextId = 0; const pending = new Map< number, - { resolve: (v: any) => void; reject: (e: Error) => void } + { + resolve: (v: any) => void; + reject: (e: Error) => void; + executionTimer: ReturnType; + } >(); const isNodeEnv = @@ -56,6 +70,7 @@ async function createWorker(): Promise { if (!entry) { return; } + clearTimeout(entry.executionTimer); pending.delete(response.id); if (!response.ok) { entry.reject(response.error); @@ -86,8 +101,16 @@ async function createWorker(): Promise { export async function callWorker(args: unknown[]): Promise { const activeWorker = await createWorker(); const id = nextId++; + const executionTimeoutMs = getExecutionTimeoutMs(); const promise = new Promise((resolve, reject) => { - pending.set(id, { resolve, reject }); + const executionTimer = setTimeout(() => { + // Terminate the worker is this message is taking too long to execute, + // this means all the other pending requests will also be terminated. + terminateWorker( + new Error(`Worker execution timed out after ${executionTimeoutMs}ms`), + ); + }, executionTimeoutMs); + pending.set(id, { resolve, reject, executionTimer }); }); try { activeWorker.postMessage({ @@ -95,6 +118,8 @@ export async function callWorker(args: unknown[]): Promise { args: trackBSON(args), }); } catch (err) { + const entry = pending.get(id); + if (entry) clearTimeout(entry.executionTimer); pending.get(id)?.reject(err as Error); pending.delete(id); } @@ -111,6 +136,7 @@ export function terminateWorker( blobUrl = null; for (const [id, entry] of pending) { + clearTimeout(entry.executionTimer); entry.reject(reason); pending.delete(id); } diff --git a/packages/shell-bson-parser/test/fixtures/slow-worker.mjs b/packages/shell-bson-parser/test/fixtures/slow-worker.mjs new file mode 100644 index 000000000..c955c1f9c --- /dev/null +++ b/packages/shell-bson-parser/test/fixtures/slow-worker.mjs @@ -0,0 +1,9 @@ +self.onmessage = (event) => { + const { id, args } = event.data; + const [delayMs] = args; + const start = Date.now(); + while (Date.now() - start < delayMs) { + // noop + } + self.postMessage({ id, ok: true, result: 'done' }); +}; diff --git a/packages/shell-bson-parser/tsconfig-lint.json b/packages/shell-bson-parser/tsconfig-lint.json index 6bdef84f3..5b09165f8 100644 --- a/packages/shell-bson-parser/tsconfig-lint.json +++ b/packages/shell-bson-parser/tsconfig-lint.json @@ -1,5 +1,5 @@ { "extends": "./tsconfig.json", "include": ["**/*"], - "exclude": ["node_modules", "dist"] + "exclude": ["node_modules", "dist", "test/fixtures"] } From 219147c6d280fdc2df080154c2fdca82e85e840b Mon Sep 17 00:00:00 2001 From: Basit Chonka Date: Wed, 23 Sep 2026 13:01:02 +0200 Subject: [PATCH 2/6] fix tests --- .../test/fixtures/slow-worker.mjs | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/packages/shell-bson-parser/test/fixtures/slow-worker.mjs b/packages/shell-bson-parser/test/fixtures/slow-worker.mjs index c955c1f9c..fc78c4569 100644 --- a/packages/shell-bson-parser/test/fixtures/slow-worker.mjs +++ b/packages/shell-bson-parser/test/fixtures/slow-worker.mjs @@ -1,9 +1,21 @@ +// Avoiding the import of these two functions, +// adding them here to keep this file self-contained +function unmarkBSON(value) { + return value.data; +} +function markBSON(value) { + return { data: value, bsonTypes: new Map() }; +} self.onmessage = (event) => { const { id, args } = event.data; - const [delayMs] = args; + const [delayMs] = unmarkBSON(args); const start = Date.now(); while (Date.now() - start < delayMs) { // noop } - self.postMessage({ id, ok: true, result: 'done' }); + self.postMessage({ + id, + ok: true, + result: markBSON('done'), + }); }; From f3d8b28e19929fc176819fc2a75227700fb2652a Mon Sep 17 00:00:00 2001 From: Basit Chonka Date: Wed, 23 Sep 2026 13:18:57 +0200 Subject: [PATCH 3/6] accept optional timeout --- packages/shell-bson-parser/src/index.ts | 13 +++++++++-- .../shell-bson-parser/src/worker-client.ts | 22 ++++++++++++++----- 2 files changed, 27 insertions(+), 8 deletions(-) diff --git a/packages/shell-bson-parser/src/index.ts b/packages/shell-bson-parser/src/index.ts index f639f256b..3d8e4214e 100644 --- a/packages/shell-bson-parser/src/index.ts +++ b/packages/shell-bson-parser/src/index.ts @@ -1,12 +1,21 @@ import type { parse as parseSync } from './parse.js'; import { toJSString } from './stringify.js'; import { ParseMode } from './options.js'; +import type { Options } from './options.js'; import { callWorker, terminateWorker } from './worker-client.js'; +import type { ExecutionOptions } from './worker-client.js'; export const parse = ( - ...args: Parameters -): Promise> => callWorker(args); + input: string, + { + executionTimeoutMs, + ...parseOptions + }: Partial = {}, +): Promise> => { + return callWorker([input, parseOptions], { executionTimeoutMs }); +}; export { ParseMode, toJSString, terminateWorker }; +export type { ExecutionOptions, Options as ParseOptions }; export default parse; diff --git a/packages/shell-bson-parser/src/worker-client.ts b/packages/shell-bson-parser/src/worker-client.ts index f88b54a77..32fad583e 100644 --- a/packages/shell-bson-parser/src/worker-client.ts +++ b/packages/shell-bson-parser/src/worker-client.ts @@ -7,12 +7,17 @@ import type { WorkerResponse } from './worker-types.js'; /** Default execution timeout for worker requests */ const DEFAULT_EXECUTION_TIMEOUT_MS = 120_000; -function getExecutionTimeoutMs(): number { - return process.env.TEST_EXECUTION_TIMEOUT_MS - ? Number(process.env.TEST_EXECUTION_TIMEOUT_MS) - : DEFAULT_EXECUTION_TIMEOUT_MS; +function getExecutionTimeoutMs(initialExecutionMs?: number): number { + if (process.env.TEST_EXECUTION_TIMEOUT_MS) { + return Number(process.env.TEST_EXECUTION_TIMEOUT_MS); + } + return initialExecutionMs ?? DEFAULT_EXECUTION_TIMEOUT_MS; } +export type ExecutionOptions = { + /** Defaults to `120_000` (2 minutes). */ + executionTimeoutMs?: number; +}; let worker: Worker | null = null; let blobUrl: string | null = null; @@ -98,10 +103,15 @@ async function createWorker(): Promise { return worker; } -export async function callWorker(args: unknown[]): Promise { +export async function callWorker( + args: unknown[], + executionOptions?: ExecutionOptions, +): Promise { const activeWorker = await createWorker(); const id = nextId++; - const executionTimeoutMs = getExecutionTimeoutMs(); + const executionTimeoutMs = getExecutionTimeoutMs( + executionOptions?.executionTimeoutMs, + ); const promise = new Promise((resolve, reject) => { const executionTimer = setTimeout(() => { // Terminate the worker is this message is taking too long to execute, From 0ec2753f0959d58d3423f70ff841719947e3bab2 Mon Sep 17 00:00:00 2001 From: Basit Chonka Date: Wed, 23 Sep 2026 16:47:00 +0200 Subject: [PATCH 4/6] clean up --- packages/shell-bson-parser/src/index.spec.ts | 4 ++-- packages/shell-bson-parser/src/worker-client.ts | 12 ++---------- 2 files changed, 4 insertions(+), 12 deletions(-) diff --git a/packages/shell-bson-parser/src/index.spec.ts b/packages/shell-bson-parser/src/index.spec.ts index 0de9493b6..d7ed218cc 100644 --- a/packages/shell-bson-parser/src/index.spec.ts +++ b/packages/shell-bson-parser/src/index.spec.ts @@ -220,13 +220,13 @@ describe('shell-bson-parser with webworker processing', function () { if (initialWorkerScriptUrl) { process.env.TEST_WORKER_SCRIPT_URL = initialWorkerScriptUrl; } else { - delete process.env.TEST_WORKER_SCRIPT_URL + delete process.env.TEST_WORKER_SCRIPT_URL; } }); it('rejects a request whose worker thread is wedged past the timeout', async function () { try { - await callWorker([1000]); + await callWorker([1000], { executionTimeoutMs: 500 }); expect.fail('Expected callWorker to throw an error due to timeout'); } catch (err) { expect((err as Error)?.message).to.equal( diff --git a/packages/shell-bson-parser/src/worker-client.ts b/packages/shell-bson-parser/src/worker-client.ts index 32fad583e..f79659e88 100644 --- a/packages/shell-bson-parser/src/worker-client.ts +++ b/packages/shell-bson-parser/src/worker-client.ts @@ -7,13 +7,6 @@ import type { WorkerResponse } from './worker-types.js'; /** Default execution timeout for worker requests */ const DEFAULT_EXECUTION_TIMEOUT_MS = 120_000; -function getExecutionTimeoutMs(initialExecutionMs?: number): number { - if (process.env.TEST_EXECUTION_TIMEOUT_MS) { - return Number(process.env.TEST_EXECUTION_TIMEOUT_MS); - } - return initialExecutionMs ?? DEFAULT_EXECUTION_TIMEOUT_MS; -} - export type ExecutionOptions = { /** Defaults to `120_000` (2 minutes). */ executionTimeoutMs?: number; @@ -109,9 +102,8 @@ export async function callWorker( ): Promise { const activeWorker = await createWorker(); const id = nextId++; - const executionTimeoutMs = getExecutionTimeoutMs( - executionOptions?.executionTimeoutMs, - ); + const executionTimeoutMs = + executionOptions?.executionTimeoutMs ?? DEFAULT_EXECUTION_TIMEOUT_MS; const promise = new Promise((resolve, reject) => { const executionTimer = setTimeout(() => { // Terminate the worker is this message is taking too long to execute, From 0368c0adfd0fea8864511cf8414e6955d8ae0e42 Mon Sep 17 00:00:00 2001 From: Basit Chonka Date: Mon, 28 Sep 2026 14:41:44 +0200 Subject: [PATCH 5/6] fix typos --- packages/shell-bson-parser/src/index.spec.ts | 2 +- packages/shell-bson-parser/src/worker-client.ts | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/shell-bson-parser/src/index.spec.ts b/packages/shell-bson-parser/src/index.spec.ts index d7ed218cc..545e33726 100644 --- a/packages/shell-bson-parser/src/index.spec.ts +++ b/packages/shell-bson-parser/src/index.spec.ts @@ -236,7 +236,7 @@ describe('shell-bson-parser with webworker processing', function () { }); it('spins up a fresh worker for the next call after a timeout kill', async function () { - await callWorker([1000]).catch(() => {}); // timeouts out + await callWorker([1000]).catch(() => {}); // times out const result = await callWorker([0]); expect(result).to.equal('done'); }); diff --git a/packages/shell-bson-parser/src/worker-client.ts b/packages/shell-bson-parser/src/worker-client.ts index f79659e88..ae2726086 100644 --- a/packages/shell-bson-parser/src/worker-client.ts +++ b/packages/shell-bson-parser/src/worker-client.ts @@ -106,7 +106,7 @@ export async function callWorker( executionOptions?.executionTimeoutMs ?? DEFAULT_EXECUTION_TIMEOUT_MS; const promise = new Promise((resolve, reject) => { const executionTimer = setTimeout(() => { - // Terminate the worker is this message is taking too long to execute, + // Terminate the worker if this message is taking too long to execute, // this means all the other pending requests will also be terminated. terminateWorker( new Error(`Worker execution timed out after ${executionTimeoutMs}ms`), From a2f74b5e2ff71379a7b277235c4149c8dfbd9fc8 Mon Sep 17 00:00:00 2001 From: Basit Chonka Date: Mon, 28 Sep 2026 15:01:17 +0200 Subject: [PATCH 6/6] handle worker init --- .../shell-bson-parser/src/worker-client.ts | 85 +++++++++++-------- 1 file changed, 49 insertions(+), 36 deletions(-) diff --git a/packages/shell-bson-parser/src/worker-client.ts b/packages/shell-bson-parser/src/worker-client.ts index ae2726086..b6eee1ebc 100644 --- a/packages/shell-bson-parser/src/worker-client.ts +++ b/packages/shell-bson-parser/src/worker-client.ts @@ -13,6 +13,7 @@ export type ExecutionOptions = { }; let worker: Worker | null = null; +let workerPromise: Promise | null = null; let blobUrl: string | null = null; let nextId = 0; const pending = new Map< @@ -54,46 +55,57 @@ async function getWorkerScriptUrl(): Promise { return blobUrl; } -async function createWorker(): Promise { +function createWorker(): Promise { if (worker) { - return worker; + return Promise.resolve(worker); + } + if (workerPromise) { + return workerPromise; } - const scriptUrl = await getWorkerScriptUrl(); - worker = new WebWorker(scriptUrl, { type: 'module' }); + workerPromise = (async () => { + const scriptUrl = await getWorkerScriptUrl(); + const newWorker = new WebWorker(scriptUrl, { type: 'module' }); + const onMessageHandler = (event: MessageEvent) => { + const response = event.data; + const entry = pending.get(response.id); + if (!entry) { + return; + } + clearTimeout(entry.executionTimer); + pending.delete(response.id); + if (!response.ok) { + entry.reject(response.error); + return; + } + try { + entry.resolve(untrackBSON(response.result)); + } catch (err) { + entry.reject(err as Error); + } + }; + + const onErrorHandler = (event: ErrorEvent) => { + terminateWorker(new Error(event.message || 'Worker error')); + }; + + const onMessageErrorHandler = () => { + terminateWorker(new Error('Worker message could not be deserialized')); + }; + + newWorker.addEventListener('message', onMessageHandler); + newWorker.addEventListener('error', onErrorHandler); + newWorker.addEventListener('messageerror', onMessageErrorHandler); + + worker = newWorker; + return newWorker; + })(); + + workerPromise.catch(() => { + workerPromise = null; + }); - const onMessageHandler = (event: MessageEvent) => { - const response = event.data; - const entry = pending.get(response.id); - if (!entry) { - return; - } - clearTimeout(entry.executionTimer); - pending.delete(response.id); - if (!response.ok) { - entry.reject(response.error); - return; - } - try { - entry.resolve(untrackBSON(response.result)); - } catch (err) { - entry.reject(err as Error); - } - }; - - const onErrorHandler = (event: ErrorEvent) => { - terminateWorker(new Error(event.message || 'Worker error')); - }; - - const onMessageErrorHandler = () => { - terminateWorker(new Error('Worker message could not be deserialized')); - }; - - worker.addEventListener('message', onMessageHandler); - worker.addEventListener('error', onErrorHandler); - worker.addEventListener('messageerror', onMessageErrorHandler); - - return worker; + return workerPromise; } export async function callWorker( @@ -135,6 +147,7 @@ export function terminateWorker( if (blobUrl) URL.revokeObjectURL(blobUrl); worker = null; + workerPromise = null; blobUrl = null; for (const [id, entry] of pending) {