diff --git a/.changeset/channel-fairness.md b/.changeset/channel-fairness.md new file mode 100644 index 0000000..9312913 --- /dev/null +++ b/.changeset/channel-fairness.md @@ -0,0 +1,7 @@ +--- +'moroutine': patch +--- + +`channel()` fan-out now rotates consumer selection round-robin instead of always preferring the lowest-index consumer. Previously, when several consumers were below `highWaterMark`, the distributor would pick `consumers[0]` every time, skewing work toward the first worker. The scan now starts from a rotating cursor, so ties distribute evenly. + +Skip-based RR semantics are unchanged — consumers at cap are still skipped so a saturated worker doesn't stall the pipeline. Fairness is best at higher volumes (≥10K items) and under any real backpressure, where initial warm-up asymmetry washes out. See `examples/benchmark-dispatch/fanout.ts` for measured distributions. diff --git a/examples/benchmark-dispatch/fanout.ts b/examples/benchmark-dispatch/fanout.ts new file mode 100644 index 0000000..56fb488 --- /dev/null +++ b/examples/benchmark-dispatch/fanout.ts @@ -0,0 +1,83 @@ +// Measures channel() fan-out distribution and throughput under various +// load profiles: +// - No backpressure: consumers are trivial, producer races ahead +// - Backpressure: consumers burn CPU per item, producer stalls on caps +// +// Skip-based RR skips consumers at cap, so distribution skew is highest +// when the producer outpaces consumers AND initial warm-up timing varies. +// As N grows the transient fades and counts even out. +// +// Requires Node v24+. +// Run: node examples/benchmark-dispatch/fanout.ts + +import { availableParallelism } from 'node:os'; +import { workers, channel, assign } from '../../src/index.ts'; +import { noopConsume, noopConsumeBusy } from './noop.ts'; + +const maxWorkers = availableParallelism(); + +interface FanoutResult { + counts: number[]; + throughput: number; + elapsedMs: number; +} + +async function fanout(n: number, numWorkers: number): Promise { + using run = workers(numWorkers); + async function* source() { + for (let i = 0; i < n; i++) yield i; + } + const ch = channel(source()); + const tasks = run.workers.map((w) => assign(w, noopConsume(ch))); + const t0 = performance.now(); + const counts = (await run(tasks)) as number[]; + const elapsedMs = performance.now() - t0; + const total = counts.reduce((a, b) => a + b, 0); + if (total !== n) throw new Error(`expected ${n} items, got ${total}`); + return { counts, throughput: n / (elapsedMs / 1000), elapsedMs }; +} + +async function fanoutBusy(n: number, numWorkers: number, busyUs: number): Promise { + using run = workers(numWorkers); + async function* source() { + for (let i = 0; i < n; i++) yield i; + } + const ch = channel(source()); + const tasks = run.workers.map((w) => assign(w, noopConsumeBusy(ch, busyUs))); + const t0 = performance.now(); + const counts = (await run(tasks)) as number[]; + const elapsedMs = performance.now() - t0; + const total = counts.reduce((a, b) => a + b, 0); + if (total !== n) throw new Error(`expected ${n} items, got ${total}`); + return { counts, throughput: n / (elapsedMs / 1000), elapsedMs }; +} + +function fmtCounts(counts: number[]): string { + const mean = counts.reduce((a, b) => a + b, 0) / counts.length; + const min = Math.min(...counts); + const max = Math.max(...counts); + const spread = ((max - min) / mean) * 100; + return `[${counts.map((c) => c.toString().padStart(6)).join(', ')}] spread ${spread.toFixed(1)}% (min ${min}, max ${max}, mean ${Math.round(mean)})`; +} + +function fmtThroughput(r: FanoutResult): string { + return `${Math.round(r.throughput).toLocaleString().padStart(10)} items/s (${r.elapsedMs.toFixed(1)}ms)`; +} + +const W = 4; +console.log(`workers: ${W} (system has ${maxWorkers} cpus)\n`); + +console.log('no-backpressure fan-out (trivial consumer, producer races ahead):'); +for (const n of [400, 4_000, 40_000, 400_000]) { + const r = await fanout(n, W); + console.log(` n=${n.toString().padStart(7)} ${fmtThroughput(r)} ${fmtCounts(r.counts)}`); +} + +console.log('\nbackpressure fan-out (consumer burns CPU per item):'); +for (const busyUs of [10, 50, 200]) { + const n = busyUs >= 200 ? 4_000 : 20_000; + const r = await fanoutBusy(n, W, busyUs); + console.log( + ` busy=${String(busyUs).padStart(3)}µs n=${n.toString().padStart(6)} ${fmtThroughput(r)} ${fmtCounts(r.counts)}`, + ); +} diff --git a/examples/benchmark-dispatch/noop.ts b/examples/benchmark-dispatch/noop.ts index fe41034..02cb1bd 100644 --- a/examples/benchmark-dispatch/noop.ts +++ b/examples/benchmark-dispatch/noop.ts @@ -24,3 +24,22 @@ export const noopConsume = mo(import.meta, async (items: AsyncIterable): for await (const _ of items) count++; return count; }); + +// Consumer that burns CPU for a target number of microseconds per item — +// simulates real work and creates backpressure when the producer can outpace +// it. +export const noopConsumeBusy = mo( + import.meta, + async (items: AsyncIterable, busyUs: number): Promise => { + let count = 0; + const spinNs = BigInt(Math.round(busyUs * 1000)); + for await (const _ of items) { + count++; + const start = process.hrtime.bigint(); + while (process.hrtime.bigint() - start < spinNs) { + /* spin */ + } + } + return count; + }, +); diff --git a/examples/benchmark/main.ts b/examples/benchmark/main.ts index 7348fea..29b9b7a 100644 --- a/examples/benchmark/main.ts +++ b/examples/benchmark/main.ts @@ -6,7 +6,7 @@ // Run: node examples/benchmark/main.ts import { availableParallelism } from 'node:os'; -import { workers, channel } from '../../src/index.ts'; +import { workers, channel, assign } from '../../src/index.ts'; import { passthrough } from './work.ts'; const ITEMS = 100_000; @@ -20,8 +20,7 @@ async function bench(numWorkers: number): Promise { const data = channel(source()); using run = workers(numWorkers); - const tasks = Array.from({ length: numWorkers }, () => passthrough(data)); - const streams = tasks.map((t) => run(t)); + const streams = run.workers.map((w) => run(assign(w, passthrough(data)))); const start = performance.now(); let count = 0; diff --git a/src/channel.ts b/src/channel.ts index 75a5b49..ac14a70 100644 --- a/src/channel.ts +++ b/src/channel.ts @@ -32,16 +32,17 @@ interface Consumer { * One-producer, many-consumer source with atomics-based backpressure. * * A single `tryPull` loop iterates the shared source and dispatches each - * item to the first consumer whose `inflight` atomic is below highWater. - * When every consumer is at cap, the loop parks on a shared `readySignal` - * atomic until any consumer pulls from its port (which decrements its - * inflight and bumps readySignal, waking the loop). + * item to the first below-cap consumer starting from `cursor` (skip-based + * round-robin). When every consumer is at cap, the loop parks on a shared + * `readySignal` atomic until any consumer pulls from its port (which + * decrements its inflight and bumps readySignal, waking the loop). */ export class Channel { private readonly iter: AsyncIterator; private readonly consumers: Consumer[] = []; private readonly readySignal: Int32Atomic = new Int32Atomic(); private readonly highWater: number; + private cursor = 0; private pulling = false; private done = false; private error: Error | null = null; @@ -102,12 +103,23 @@ export class Channel { return handle; } + // Skip-based round-robin: scan from `cursor`, pick the first non-cancelled + // consumer below cap. Cursor advances past the picked consumer so ties + // rotate across the below-cap set. Consumers at cap are skipped — this + // lets faster consumers keep pulling when a peer is saturated, trading + // strict fairness for resilience and throughput. private findReadyConsumer(): Consumer | null { const hw = this.highWater; - for (let i = 0; i < this.consumers.length; i++) { + const n = this.consumers.length; + if (n === 0) return null; + for (let k = 0; k < n; k++) { + const i = (this.cursor + k) % n; const c = this.consumers[i]; if (c.flags.fields.state.load() === CANCEL) continue; - if (c.flags.fields.inflight.load() < hw) return c; + if (c.flags.fields.inflight.load() < hw) { + this.cursor = (i + 1) % n; + return c; + } } return null; } diff --git a/test/channel-fanout.test.ts b/test/channel-fanout.test.ts index 4c19c7b..8d9486d 100644 --- a/test/channel-fanout.test.ts +++ b/test/channel-fanout.test.ts @@ -1,37 +1,31 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; -import { workers, channel } from 'moroutine'; +import { workers, channel, assign } from 'moroutine'; import { collectItems, generate } from './fixtures/channel-fanout.ts'; describe('channel fan-out', () => { it('distributes items across multiple consumers', async () => { const data = channel(generate(20)); - const run = workers(4); - try { - const [a, b, c, d] = await run([collectItems(data), collectItems(data), collectItems(data), collectItems(data)]); - const all = [...a, ...b, ...c, ...d].sort((x, y) => x - y); - assert.deepEqual( - all, - Array.from({ length: 20 }, (_, i) => i), - ); - } finally { - run[Symbol.dispose](); - } + using run = workers(4); + const fanout = run.workers.map((w) => assign(w, collectItems(data))); + const [a, b, c, d] = await run(fanout); + const all = [...a, ...b, ...c, ...d].sort((x, y) => x - y); + assert.deepEqual( + all, + Array.from({ length: 20 }, (_, i) => i), + ); }); it('each item goes to exactly one consumer', async () => { const data = channel(generate(100)); - const run = workers(2); - try { - const [a, b] = await run([collectItems(data), collectItems(data)]); - const setA = new Set(a); - for (const item of b) { - assert.ok(!setA.has(item), `Item ${item} appeared in both consumers`); - } - assert.equal(a.length + b.length, 100); - } finally { - run[Symbol.dispose](); + using run = workers(2); + const fanout = run.workers.map((w) => assign(w, collectItems(data))); + const [a, b] = await run(fanout); + const setA = new Set(a); + for (const item of b) { + assert.ok(!setA.has(item), `Item ${item} appeared in both consumers`); } + assert.equal(a.length + b.length, 100); }); it('single consumer via channel() still works', async () => { @@ -40,13 +34,9 @@ describe('channel fan-out', () => { yield 2; yield 3; } - const run = workers(1); - try { - const result = await run(collectItems(channel(numbers()))); - assert.deepEqual(result, [1, 2, 3]); - } finally { - run[Symbol.dispose](); - } + using run = workers(1); + const result = await run(collectItems(channel(numbers()))); + assert.deepEqual(result, [1, 2, 3]); }); it('channel with local AsyncIterable fan-out', async () => { @@ -54,16 +44,32 @@ describe('channel fan-out', () => { for (let i = 0; i < 10; i++) yield i; } const data = channel(localGen()); - const run = workers(2); - try { - const [a, b] = await run([collectItems(data), collectItems(data)]); - const all = [...a, ...b].sort((x, y) => x - y); - assert.deepEqual( - all, - Array.from({ length: 10 }, (_, i) => i), - ); - } finally { - run[Symbol.dispose](); + using run = workers(2); + const fanout = run.workers.map((w) => assign(w, collectItems(data))); + const [a, b] = await run(fanout); + const all = [...a, ...b].sort((x, y) => x - y); + assert.deepEqual( + all, + Array.from({ length: 10 }, (_, i) => i), + ); + }); + + it('does not starve any consumer (skip-based round-robin)', async () => { + async function* localGen() { + for (let i = 0; i < 4000; i++) yield i; } + const data = channel(localGen()); + using run = workers(4); + const fanout = run.workers.map((w) => assign(w, collectItems(data))); + const counts = await run(fanout); + const lengths = counts.map((c) => c.length); + assert.equal( + lengths.reduce((a, b) => a + b, 0), + 4000, + ); + // Skip-based RR doesn't guarantee even splits under heterogeneous drain + // rates, but no consumer should be starved to just its initial fill. + const min = Math.min(...lengths); + assert.ok(min >= 100, `expected each consumer ≥ 100 items, got ${lengths.join(',')}`); }); });