Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 64 additions & 7 deletions fibers_async.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
18 changes: 18 additions & 0 deletions src/fibers.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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<String> property, const uni::GetterCallbackInfo& info) {
return uni::Return(uni::NewNumber(Isolate::GetCurrent(), microtasks_dispatched), info);
}
Expand Down Expand Up @@ -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

Expand Down