diff --git a/examples/benchmark-dispatch/main.ts b/examples/benchmark-dispatch/main.ts new file mode 100644 index 0000000..68110da --- /dev/null +++ b/examples/benchmark-dispatch/main.ts @@ -0,0 +1,121 @@ +// Measures pure round-trip dispatch overhead for value tasks — a noop +// function is dispatched N times and timed. Surfaces the cost of +// postMessage + event loop re-entry per task. +// +// Three benchmarks: +// - Ping-pong latency: strict await-each-dispatch, 1 worker +// - Pipelined throughput: N in-flight at once, 1 worker +// - Parallel throughput: N in-flight spread across pool +// +// Requires Node v24+. +// Run: node examples/benchmark-dispatch/main.ts [--atomics] + +import { availableParallelism } from 'node:os'; +import { workers } from '../../src/index.ts'; +import { noop, noopBatch, noopStream } from './noop.ts'; + +const atomicsFlag = process.argv.includes('--atomics'); +const maxWorkers = availableParallelism(); + +// Passing `atomics: true` is forward-compatible — the option is ignored +// until the feature lands, so the same bench file times both modes. +const poolOpts = atomicsFlag ? ({ atomics: true } as any) : {}; + +async function warmup(run: ReturnType, n: number) { + for (let i = 0; i < n; i++) await run(noop(i)); +} + +async function pingPong(iters: number): Promise { + using run = workers(1, poolOpts); + await warmup(run, 100); + const t0 = performance.now(); + for (let i = 0; i < iters; i++) await run(noop(i)); + return iters / ((performance.now() - t0) / 1000); +} + +async function pipelined(iters: number, inflight: number): Promise { + using run = workers(1, poolOpts); + await warmup(run, 100); + const t0 = performance.now(); + let next = 0; + async function slot() { + while (next < iters) { + const i = next++; + await run(noop(i)); + } + } + await Promise.all(Array.from({ length: inflight }, () => slot())); + return iters / ((performance.now() - t0) / 1000); +} + +async function parallel(iters: number, numWorkers: number): Promise { + using run = workers(numWorkers, poolOpts); + await warmup(run, 100); + const t0 = performance.now(); + const promises = Array.from({ length: iters }, (_, i) => run(noop(i))); + await Promise.all(promises); + return iters / ((performance.now() - t0) / 1000); +} + +// One task carries the whole batch — items are not dispatched individually. +async function batched(iters: number): Promise { + using run = workers(1, poolOpts); + await run(noopBatch([1, 2, 3])); // warm the handler + const items = Array.from({ length: iters }, (_, i) => i); + const t0 = performance.now(); + await run(noopBatch(items)); + return iters / ((performance.now() - t0) / 1000); +} + +// Streaming: items flow over a MessageChannel to the worker-side async +// generator. Worker drains at its own pace — no per-item task envelope. +async function streaming(iters: number): Promise { + using run = workers(1, poolOpts); + async function* probe() { + for (let i = 0; i < 3; i++) yield i; + } + for await (const _ of run(noopStream(probe()))) void _; // warm + + async function* gen() { + for (let i = 0; i < iters; i++) yield i; + } + const t0 = performance.now(); + let count = 0; + for await (const _ of run(noopStream(gen()))) { + count++; + void _; + } + if (count !== iters) throw new Error(`expected ${iters} items, got ${count}`); + return iters / ((performance.now() - t0) / 1000); +} + +const SEQ = 20_000; +const PIPE = 50_000; +const PAR = 100_000; + +console.log(`mode: ${atomicsFlag ? 'atomics' : 'default (postMessage events)'}`); +console.log(`cpus: ${maxWorkers}\n`); + +console.log('ping-pong latency (strict await-each, 1 worker):'); +const pp = await pingPong(SEQ); +console.log(` ${SEQ.toLocaleString()} tasks ${Math.round(pp).toLocaleString().padStart(9)} ops/s (${(1e6 / pp).toFixed(1)} µs/roundtrip)`); + +console.log('\npipelined throughput (1 worker, in-flight window):'); +for (const window of [4, 16, 64, 256]) { + const ips = await pipelined(PIPE, window); + console.log(` window=${String(window).padStart(3)} ${Math.round(ips).toLocaleString().padStart(9)} ops/s`); +} + +console.log('\nparallel throughput (N workers, 100K tasks):'); +for (const n of [1, 2, 4, maxWorkers]) { + const ips = await parallel(PAR, n); + console.log(` ${String(n).padStart(2)} workers ${Math.round(ips).toLocaleString().padStart(9)} ops/s`); +} + +console.log('\nbatch vs stream vs per-task (1 worker, 100K items):'); +const perTask = await parallel(PAR, 1); +const bat = await batched(PAR); +const str = await streaming(PAR); +console.log(` per-task (run(noop(i)) ×N) ${Math.round(perTask).toLocaleString().padStart(9)} items/s`); +console.log(` batch (run(noopBatch([...N]))) ${Math.round(bat).toLocaleString().padStart(9)} items/s`); +console.log(` stream (run(noopStream(gen()))) ${Math.round(str).toLocaleString().padStart(9)} items/s`); diff --git a/examples/benchmark-dispatch/noop.ts b/examples/benchmark-dispatch/noop.ts new file mode 100644 index 0000000..368acdd --- /dev/null +++ b/examples/benchmark-dispatch/noop.ts @@ -0,0 +1,15 @@ +import { mo } from '../../src/index.ts'; + +// Tasks that do essentially zero work. Any non-trivial measurement is +// dominated by dispatch overhead (postMessage + event loop re-entry). + +// One task per item — full per-task round trip. +export const noop = mo(import.meta, (x: number): number => x); + +// One task, N items bundled as an array arg. Worker loops locally, returns once. +export const noopBatch = mo(import.meta, (items: number[]): number[] => items); + +// One streaming task; items arrive over a MessageChannel, no per-item task envelope. +export const noopStream = mo(import.meta, async function* (items: AsyncIterable) { + for await (const item of items) yield item; +});