From 2190ecfc340b6478fd54071f4509a30e191c5d55 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Fri, 10 Apr 2026 16:58:40 -0400 Subject: [PATCH] feat: channel-fanout and benchmark examples, fix Moroutine Awaited type - Add channel-fanout example demonstrating work-stealing fan-out - Add benchmark example measuring roundtrip channel throughput - Remove channel() from pipeline example (auto-detection handles it) - Fix Moroutine type to use Awaited so async tasks resolve correctly - Fix channel-fanout test types to use destructuring instead of any[] - Add new examples to README Co-Authored-By: Claude Opus 4.6 (1M context) --- README.md | 3 +++ examples/benchmark/main.ts | 46 +++++++++++++++++++++++++++++++++ examples/benchmark/work.ts | 10 +++++++ examples/channel-fanout/main.ts | 21 +++++++++++++++ examples/channel-fanout/work.ts | 17 ++++++++++++ examples/pipeline/main.ts | 7 +++-- src/mo.ts | 4 +-- test/channel-fanout.test.ts | 12 ++++----- 8 files changed, 107 insertions(+), 13 deletions(-) create mode 100644 examples/benchmark/main.ts create mode 100644 examples/benchmark/work.ts create mode 100644 examples/channel-fanout/main.ts create mode 100644 examples/channel-fanout/work.ts diff --git a/README.md b/README.md index f1f75ef..d0ac526 100644 --- a/README.md +++ b/README.md @@ -245,3 +245,6 @@ All examples require Node v24+ and can be run directly, e.g. `node examples/prim - [`examples/multi-module`](examples/multi-module) -- moroutines from multiple modules on one worker - [`examples/transfer`](examples/transfer) -- zero-copy buffer transfer to and from a worker - [`examples/sqlite`](examples/sqlite) -- shared SQLite database on a worker via task-arg caching +- [`examples/pipeline`](examples/pipeline) -- streaming pipeline across dedicated workers +- [`examples/channel-fanout`](examples/channel-fanout) -- fan-out a channel to multiple workers via work stealing +- [`examples/benchmark`](examples/benchmark) -- roundtrip channel throughput with 1–N workers diff --git a/examples/benchmark/main.ts b/examples/benchmark/main.ts new file mode 100644 index 0000000..7348fea --- /dev/null +++ b/examples/benchmark/main.ts @@ -0,0 +1,46 @@ +// Measures roundtrip channel throughput (main → worker → main) with +// 1 up to availableParallelism() workers. The worker is a passthrough +// streaming moroutine that yields each item back unchanged. +// Requires Node v24+. +// +// Run: node examples/benchmark/main.ts + +import { availableParallelism } from 'node:os'; +import { workers, channel } from '../../src/index.ts'; +import { passthrough } from './work.ts'; + +const ITEMS = 100_000; +const MAX_WORKERS = availableParallelism(); + +async function bench(numWorkers: number): Promise { + async function* source() { + for (let i = 0; i < ITEMS; i++) yield i; + } + + const data = channel(source()); + + using run = workers(numWorkers); + const tasks = Array.from({ length: numWorkers }, () => passthrough(data)); + const streams = tasks.map((t) => run(t)); + + const start = performance.now(); + let count = 0; + await Promise.all( + streams.map(async (s) => { + for await (const _ of s) count++; + }), + ); + const elapsed = performance.now() - start; + + return count / (elapsed / 1000); +} + +console.log(`Roundtrip throughput: ${ITEMS.toLocaleString()} items, 1–${MAX_WORKERS} workers\n`); + +for (let n = 1; n <= MAX_WORKERS; n++) { + const ips = await bench(n); + const bar = '#'.repeat(Math.round(ips / 1000)); + console.log( + ` ${String(n).padStart(2)} worker${n === 1 ? ' ' : 's'} ${Math.round(ips).toLocaleString().padStart(8)} items/s ${bar}`, + ); +} diff --git a/examples/benchmark/work.ts b/examples/benchmark/work.ts new file mode 100644 index 0000000..dd91ccd --- /dev/null +++ b/examples/benchmark/work.ts @@ -0,0 +1,10 @@ +import { mo } from '../../src/index.ts'; + +export const passthrough = mo(import.meta, async function* (input: AsyncIterable) { + for await (const n of input) { + // Small CPU cost per item so workers become the bottleneck + let x = n; + for (let i = 0; i < 10_000; i++) x = (x * 1103515245 + 12345) & 0x7fffffff; + yield x; + } +}); diff --git a/examples/channel-fanout/main.ts b/examples/channel-fanout/main.ts new file mode 100644 index 0000000..1a56afb --- /dev/null +++ b/examples/channel-fanout/main.ts @@ -0,0 +1,21 @@ +// Fan-out a single channel to multiple workers using work stealing. +// Each item goes to whichever worker is ready first. +// Requires Node v24+. +// +// Run: node examples/channel-fanout/main.ts + +import { workers, channel } from '../../src/index.ts'; +import { generate, process } from './work.ts'; + +{ + using run = workers(4); + const data = channel(generate(200)); + const results: number[][] = await run([process(data), process(data), process(data), process(data)]); + + for (let i = 0; i < results.length; i++) { + console.log(`Worker ${i}: processed ${results[i].length} items`); + } + + const all = results.flat().sort((a, b) => a - b); + console.log(`\nTotal: ${all.length} items, none lost, none duplicated`); +} diff --git a/examples/channel-fanout/work.ts b/examples/channel-fanout/work.ts new file mode 100644 index 0000000..871cfb0 --- /dev/null +++ b/examples/channel-fanout/work.ts @@ -0,0 +1,17 @@ +import { setTimeout } from 'node:timers/promises'; +import { mo } from '../../src/index.ts'; + +export const generate = mo(import.meta, async function* (n: number) { + for (let i = 0; i < n; i++) { + yield i; + } +}); + +export const process = mo(import.meta, async (input: AsyncIterable): Promise => { + const results: number[] = []; + for await (const n of input) { + await setTimeout(10); // simulate async work + results.push(n); + } + return results; +}); diff --git a/examples/pipeline/main.ts b/examples/pipeline/main.ts index cc37475..43e46ed 100644 --- a/examples/pipeline/main.ts +++ b/examples/pipeline/main.ts @@ -4,13 +4,12 @@ // // Run: node examples/pipeline/main.ts -import { channel } from '../../src/index.ts'; import { generate, double, square, toString } from './steps.ts'; const numbers = generate(5); -const doubled = double(channel(numbers)); -const squared = square(channel(doubled)); -const labels = toString(channel(squared)); +const doubled = double(numbers); +const squared = square(doubled); +const labels = toString(squared); for await (const label of labels) { console.log(label); diff --git a/src/mo.ts b/src/mo.ts index 4395d47..2f493da 100644 --- a/src/mo.ts +++ b/src/mo.ts @@ -16,8 +16,8 @@ export type Arg = T | Task; type TaskableArgs = { [K in keyof A]: Arg }; type Moroutine = { - (...args: A): Task; - (...args: TaskableArgs): Task; + (...args: A): Task>; + (...args: TaskableArgs): Task>; }; type StreamMoroutine = { diff --git a/test/channel-fanout.test.ts b/test/channel-fanout.test.ts index 0be9878..4eb383f 100644 --- a/test/channel-fanout.test.ts +++ b/test/channel-fanout.test.ts @@ -8,13 +8,13 @@ describe('channel fan-out', () => { const data = channel(generate(20)); const run = workers(4); try { - const results: any[] = await run([ + const [a, b, c, d] = await run([ collectItems(data), collectItems(data), collectItems(data), collectItems(data), ]); - const all = [...results[0], ...results[1], ...results[2], ...results[3]].sort((x: number, y: number) => x - y); + const all = [...a, ...b, ...c, ...d].sort((x, y) => x - y); assert.deepEqual(all, Array.from({ length: 20 }, (_, i) => i)); } finally { run[Symbol.dispose](); @@ -25,12 +25,10 @@ describe('channel fan-out', () => { const data = channel(generate(100)); const run = workers(2); try { - const results: any[] = await run([ + const [a, b] = await run([ collectItems(data), collectItems(data), ]); - const a: number[] = results[0]; - const b: number[] = results[1]; const setA = new Set(a); for (const item of b) { assert.ok(!setA.has(item), `Item ${item} appeared in both consumers`); @@ -59,11 +57,11 @@ describe('channel fan-out', () => { const data = channel(localGen()); const run = workers(2); try { - const results: any[] = await run([ + const [a, b] = await run([ collectItems(data), collectItems(data), ]); - const all = [...results[0], ...results[1]].sort((x: number, y: number) => x - y); + const all = [...a, ...b].sort((x, y) => x - y); assert.deepEqual(all, Array.from({ length: 10 }, (_, i) => i)); } finally { run[Symbol.dispose](); -- 2.51.2