Skip to content
Closed
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
18 changes: 13 additions & 5 deletions binding.gyp
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,15 @@
['OS == "linux"',
{
'cflags_c': [ '-std=gnu11' ],
'defines': ['CORO_PTHREAD'],
'variables': {
'USE_MUSL': '<!(ldd --version 2>&1 | head -n1 | grep "musl" | wc -l)',
},
'conditions': [
['<(USE_MUSL) == 1',
{'defines': ['CORO_ASM', '__MUSL__']},
{'defines': ['CORO_UCONTEXT']}
],
],
},
],
['OS == "solaris" or OS == "sunos" or OS == "freebsd" or OS == "aix"', {'defines': ['CORO_UCONTEXT']}],
Expand All @@ -50,15 +58,15 @@
['target_arch == "arm"',
{
# There's been problems getting real fibers working on arm
'defines': ['CORO_PTHREAD'],
'defines!': ['CORO_UCONTEXT', 'CORO_SJLJ', 'CORO_ASM'],
'defines': ['CORO_UCONTEXT', '_XOPEN_SOURCE'],
'defines!': ['CORO_PTHREAD', 'CORO_SJLJ', 'CORO_ASM'],
},
],
['target_arch == "arm64"',
{
# There's been problems getting real fibers working on arm
'defines': ['CORO_PTHREAD'],
'defines!': ['CORO_UCONTEXT', 'CORO_SJLJ', 'CORO_ASM'],
'defines': ['CORO_UCONTEXT', '_XOPEN_SOURCE'],
'defines!': ['CORO_PTHREAD', 'CORO_SJLJ', 'CORO_ASM'],
},
],
],
Expand Down
3 changes: 2 additions & 1 deletion fibers_async.js
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,8 @@ function Fiber(fn, ...args) {
const ar = new AsyncResource('Fiber');
const actualFn = (...args1) => ar.runInAsyncScope(() => {
Fiber.current._meteor_dynamics = undefined;
fn(...args1);
// return the fiber function's value so run() resolves to it when the fiber finishes (fibers README semantics)
return fn(...args1);
});
const _fiber = _Fiber(actualFn, ...args);
asyncResourceWeakMap.set(_fiber, ar);
Expand Down
204 changes: 156 additions & 48 deletions src/coroutine.cc
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,11 @@
#endif

#include <stdexcept>
#include <cstdio>
#include <cstdlib>
#ifndef WINDOWS
#include <dlfcn.h>
#endif
#include <stack>
#include <vector>
using namespace std;
Expand All @@ -28,6 +33,18 @@ static pthread_key_t isolate_key = 0x7777;
static pthread_key_t thread_id_key = 0x7777;
static pthread_key_t thread_data_key = 0x7777;

/**
* Qualia's node build exports `v8_qualia_set_thread_stack_start(void*)`. V8 >= 12 records the
* start of the OS thread's stack when an isolate is entered and cppgc conservatively scans from
* the current stack pointer up to it during GC; on a coroutine stack that range is garbage. When
* the symbol exists we tell V8 which stack is running on every switch (nullptr = the thread's own
* stack). Resolved with dlsym so this addon still loads on a node without the patch.
*/
typedef void (*set_thread_stack_start_t)(void*);
static set_thread_stack_start_t set_thread_stack_start = NULL;

static void notify_stack_start(const Coroutine& next);

static size_t stack_size = 0;
static size_t coroutines_created_ = 0;
static vector<Coroutine*> fiber_pool;
Expand Down Expand Up @@ -75,66 +92,129 @@ namespace v8 {
}
#endif

struct tls_snapshot_t {
v8::Isolate* isolate;
std::vector<void*> values;
};

/**
* Runs on a fresh helper thread: enter the isolate (which makes V8 assign this thread a new
* ThreadId and thread data) and copy every pthread TLS slot below `coro_thread_key`.
*/
#ifndef WINDOWS
static void* find_thread_id_key(void* arg)
static void* snapshot_tls(void* arg)
#else
static DWORD __stdcall find_thread_id_key(LPVOID arg)
static DWORD __stdcall snapshot_tls(LPVOID arg)
#endif
{
v8::Isolate* isolate = static_cast<v8::Isolate*>(arg);
assert(isolate != NULL);
v8::Locker locker(isolate);
isolate->Enter();

// First pass-- find isolate thread key
tls_snapshot_t* snap = static_cast<tls_snapshot_t*>(arg);
assert(snap->isolate != NULL);
v8::Locker locker(snap->isolate);
snap->isolate->Enter();
#ifdef __MUSL__
// 128 is default max key in musl
for (pthread_key_t ii = 1; ii < 128; ++ii) {
const pthread_key_t key_count = 128;
#else
for (pthread_key_t ii = coro_thread_key; ii > 0; --ii) {
const pthread_key_t key_count = coro_thread_key;
#endif
void* tls = pthread_getspecific(ii - 1);
if (tls == isolate) {
snap->values.assign(key_count, NULL);
for (pthread_key_t ii = 0; ii < key_count; ++ii) {
snap->values[ii] = pthread_getspecific(ii);
}
snap->isolate->Exit();
return NULL;
}

static void take_tls_snapshot(v8::Isolate* isolate, tls_snapshot_t& snap) {
snap.isolate = isolate;
pthread_t thread;
pthread_create(&thread, NULL, snapshot_tls, &snap);
pthread_join(thread, NULL);
}

/**
* Locate the pthread TLS keys V8 uses for the current isolate, per-isolate thread data and
* ThreadId. Up to V8 12 all three were pthread keys and can be found by value (the isolate
* pointer, a struct whose first word is the isolate pointer, and the int thread id read from
* that struct). Newer V8 keeps the isolate and thread data in compiler thread_local storage
* (which Locker/Unlocker maintain for us), so only the ThreadId key matters: it is the slot
* holding a small positive int that increments between two consecutively created threads.
*/
static void find_thread_id_key(v8::Isolate* isolate) {
tls_snapshot_t a;
take_tls_snapshot(isolate, a);
const pthread_key_t key_count = a.values.size();

// First pass-- find isolate thread key
for (pthread_key_t ii = key_count; ii > 0; --ii) {
if (a.values[ii - 1] == isolate) {
isolate_key = ii - 1;
break;
}
}
assert(isolate_key != 0x7777);

// Second pass-- find data key
int thread_id = 0;
#ifdef __MUSL__
for (pthread_key_t ii = 0; ii < 128; ++ii) {
#else
for (pthread_key_t ii = isolate_key + 1; ii < coro_thread_key; ++ii) {
#endif
void* tls = pthread_getspecific(ii);
if (can_poke(tls) && *(void**)tls == isolate) {
// First member of per-thread data is the isolate
thread_data_key = ii;
// Second member is the thread id
thread_id = *(int*)((void**)tls + 1);
break;
if (isolate_key != 0x7777) {
for (pthread_key_t ii = isolate_key + 1; ii < key_count; ++ii) {
void* tls = a.values[ii];
if (can_poke(tls) && *(void**)tls == isolate) {
// First member of per-thread data is the isolate
thread_data_key = ii;
// Second member is the thread id
thread_id = *(int*)((void**)tls + 1);
break;
}
}
}
assert(thread_data_key != 0x7777);

// Third pass-- find thread id key
#ifdef __MUSL__
for (pthread_key_t ii = 0; ii < 128; ++ii) {
#else
for (pthread_key_t ii = isolate_key + 1; ii < coro_thread_key; ++ii) {
#endif
int tls = static_cast<int>(reinterpret_cast<intptr_t>(pthread_getspecific(ii)));
if (tls == thread_id) {
thread_id_key = ii;
break;
if (thread_data_key != 0x7777) {
for (pthread_key_t ii = isolate_key + 1; ii < key_count; ++ii) {
int tls = static_cast<int>(reinterpret_cast<intptr_t>(a.values[ii]));
if (tls == thread_id) {
thread_id_key = ii;
break;
}
}
}
assert(thread_id_key != 0x7777);

isolate->Exit();
return NULL;
if (thread_id_key == 0x7777) {
// Fallback for V8 >= 13: compare against further helper threads. V8 hands out thread ids
// from an atomic counter, so the ThreadId slot is the one whose small positive value grows
// between two consecutively created threads. Other threads (V8 platform workers) may be
// created concurrently and consume ids, so allow a small gap, require the match to be
// unique, and retry with a fresh snapshot pair when it is not.
for (int attempt = 0; attempt < 8 && thread_id_key == 0x7777; ++attempt) {
tls_snapshot_t b;
take_tls_snapshot(isolate, b);
pthread_key_t candidate = 0x7777;
int candidates = 0;
for (pthread_key_t ii = 0; ii < key_count && ii < b.values.size(); ++ii) {
intptr_t va = reinterpret_cast<intptr_t>(a.values[ii]);
intptr_t vb = reinterpret_cast<intptr_t>(b.values[ii]);
if (va > 0 && va < (1 << 24) && vb > va && vb - va <= 64) {
candidate = ii;
++candidates;
}
}
if (candidates == 1) {
thread_id_key = candidate;
}
a = b;
}
}
if (thread_id_key == 0x7777) {
// Without this key every coroutine shares the OS thread's V8 ThreadId and Locker/Unlocker
// archiving silently corrupts JS stacks. Refuse to run rather than fail intermittently
// later (an assert would be compiled out of Release builds).
fprintf(stderr, "fibers: could not locate V8's ThreadId thread-local key; this node/V8 build is not supported (set FIBERS_DEBUG_TLS=1 for details)\n");
abort();
}
if (getenv("FIBERS_DEBUG_TLS")) {
fprintf(stderr, "fibers: v8 tls keys isolate=%d thread_data=%d thread_id=%d (0x7777 = not found)\n",
(int)isolate_key, (int)thread_data_key, (int)thread_id_key);
}
}

/**
Expand All @@ -149,9 +229,13 @@ void Coroutine::init(v8::Isolate* isolate) {
thread_data_key = v8::internal::Isolate::per_isolate_thread_data_key_;
thread_id_key = v8::internal::Isolate::thread_id_key_;
#elif !defined(CORO_PTHREAD)
pthread_t thread;
pthread_create(&thread, NULL, find_thread_id_key, isolate);
pthread_join(thread, NULL);
find_thread_id_key(isolate);
#endif
#if !defined(WINDOWS) && !defined(CORO_PTHREAD)
set_thread_stack_start = reinterpret_cast<set_thread_stack_start_t>(dlsym(RTLD_DEFAULT, "v8_qualia_set_thread_stack_start"));
if (getenv("FIBERS_DEBUG_TLS")) {
fprintf(stderr, "fibers: v8_qualia_set_thread_stack_start %s\n", set_thread_stack_start ? "found" : "not found (GC may scan the wrong stack)");
}
#endif
}

Expand Down Expand Up @@ -253,18 +337,34 @@ void Coroutine::reset(entry_t* entry, void* arg) {
this->arg = arg;
}

static void notify_stack_start(const Coroutine& next) {
if (!set_thread_stack_start) {
return;
}
void* bottom = next.bottom();
// The original thread's Coroutine has no stack of its own: fall back to the OS thread's stack.
set_thread_stack_start(bottom ? static_cast<char*>(bottom) + next.stack_bytes() : NULL);
}

void Coroutine::transfer(Coroutine& next) {
assert(this != &next);
#ifndef CORO_PTHREAD
fls_data[0] = pthread_getspecific(isolate_key);
fls_data[1] = pthread_getspecific(thread_id_key);
fls_data[2] = pthread_getspecific(thread_data_key);

pthread_setspecific(isolate_key, next.fls_data[0]);
pthread_setspecific(thread_id_key, next.fls_data[1]);
pthread_setspecific(thread_data_key, next.fls_data[2]);
// Keys V8 no longer keeps in pthread TLS stay at 0x7777 and are skipped.
if (isolate_key != 0x7777) {
fls_data[0] = pthread_getspecific(isolate_key);
pthread_setspecific(isolate_key, next.fls_data[0]);
}
if (thread_id_key != 0x7777) {
fls_data[1] = pthread_getspecific(thread_id_key);
pthread_setspecific(thread_id_key, next.fls_data[1]);
}
if (thread_data_key != 0x7777) {
fls_data[2] = pthread_getspecific(thread_data_key);
pthread_setspecific(thread_data_key, next.fls_data[2]);
}

pthread_setspecific(coro_thread_key, &next);
notify_stack_start(next);
#endif
coro_transfer(&context, &next.context);
#ifndef CORO_PTHREAD
Expand Down Expand Up @@ -322,6 +422,14 @@ void* Coroutine::bottom() const {
#endif
}

size_t Coroutine::stack_bytes() const {
#ifdef CORO_FIBER
return stack_size * sizeof(void*);
#else
return stack.ssze;
#endif
}

size_t Coroutine::size() const {
return sizeof(Coroutine) + stack_size * sizeof(void*);
}
5 changes: 5 additions & 0 deletions src/coroutine.h
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,11 @@ class Coroutine {
*/
void* bottom() const;

/**
* Size in bytes of this coroutine's stack (0 for the original thread).
*/
size_t stack_bytes() const;

/**
* Returns the size this Coroutine takes up in the heap.
*/
Expand Down
Loading