From 253d1e7c111a6054ad928444a5b1bae0aecdcc11 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Sun, 19 Apr 2026 22:24:09 -0400 Subject: [PATCH] feat(channel): respect highWaterMark option MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit channel(source, { highWaterMark }) now plumbs the threshold through: - Channel stores it from opts (default 16) - findReadyConsumer uses it for the cap check - addConsumer returns it in the per-consumer handle - prepareArg forwards it in the {__stream__, ...} arg - Worker portToAsyncIterable receives it; uses it in the cap→below-cap transition check that signals readySignal - Legacy message-based path uses it for the pause-threshold too Two sides both need it (producer scans against it; consumer detects the transition to signal readiness), so it travels with the arg. Co-Authored-By: Claude Opus 4.7 (1M context) --- src/channel.ts | 34 +++++++++++++++++++++++----------- src/execute.ts | 4 ++-- src/worker-entry.ts | 24 ++++++++++++++++-------- 3 files changed, 41 insertions(+), 21 deletions(-) diff --git a/src/channel.ts b/src/channel.ts index b04aef1..38a6080 100644 --- a/src/channel.ts +++ b/src/channel.ts @@ -9,11 +9,13 @@ import { serializeArg } from './shared/reconstruct.ts'; import { collectTransferables } from './transfer.ts'; const CHANNEL = Symbol.for('moroutine.channel'); -const HIGH_WATER = 16; +const DEFAULT_HIGH_WATER = 16; /** Options for configuring a channel. */ export interface ChannelOptions { - /** Maximum number of items buffered before signaling backpressure. Defaults to 16. */ + /** Max number of items any one consumer can have in-flight before the + * distributor skips it and gives items to other ready consumers. + * Defaults to 16. */ highWaterMark?: number; } @@ -36,12 +38,14 @@ export class Channel { private readonly consumers: Consumer[] = []; private readonly readySignalBuf: SharedArrayBuffer = new SharedArrayBuffer(4); private readonly readySignal: Int32Atomic = new Int32Atomic(this.readySignalBuf, 0); + private readonly highWater: number; private pulling = false; private done = false; private error: Error | null = null; readonly [CHANNEL] = true; - constructor(source: AsyncIterable) { + constructor(source: AsyncIterable, opts?: ChannelOptions) { + this.highWater = opts?.highWaterMark ?? DEFAULT_HIGH_WATER; const src = source instanceof AsyncIterableTask ? (runStreamOnDedicated((source as any).id, (source as any).args) as AsyncIterable) @@ -50,8 +54,15 @@ export class Channel { } /** Per-consumer handle. Returns the worker-facing port, the consumer's - * private flags (inflight + state) and the channel-wide readySignal. */ - addConsumer(): { port: MessagePort; flags: SharedArrayBuffer; readySignal: SharedArrayBuffer } { + * private flags (inflight + state), the channel-wide readySignal, and + * the channel's highWaterMark (so the worker knows when to signal the + * cap→below-cap transition). */ + addConsumer(): { + port: MessagePort; + flags: SharedArrayBuffer; + readySignal: SharedArrayBuffer; + highWater: number; + } { const { port1, port2 } = new MessageChannel(); port1.unref(); const flags = newPipeFlags(); @@ -65,7 +76,7 @@ export class Channel { try { port1.close(); } catch {} - return { port: port2, flags: flags.buffer, readySignal: this.readySignalBuf }; + return { port: port2, flags: flags.buffer, readySignal: this.readySignalBuf, highWater: this.highWater }; } const consumer: Consumer = { port: port1, flags }; @@ -85,14 +96,15 @@ export class Channel { this.readySignal.notify(); void this.tryPull(); - return { port: port2, flags: flags.buffer, readySignal: this.readySignalBuf }; + return { port: port2, flags: flags.buffer, readySignal: this.readySignalBuf, highWater: this.highWater }; } private findReadyConsumer(): Consumer | null { + const hw = this.highWater; for (let i = 0; i < this.consumers.length; i++) { const c = this.consumers[i]; if (c.flags.state.load() === CANCEL) continue; - if (c.flags.inflight.load() < HIGH_WATER) return c; + if (c.flags.inflight.load() < hw) return c; } return null; } @@ -167,11 +179,11 @@ export class Channel { * Creates a channel for streaming values to workers. Supports fan-out when * the same channel is passed to multiple tasks. * @param iterable - The async iterable or StreamTask to distribute. - * @param _opts - Optional. Reserved for future configuration. + * @param opts - Optional. `highWaterMark` bounds per-consumer in-flight items. * @returns A Channel, typed as AsyncIterable for transparent moroutine arg use. */ -export function channel(iterable: AsyncIterable, _opts?: ChannelOptions): AsyncIterable { - return new Channel(iterable) as unknown as AsyncIterable; +export function channel(iterable: AsyncIterable, opts?: ChannelOptions): AsyncIterable { + return new Channel(iterable, opts) as unknown as AsyncIterable; } export { CHANNEL }; diff --git a/src/execute.ts b/src/execute.ts index 58cc0f4..b9935e5 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -66,9 +66,9 @@ function prepareArg(arg: unknown): unknown { } // Channel wrapper — single distributor loop + per-consumer atomics backpressure. if (arg instanceof Channel) { - const { port, flags, readySignal } = arg.addConsumer(); + const { port, flags, readySignal, highWater } = arg.addConsumer(); streamPortStack[streamPortStack.length - 1].push(port); - return { __stream__: true, port, flags, readySignal }; + return { __stream__: true, port, flags, readySignal, highWater }; } if (arg instanceof PromiseLikeTask) { return { __task__: arg.uid, id: arg.id, args: arg.args.map(prepareArg) }; diff --git a/src/worker-entry.ts b/src/worker-entry.ts index 160d3be..db836f7 100644 --- a/src/worker-entry.ts +++ b/src/worker-entry.ts @@ -17,9 +17,13 @@ function isTaskArg(arg: unknown): arg is { __task__: number; id: string; args: u return typeof arg === 'object' && arg !== null && '__task__' in arg; } -function isStreamArg( - arg: unknown, -): arg is { __stream__: true; port: MessagePort; flags: SharedArrayBuffer; readySignal?: SharedArrayBuffer } { +function isStreamArg(arg: unknown): arg is { + __stream__: true; + port: MessagePort; + flags: SharedArrayBuffer; + readySignal?: SharedArrayBuffer; + highWater?: number; +} { return typeof arg === 'object' && arg !== null && (arg as any).__stream__ === true; } @@ -47,7 +51,12 @@ async function resolveFnSlow(id: string): Promise { return fn; } -function portToAsyncIterable(port: MessagePort, flags?: PipeFlags, readySignal?: Int32Atomic): AsyncIterable { +function portToAsyncIterable( + port: MessagePort, + flags?: PipeFlags, + readySignal?: Int32Atomic, + highWater: number = 16, +): AsyncIterable { const queue: T[] = []; let done = false; let error: Error | null = null; @@ -55,7 +64,6 @@ function portToAsyncIterable(port: MessagePort, flags?: PipeFlags, readySigna // Legacy message-based backpressure (when flags is absent). let paused = false; - const HIGH_WATER = 16; port.on('message', (msg: { value?: unknown; done?: boolean; error?: Error }) => { if (msg.error) { @@ -80,7 +88,7 @@ function portToAsyncIterable(port: MessagePort, flags?: PipeFlags, readySigna waiting(); waiting = null; } - if (!flags && !paused && queue.length >= HIGH_WATER) { + if (!flags && !paused && queue.length >= highWater) { paused = true; port.postMessage('pause'); } @@ -100,7 +108,7 @@ function portToAsyncIterable(port: MessagePort, flags?: PipeFlags, readySigna // we only wake it on the cap→below-cap transition. const prev = flags.inflight.sub(1); if (readySignal) { - if (prev === 16) { + if (prev === highWater) { readySignal.add(1); readySignal.notify(); } @@ -141,7 +149,7 @@ async function resolveArg(arg: unknown): Promise { if (isStreamArg(arg)) { const flags = pipeFlagsFromBuffer(arg.flags); const readySignal = arg.readySignal ? new Int32Atomic(arg.readySignal, 0) : undefined; - return portToAsyncIterable(arg.port, flags, readySignal); + return portToAsyncIterable(arg.port, flags, readySignal, arg.highWater); } if (arg instanceof MessagePort) { return portToAsyncIterable(arg); -- 2.51.2