diff --git a/src/execute.ts b/src/execute.ts index 150c4d7..5e84e1c 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -1,4 +1,3 @@ -import type { Worker } from 'node:worker_threads'; import { MessageChannel } from 'node:worker_threads'; import type { MessagePort, Transferable } from 'node:worker_threads'; import { transferableAbortSignal } from 'node:util'; @@ -18,8 +17,15 @@ let nextCallId = 0; const pending = new Map void; reject: (reason: any) => void }>(); const streamPortStack: MessagePort[][] = []; -export function setupWorker(worker: Worker): void { - worker.on('message', (msg: { callId: number; value?: unknown; error?: Error }) => { +/** Anything dispatch can talk over: a Worker from the main thread, or a + * MessagePort between peers. Structurally satisfied by both. */ +export interface Endpoint { + postMessage(value: unknown, transferList?: readonly Transferable[]): void; + 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); @@ -85,7 +91,7 @@ function prepareArg(arg: unknown): unknown { return serializeArg(arg); } -export function execute(worker: Worker, id: string, args: unknown[]): Promise { +export function execute(endpoint: Endpoint, id: string, args: unknown[]): Promise { const url = id.slice(0, id.lastIndexOf('#')); freezeModule(url); const callId = nextCallId++; @@ -96,7 +102,7 @@ export function execute(worker: Worker, id: string, args: unknown[]): Promise const preparedArgs = extracted.args.map(prepareArg); const ports = streamPortStack.pop()!; const msg = { callId, id, args: preparedArgs }; - worker.postMessage(msg, [...extracted.transfer, ...ports] as any[]); + endpoint.postMessage(msg, [...extracted.transfer, ...ports] as any[]); }); } @@ -106,7 +112,7 @@ export interface StreamDispatch { } export function dispatchStream( - worker: Worker, + endpoint: Endpoint, id: string, args: unknown[], opts?: ChannelOptions, @@ -128,7 +134,7 @@ export function dispatchStream( const preparedArgs = extracted.args.map(prepareArg); const ports = streamPortStack.pop()!; const msg = { id, args: preparedArgs, stream: serializeStreamHandle(handle) }; - worker.postMessage(msg, [...extracted.transfer, ...ports, port2] as any[]); + endpoint.postMessage(msg, [...extracted.transfer, ...ports, port2] as any[]); let resolveDone: () => void; const donePromise = new Promise((r) => {