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
84 changes: 83 additions & 1 deletion fibers_async.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ const asyncResourceWeakMap = new WeakMap();
function Fiber(fn, ...args) {
const ar = new AsyncResource('Fiber');
const actualFn = (...args1) => ar.runInAsyncScope(() => {
Fiber.current._meteor_dynamics = undefined;
_Fiber.current._meteor_dynamics = undefined;
fn(...args1);
});
const _fiber = _Fiber(actualFn, ...args);
Expand All @@ -21,7 +21,20 @@ function Fiber(fn, ...args) {
Fiber.__proto__ = _Fiber;
Fiber.prototype = _Fiber.prototype;

// Fibers that run promise reactions (see below) are hidden: `Fiber.current` reports them as
// undefined, exactly what code after an `await` sees today, so `Fiber.current`-based mode switches
// ("in a fiber? return sync, else return a Promise") keep behaving as before. Blocking primitives
// (Future.wait, Promise.await) use `Fiber.currentIncludingHidden` so they can still yield.
const kHidden = Symbol.for('fibers.hiddenAwaitContinuation');

Object.defineProperty(Fiber, 'current', {
get() {
const current = _Fiber.current;
return current && current[kHidden] ? undefined : current;
}
})

Object.defineProperty(Fiber, 'currentIncludingHidden', {
get() {
return _Fiber.current;
}
Expand All @@ -37,4 +50,73 @@ _Fiber[Symbol.hasInstance] = function(obj) {
return obj instanceof Fiber || obj.run;
};

// 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 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') {
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;
3 changes: 2 additions & 1 deletion future.js
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,8 @@ Future.wait = function wait(/* ... */) {
}

// Resumes current fiber
var fiber = Fiber.current;
// Includes the hidden fibers that run code after an `await` (see fibers_async.js).
var fiber = Fiber.currentIncludingHidden || Fiber.current;
if (!fiber) {
throw new Error('Can\'t wait without a fiber. Most likely you called `Promise.await()` after calling `await somePromiseAPI()`. Read https://engdocs.qualia.io/qualia/advanced-topics/fibers');
}
Expand Down
157 changes: 157 additions & 0 deletions src/fibers.cc
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,26 @@
#include <vector>
#include <iostream>

// Prototype (Qualia): patched node exposes v8::SetMicrotaskDispatchCallback so the code after an
// `await` (or a `.then` callback) registered inside a fiber also runs inside a fiber.
#if defined(__has_include)
#if __has_include(<v8-microtask-dispatch.h>)
#include <v8-microtask-dispatch.h>
#include <dlfcn.h>
#define FIBERS_AWAIT_DISPATCH 1
#endif
#endif

#ifdef FIBERS_AWAIT_DISPATCH
// 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)

using namespace std;
Expand Down Expand Up @@ -415,6 +435,13 @@ class Fiber {
static Fiber* current;
static vector<Fiber*> orphaned_fibers;
static Persistent<Value> fatal_stack;
#ifdef FIBERS_AWAIT_DISPATCH
static Persistent<Object> cped_marker;
static Persistent<Function> microtask_dispatcher;
static Persistent<Function> run_dispatched;
static Persistent<ObjectTemplate> token_template;
static double microtasks_dispatched;
#endif

Isolate* isolate;
Persistent<Object> handle;
Expand Down Expand Up @@ -689,6 +716,21 @@ class Fiber {
Fiber* last_fiber = current;
current = this;

#ifdef FIBERS_AWAIT_DISPATCH
// While this fiber runs, the context's continuation-preserved embedder data (CPED) is the
// marker. V8 copies CPED into every promise reaction registered meanwhile, and offers
// those reactions to DispatchMicrotask() instead of running them on whatever stack is
// draining the microtask queue.
bool tag_reactions = !microtask_dispatcher.IsEmpty();
Local<Context> cped_context;
Local<Value> previous_cped;
if (tag_reactions) {
cped_context = Local<Context>::New(isolate, v8_context);
previous_cped = cped_context->GetContinuationPreservedEmbedderData();
cped_context->SetContinuationPreservedEmbedderData(Local<Object>::New(isolate, cped_marker));
}
#endif

// This will jump into either `RunFiber()` or `Yield()`, depending on if the fiber was
// already running.
{
Expand All @@ -697,6 +739,12 @@ class Fiber {
this_fiber->run();
}

#ifdef FIBERS_AWAIT_DISPATCH
if (tag_reactions) {
cped_context->SetContinuationPreservedEmbedderData(previous_cped);
}
#endif

// At this point the fiber either returned or called `yield()`.
current = last_fiber;
}
Expand Down Expand Up @@ -858,6 +906,89 @@ class Fiber {
return uni::Return(uni::NewNumber(Isolate::GetCurrent(), Coroutine::coroutines_created()), info);
}

#ifdef FIBERS_AWAIT_DISPATCH
/**
* Called by V8 for a promise reaction (e.g. the code after an `await`) that was registered
* while a fiber was running. Hands the job to the JS dispatcher, which runs it in a fiber.
* Returns false (V8 runs the job inline) if the dispatcher didn't start it.
*/
static bool DispatchMicrotask(Isolate* isolate, v8::DispatchedMicrotask* task, void* data) {
if (microtask_dispatcher.IsEmpty()) {
return false;
}
HandleScope scope(isolate);
Local<Context> context = isolate->GetCurrentContext();
Local<Object> token;
if (!Local<ObjectTemplate>::New(isolate, token_template)->NewInstance(context).ToLocal(&token)) {
return false;
}
token->SetAlignedPointerInInternalField(0, task);
Local<Value> argv[2] = { Local<Function>::New(isolate, run_dispatched), token };
{
uni::TryCatch try_catch(isolate);
Local<Value> ignored;
if (!Local<Function>::New(isolate, microtask_dispatcher)->Call(context, Undefined(isolate), 2, argv).ToLocal(&ignored) && try_catch.HasCaught()) {
String::Utf8Value message(isolate, try_catch.Exception());
cerr << "fibers: microtask dispatcher threw: " << (*message ? *message : "<unknown>") << "\n";
}
}
if (token->GetAlignedPointerFromInternalField(0) != NULL) {
// Never started: make the token inert and let V8 run the job inline.
token->SetAlignedPointerInInternalField(0, NULL);
return false;
}
++microtasks_dispatched;
return true;
}

/**
* runDispatched(token): runs the job behind `token` on the current stack. Only the first call
* does anything.
*/
static uni::FunctionType RunDispatched(const uni::Arguments& args) {
Isolate* isolate = args.GetIsolate();
if (args.Length() != 1 || !args[0]->IsObject() || args[0].As<Object>()->InternalFieldCount() != 1) {
THROW(Exception::TypeError, "runDispatched expects a dispatch token");
}
Local<Object> token = args[0].As<Object>();
v8::DispatchedMicrotask* task = static_cast<v8::DispatchedMicrotask*>(token->GetAlignedPointerFromInternalField(0));
if (task != NULL) {
token->SetAlignedPointerInInternalField(0, NULL);
run_dispatched_microtask(isolate, task);
}
return uni::Return(uni::Undefined(isolate), args);
}

/**
* Fiber.__setMicrotaskDispatcher(fn): fn(runDispatched, token) must call
* runDispatched(token) synchronously, normally inside a new fiber.
*/
static uni::FunctionType SetMicrotaskDispatcher(const uni::Arguments& args) {
Isolate* isolate = args.GetIsolate();
if (args.Length() != 1 || !args[0]->IsFunction()) {
THROW(Exception::TypeError, "__setMicrotaskDispatcher expects a function");
}
microtask_dispatcher.Reset(isolate, args[0].As<Function>());
set_microtask_dispatch_callback(isolate, DispatchMicrotask, NULL);
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);
}
#endif

public:
/**
* Initialize the Fiber library.
Expand Down Expand Up @@ -904,6 +1035,25 @@ class Fiber {
uni::SetAccessor(isolate, fn, uni::NewLatin1Symbol(isolate, "current"), GetCurrent);
uni::SetAccessor(isolate, fn, uni::NewLatin1Symbol(isolate, "poolSize"), GetPoolSize, SetPoolSize);
uni::SetAccessor(isolate, fn, uni::NewLatin1Symbol(isolate, "fibersCreated"), GetFibersCreated);
#ifdef FIBERS_AWAIT_DISPATCH
set_microtask_dispatch_callback = (SetMicrotaskDispatchCallbackFn)dlsym(RTLD_DEFAULT, "v8_qualia_SetMicrotaskDispatchCallback");
run_dispatched_microtask = (RunDispatchedMicrotaskFn)dlsym(RTLD_DEFAULT, "v8_qualia_RunDispatchedMicrotask");
if (set_microtask_dispatch_callback && run_dispatched_microtask) {
Local<ObjectTemplate> token_tmpl = ObjectTemplate::New(isolate);
token_tmpl->SetInternalFieldCount(1);
token_template.Reset(isolate, token_tmpl);
cped_marker.Reset(isolate, Object::New(isolate));
run_dispatched.Reset(isolate, uni::GetFunction(uni::NewFunctionTemplate(isolate, RunDispatched)));
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

// Global Fiber
target->Set(context, uni::NewLatin1Symbol(isolate, "Fiber"), fn).FromJust();
Expand All @@ -917,6 +1067,13 @@ Locker* Fiber::global_locker;
Fiber* Fiber::current = NULL;
vector<Fiber*> Fiber::orphaned_fibers;
Persistent<Value> Fiber::fatal_stack;
#ifdef FIBERS_AWAIT_DISPATCH
Persistent<Object> Fiber::cped_marker;
Persistent<Function> Fiber::microtask_dispatcher;
Persistent<Function> Fiber::run_dispatched;
Persistent<ObjectTemplate> Fiber::token_template;
double Fiber::microtasks_dispatched = 0;
#endif
bool did_init = false;

#if !NODE_VERSION_AT_LEAST(0,10,0)
Expand Down