diff --git a/src/execute.ts b/src/execute.ts index c46a572..c91b36b 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -9,7 +9,7 @@ import { PromiseLikeTask } from './task.ts'; import { AsyncIterableTask } from './stream-task.ts'; import { runStreamOnDedicated } from './dedicated-runner.ts'; import { CHANNEL, Channel } from './channel.ts'; -import { pipeIterable } from './pipe.ts'; +import { pipeIterable, FLAGS_BYTE_LEN, INFLIGHT, STATE, CANCEL } from './pipe.ts'; import type { ChannelOptions } from './channel.ts'; let nextCallId = 0; @@ -107,12 +107,18 @@ export function dispatchStream( const highWater = opts?.highWaterMark ?? DEFAULT_HIGH_WATER; const { port1, port2 } = new MessageChannel(); + // Atomics-based backpressure: shared inflight count + cancel flag. + // Worker parks on INFLIGHT when it hits highWater; parent decrements + // on pull and notifies. Replaces pause/resume messages and the + // adaptive-yield setImmediate dance. + const flagsBuf = new SharedArrayBuffer(FLAGS_BYTE_LEN); + const flags = new Int32Array(flagsBuf); const extracted = extractTransferables(args); streamPortStack.push([]); const preparedArgs = extracted.args.map(prepareArg); const ports = streamPortStack.pop()!; - const msg = { id, args: preparedArgs, port: port2 }; + const msg = { id, args: preparedArgs, port: port2, flags: flagsBuf }; worker.postMessage(msg, [...extracted.transfer, ...ports, port2] as any[]); let resolveDone: () => void; @@ -123,7 +129,6 @@ export function dispatchStream( const queue: T[] = []; let done = false; let error: Error | null = null; - let paused = false; let waiting: (() => void) | null = null; port1.on('message', (msg: { value?: unknown; done?: boolean; error?: Error }) => { @@ -153,10 +158,6 @@ export function dispatchStream( waiting(); waiting = null; } - if (!paused && queue.length >= highWater) { - paused = true; - port1.postMessage('pause'); - } }); port1.unref(); @@ -168,10 +169,9 @@ export function dispatchStream( while (true) { if (queue.length > 0) { const value = queue.shift()!; - if (paused && queue.length <= LOW_WATER) { - paused = false; - port1.postMessage('resume'); - } + // Signal consumption to producer + Atomics.sub(flags, INFLIGHT, 1); + Atomics.notify(flags, INFLIGHT); return { done: false, value }; } if (error) throw error; @@ -182,6 +182,8 @@ export function dispatchStream( } }, async return(): Promise> { + Atomics.store(flags, STATE, CANCEL); + Atomics.notify(flags, INFLIGHT); port1.close(); resolveDone!(); return { done: true, value: undefined }; diff --git a/src/pipe.ts b/src/pipe.ts index 1add5b3..df00fdc 100644 --- a/src/pipe.ts +++ b/src/pipe.ts @@ -2,9 +2,17 @@ import type { MessagePort, Transferable } from 'node:worker_threads'; import { serializeArg } from './shared/reconstruct.ts'; import { collectTransferables, extractTransferables } from './transfer.ts'; +export const INFLIGHT = 0; +export const STATE = 1; +export const RUN = 0; +export const CANCEL = 1; +export const FLAGS_BYTE_LEN = 2 * Int32Array.BYTES_PER_ELEMENT; + export interface PipeOptions { - /** How many emits without a natural event-loop tick before forcing a yield. Defaults to 16. */ + /** How many emits without a natural event-loop tick before forcing a yield. Defaults to 16. Unused when `flags` is provided. */ yieldEvery?: number; + /** Target in-flight count before the producer parks on Atomics. Defaults to 16. Used only with `flags`. */ + highWater?: number; /** * When true, call extractTransferables on each value to pull out any * `transfer(buf)` markers into the transfer list. Set on the parent side @@ -12,16 +20,78 @@ export interface PipeOptions { * come from a user generator and don't need this processing. */ extractTransfers?: boolean; + /** + * When provided, pipe uses an atomics-based backpressure protocol: + * the producer parks on `INFLIGHT` when it reaches `highWater`, the + * consumer increments/decrements via `Atomics.sub`+`notify` on pull, + * and cancellation is signalled by writing `CANCEL` to the `STATE` slot. + * Eliminates pause/resume messages and the adaptive-yield overhead. + */ + flags?: Int32Array; } /** * Pumps values from an async iterable to a MessagePort. Owns pause/resume - * handling, cancellation on port close, serialization + transferable - * collection, and an adaptive yield policy that ticks the event loop only - * when the producer is starving it (pure-CPU generators) and stays out of - * the way when natural awaits already tick it (I/O-backed generators). + * wiring, cancellation, serialization, and transferable collection. With + * `flags` provided, uses atomics-based backpressure (tight, no per-item + * yield); without, falls back to pause/resume messages with an adaptive + * yield policy. */ export async function pipeIterable(source: AsyncIterable, port: MessagePort, opts: PipeOptions = {}): Promise { + if (opts.flags) return pipeWithAtomics(source, port, opts.flags, opts); + return pipeWithMessages(source, port, opts); +} + +async function pipeWithAtomics( + source: AsyncIterable, + port: MessagePort, + flags: Int32Array, + opts: PipeOptions, +): Promise { + const highWater = opts.highWater ?? 16; + const extract = opts.extractTransfers ?? false; + + try { + for await (const value of source) { + if (Atomics.load(flags, STATE) === CANCEL) return; + + let serialized: unknown; + const transferList: Transferable[] = []; + if (extract) { + const extracted = extractTransferables([value]); + serialized = serializeArg(extracted.args[0]); + transferList.push(...extracted.transfer); + collectTransferables(extracted.args[0], transferList); + } else { + serialized = serializeArg(value); + collectTransferables(value, transferList); + } + port.postMessage({ value: serialized, done: false }, transferList); + + const newInflight = Atomics.add(flags, INFLIGHT, 1) + 1; + if (newInflight >= highWater) { + // Park until parent drains below highWater or we're cancelled. + while (true) { + if (Atomics.load(flags, STATE) === CANCEL) return; + const current = Atomics.load(flags, INFLIGHT); + if (current < highWater) break; + const { async, value: p } = Atomics.waitAsync(flags, INFLIGHT, current); + if (async) await p; + } + } + } + if (Atomics.load(flags, STATE) !== CANCEL) port.postMessage({ done: true }); + } catch (err) { + if (Atomics.load(flags, STATE) !== CANCEL) { + port.postMessage({ done: true, error: err instanceof Error ? err : new Error(String(err)) }); + } + } + try { + port.close(); + } catch {} +} + +async function pipeWithMessages(source: AsyncIterable, port: MessagePort, opts: PipeOptions): Promise { const yieldEvery = opts.yieldEvery ?? 16; const extract = opts.extractTransfers ?? false; diff --git a/src/worker-entry.ts b/src/worker-entry.ts index ed7cb51..947a5e9 100644 --- a/src/worker-entry.ts +++ b/src/worker-entry.ts @@ -130,7 +130,7 @@ async function resolveArg(arg: unknown): Promise { return deserializeArg(arg); } -type TaskMsg = { callId?: number; id: string; args: unknown[]; port?: MessagePort }; +type TaskMsg = { callId?: number; id: string; args: unknown[]; port?: MessagePort; flags?: SharedArrayBuffer }; function postResult(callId: number, value: unknown): void { const returnValue = serializeArg(value); @@ -210,7 +210,7 @@ function invokeAndRespond(callId: number, fn: Fn, args: unknown[]): void { } async function handleStreamTask(msg: TaskMsg): Promise { - const { id, args, port } = msg; + const { id, args, port, flags } = msg; const callId = msg.callId; try { const fnM = resolveFn(id); @@ -221,7 +221,8 @@ async function handleStreamTask(msg: TaskMsg): Promise { resolvedArgs[i] = needsAsyncResolve(a) ? await resolveArg(a) : deserializeArg(a); } const gen = fn(...resolvedArgs) as AsyncGenerator; - await pipeIterable(gen, port!); + const flagsView = flags ? new Int32Array(flags) : undefined; + await pipeIterable(gen, port!, { flags: flagsView }); } catch (err) { if (callId != null) postError(callId, err); }