From 6946b8dc6397446ed6abd4a70754791611ac5871 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Mon, 20 Jul 2026 12:41:15 -0700 Subject: [PATCH] feat: validate runtime balancer state is thread-safe at boot --- src/runner.ts | 2 ++ src/runtime.ts | 19 +++++++++++++++++++ src/shared/reconstruct.ts | 7 ++++++- src/worker-pool.ts | 3 ++- test/fixtures/runtime-bad-balancer-def.ts | 12 ++++++++++++ test/fixtures/runtime-bad-balancer.ts | 11 +++++++++++ test/runtime-lifecycle.test.ts | 7 +++++++ 7 files changed, 59 insertions(+), 2 deletions(-) create mode 100644 test/fixtures/runtime-bad-balancer-def.ts create mode 100644 test/fixtures/runtime-bad-balancer.ts diff --git a/src/runner.ts b/src/runner.ts index 3dc600f..4b23d61 100644 --- a/src/runner.ts +++ b/src/runner.ts @@ -53,6 +53,8 @@ export interface WorkerOptions { * used by the global runtime. Default 'held': workers keep the event loop alive * until disposal. */ refMode?: 'held' | 'lazy'; + /** @internal Pre-computed balancer state (already validated by the runtime). */ + balancerState?: unknown; } /** A handle to a specific worker in a pool. Uniform across threads. */ diff --git a/src/runtime.ts b/src/runtime.ts index 9f92d9e..0710e9b 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -1,5 +1,6 @@ import { isMainThread } from 'node:worker_threads'; import { workers } from './worker-pool.ts'; +import { isSharedValue } from './shared/reconstruct.ts'; import { getDefineId } from './define.ts'; import { defaultRuntime } from './runtime-definition.ts'; import type { RuntimeDefinition } from './runtime-definition.ts'; @@ -55,22 +56,40 @@ export function getRuntime(): Runner { ); } const def = registered ?? defaultRuntime; + // Balancer state crosses the future worker handshake — validate at boot + // that it can survive that trip intact (crash early, not corrupt later). + const state = def.balance?.initialState?.(); + if (state !== undefined && !isSharedValue(state) && !isPlainSerializable(state)) { + throw new Error( + 'Global runtime balancer state must be a single shared value (e.g. uint32atomic() ' + + 'or shared({...})) or a plain serializable value — got a non-shared object, which ' + + 'cannot stay consistent across scheduling threads.', + ); + } pool = def.size !== undefined ? workers(def.size, { balance: def.balance, shutdownTimeout: def.shutdownTimeout, refMode: 'lazy', + balancerState: state, }) : workers({ balance: def.balance, shutdownTimeout: def.shutdownTimeout, refMode: 'lazy', + balancerState: state, }); } return pool; } +function isPlainSerializable(value: unknown): boolean { + // Primitives serialize fine; objects/functions (other than shared values, + // checked by the caller) do not survive the worker handshake intact. + return value === null || (typeof value !== 'object' && typeof value !== 'function'); +} + /** Returns the booted pool or throws — for surface that must not boot. */ function bootedRuntime(): Runner { if (shutdownStarted) { diff --git a/src/shared/reconstruct.ts b/src/shared/reconstruct.ts index d482741..5bf64a0 100644 --- a/src/shared/reconstruct.ts +++ b/src/shared/reconstruct.ts @@ -11,8 +11,13 @@ export function registerSync(tag: string, ctor: new (buffer: SharedArrayBuffer, registry.set(tag, ctor); } +/** True when `value` is a shared-branded value (e.g. `uint32atomic()`, `shared({...})`). */ +export function isSharedValue(value: unknown): boolean { + return typeof value === 'object' && value !== null && SHARED in value; +} + export function serializeArg(arg: unknown): unknown { - if (typeof arg === 'object' && arg !== null && SHARED in arg) { + if (isSharedValue(arg)) { const data = (arg as any)[SHARED](); if (data.tag === 'SharedStruct') { const serializedFields: Record = {}; diff --git a/src/worker-pool.ts b/src/worker-pool.ts index ce54599..280ca59 100644 --- a/src/worker-pool.ts +++ b/src/worker-pool.ts @@ -32,7 +32,8 @@ export function workers(sizeOrOpts?: number | WorkerOptions, opts?: WorkerOption } const balancer: Balancer = opts?.balance ?? leastBusy(); - const balancerState: unknown = balancer.initialState?.(); + const balancerState: unknown = + opts != null && 'balancerState' in opts ? opts.balancerState : balancer.initialState?.(); const pool: Worker[] = []; for (let i = 0; i < size; i++) { diff --git a/test/fixtures/runtime-bad-balancer-def.ts b/test/fixtures/runtime-bad-balancer-def.ts new file mode 100644 index 0000000..7abc99e --- /dev/null +++ b/test/fixtures/runtime-bad-balancer-def.ts @@ -0,0 +1,12 @@ +import { define } from 'moroutine'; +import type { Balancer, RuntimeDefinition } from 'moroutine'; + +// A balancer whose state is a plain object wrapping nothing shared — legal in +// a workers() pool, but the global runtime must reject it at boot because the +// state cannot cross the future worker handshake intact. +const badBalancer: Balancer<{ n: number }> = { + initialState: () => ({ n: 0 }), + select: (w, _t, _s) => w[0], +}; + +export const badDef: RuntimeDefinition = define(import.meta, { size: 1, balance: badBalancer }); diff --git a/test/fixtures/runtime-bad-balancer.ts b/test/fixtures/runtime-bad-balancer.ts new file mode 100644 index 0000000..f0d8d53 --- /dev/null +++ b/test/fixtures/runtime-bad-balancer.ts @@ -0,0 +1,11 @@ +import { runtime, registerRuntime } from 'moroutine'; +import { badDef } from './runtime-bad-balancer-def.ts'; +import { double } from './math.ts'; + +registerRuntime(badDef); +try { + await runtime.run(double(1)); + console.log('BOOT no-throw'); +} catch (err) { + console.log('BOOT threw: ' + (err as Error).message); +} diff --git a/test/runtime-lifecycle.test.ts b/test/runtime-lifecycle.test.ts index 6f765bf..dd87854 100644 --- a/test/runtime-lifecycle.test.ts +++ b/test/runtime-lifecycle.test.ts @@ -35,6 +35,13 @@ describe('global runtime lifecycle', () => { assert.ok(stdout.includes('RUNTIME-BOOTED true')); }); + it('boot rejects balancer state that cannot cross threads', async () => { + const { stdout } = await exec(process.execPath, ['--no-warnings', join(fixturesDir, 'runtime-bad-balancer.ts')], { + timeout: 15000, + }); + assert.match(stdout, /BOOT threw: .*balancer state/i); + }); + it('shutdown() fires signal, drains, and locks out further dispatch', async () => { const { stdout } = await exec(process.execPath, ['--no-warnings', join(fixturesDir, 'runtime-shutdown.ts')], { timeout: 15000, -- 2.51.2