From 12124c8969eec3e6fb651e9798a2e21e60f86be3 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Thu, 23 Jul 2026 00:20:18 -0400 Subject: [PATCH] docs: add benchmark-peers example for worker-originated dispatch --- README.md | 1 + examples/benchmark-peers/bench-runtime.ts | 6 +++ examples/benchmark-peers/main.ts | 45 +++++++++++++++++++++++ examples/benchmark-peers/work.ts | 39 ++++++++++++++++++++ 4 files changed, 91 insertions(+) create mode 100644 examples/benchmark-peers/bench-runtime.ts create mode 100644 examples/benchmark-peers/main.ts create mode 100644 examples/benchmark-peers/work.ts diff --git a/README.md b/README.md index 093e747..0192b18 100644 --- a/README.md +++ b/README.md @@ -728,3 +728,4 @@ All examples require Node v24+ and can be run directly, e.g. `node examples/prim - [`examples/server-threads`](examples/server-threads) - scale an HTTP server across worker threads via fd-fanout - [`examples/benchmark`](examples/benchmark) - roundtrip channel throughput with 1–N workers - [`examples/benchmark-dispatch`](examples/benchmark-dispatch) - task dispatch overhead measurement +- [`examples/benchmark-peers`](examples/benchmark-peers) - worker-originated dispatch: peer channels, loopback, nested fork-join diff --git a/examples/benchmark-peers/bench-runtime.ts b/examples/benchmark-peers/bench-runtime.ts new file mode 100644 index 0000000..1bd0aa3 --- /dev/null +++ b/examples/benchmark-peers/bench-runtime.ts @@ -0,0 +1,6 @@ +import { define } from '../../src/index.ts'; +import type { RuntimeDefinition } from '../../src/index.ts'; + +export default define(import.meta, { + size: 4, +}); diff --git a/examples/benchmark-peers/main.ts b/examples/benchmark-peers/main.ts new file mode 100644 index 0000000..4310c5c --- /dev/null +++ b/examples/benchmark-peers/main.ts @@ -0,0 +1,45 @@ +// Measures worker-originated dispatch: peer-to-peer round-trips over +// factory channels, loopback self-dispatch, and nested fork-join — the +// paths added by workers-as-peers. Compare against main→worker numbers +// from examples/benchmark-dispatch. +// +// Requires Node v24+. +// Run: node examples/benchmark-peers/main.ts + +import { registerRuntime, runtime, assign } from '../../src/index.ts'; +import { noop, pingPeer, pipePeer, treeSum } from './work.ts'; +import benchRuntime from './bench-runtime.ts'; + +registerRuntime(benchRuntime); +await runtime.run(noop(0)); // boot the runtime so runtime.workers resolves + +const SEQ = 20_000; +const PIPE = 50_000; + +function fmt(ops: number): string { + return `${Math.round(ops).toLocaleString().padStart(9)} ops/s (${(1e6 / ops).toFixed(1)} µs/roundtrip)`; +} + +// Worker 0 dispatches to worker 1 (peer channel) vs to itself (loopback). +console.log('worker-originated ping-pong latency (strict await-each):'); +const peer = await runtime.run(assign(runtime.workers[0], pingPeer(1, SEQ))); +console.log(` worker 0 → worker 1 (peer channel) ${fmt(peer)}`); +const self = await runtime.run(assign(runtime.workers[0], pingPeer(0, SEQ))); +console.log(` worker 0 → worker 0 (loopback) ${fmt(self)}`); + +console.log('\nworker-originated pipelined throughput (worker 0 → worker 1):'); +for (const window of [4, 16, 64]) { + const ips = await runtime.run(assign(runtime.workers[0], pipePeer(1, PIPE, window))); + console.log(` window=${String(window).padStart(2)} ${Math.round(ips).toLocaleString().padStart(9)} ops/s`); +} + +console.log('\nnested fork-join (treeSum over 1M values, 1024-leaf):'); +const values = Array.from({ length: 1 << 20 }, (_, i) => i); +const expected = ((1 << 20) / 2) * ((1 << 20) - 1); +{ + const t0 = performance.now(); + const sum = await treeSum(values); + const ms = performance.now() - t0; + if (sum !== expected) throw new Error(`expected ${expected}, got ${sum}`); + console.log(` ${ms.toFixed(0)} ms (${Math.round((1 << 20) / (ms / 1000)).toLocaleString()} values/s)`); +} diff --git a/examples/benchmark-peers/work.ts b/examples/benchmark-peers/work.ts new file mode 100644 index 0000000..7a3c9bc --- /dev/null +++ b/examples/benchmark-peers/work.ts @@ -0,0 +1,39 @@ +import { runtime, assign } from '../../src/index.ts'; +import { mo } from '../../src/index.ts'; + +export const noop = mo(import.meta, (i: number): number => i); + +/** Strict await-each dispatch from THIS worker to a target runtime worker. + * Measures peer-channel (or loopback, when target === self) round-trips. */ +export const pingPeer = mo(import.meta, async (target: number, iters: number): Promise => { + // Warm the channel + handler. + for (let i = 0; i < 100; i++) await runtime.run(assign(runtime.workers[target], noop(i))); + const t0 = performance.now(); + for (let i = 0; i < iters; i++) { + await runtime.run(assign(runtime.workers[target], noop(i))); + } + return iters / ((performance.now() - t0) / 1000); +}); + +/** Pipelined dispatch from THIS worker to a target: `inflight` concurrent slots. */ +export const pipePeer = mo(import.meta, async (target: number, iters: number, inflight: number): Promise => { + for (let i = 0; i < 100; i++) await runtime.run(assign(runtime.workers[target], noop(i))); + const t0 = performance.now(); + let next = 0; + async function slot(): Promise { + while (next < iters) { + const i = next++; + await runtime.run(assign(runtime.workers[target], noop(i))); + } + } + await Promise.all(Array.from({ length: inflight }, () => slot())); + return iters / ((performance.now() - t0) / 1000); +}); + +/** Recursive fork-join: splits across the pool via balancer placement. */ +export const treeSum = mo(import.meta, async (values: number[]): Promise => { + if (values.length <= 1024) return values.reduce((a, b) => a + b, 0); + const mid = values.length / 2; + const [left, right] = await runtime.run([treeSum(values.slice(0, mid)), treeSum(values.slice(mid))]); + return left + right; +}); -- 2.51.2