diff --git a/fibers_async.js b/fibers_async.js index 0fbf1cc..b58f09a 100644 --- a/fibers_async.js +++ b/fibers_async.js @@ -51,15 +51,72 @@ _Fiber[Symbol.hasInstance] = function(obj) { }; // Prototype (Qualia): with a patched node, V8 offers us every promise reaction (e.g. the code after -// an `await`) that was registered inside a fiber. Run it in a new, hidden fiber so it can still -// block. The job's own promise hooks restore its async context inside the fiber. +// an `await`) that was registered inside a fiber. Run it in a hidden fiber so it can still block. +// The job's own promise hooks restore its async context inside the fiber. // FIBERS_AWAIT_DISPATCH=0 turns this off. if (typeof _Fiber.__setMicrotaskDispatcher === 'function' && process.env.FIBERS_AWAIT_DISPATCH !== '0') { - _Fiber.__setMicrotaskDispatcher(function dispatchMicrotask(runDispatched, token) { - const fiber = Fiber(() => runDispatched(token)); - fiber[kHidden] = true; - fiber.run(); - }); + const runNextDispatchable = process.env.FIBERS_AWAIT_REUSE !== '0' ? _Fiber.__runNextDispatchable : undefined; + _Fiber.__setMicrotaskDispatcher(typeof runNextDispatchable === 'function' + ? reusingDispatcher(runNextDispatchable) + : dispatchOnNewFiber); +} + +function dispatchOnNewFiber(runDispatched, token) { + const fiber = Fiber(() => runDispatched(token)); + fiber[kHidden] = true; + fiber.run(); +} + +// Runs jobs on one long-lived hidden fiber instead of a new fiber per job. After each job it keeps +// running the dispatchable jobs queued right behind it, so a run of `await`s in fibered code costs +// 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; +function reusingDispatcher(runNextDispatchable) { + const kIdle = {}; + let dispatchFiber = null; + let runDispatched; + + function dispatchLoop(token) { + const self = _Fiber.current; + for (;;) { + runDispatched(token); + // 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++; + } + if (self !== dispatchFiber) { + return; + } + token = _Fiber.yield(kIdle); + } + } + + return function dispatchMicrotask(run, token) { + runDispatched = run; + if (dispatchFiber && dispatchFiber === _Fiber.current) { + // A nested microtask drain on the dispatch fiber itself; we can't switch into a running fiber. + dispatchOnNewFiber(run, token); + return; + } + if (!dispatchFiber) { + dispatchFiber = Fiber(dispatchLoop); + dispatchFiber[kHidden] = true; + } + const fiber = dispatchFiber; + if (fiber.run(token) !== kIdle && dispatchFiber === fiber) { + // A job blocked on this fiber; it belongs to that job now. + dispatchFiber = null; + } + }; } +Object.defineProperty(Fiber, 'microtasksBatched', { + get() { + return microtasksBatched; + } +}) + module.exports = _Fiber.Fiber = Fiber; diff --git a/src/fibers.cc b/src/fibers.cc index ca86aaa..0fc8350 100644 --- a/src/fibers.cc +++ b/src/fibers.cc @@ -21,8 +21,10 @@ // Looked up at load time so this binary still works (with the feature off) on unpatched node. typedef void (*SetMicrotaskDispatchCallbackFn)(v8::Isolate*, v8::MicrotaskDispatchCallback, void*); typedef void (*RunDispatchedMicrotaskFn)(v8::Isolate*, v8::DispatchedMicrotask*); +typedef bool (*RunNextDispatchableMicrotaskFn)(v8::Isolate*); static SetMicrotaskDispatchCallbackFn set_microtask_dispatch_callback = NULL; static RunDispatchedMicrotaskFn run_dispatched_microtask = NULL; +static RunNextDispatchableMicrotaskFn run_next_dispatchable_microtask = NULL; #endif #define THROW(x, m) return uni::Return(uni::ThrowException(Isolate::GetCurrent(), x(uni::NewLatin1String(Isolate::GetCurrent(), m))), args) @@ -971,6 +973,17 @@ class Fiber { return uni::Return(uni::Undefined(isolate), args); } + /** + * Fiber.__runNextDispatchable(): if the next queued microtask is also one that would be + * dispatched, run it on the current stack and return true. Lets the dispatch fiber run a + * run of such jobs without switching stacks for each one. + */ + static uni::FunctionType RunNextDispatchable(const uni::Arguments& args) { + Isolate* isolate = args.GetIsolate(); + bool ran = run_next_dispatchable_microtask(isolate); + return uni::Return(uni::NewBoolean(isolate, ran), args); + } + static uni::FunctionType GetMicrotasksDispatched(Local property, const uni::GetterCallbackInfo& info) { return uni::Return(uni::NewNumber(Isolate::GetCurrent(), microtasks_dispatched), info); } @@ -1034,6 +1047,11 @@ 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); + run_next_dispatchable_microtask = (RunNextDispatchableMicrotaskFn)dlsym(RTLD_DEFAULT, "v8_qualia_RunNextDispatchableMicrotask"); + if (run_next_dispatchable_microtask) { + fn->Set(context, uni::NewLatin1Symbol(isolate, "__runNextDispatchable"), + uni::GetFunction(uni::NewFunctionTemplate(isolate, RunNextDispatchable))).FromJust(); + } } #endif