diff --git a/fiber-await/README.md b/fiber-await/README.md new file mode 100644 index 0000000..aa5cd43 --- /dev/null +++ b/fiber-await/README.md @@ -0,0 +1,70 @@ +# Measuring the fiber-await prototype + +Measurement tooling for running `await` continuations on fibers ([node#7](https://github.com/qualialabs/node/pull/7) + [node#9](https://github.com/qualialabs/node/pull/9), [node-fibers#6](https://github.com/qualialabs/node-fibers/pull/6) + [#8](https://github.com/qualialabs/node-fibers/pull/8)). It lives on its own branch so it never ships: build fibers from this branch only in environments you want to measure. + +This branch adds, on top of #8: +- `Fiber.microtaskDrains`, `Fiber.drainsWithDispatch`, `Fiber.redispatches` (C++ counters; `Fiber.microtasksDispatched` and `Fiber.microtasksBatched` already exist) +- `fiber-await/test.js`, `fiber-await/bench.js`, `fiber-await/app_stats.mjs` + +## What gets measured + +A promise callback (the code after an `await`, a `.then` callback) is **tagged** if it was registered while a fiber was running. Tagged callbacks are run on a hidden *dispatch fiber*; untagged ones run on the main stack as on stock node. + +| Counter | Meaning | +|---|---| +| `microtasksDispatched` | Switches onto a dispatch fiber | +| `microtasksBatched` | Tagged callbacks run on the dispatch fiber without a switch (they were queued right behind another tagged one) | +| `microtaskDrains` | Microtask checkpoints on the default queue, including empty ones (node checkpoints about twice per event-loop turn) | +| `drainsWithDispatch` | Drains that switched onto a dispatch fiber at least once | +| `redispatches` | Further switches within the same drain: an untagged callback between two tagged ones sent us back to the main stack | + +Each switch costs about one fiber `run()` + `yield()` round trip. On the `CORO_PTHREAD` backend (Linux, arm64) every fiber is an OS thread, so that's a thread handoff. + +## Behaviour checks + +```sh +node --no-deprecation fiber-await/test.js # expect: all passed +FIBERS_AWAIT_REUSE=0 node --no-deprecation fiber-await/test.js # one new fiber per job; all passed +FIBERS_AWAIT_DISPATCH=0 node --no-deprecation fiber-await/test.js # baseline: 7 fail, as on stock node +``` + +znewsham's scenarios from [node-fibers#7](https://github.com/qualialabs/node-fibers/pull/7) (`test-microtask-fibers.js`) also run here once `Fiber.current` is replaced with `(Fiber.currentIncludingHidden || Fiber.current)`, since this design hides the fiber from `Fiber.current`. With `Fiber.poolSize = 1e9`, 7/7 behaviour tests pass. Its 100k park/resume test takes about 10 s here, just over its 10 s timeout. + +## Microbenchmarks + +```sh +node --no-deprecation fiber-await/bench.js +BENCH=switch node --no-deprecation fiber-await/bench.js # one benchmark; also works on stock node + fibers +``` + +Run it with the same binary and `FIBERS_AWAIT_REUSE=0` / `FIBERS_AWAIT_DISPATCH=0`, and with stock node, to compare. `switch` needs no patch, so it can measure the thread-handoff floor on any environment before deploying anything. + +Reference numbers from one `bench.js` run per column (arm64 Docker Desktop VM on an M-series Mac, node 18.16.1, 2026-10-02; expect run-to-run variance of 10-30%): + +| Benchmark | stock node + fibers | patched, `FIBERS_AWAIT_REUSE=0` | patched (reuse on) | +|---|---|---|---| +| `switch` (run+yield round trip) | 20.4 us | 24.7 us | 20.6 us | +| `chain` | 0.04 us | 24.8 us | 0.13 us | +| `io` | 0.61 us | 31.8 us | 20.8 us | +| `interleaved` | 0.09 us | 28.3 us | 21.8 us | +| `fiberless` | 0.04 us | 0.04 us | 0.05 us | +| `drains` | 0.57 us | 0.70 us | 0.74 us | +| `park` | skipped | 158.9 us | 145.9 us | + +The harness runs inside one long-lived async loop, so its counter deltas can show a `redispatches` count where a real app would see a new drain. Use the per-op times from `bench.js` and the counters from `app_stats.mjs`. + +## Sampling a running app + +From a Meteor shell (or another REPL in the app process). If `fibers/...` doesn't resolve from the shell, import the file by absolute path instead: + +```js +const s = await import('fibers/fiber-await/app_stats.mjs'); +s.start(); // then use the app +s.stop(); // counts and percentages for the window +s.measureSwitch(); // { roundTripUs } in this process; blocks the event loop briefly +s.stop({ roundTripUs }) // after another start(): also estimates the time spent switching +``` + +`start()` counts every promise callback with a `v8.promiseHooks` `before` hook (one JS call per promise job), so stop it when you're done. Non-promise `queueMicrotask` callbacks are not counted; they're never tagged. + +Reference (local qualia, one user clicking for 88 s, 2026-10-02): 13.7% of promise callbacks tagged, 3.2% of untagged callbacks caused a switch, 1.9% of drains used the dispatch fiber, 71% of tagged callbacks batched; 4,738 switches ≈ 0.1–0.25% of a core at 20–45 us each. Idle: 6.9% tagged, 0.29% of untagged caused a switch. diff --git a/fiber-await/app_stats.mjs b/fiber-await/app_stats.mjs new file mode 100644 index 0000000..c929222 --- /dev/null +++ b/fiber-await/app_stats.mjs @@ -0,0 +1,84 @@ +// Sample await-dispatch activity in a running app. See README.md. +// +// From a Meteor shell (or any REPL inside the app process): +// const s = await import('fibers/fiber-await/app_stats.mjs'); +// s.start(); // ...use the app... +// s.stop(); // counts and percentages for the window +// s.measureSwitch(); // optional: fiber switch cost in this process (blocks ~0.1-0.5s) +// +// start() adds a promise `before` hook that counts every promise callback; it costs a JS call per +// promise job while armed, so leave it off when you're not sampling. +import { createRequire } from 'node:module'; +import { promiseHooks } from 'node:v8'; + +// The native fibers module every copy of fibers in the process shares. Read it rather than loading +// fibers_async.js, which would register a dispatcher of its own. +const Fiber = process.fiberLib || createRequire(import.meta.url)('../fibers_sync.js'); + +let armed = null; + +function snapshot() { + return { + at: process.hrtime.bigint(), + drains: Fiber.microtaskDrains, + drainsWithDispatch: Fiber.drainsWithDispatch, + dispatched: Fiber.microtasksDispatched, + redispatches: Fiber.redispatches, + batched: Fiber.__microtasksBatched || 0, + }; +} + +export function start() { + if (Fiber.microtasksDispatched === undefined || Fiber.redispatches === undefined) { + throw new Error('this fibers build has no dispatch counters (needs the fiber-await measurement branch)'); + } + if (armed) armed.stopHook(); + const state = { callbacks: 0 }; + state.stopHook = promiseHooks.onBefore(() => { state.callbacks++; }); + state.start = snapshot(); + armed = state; + return 'armed'; +} + +export function stop({ roundTripUs } = {}) { + if (!armed) throw new Error('call start() first'); + const end = snapshot(); + armed.stopHook(); + const { start: begin, callbacks } = armed; + armed = null; + const d = (k) => end[k] - begin[k]; + const pct = (part, whole) => (whole ? `${(100 * part / whole).toFixed(2)}%` : 'n/a'); + const tagged = d('dispatched') + d('batched'); + const untagged = callbacks - tagged; + const result = { + seconds: Number(end.at - begin.at) / 1e9, + promiseCallbacks: callbacks, + tagged, + taggedShare: pct(tagged, callbacks), + untagged, + // Each redispatch is a gap of untagged callbacks between two tagged ones in the same drain. + untaggedThatCausedASwitch: d('redispatches'), + untaggedThatCausedASwitchShare: pct(d('redispatches'), untagged), + drains: d('drains'), + drainsWithDispatch: d('drainsWithDispatch'), + drainsWithDispatchShare: pct(d('drainsWithDispatch'), d('drains')), + switchesOntoDispatchFiber: d('dispatched'), + batchedWithoutSwitch: d('batched'), + batchedShareOfTagged: pct(d('batched'), tagged), + }; + if (roundTripUs) { + result.estimatedSwitchMs = Number((d('dispatched') * roundTripUs / 1000).toFixed(1)); + result.estimatedRedispatchMs = Number((d('redispatches') * roundTripUs / 1000).toFixed(1)); + } + return result; +} + +// Cost of one fiber run+yield round trip in this process (the cost of each switch onto the dispatch +// fiber and back). Blocks the event loop while it runs. +export function measureSwitch(n = 5000) { + const fiber = Fiber(() => { for (;;) Fiber.yield(); }); + fiber.run(); + const begin = process.hrtime.bigint(); + for (let i = 0; i < n; i++) fiber.run(); + return { roundTripUs: Number((Number(process.hrtime.bigint() - begin) / 1e3 / n).toFixed(2)) }; +} diff --git a/fiber-await/bench.js b/fiber-await/bench.js new file mode 100644 index 0000000..d8e4ff4 --- /dev/null +++ b/fiber-await/bench.js @@ -0,0 +1,133 @@ +// Microbenchmarks for running await continuations on fibers. See README.md. +// Run: node --no-deprecation fiber-await/bench.js (all) +// BENCH=switch,io node --no-deprecation fiber-await/bench.js +// Compare the same binary with FIBERS_AWAIT_REUSE=0 or FIBERS_AWAIT_DISPATCH=0, and with stock node. +const path = require('path'); +const Fiber = require(path.join(__dirname, '..')); + +Fiber.poolSize = 1e9; // production setting; also avoids the pthread coro_destroy crash in test/pool.js + +// The fiber code after an await runs on, whether or not Fiber.current shows it. +const runningFiber = () => Fiber.currentIncludingHidden || Fiber.current; +const dispatchAvailable = typeof Fiber.__setMicrotaskDispatcher === 'function'; +const now = () => process.hrtime.bigint(); +const msSince = (start) => Number(now() - start) / 1e6; +const nextTurn = () => new Promise((resolve) => setImmediate(resolve)); + +function inFiber(fn) { + return new Promise((resolve, reject) => Fiber(() => { fn().then(resolve, reject); }).run()); +} + +function counters() { + return { + dispatched: Fiber.microtasksDispatched || 0, + batched: Fiber.microtasksBatched || 0, + redispatches: Fiber.redispatches || 0, + }; +} + +const benches = { + // Two fiber switches (run + yield), with fibers' async-hooks stack save/restore. The floor for any + // design that moves work onto a fiber; works on stock node + fibers. + async switch() { + const n = 50000; + const fiber = Fiber(() => { for (;;) Fiber.yield(); }); + fiber.run(); + const start = now(); + for (let i = 0; i < n; i++) fiber.run(); + return { n, ms: msSince(start), unit: 'run+yield round trip' }; + }, + + // A chain of awaits in a fiber with nothing else queued: batching's best case. + async chain() { + const n = 100000; + let ms; + await inFiber(async () => { + const start = now(); + for (let i = 0; i < n; i++) await null; + ms = msSince(start); + }); + return { n, ms, unit: 'await in a fiber' }; + }, + + // Each await resumes from its own macrotask, like a DB call: one dispatch per await. + async io() { + const n = 20000; + let ms; + await inFiber(async () => { + const start = now(); + for (let i = 0; i < n; i++) await nextTurn(); + ms = msSince(start); + }); + return { n, ms, unit: 'I/O-style await in a fiber' }; + }, + + // A fibered chain alongside a fiberless one, so tagged and untagged jobs alternate. + async interleaved() { + const n = 50000; + const chain = async () => { for (let i = 0; i < n; i++) await null; }; + const start = now(); + await Promise.all([inFiber(chain), chain()]); + return { n, ms: msSince(start), unit: 'await in a fiber, interleaved with fiberless' }; + }, + + // Awaits in code that never had a fiber: should cost the same as stock. + async fiberless() { + const n = 1000000; + const start = now(); + for (let i = 0; i < n; i++) await null; + return { n, ms: msSince(start), unit: 'await outside any fiber' }; + }, + + // One non-empty microtask drain per macrotask, in fiberless code (znewsham's drain benchmark). + async drains() { + const n = 200000; + const start = now(); + await new Promise((resolve) => { + let i = 0; + (function next() { + if (++i === n) return resolve(); + Promise.resolve().then(() => setImmediate(next)); + })(); + }); + return { n, ms: msSince(start), unit: 'macrotask with a microtask drain' }; + }, + + // Park the fiber after an await and resume it from the next macrotask, 1000 at a time. + async park() { + if (!dispatchAvailable) return { skipped: 'needs the fiber-await node + fibers' }; + const n = 20000; + const batch = 1000; + const start = now(); + for (let i = 0; i < n; i += batch) { + await Promise.all(Array.from({ length: batch }, () => inFiber(async () => { + await null; + const fiber = runningFiber(); + setImmediate(() => fiber.run()); + Fiber.yield(); + }))); + } + return { n, ms: msSince(start), unit: 'park/resume cycle after await' }; + }, +}; + +(async () => { + const selected = process.env.BENCH ? process.env.BENCH.split(',') : Object.keys(benches); + console.log(`node ${process.version}; dispatch API: ${dispatchAvailable}; ` + + `FIBERS_AWAIT_DISPATCH=${process.env.FIBERS_AWAIT_DISPATCH ?? '(unset)'} ` + + `FIBERS_AWAIT_REUSE=${process.env.FIBERS_AWAIT_REUSE ?? '(unset)'}`); + for (const name of selected) { + await benches[name](); // warm up + const before = counters(); + const result = await benches[name](); + const after = counters(); + if (result.skipped) { + console.log(`${name.padEnd(12)} skipped: ${result.skipped}`); + continue; + } + const delta = Object.fromEntries(Object.keys(after).map((k) => [k, after[k] - before[k]])); + const us = (result.ms * 1000 / result.n).toFixed(2); + console.log(`${name.padEnd(12)} ${us.padStart(8)} us per ${result.unit}` + + (dispatchAvailable ? ` ${JSON.stringify(delta)}` : '')); + } +})(); diff --git a/fiber-await/test.js b/fiber-await/test.js new file mode 100644 index 0000000..b31fd67 --- /dev/null +++ b/fiber-await/test.js @@ -0,0 +1,116 @@ +// Behaviour checks for running await continuations on fibers. See README.md. +// Run: node fiber-await/test.js (dispatch on) +// FIBERS_AWAIT_REUSE=0 node fiber-await/test.js (one new fiber per job) +// FIBERS_AWAIT_DISPATCH=0 node fiber-await/test.js (baseline: the fiber checks fail, as on stock) +const assert = require('assert'); +const { AsyncLocalStorage } = require('async_hooks'); +const FIBERS = require('path').join(__dirname, '..'); +const Fiber = require(FIBERS); +const Future = require(`${FIBERS}/future`); + +const als = new AsyncLocalStorage(); + +// A fiber-only API, like Orders.find().count(): blocks the current fiber. +function sleepSync(ms) { + const fut = new Future(); + setTimeout(() => fut.return(ms), ms); + return fut.wait(); +} +const tick = (ms = 2) => new Promise((resolve) => setTimeout(resolve, ms)); +// The fiber code after an await runs on, whether or not Fiber.current shows it. +const anyFiber = () => Fiber.currentIncludingHidden || Fiber.current; + +// A dual-mode API, like collection-hooks' cursor.observeChanges: sync in a fiber, Promise otherwise. +function dualMode() { + const p = tick().then(() => 'handle'); + return Fiber.current ? Future.fromPromise(p).wait() : p; +} + +// Run fn in a fiber; resolve with what its returned promise resolves to. +function inFiber(fn) { + return new Promise((resolve, reject) => { + Fiber(() => { Promise.resolve(fn()).then(resolve, reject); }).run(); + }); +} + +const tests = [ + ['JordanTest shape: await, then fiber-only call', () => inFiber(async () => { + await tick(); + assert.ok(anyFiber(), 'no fiber after await'); + assert.strictEqual(sleepSync(3), 3); + })], + ['several awaits, fiber-only call after each', () => inFiber(async () => { + for (let i = 0; i < 5; i++) { + await tick(1); + assert.ok(anyFiber(), `no fiber after await #${i}`); + sleepSync(1); + } + })], + ['nested async function called from a fiber', () => inFiber(async () => { + const inner = async () => { await tick(); sleepSync(1); return 'inner'; }; + assert.strictEqual(await inner(), 'inner'); + sleepSync(1); + })], + ['.then callback registered in a fiber runs in a fiber', () => inFiber(async () => { + const seen = await tick().then(() => !!anyFiber()); + assert.ok(seen, '.then callback ran without a fiber'); + })], + ['AsyncLocalStorage store survives the hop', () => als.run({ user: 'u1' }, () => inFiber(async () => { + await tick(); + assert.deepStrictEqual(als.getStore(), { user: 'u1' }); + sleepSync(1); + assert.deepStrictEqual(als.getStore(), { user: 'u1' }); + }))], + ['concurrent async calls still overlap (Promise.all)', () => inFiber(async () => { + const start = Date.now(); + const work = async () => { await tick(1); sleepSync(40); }; + await Promise.all([work(), work(), work()]); + const elapsed = Date.now() - start; + assert.ok(elapsed < 100, `expected overlap, took ${elapsed}ms`); + })], + ['errors after await reject normally', () => inFiber(async () => { + const bad = async () => { await tick(); sleepSync(1); throw new Error('boom'); }; + await assert.rejects(bad(), /boom/); + })], + ['Fiber.current still reads undefined after await (hidden fiber)', () => inFiber(async () => { + assert.ok(Fiber.current, 'expected a visible fiber before the first await'); + await tick(); + assert.strictEqual(Fiber.current, undefined); + })], + ['dual-mode API still returns a Promise after await', () => inFiber(async () => { + assert.strictEqual(dualMode(), 'handle', 'expected sync result before the first await'); + await tick(); + const result = dualMode(); + assert.strictEqual(typeof result.catch, 'function', 'got a sync result after await'); + assert.strictEqual(await result, 'handle'); + })], + ['fiberless code is untouched', async () => { + assert.ok(!Fiber.current); + await tick(); + assert.ok(!anyFiber(), 'fiberless continuation got a fiber'); + }], + ['microtask order unchanged for non-yielding continuations', () => inFiber(async () => { + const order = []; + const a = (async () => { await null; order.push('a1'); await null; order.push('a2'); })(); + const b = (async () => { await null; order.push('b1'); await null; order.push('b2'); })(); + await Promise.all([a, b]); + assert.deepStrictEqual(order, ['a1', 'b1', 'a2', 'b2']); + })], +]; + +(async () => { + console.log(`node ${process.version}, dispatch API: ${typeof Fiber.__setMicrotaskDispatcher === 'function'}, FIBERS_AWAIT_DISPATCH=${process.env.FIBERS_AWAIT_DISPATCH ?? '(unset)'}`); + let failed = 0; + for (const [name, fn] of tests) { + try { + await fn(); + console.log(`PASS ${name}`); + } catch (e) { + failed++; + console.log(`FAIL ${name}: ${e.message.split('\n')[0]}`); + } + } + console.log(`microtasksDispatched=${Fiber.microtasksDispatched}, fibersCreated=${Fiber.fibersCreated}`); + console.log(failed ? `${failed} failed` : 'all passed'); + process.exitCode = failed ? 1 : 0; +})(); diff --git a/fibers_async.js b/fibers_async.js index b58f09a..f0c0541 100644 --- a/fibers_async.js +++ b/fibers_async.js @@ -72,7 +72,8 @@ function dispatchOnNewFiber(runDispatched, token) { // one switch onto the fiber and one back, not two per `await`. A job that blocks keeps the fiber // (it's parked inside that job), and the next dispatch starts a new one. FIBERS_AWAIT_REUSE=0 turns // this off. -let microtasksBatched = 0; +// Kept on the native module, which every copy of this file shares (an app can load more than one). +_Fiber.__microtasksBatched = _Fiber.__microtasksBatched || 0; function reusingDispatcher(runNextDispatchable) { const kIdle = {}; let dispatchFiber = null; @@ -85,7 +86,7 @@ function reusingDispatcher(runNextDispatchable) { // Only while this is still the dispatch fiber, i.e. we're nested inside the microtask loop. // A fiber that was parked by a job and later resumed must not pull in unrelated jobs. while (self === dispatchFiber && runNextDispatchable()) { - microtasksBatched++; + _Fiber.__microtasksBatched++; } if (self !== dispatchFiber) { return; @@ -115,7 +116,7 @@ function reusingDispatcher(runNextDispatchable) { Object.defineProperty(Fiber, 'microtasksBatched', { get() { - return microtasksBatched; + return _Fiber.__microtasksBatched || 0; } }) diff --git a/src/fibers.cc b/src/fibers.cc index 0fc8350..3e7274c 100644 --- a/src/fibers.cc +++ b/src/fibers.cc @@ -441,6 +441,12 @@ class Fiber { static Persistent run_dispatched; static Persistent token_template; static double microtasks_dispatched; + // Diagnostics: how often a drain switches onto a dispatch fiber more than once (an untagged + // job between two tagged ones sends us back to the main stack). + static double microtask_drains; + static double drains_with_dispatch; + static double redispatches; + static double last_dispatch_drain; #endif Isolate* isolate; @@ -938,9 +944,32 @@ class Fiber { return false; } ++microtasks_dispatched; + if (last_dispatch_drain == microtask_drains) { + ++redispatches; + } else { + ++drains_with_dispatch; + last_dispatch_drain = microtask_drains; + } return true; } + // Runs after every checkpoint of the isolate's default microtask queue. + static void OnMicrotasksCompleted(Isolate* isolate, void* data) { + ++microtask_drains; + } + + static uni::FunctionType GetMicrotaskDrains(Local property, const uni::GetterCallbackInfo& info) { + return uni::Return(uni::NewNumber(Isolate::GetCurrent(), microtask_drains), info); + } + + static uni::FunctionType GetDrainsWithDispatch(Local property, const uni::GetterCallbackInfo& info) { + return uni::Return(uni::NewNumber(Isolate::GetCurrent(), drains_with_dispatch), info); + } + + static uni::FunctionType GetRedispatches(Local property, const uni::GetterCallbackInfo& info) { + return uni::Return(uni::NewNumber(Isolate::GetCurrent(), redispatches), info); + } + /** * runDispatched(token): runs the job behind `token` on the current stack. Only the first call * does anything. @@ -970,6 +999,7 @@ class Fiber { } microtask_dispatcher.Reset(isolate, args[0].As()); set_microtask_dispatch_callback(isolate, DispatchMicrotask, NULL); + isolate->AddMicrotasksCompletedCallback(OnMicrotasksCompleted, NULL); return uni::Return(uni::Undefined(isolate), args); } @@ -1047,6 +1077,9 @@ class Fiber { fn->Set(context, uni::NewLatin1Symbol(isolate, "__setMicrotaskDispatcher"), uni::GetFunction(uni::NewFunctionTemplate(isolate, SetMicrotaskDispatcher))).FromJust(); uni::SetAccessor(isolate, fn, uni::NewLatin1Symbol(isolate, "microtasksDispatched"), GetMicrotasksDispatched); + uni::SetAccessor(isolate, fn, uni::NewLatin1Symbol(isolate, "microtaskDrains"), GetMicrotaskDrains); + uni::SetAccessor(isolate, fn, uni::NewLatin1Symbol(isolate, "drainsWithDispatch"), GetDrainsWithDispatch); + uni::SetAccessor(isolate, fn, uni::NewLatin1Symbol(isolate, "redispatches"), GetRedispatches); run_next_dispatchable_microtask = (RunNextDispatchableMicrotaskFn)dlsym(RTLD_DEFAULT, "v8_qualia_RunNextDispatchableMicrotask"); if (run_next_dispatchable_microtask) { fn->Set(context, uni::NewLatin1Symbol(isolate, "__runNextDispatchable"), @@ -1073,6 +1106,10 @@ Persistent Fiber::microtask_dispatcher; Persistent Fiber::run_dispatched; Persistent Fiber::token_template; double Fiber::microtasks_dispatched = 0; +double Fiber::microtask_drains = 0; +double Fiber::drains_with_dispatch = 0; +double Fiber::redispatches = 0; +double Fiber::last_dispatch_drain = -1; #endif bool did_init = false;