From 53c296bbe7ce661e6523cfd92c743c6a908bba46 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Tue, 21 Jul 2026 17:17:15 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20worker-side=20runtime=20view=20?= =?UTF-8?q?=E2=80=94=20nested=20dispatch,=20bare=20await,=20loopback=20sel?= =?UTF-8?q?f-dispatch?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/peers.ts | 12 +- src/runtime-worker.ts | 206 +++++++++++++++++++++++++++ src/runtime.ts | 22 ++- src/worker-entry.ts | 11 +- test/fixtures/runtime-nested-main.ts | 7 + test/fixtures/runtime-nested.ts | 20 +++ test/fixtures/runtime-pinned-main.ts | 8 ++ test/handshake.test.ts | 9 +- test/runtime-lifecycle.test.ts | 14 ++ 9 files changed, 299 insertions(+), 10 deletions(-) create mode 100644 src/runtime-worker.ts create mode 100644 test/fixtures/runtime-nested-main.ts create mode 100644 test/fixtures/runtime-nested.ts create mode 100644 test/fixtures/runtime-pinned-main.ts diff --git a/src/peers.ts b/src/peers.ts index f0c43cb..4e22893 100644 --- a/src/peers.ts +++ b/src/peers.ts @@ -1,6 +1,6 @@ import { parentPort } from 'node:worker_threads'; import type { MessagePort } from 'node:worker_threads'; -import { attachEndpointHandler } from './worker-entry.ts'; +import type { Endpoint } from './execute.ts'; /** * Worker-side peer table: one full-duplex MessagePort per sibling worker, @@ -15,17 +15,23 @@ const waiting = new Map void>(); let poolSize = 0; let selfIndex = -1; +// worker-entry injects its combined handler here — a static import of +// worker-entry would put it in the `moroutine` import graph (via +// runtime-worker), deadlocking worker-entry's top-level await when a worker +// dynamically imports a module that imports `moroutine`. +let attachHandler: ((endpoint: Endpoint) => void) | null = null; /** Called once by worker-entry from the handshake. */ -export function initPeers(index: number, size: number): void { +export function initPeers(index: number, size: number, attach: (endpoint: Endpoint) => void): void { selfIndex = index; poolSize = size; + attachHandler = attach; } function adoptPeerPort(peer: number, port: MessagePort): void { // Same-tick wiring — never let a peer port sit unlistened: 'message' events // emitted without a listener are lost (only undelivered messages buffer). - attachEndpointHandler(port); + attachHandler!(port); port.unref(); peers.set(peer, port); const waiter = waiting.get(peer); diff --git a/src/runtime-worker.ts b/src/runtime-worker.ts new file mode 100644 index 0000000..810e70d --- /dev/null +++ b/src/runtime-worker.ts @@ -0,0 +1,206 @@ +import { MessageChannel } from 'node:worker_threads'; +import type { MessagePort } from 'node:worker_threads'; +import { execute, dispatchStream } from './execute.ts'; +import type { Endpoint } from './execute.ts'; +import { getPeerPort } from './peers.ts'; +import { ActiveCounts } from './active-counts.ts'; +import { deserializeArg } from './shared/reconstruct.ts'; +import { defineRegistry } from './define.ts'; +import { leastBusy } from './balancers.ts'; +import { AsyncIterableTask } from './stream-task.ts'; +import type { Int32Atomic } from './shared/int32-atomic.ts'; +import type { WorkerHandshake } from './handshake.ts'; +import type { RuntimeDefinition } from './runtime-definition.ts'; +import type { Task, Balancer, Runner, WorkerHandle } from './runner.ts'; +import type { ChannelOptions } from './channel.ts'; + +/** + * Worker-side view of the global runtime, bound from the boot handshake. + * Dispatch resolves a pin or balancer pick to an endpoint: a factory-created + * peer port, or a loopback channel for self-dispatch (identical protocol, so + * the "shortcut" is unobservable by construction). Shutdown observation is + * the transferred AbortSignal from the handshake — no frames. + */ + +let view: Runner | null = null; +let shutdownSignal: AbortSignal | null = null; + +// worker-entry injects its combined handler at init to avoid a static cycle +// (worker-entry imports runtime-worker; the handler closes over worker-entry's +// task machinery). NOTE: require() is unavailable in ESM — injection is the +// cycle-free wiring. +let attachHandler: ((endpoint: Endpoint) => void) | null = null; + +/** Loopback: a port pair whose far side runs THIS worker's own task handler. + * Both ends get the same combined handler wiring a peer port would. */ +let selfPort: MessagePort | null = null; +function getSelfPort(): MessagePort { + if (selfPort === null) { + const { port1, port2 } = new MessageChannel(); + // With the combined shape-router on BOTH ends, port1 receives only + // responses (routed to handleResponse) and port2 receives only requests + // (routed to the task handlers) — the full-duplex demux makes the + // loopback symmetric for free. + attachHandler!(port2); // far side: executes our dispatches + attachHandler!(port1); // near side: routes responses back to our pending map + port1.unref(); + port2.unref(); + selfPort = port1; + } + return selfPort; +} + +/** Initializes the worker-side runtime view. Called once by worker-entry when + * the handshake carries runtime fields. */ +export async function initRuntimeWorker( + handshake: WorkerHandshake, + attach: (endpoint: Endpoint) => void, +): Promise { + attachHandler = attach; + shutdownSignal = handshake.shutdownSignal!; // transferred copy of main's runtime signal + const counts = new ActiveCounts(handshake.countsBuffer); + const selfIndex = handshake.index; + const size = handshake.size; + + // Reconstruct the definition by define() identity (module evaluation), and + // the balancer state from the handshake (NEVER call initialState() here — + // main allocated the single instance). + const defId = handshake.runtimeId!; + const url = defId.slice(0, defId.lastIndexOf('#')); + await import(url); + const def = defineRegistry.get(defId) as RuntimeDefinition | undefined; + if (def === undefined) throw new Error(`Runtime definition not found: ${defId}`); + const balancer: Balancer = def.balance ?? leastBusy(); + const balancerState: unknown = + handshake.balancerState !== undefined ? deserializeArg(handshake.balancerState) : undefined; + const busyPeers = deserializeArg(handshake.busyPeers) as Int32Atomic; + + async function endpointFor(index: number): Promise { + if (index === selfIndex) return getSelfPort(); + return await getPeerPort(index); + } + + // Edge-counted liveness: the shared busyPeers counter tracks WORKERS with + // outstanding outgoing dispatches, not dispatches. localOutgoing is plain — + // this thread only. CAUSALITY INVARIANT: the add(1) below is synchronous + // inside dispatch, before this worker's covering task can post its result + // to main — so main can never observe busyPeers 0 while our work exists. + let localOutgoing = 0; + function track(): void { + if (localOutgoing === 0) busyPeers.add(1); + localOutgoing++; + } + function untrack(): void { + localOutgoing--; + if (localOutgoing === 0) { + // sub() returns the previous value: 1 means we were the last busy + // worker — wake main's drain wait. + if (busyPeers.sub(1) === 1) busyPeers.notify(); + } + } + + const handles: WorkerHandle[] = []; + for (let i = 0; i < size; i++) { + const idx = i; + handles.push({ + exec(task: Task, channelOpts?: ChannelOptions): any { + return dispatchTo(idx, task, channelOpts); + }, + get index() { + return idx; + }, + get activeCount() { + return counts.get(idx); + }, + }); + } + const workerHandles: readonly WorkerHandle[] = Object.freeze(handles); + + function resolveIndex(task: Task): number { + if (task.worker != null) { + const idx = workerHandles.indexOf(task.worker); + if (idx === -1) throw new Error('Task is pinned to a worker from another pool'); + return idx; + } + const handle = balancer.select(workerHandles, task, balancerState); + const idx = workerHandles.indexOf(handle); + if (idx === -1) throw new Error('Balancer returned a handle that is not in this pool'); + return idx; + } + + function dispatchTo(idx: number, task: Task, channelOpts?: ChannelOptions): any { + // The transferred signal IS main's runtime.signal — post-shutdown dispatch + // throws here without any propagation machinery. + if (shutdownSignal!.aborted) { + if (task instanceof AsyncIterableTask) throw new Error('The global runtime has been shut down.'); + return Promise.reject(new Error('The global runtime has been shut down.')); + } + if (task instanceof AsyncIterableTask) { + // Streaming: resolve endpoint, dispatch, count + edge-track for the + // stream's life. done.then and the catch below can both fire on error + // paths — the settled guard makes counts/track settle exactly once. + counts.inc(idx); + track(); + let settled = false; + const settle = (): void => { + if (settled) return; + settled = true; + counts.dec(idx); + untrack(); + }; + const iterable = (async function* () { + try { + const endpoint = await endpointFor(idx); + const { iterable: inner, done } = dispatchStream(endpoint, task.id, task.args, channelOpts); + void done.then(settle); + yield* inner as AsyncIterable; + } catch (err) { + settle(); + throw err; + } + })(); + return iterable; + } + counts.inc(idx); + track(); + return (async () => { + try { + const endpoint = await endpointFor(idx); + return await execute(endpoint, task.id, task.args); + } finally { + counts.dec(idx); + untrack(); + } + })(); + } + + view = Object.assign( + (taskOrTasks: Task | Task[], channelOpts?: ChannelOptions): any => { + if (Array.isArray(taskOrTasks)) { + return Promise.all(taskOrTasks.map((t) => dispatchTo(resolveIndex(t), t))); + } + return dispatchTo(resolveIndex(taskOrTasks), taskOrTasks, channelOpts); + }, + { + get signal() { + // The handshake signal is a transferred copy of main's runtime.signal: + // same abort event, same reason. No worker-side controller exists. + return shutdownSignal!; + }, + get workers() { + return workerHandles; + }, + [Symbol.dispose]() { + throw new Error('The global runtime cannot be disposed from a worker thread.'); + }, + async [Symbol.asyncDispose]() { + throw new Error('The global runtime cannot be disposed from a worker thread.'); + }, + }, + ) as Runner; +} + +/** The bound view, or null on non-runtime threads. @internal */ +export function getWorkerRuntime(): Runner | null { + return view; +} diff --git a/src/runtime.ts b/src/runtime.ts index 2bd4324..27f75ee 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 { getWorkerRuntime } from './runtime-worker.ts'; import { isSharedValue, serializeArg } from './shared/reconstruct.ts'; import { getDefineId } from './define.ts'; import { defaultRuntime } from './runtime-definition.ts'; @@ -50,10 +51,15 @@ export function getRuntime(): Runner { } if (pool === null) { if (!isMainThread) { - throw new Error( - 'The global runtime cannot boot on a worker thread yet — dispatch from the main thread. ' + - '(Worker-side runtime participation is planned.)', - ); + const workerView = getWorkerRuntime(); + if (workerView === null) { + throw new Error( + 'This worker thread is not part of the global runtime — bare dispatch and ' + + 'runtime.run() require runtime membership (dispatch from the main thread ' + + 'or from a runtime worker).', + ); + } + return workerView; } const def = registered ?? defaultRuntime; // Balancer state crosses the future worker handshake — validate at boot @@ -96,6 +102,11 @@ function isPlainSerializable(value: unknown): boolean { /** Returns the booted pool or throws — for surface that must not boot. */ function bootedRuntime(): Runner { + if (!isMainThread) { + const workerView = getWorkerRuntime(); + if (workerView !== null) return workerView; + throw new Error('This worker thread is not part of the global runtime.'); + } if (shutdownStarted) { throw new Error('The global runtime has been shut down.'); } @@ -139,6 +150,9 @@ export const runtime = { * re-booted (one runtime per process lifetime). */ async shutdown(): Promise { + if (!isMainThread) { + throw new Error('runtime.shutdown() must be called from the main thread.'); + } shutdownStarted = true; if (pool === null) return; // never booted: nothing to drain, lock the door shutdownPromise ??= pool[Symbol.asyncDispose](); diff --git a/src/worker-entry.ts b/src/worker-entry.ts index 2627def..4309f35 100644 --- a/src/worker-entry.ts +++ b/src/worker-entry.ts @@ -11,6 +11,7 @@ import type { CtrlFrame, WorkerHandshake } from './handshake.ts'; import { receivePeer, initPeers } from './peers.ts'; import { handleResponse } from './execute.ts'; import type { Endpoint } from './execute.ts'; +import { initRuntimeWorker } from './runtime-worker.ts'; const imported = new Set(); const taskCache = new Map(); @@ -307,10 +308,18 @@ export function attachEndpointHandler(endpoint: Endpoint): void { // channels are pool-agnostic even though only runtime workers use them via // the runtime view. const handshake = workerData as WorkerHandshake | null; -if (handshake != null) initPeers(handshake.index, handshake.size); +if (handshake != null) initPeers(handshake.index, handshake.size, attachEndpointHandler); // parentPort is null when this module evaluates on the MAIN thread — user // code importing peers.ts (e.g. a moroutine that dispatches to a peer) pulls // this module in transitively. There is no endpoint to attach there; the // handshake above is null on main too, so evaluation is inert by design. if (parentPort !== null) attachEndpointHandler(parentPort); + +// Runtime worker: bind the ambient runtime view before any task can run. +// Top-level await — worker-entry is ESM; tasks queue behind module eval. +// Main-thread evaluation never reaches the await: parentPort is null there +// and the handshake carries no runtimeId for plain pools. +if (parentPort !== null && handshake?.runtimeId !== undefined) { + await initRuntimeWorker(handshake, attachEndpointHandler); +} diff --git a/test/fixtures/runtime-nested-main.ts b/test/fixtures/runtime-nested-main.ts new file mode 100644 index 0000000..8cee286 --- /dev/null +++ b/test/fixtures/runtime-nested-main.ts @@ -0,0 +1,7 @@ +import { runtime, registerRuntime } from 'moroutine'; +import { tinyRuntime } from './runtime-def.ts'; +import { fanOut } from './runtime-nested.ts'; + +registerRuntime(tinyRuntime); // size 1: nested fan-out MUST work at pool size 1 (liveness invariant) +const r = await runtime.run(fanOut(10)); +console.log('NESTED ' + JSON.stringify(r)); diff --git a/test/fixtures/runtime-nested.ts b/test/fixtures/runtime-nested.ts new file mode 100644 index 0000000..a431ae0 --- /dev/null +++ b/test/fixtures/runtime-nested.ts @@ -0,0 +1,20 @@ +import { threadId } from 'node:worker_threads'; +import { mo, runtime, assign } from 'moroutine'; +import { double } from './math.ts'; + +// Nested dispatch: runs on one runtime worker, fans out to the pool. +export const fanOut = mo(import.meta, async (n: number): Promise => { + // bare await inside a runtime worker + const a = await double(n); + // explicit runtime.run with a batch + const [b, c] = await runtime.run([double(n + 1), double(n + 2)]); + return [a, b, c]; +}); + +// Pinned nested dispatch: prove worker-side assign() + same-worker loopback. +export const pinnedSelf = mo(import.meta, async (selfIndex: number, n: number): Promise<[number, number]> => { + const viaSelf = await runtime.run(assign(runtime.workers[selfIndex], double(n))); + return [viaSelf, threadId]; +}); + +export const whichThread = mo(import.meta, (): number => threadId); diff --git a/test/fixtures/runtime-pinned-main.ts b/test/fixtures/runtime-pinned-main.ts new file mode 100644 index 0000000..4112999 --- /dev/null +++ b/test/fixtures/runtime-pinned-main.ts @@ -0,0 +1,8 @@ +import { runtime, registerRuntime } from 'moroutine'; +import { tinyRuntime } from './runtime-def.ts'; +import { pinnedSelf, whichThread } from './runtime-nested.ts'; + +registerRuntime(tinyRuntime); +const workerThread = await runtime.run(whichThread()); +const [viaSelf, taskThread] = await runtime.run(pinnedSelf(0, 21)); +console.log('PINNED ' + viaSelf + ' SAME-THREAD ' + (workerThread === taskThread)); diff --git a/test/handshake.test.ts b/test/handshake.test.ts index a669498..32aacbf 100644 --- a/test/handshake.test.ts +++ b/test/handshake.test.ts @@ -1,16 +1,21 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; import { workers } from 'moroutine'; +import { getDefineId } from '../src/define.ts'; import { readHandshake } from './fixtures/handshake-probe.ts'; +import { tinyRuntime } from './fixtures/runtime-def.ts'; describe('worker handshake', () => { it('runtime handshake fields arrive in workerData', async () => { + // A real define()-branded id: workers dereference runtimeId at boot + // (module import), so a fake id would crash worker module evaluation. + const runtimeId = getDefineId(tinyRuntime)!; const run = workers(2, { - runtimeHandshake: { runtimeId: 'test://fake#0', balancerState: 42 }, + runtimeHandshake: { runtimeId, balancerState: 42 }, } as any); try { const h = await run(readHandshake()); - assert.equal(h.runtimeId, 'test://fake#0'); + assert.equal(h.runtimeId, runtimeId); assert.equal(h.balancerState, 42); assert.equal(h.size, 2); assert.equal(typeof h.index, 'number'); diff --git a/test/runtime-lifecycle.test.ts b/test/runtime-lifecycle.test.ts index dd87854..62ded7c 100644 --- a/test/runtime-lifecycle.test.ts +++ b/test/runtime-lifecycle.test.ts @@ -49,4 +49,18 @@ describe('global runtime lifecycle', () => { assert.ok(stdout.includes('SIGNAL true')); assert.ok(stdout.includes('POST-SHUTDOWN-DISPATCH threw')); }); + + it('nested dispatch from a runtime worker (fork-join at pool size 1 stays live)', async () => { + const { stdout } = await exec(process.execPath, ['--no-warnings', join(fixturesDir, 'runtime-nested-main.ts')], { + timeout: 20000, + }); + assert.ok(stdout.includes('NESTED [20,22,24]')); + }); + + it('worker-side pin to self dispatches via loopback on the same thread', async () => { + const { stdout } = await exec(process.execPath, ['--no-warnings', join(fixturesDir, 'runtime-pinned-main.ts')], { + timeout: 20000, + }); + assert.ok(stdout.includes('PINNED 42 SAME-THREAD true')); + }); }); -- 2.51.2