From 9a58d315f24b35a7fd6c9b1c7b5c647ca339123e Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Tue, 21 Jul 2026 15:53:47 -0700 Subject: [PATCH] feat: control-frame discriminator, crash on unknown frames --- src/execute.ts | 36 +++++++++++++++++++++-------- src/worker-entry.ts | 37 ++++++++++++++++++++++++++---- test/ctrl-frames.test.ts | 25 ++++++++++++++++++++ test/fixtures/ctrl-unknown-main.ts | 13 +++++++++++ test/fixtures/ctrl-unknown.ts | 12 ++++++++++ 5 files changed, 109 insertions(+), 14 deletions(-) create mode 100644 test/ctrl-frames.test.ts create mode 100644 test/fixtures/ctrl-unknown-main.ts create mode 100644 test/fixtures/ctrl-unknown.ts diff --git a/src/execute.ts b/src/execute.ts index cbd3354..bdee045 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -13,10 +13,15 @@ import { Channel } from './channel.ts'; import { pipeIterable, newPipeFlags, CANCEL, DEFAULT_HIGH_WATER, serializeStreamHandle } from './pipe.ts'; import type { StreamHandle, SerializedStreamHandle } from './pipe.ts'; import type { ChannelOptions } from './channel.ts'; +import { isCtrlFrame } from './handshake.ts'; +import type { CtrlFrame } from './handshake.ts'; let nextCallId = 0; const pending = new Map void; reject: (reason: any) => void }>(); const streamPortStack: MessagePort[][] = []; +// Keyed by endpoint (per Task 8: endpoint-death rejection). Reads are no-ops +// until that task populates it. +const endpointCalls = new WeakMap>(); /** Anything dispatch can talk over: a Worker from the main thread, or a * MessagePort between peers. Structurally satisfied by both. */ @@ -25,16 +30,29 @@ export interface Endpoint { on(event: 'message', listener: (value: any) => void): unknown; } -export function setupWorker(endpoint: Endpoint): void { - endpoint.on('message', (msg: { callId: number; value?: unknown; error?: Error }) => { - const call = pending.get(msg.callId); - if (!call) return; - pending.delete(msg.callId); - if (msg.error !== undefined) { - call.reject(msg.error); - } else { - call.resolve(deserializeArg(msg.value)); +/** Settles the pending call a response message answers. @internal */ +export function handleResponse(endpoint: Endpoint, msg: { callId: number; value?: unknown; error?: Error }): void { + const call = pending.get(msg.callId); + if (!call) return; + pending.delete(msg.callId); + endpointCalls.get(endpoint)?.delete(msg.callId); // no-op until Task 8 adds the map + if (msg.error !== undefined) { + call.reject(msg.error); + } else { + call.resolve(deserializeArg(msg.value)); + } +} + +export function setupWorker(endpoint: Endpoint, onCtrl?: (frame: CtrlFrame) => void): void { + endpoint.on('message', (msg: { callId: number; value?: unknown; error?: Error } | CtrlFrame) => { + if (isCtrlFrame(msg)) { + if (onCtrl === undefined) { + throw new Error(`Unknown control frame: ${String(msg.__ctrl__)} (no control handler on this endpoint)`); + } + onCtrl(msg); + return; } + handleResponse(endpoint, msg); }); } diff --git a/src/worker-entry.ts b/src/worker-entry.ts index 287daa4..e09a4d0 100644 --- a/src/worker-entry.ts +++ b/src/worker-entry.ts @@ -6,6 +6,9 @@ import { deserializeArg, serializeArg } from './shared/reconstruct.ts'; import { collectTransferables } from './transfer.ts'; import { pipeIterable, CANCEL, DEFAULT_HIGH_WATER, deserializeStreamHandle, isSerializedStreamHandle } from './pipe.ts'; import type { StreamHandle, SerializedStreamHandle } from './pipe.ts'; +import { isCtrlFrame } from './handshake.ts'; +import type { CtrlFrame } from './handshake.ts'; +import { handleResponse } from './execute.ts'; const imported = new Set(); const taskCache = new Map(); @@ -267,10 +270,34 @@ async function handleStreamTask(msg: TaskMsg): Promise { } } -parentPort!.on('message', (msg: TaskMsg) => { - if (msg.stream) { - void handleStreamTask(msg); - } else { - handleValueTask(msg); +function handleCtrlFrame(frame: CtrlFrame): void { + switch (frame.__ctrl__) { + // 'peer' handled in Task 5; the exhaustive default keeps unknowns loud. + default: + throw new Error(`Unknown control frame: ${String((frame as { __ctrl__: unknown }).__ctrl__)}`); } +} + +// ONE combined handler: ctrl / incoming task / response, discriminated by +// shape IN THAT ORDER. At this task the handlers are still parentPort-bound +// (Task 3 threads the endpoint parameter through and extracts the reusable +// attachEndpointHandler(endpoint) for peer ports); the ROUTING SHAPE lands +// here and does not change again. Full-duplex-safe by construction: an +// incoming request {callId, id, args} carries the sender's callId and must +// hit the task branch, never the response branch. +parentPort!.on('message', (msg: TaskMsg | CtrlFrame | { callId: number }) => { + if (isCtrlFrame(msg)) { + handleCtrlFrame(msg); + return; + } + if ('id' in msg && msg.id !== undefined) { + const task = msg as TaskMsg; + if (task.stream) { + void handleStreamTask(task); + } else { + handleValueTask(task); + } + return; + } + handleResponse(parentPort!, msg as { callId: number }); }); diff --git a/test/ctrl-frames.test.ts b/test/ctrl-frames.test.ts new file mode 100644 index 0000000..f131707 --- /dev/null +++ b/test/ctrl-frames.test.ts @@ -0,0 +1,25 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { execFile } from 'node:child_process'; +import { promisify } from 'node:util'; +import { fileURLToPath } from 'node:url'; +import { join } from 'node:path'; + +const exec = promisify(execFile); +const fixturesDir = join(fileURLToPath(import.meta.url), '..', 'fixtures'); + +describe('control frames', () => { + it('worker crashes loudly on an unknown ctrl frame', async () => { + const { stdout } = await exec(process.execPath, ['--no-warnings', join(fixturesDir, 'ctrl-unknown.ts')], { + timeout: 15000, + }); + assert.match(stdout, /WORKER-ERROR .*Unknown control frame/); + }); + + it('main crashes loudly on an unknown ctrl frame from a worker', async () => { + const { stdout } = await exec(process.execPath, ['--no-warnings', join(fixturesDir, 'ctrl-unknown-main.ts')], { + timeout: 15000, + }); + assert.match(stdout, /MAIN-ERROR .*Unknown control frame/); + }); +}); diff --git a/test/fixtures/ctrl-unknown-main.ts b/test/fixtures/ctrl-unknown-main.ts new file mode 100644 index 0000000..628444a --- /dev/null +++ b/test/fixtures/ctrl-unknown-main.ts @@ -0,0 +1,13 @@ +// Simulate a worker sending main an unknown ctrl frame: setupWorker's handler +// must throw. We attach setupWorker to a MessagePort pair and inject. +import { MessageChannel } from 'node:worker_threads'; +import { setupWorker } from '../../src/execute.ts'; + +const { port1, port2 } = new MessageChannel(); +setupWorker(port1); +process.on('uncaughtException', (err) => { + console.log('MAIN-ERROR ' + err.message); + port1.close(); + port2.close(); +}); +port2.postMessage({ __ctrl__: 'bogus' }); diff --git a/test/fixtures/ctrl-unknown.ts b/test/fixtures/ctrl-unknown.ts new file mode 100644 index 0000000..f4e1472 --- /dev/null +++ b/test/fixtures/ctrl-unknown.ts @@ -0,0 +1,12 @@ +// Send a bogus ctrl frame straight at a pool worker's parentPort and observe +// the worker error. Reaches the raw Worker via an internal probe: the pool +// does not expose threads, so spawn the entry directly. +import { Worker } from 'node:worker_threads'; + +const entryUrl = new URL('../../src/worker-entry.ts', import.meta.url); +const w = new Worker(entryUrl, { workerData: { index: 0, size: 1, countsBuffer: new SharedArrayBuffer(64) } }); +w.on('error', (err: Error) => { + console.log('WORKER-ERROR ' + err.message); + void w.terminate(); +}); +w.postMessage({ __ctrl__: 'bogus' }); -- 2.51.2