From 8fca2352f53be9ba8c7d9dc61c508f1e0d8dcdb4 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Sun, 19 Apr 2026 00:33:12 -0400 Subject: [PATCH] refactor: use Int32Atomic for pipe backpressure; add waitAsync/notify MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Dogfood moroutine's own shared-memory primitives for the stream backpressure path instead of reinventing the Atomics layer. - Add `waitAsync(expected, timeoutMs?)` and `notify(count?)` methods to Int32Atomic, wrapping `Atomics.waitAsync` / `Atomics.notify`. Returns 'ok' | 'not-equal' | 'timed-out' from the wait. 5 new tests. - Replace raw Int32Array + Atomics.* calls in pipe.ts with Int32Atomic method calls. A PipeFlags struct holds two Int32Atomic views (inflight at offset 0, state at offset 4) over a single 8-byte SAB. - Expose newPipeFlags() / pipeFlagsFromBuffer() so callers only work with typed wrappers; the raw SharedArrayBuffer is passed over the dispatch message edge but both sides reconstruct PipeFlags from it. Behavior unchanged — same throughput (~450K/s) and backpressure tightness (~89ms steady latency with 5ms-per-item consumer). Co-Authored-By: Claude Opus 4.7 (1M context) --- src/execute.ts | 19 ++++----- src/pipe.ts | 71 ++++++++++++++++++++++---------- src/shared/int32-atomic.ts | 20 +++++++++ src/worker-entry.ts | 5 +-- test/shared/int32-atomic.test.ts | 35 ++++++++++++++++ 5 files changed, 115 insertions(+), 35 deletions(-) diff --git a/src/execute.ts b/src/execute.ts index c91b36b..41a30fb 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, FLAGS_BYTE_LEN, INFLIGHT, STATE, CANCEL } from './pipe.ts'; +import { pipeIterable, newPipeFlags, CANCEL } from './pipe.ts'; import type { ChannelOptions } from './channel.ts'; let nextCallId = 0; @@ -108,17 +108,16 @@ export function dispatchStream( 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 + // Worker parks on flags.inflight when it hits highWater; parent decrements + // + notifies on pull. Replaces pause/resume messages and the // adaptive-yield setImmediate dance. - const flagsBuf = new SharedArrayBuffer(FLAGS_BYTE_LEN); - const flags = new Int32Array(flagsBuf); + const flags = newPipeFlags(); const extracted = extractTransferables(args); streamPortStack.push([]); const preparedArgs = extracted.args.map(prepareArg); const ports = streamPortStack.pop()!; - const msg = { id, args: preparedArgs, port: port2, flags: flagsBuf }; + const msg = { id, args: preparedArgs, port: port2, flags: flags.buffer }; worker.postMessage(msg, [...extracted.transfer, ...ports, port2] as any[]); let resolveDone: () => void; @@ -170,8 +169,8 @@ export function dispatchStream( if (queue.length > 0) { const value = queue.shift()!; // Signal consumption to producer - Atomics.sub(flags, INFLIGHT, 1); - Atomics.notify(flags, INFLIGHT); + flags.inflight.sub(1); + flags.inflight.notify(); return { done: false, value }; } if (error) throw error; @@ -182,8 +181,8 @@ export function dispatchStream( } }, async return(): Promise> { - Atomics.store(flags, STATE, CANCEL); - Atomics.notify(flags, INFLIGHT); + flags.state.store(CANCEL); + flags.inflight.notify(); port1.close(); resolveDone!(); return { done: true, value: undefined }; diff --git a/src/pipe.ts b/src/pipe.ts index df00fdc..f37c217 100644 --- a/src/pipe.ts +++ b/src/pipe.ts @@ -1,17 +1,42 @@ import type { MessagePort, Transferable } from 'node:worker_threads'; import { serializeArg } from './shared/reconstruct.ts'; import { collectTransferables, extractTransferables } from './transfer.ts'; +import { Int32Atomic } from './shared/int32-atomic.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; +const RUN = 0; +const CANCEL = 1; +const FLAGS_BYTE_LEN = 8; +const INFLIGHT_OFFSET = 0; +const STATE_OFFSET = 4; + +export interface PipeFlags { + inflight: Int32Atomic; + state: Int32Atomic; + /** Backing SharedArrayBuffer — transport via postMessage, reconstruct with {@link pipeFlagsFromBuffer} on the other side. */ + buffer: SharedArrayBuffer; +} + +export function newPipeFlags(): PipeFlags { + const buffer = new SharedArrayBuffer(FLAGS_BYTE_LEN); + return { + inflight: new Int32Atomic(buffer, INFLIGHT_OFFSET), + state: new Int32Atomic(buffer, STATE_OFFSET), + buffer, + }; +} + +export function pipeFlagsFromBuffer(buffer: SharedArrayBuffer): PipeFlags { + return { + inflight: new Int32Atomic(buffer, INFLIGHT_OFFSET), + state: new Int32Atomic(buffer, STATE_OFFSET), + buffer, + }; +} export interface PipeOptions { /** 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`. */ + /** Target in-flight count before the producer parks on the `flags.inflight` atomic. Defaults to 16. */ highWater?: number; /** * When true, call extractTransferables on each value to pull out any @@ -21,13 +46,13 @@ export interface PipeOptions { */ 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. + * When provided, pipe uses atomics-based backpressure: the producer parks + * on `flags.inflight` when it reaches `highWater`, the consumer + * decrements + notifies on pull, and cancellation is signalled by writing + * `CANCEL` to `flags.state`. Eliminates pause/resume messages and the + * adaptive-yield overhead. */ - flags?: Int32Array; + flags?: PipeFlags; } /** @@ -45,15 +70,16 @@ export async function pipeIterable(source: AsyncIterable, port: MessagePor async function pipeWithAtomics( source: AsyncIterable, port: MessagePort, - flags: Int32Array, + flags: PipeFlags, opts: PipeOptions, ): Promise { const highWater = opts.highWater ?? 16; const extract = opts.extractTransfers ?? false; + const { inflight, state } = flags; try { for await (const value of source) { - if (Atomics.load(flags, STATE) === CANCEL) return; + if (state.load() === CANCEL) return; let serialized: unknown; const transferList: Transferable[] = []; @@ -68,21 +94,20 @@ async function pipeWithAtomics( } port.postMessage({ value: serialized, done: false }, transferList); - const newInflight = Atomics.add(flags, INFLIGHT, 1) + 1; + const newInflight = inflight.add(1) + 1; if (newInflight >= highWater) { - // Park until parent drains below highWater or we're cancelled. + // Park until consumer drains below highWater or we're cancelled. while (true) { - if (Atomics.load(flags, STATE) === CANCEL) return; - const current = Atomics.load(flags, INFLIGHT); + if (state.load() === CANCEL) return; + const current = inflight.load(); if (current < highWater) break; - const { async, value: p } = Atomics.waitAsync(flags, INFLIGHT, current); - if (async) await p; + await inflight.waitAsync(current); } } } - if (Atomics.load(flags, STATE) !== CANCEL) port.postMessage({ done: true }); + if (state.load() !== CANCEL) port.postMessage({ done: true }); } catch (err) { - if (Atomics.load(flags, STATE) !== CANCEL) { + if (state.load() !== CANCEL) { port.postMessage({ done: true, error: err instanceof Error ? err : new Error(String(err)) }); } } @@ -173,3 +198,5 @@ async function pipeWithMessages(source: AsyncIterable, port: MessagePort, port.close(); } catch {} } + +export { CANCEL, RUN }; diff --git a/src/shared/int32-atomic.ts b/src/shared/int32-atomic.ts index d031d47..433b8f0 100644 --- a/src/shared/int32-atomic.ts +++ b/src/shared/int32-atomic.ts @@ -47,6 +47,26 @@ export class Int32Atomic { return Atomics.compareExchange(this.view, 0, expected, replacement); } + /** + * Parks asynchronously until the slot's value differs from `expected` or the optional timeout elapses. + * Returns `'not-equal'` synchronously (no await) if the slot already holds a different value. + * @param expected - Value to wait against; wake when the slot's value is observed to differ. + * @param timeoutMs - Optional timeout in milliseconds. Defaults to Infinity. + * @returns `'ok'` if woken by a notify, `'not-equal'` if the value already differed, `'timed-out'` if the timeout elapsed. + */ + async waitAsync(expected: number, timeoutMs?: number): Promise<'ok' | 'not-equal' | 'timed-out'> { + const r = Atomics.waitAsync(this.view, 0, expected, timeoutMs); + return r.async ? await r.value : r.value; + } + + /** + * Wakes up to `count` waiters on this slot. Returns the number actually woken. + * @param count - Maximum number of waiters to wake. Defaults to Infinity. + */ + notify(count: number = Infinity): number { + return Atomics.notify(this.view, 0, count); + } + [Symbol.for('moroutine.shared')](): { tag: string; buffer: SharedArrayBuffer; byteOffset: number } { return { tag: 'Int32Atomic', buffer: this.view.buffer as SharedArrayBuffer, byteOffset: this.view.byteOffset }; } diff --git a/src/worker-entry.ts b/src/worker-entry.ts index 947a5e9..c911f88 100644 --- a/src/worker-entry.ts +++ b/src/worker-entry.ts @@ -3,7 +3,7 @@ import type { Transferable } from 'node:worker_threads'; import { registry } from './registry.ts'; import { deserializeArg, serializeArg } from './shared/reconstruct.ts'; import { collectTransferables } from './transfer.ts'; -import { pipeIterable } from './pipe.ts'; +import { pipeIterable, pipeFlagsFromBuffer } from './pipe.ts'; const imported = new Set(); const taskCache = new Map(); @@ -221,8 +221,7 @@ async function handleStreamTask(msg: TaskMsg): Promise { resolvedArgs[i] = needsAsyncResolve(a) ? await resolveArg(a) : deserializeArg(a); } const gen = fn(...resolvedArgs) as AsyncGenerator; - const flagsView = flags ? new Int32Array(flags) : undefined; - await pipeIterable(gen, port!, { flags: flagsView }); + await pipeIterable(gen, port!, { flags: flags ? pipeFlagsFromBuffer(flags) : undefined }); } catch (err) { if (callId != null) postError(callId, err); } diff --git a/test/shared/int32-atomic.test.ts b/test/shared/int32-atomic.test.ts index 109e077..97c4719 100644 --- a/test/shared/int32-atomic.test.ts +++ b/test/shared/int32-atomic.test.ts @@ -81,4 +81,39 @@ describe('Int32Atomic', () => { assert.equal(actual, 42); assert.equal(a.load(), 42); }); + + it('waitAsync returns "not-equal" synchronously when value already differs', async () => { + const a = int32atomic(); + a.store(5); + const result = await a.waitAsync(0); + assert.equal(result, 'not-equal'); + }); + + it('waitAsync resolves "ok" when a concurrent store + notify happens', async () => { + const a = int32atomic(); + const wait = a.waitAsync(0); + setImmediate(() => { + a.store(1); + a.notify(); + }); + assert.equal(await wait, 'ok'); + assert.equal(a.load(), 1); + }); + + it('waitAsync resolves "timed-out" when no notify arrives', async () => { + const a = int32atomic(); + const result = await a.waitAsync(0, 20); + assert.equal(result, 'timed-out'); + }); + + it('notify returns count of waiters woken', async () => { + const a = int32atomic(); + const w1 = a.waitAsync(0); + const w2 = a.waitAsync(0); + const woken = a.notify(); + a.store(1); + a.notify(); + assert.equal(woken + a.notify(), 2); + await Promise.all([w1, w2]); + }); }); -- 2.51.2