Skip to content

Commit 87c98fd

Browse files
authored
Merge pull request #23 from NativeScript/feat/worker-transfer
Transfer lists in worker postMessage, native addons in workers
2 parents fa801ab + d12793c commit 87c98fd

11 files changed

Lines changed: 733 additions & 194 deletions

File tree

‎napi-v8-shim/csrc/env_ext.cpp‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,9 @@ static_assert(sizeof(v8::Local<v8::Context>) == sizeof(void*), "v8::Local must b
1414
// From node_api_types.h, which the vendored headers don't carry.
1515
typedef napi_value (*napi_addon_register_func)(napi_env env, napi_value exports);
1616

17+
// v8-api.cpp.
18+
void ns_napi_finalize_external_arraybuffers(napi_env env);
19+
1720
extern "C" {
1821

1922
// `context` is a `v8::Local<v8::Context>` (as its underlying pointer) of the isolate current on
@@ -53,14 +56,16 @@ bool ns_napi_has_pending_finalizers(napi_env env) {
5356
return !env->pending_finalizers.empty();
5457
}
5558

56-
// Tears the env down: finalizes remaining references (running their finalizers) and frees it.
59+
// Tears the env down: finalizes remaining references and external ArrayBuffers (running their
60+
// finalizers) and frees it.
5761
void ns_napi_env_teardown(napi_env env) {
5862
v8::Isolate* isolate = env->isolate;
5963
v8::HandleScope handle_scope(isolate);
6064
// A real handle, not env->context(): that aliases the env's persistent slot, which DeleteMe
6165
// frees before this scope exits.
6266
v8::Local<v8::Context> context = v8::Local<v8::Context>::New(isolate, env->context());
6367
v8::Context::Scope context_scope(context);
68+
ns_napi_finalize_external_arraybuffers(env);
6469
env->DeleteMe();
6570
}
6671

‎packages/windows-v8/vendor/shim/v8-api.cpp‎

Lines changed: 59 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,10 @@
33
#include <cmath>
44
#include <string_view> // string_view, u16string_view
55
#include <sstream>
6+
#include <mutex>
7+
#include <unordered_map>
8+
#include <unordered_set>
9+
#include <vector>
610

711
#define NAPI_EXPERIMENTAL
812

@@ -2995,6 +2999,41 @@ const v8::ArrayBuffer *v8__ArrayBuffer__New__with_backing_store(
29952999
void std__shared_ptr__v8__BackingStore__reset(void *shared_ptr_ref);
29963000
}
29973001

3002+
namespace {
3003+
// An external ArrayBuffer's finalizer. Its backing store can outlive the env (it goes with the
3004+
// isolate's heap), so the env's teardown runs it instead and clears `cb`, as Node does.
3005+
struct AbFinalize {
3006+
napi_env env;
3007+
napi_finalize cb;
3008+
void *data;
3009+
void *hint;
3010+
};
3011+
3012+
std::mutex ab_finalizers_mutex;
3013+
std::unordered_map<napi_env, std::unordered_set<AbFinalize *>> ab_finalizers;
3014+
}
3015+
3016+
// [windows port] Called by the env's teardown (napi-v8-shim/csrc/env_ext.cpp).
3017+
void ns_napi_finalize_external_arraybuffers(napi_env env) {
3018+
// Copies: once `cb` is cleared, a deleter on another thread may free the record.
3019+
std::vector<AbFinalize> pending;
3020+
{
3021+
std::lock_guard<std::mutex> lock(ab_finalizers_mutex);
3022+
auto found = ab_finalizers.find(env);
3023+
if (found == ab_finalizers.end()) {
3024+
return;
3025+
}
3026+
for (AbFinalize *fd : found->second) {
3027+
pending.push_back(*fd);
3028+
fd->cb = nullptr;
3029+
}
3030+
ab_finalizers.erase(found);
3031+
}
3032+
for (const AbFinalize &fd : pending) {
3033+
fd.cb(env, fd.data, fd.hint);
3034+
}
3035+
}
3036+
29983037
napi_status NAPI_CDECL
29993038
napi_create_external_arraybuffer(napi_env env,
30003039
void *external_data,
@@ -3010,19 +3049,30 @@ napi_create_external_arraybuffer(napi_env env,
30103049
// [windows port] TRUE zero-copy: build a BackingStore that aliases external_data with a deleter
30113050
// that invokes the napi finalizer when the ArrayBuffer is GC'd. Goes through rusty_v8's C
30123051
// bindings (see above) to avoid the libc++ std::unique_ptr ABI boundary.
3013-
struct AbFinalize {
3014-
napi_env env;
3015-
napi_finalize cb;
3016-
void *hint;
3017-
};
30183052
void *deleter_data = nullptr;
3019-
void (*deleter)(void *, size_t, void *) = nullptr;
3053+
// V8 calls the deleter unconditionally when the backing store goes (at the latest when the
3054+
// isolate is disposed, as a terminated worker's is), so it can't be null.
3055+
void (*deleter)(void *, size_t, void *) = [](void *, size_t, void *) {};
30203056
if (finalize_cb != nullptr) {
3021-
deleter_data = new AbFinalize{env, finalize_cb, finalize_hint};
3057+
auto *fd = new AbFinalize{env, finalize_cb, external_data, finalize_hint};
3058+
{
3059+
std::lock_guard<std::mutex> lock(ab_finalizers_mutex);
3060+
ab_finalizers[env].insert(fd);
3061+
}
3062+
deleter_data = fd;
30223063
deleter = [](void *data, size_t, void *dd) {
30233064
auto *fd = static_cast<AbFinalize *>(dd);
3024-
if (fd->cb) {
3025-
fd->cb(fd->env, data, fd->hint);
3065+
napi_finalize cb;
3066+
{
3067+
std::lock_guard<std::mutex> lock(ab_finalizers_mutex);
3068+
cb = fd->cb;
3069+
auto found = ab_finalizers.find(fd->env);
3070+
if (cb && found != ab_finalizers.end()) {
3071+
found->second.erase(fd);
3072+
}
3073+
}
3074+
if (cb) {
3075+
cb(fd->env, data, fd->hint);
30263076
}
30273077
delete fd;
30283078
};

‎runtime/src/animation_frames.rs‎

Lines changed: 47 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,23 +5,66 @@
55
//! once per compositor frame from `CompositionTarget.Rendering`, outside the render walk) then
66
//! runs the queued callbacks once and drains microtasks, which is where rendering work such as
77
//! canvas presents happens. Nothing waits for vsync on the UI thread, and a continuous rAF loop
8-
//! gives the dispatcher back between frames.
8+
//! gives the dispatcher back between frames. A worker's frames follow the UI thread's: each pump
9+
//! there hands a frame to every worker that asked for one. The compositor stops raising frames
10+
//! while nothing on screen changes, so a worker not handed one in time takes its own.
911
1012
use std::cell::Cell;
13+
use std::time::{Duration, Instant};
1114

1215
use crate::DELEGATE_ISOLATE_PTR;
1316

1417
thread_local! {
1518
static REQUESTED: Cell<bool> = const { Cell::new(false) };
19+
static LAST_FRAME: Cell<Option<Instant>> = const { Cell::new(None) };
1620
}
1721

22+
#[cfg(feature = "classic")]
23+
const WORKER_FRAME_INTERVAL: Duration = Duration::from_micros(16_667);
24+
25+
#[cfg(feature = "classic")]
26+
pub(crate) fn worker_frame_wait() -> Option<Duration> {
27+
if !REQUESTED.with(|r| r.get()) {
28+
return None;
29+
}
30+
let last = LAST_FRAME.with(|l| l.get());
31+
Some(last.map_or(Duration::ZERO, |last| WORKER_FRAME_INTERVAL.saturating_sub(last.elapsed())))
32+
}
33+
34+
#[cfg(feature = "classic")]
35+
static WORKERS_REQUESTED: std::sync::Mutex<Vec<usize>> = std::sync::Mutex::new(Vec::new());
36+
1837
/// `__nsRequestFrame()`: run animation callbacks at the next pump.
1938
pub(crate) fn handle_request_frame(
2039
_scope: &mut v8::PinScope<'_, '_>,
2140
_args: v8::FunctionCallbackArguments,
2241
_retval: v8::ReturnValue,
2342
) {
2443
REQUESTED.with(|r| r.set(true));
44+
#[cfg(feature = "classic")]
45+
{
46+
let isolate = DELEGATE_ISOLATE_PTR.with(|c| c.get()) as usize;
47+
if crate::worker_threads::is_worker_isolate(isolate) {
48+
let mut workers = WORKERS_REQUESTED.lock().unwrap_or_else(|e| e.into_inner());
49+
if !workers.contains(&isolate) {
50+
workers.push(isolate);
51+
}
52+
}
53+
}
54+
}
55+
56+
#[cfg(feature = "classic")]
57+
fn hand_frames_to_workers() {
58+
let isolate = DELEGATE_ISOLATE_PTR.with(|c| c.get()) as usize;
59+
if crate::worker_threads::is_worker_isolate(isolate) {
60+
return;
61+
}
62+
let workers = std::mem::take(&mut *WORKERS_REQUESTED.lock().unwrap_or_else(|e| e.into_inner()));
63+
for worker in workers {
64+
crate::worker_threads::post_to_worker(worker, || {
65+
pump();
66+
});
67+
}
2568
}
2669

2770
/// Drops a pending request (the runtime on this thread is going away).
@@ -40,9 +83,12 @@ fn now_ms() -> f64 {
4083
/// Runs this thread's pending animation-frame callbacks, if a frame was requested, then drains
4184
/// microtasks. Returns whether callbacks ran.
4285
pub fn pump() -> bool {
86+
#[cfg(feature = "classic")]
87+
hand_frames_to_workers();
4388
if !REQUESTED.with(|r| r.replace(false)) {
4489
return false;
4590
}
91+
LAST_FRAME.with(|l| l.set(Some(Instant::now())));
4692
let isolate_ptr = DELEGATE_ISOLATE_PTR.with(|c| c.get());
4793
if isolate_ptr.is_null() {
4894
return false;

‎runtime/src/global_fns.rs‎

Lines changed: 52 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ use windows::Win32::UI::WindowsAndMessaging::{
1414

1515
use crate::type_description::build_runtime_type_descriptor;
1616
use crate::dotnet::{bin_write_str16, bin_write_str32};
17-
use crate::{normalize_js_path, proxy_manifests, throw_js_error, try_resolve_with_known_extensions, Runtime, ASYNC_PUMP_HOOK};
17+
use crate::{normalize_js_path, proxy_manifests, throw_js_error, try_resolve_with_known_extensions, ASYNC_PUMP_HOOK};
1818
use std::cell::RefCell;
1919
use std::ffi::c_void;
2020

@@ -209,23 +209,29 @@ fn default_auto_capture_path() -> PathBuf {
209209
PathBuf::from("sbg_output").join("sbg_metadata.json")
210210
}
211211

212+
/// Delivered as a `messageerror`.
213+
fn worker_error<'s>(scope: &mut v8::PinScope<'s, '_>, error: &str) -> Option<v8::Local<'s, v8::Value>> {
214+
let obj = v8::Object::new(scope);
215+
if let Some(key) = v8::String::new(scope, "__workerError") {
216+
if let Some(val) = v8::String::new(scope, error) {
217+
obj.set(scope, key.into(), val.into());
218+
}
219+
}
220+
Some(obj.into())
221+
}
222+
212223
fn polled_event_to_v8<'s>(
213224
scope: &mut v8::PinScope<'s, '_>,
214225
event: crate::worker_threads::PolledWorkerEvent,
215226
) -> Option<v8::Local<'s, v8::Value>> {
216227
match event {
217-
crate::worker_threads::PolledWorkerEvent::Message(bytes) => {
218-
Runtime::deserialize_value(scope, &bytes)
219-
}
220-
crate::worker_threads::PolledWorkerEvent::Error(error) => {
221-
let obj = v8::Object::new(scope);
222-
if let Some(key) = v8::String::new(scope, "__workerError") {
223-
if let Some(val) = v8::String::new(scope, error.as_str()) {
224-
obj.set(scope, key.into(), val.into());
225-
}
228+
crate::worker_threads::PolledWorkerEvent::Message(message) => {
229+
match crate::transfer::deserialize(scope, message) {
230+
Ok(value) => Some(value),
231+
Err(error) => worker_error(scope, &error),
226232
}
227-
Some(obj.into())
228233
}
234+
crate::worker_threads::PolledWorkerEvent::Error(error) => worker_error(scope, &error),
229235
crate::worker_threads::PolledWorkerEvent::Exited => {
230236
let obj = v8::Object::new(scope);
231237
if let Some(key) = v8::String::new(scope, "__workerExit") {
@@ -919,6 +925,11 @@ pub(crate) fn handle_resolve_module_path(
919925
throw_js_error(scope, "__nsResolveModulePath: module specifier is empty");
920926
return;
921927
}
928+
// The app directory, as NativeScript's webpack worker loader names worker chunks.
929+
let specifier = match specifier.strip_prefix("~/") {
930+
Some(rest) => rest.to_string(),
931+
None => specifier,
932+
};
922933
let parent_path = if args.length() >= 2 {
923934
value_to_string(scope, args.get(1))
924935
} else {
@@ -1061,7 +1072,7 @@ pub(crate) fn handle_worker_post_message(
10611072
if args.length() < 2 {
10621073
throw_js_error(
10631074
scope,
1064-
"__nsWorkerPostMessage(workerId, value) expects 2 arguments",
1075+
"__nsWorkerPostMessage(workerId, value, transfer?) expects 2 arguments",
10651076
);
10661077
return;
10671078
}
@@ -1070,16 +1081,38 @@ pub(crate) fn handle_worker_post_message(
10701081
throw_js_error(scope, "Invalid worker id");
10711082
return;
10721083
}
1073-
let value = args.get(1);
1074-
let Some(bytes) = Runtime::serialize_value(scope, value) else {
1075-
throw_js_error(scope, "DataCloneError: value could not be cloned.");
1084+
// On failure the exception is pending.
1085+
let Some(message) = crate::transfer::serialize(scope, args.get(1), args.get(2)) else {
10761086
return;
10771087
};
1078-
if let Err(err) = crate::worker_threads::post_message(worker_id as u64, bytes) {
1088+
if let Err(err) = crate::worker_threads::post_message(worker_id as u64, message) {
10791089
throw_js_error(scope, err.as_str());
10801090
}
10811091
}
10821092

1093+
/// A dispatcher the runtime made itself only runs when its host pumps messages: there a worker's
1094+
/// replies only arrive by polling.
1095+
pub(crate) fn handle_worker_pushes_messages(
1096+
_scope: &mut v8::PinScope<'_, '_>,
1097+
_args: v8::FunctionCallbackArguments,
1098+
mut retval: v8::ReturnValue,
1099+
) {
1100+
let isolate = crate::DELEGATE_ISOLATE_PTR.with(|c| c.get()) as usize;
1101+
let xaml_host = crate::ui_dispatcher::is_initialized() && !crate::ui_dispatcher::needs_win32_pump();
1102+
retval.set_bool(crate::worker_threads::is_worker_isolate(isolate) || xaml_host);
1103+
}
1104+
1105+
/// Cloned and transferred now, as on the web; sent once the current job is done.
1106+
pub(crate) fn handle_worker_queue_message(
1107+
scope: &mut v8::PinScope<'_, '_>,
1108+
args: v8::FunctionCallbackArguments,
1109+
_retval: v8::ReturnValue,
1110+
) {
1111+
if let Some(message) = crate::transfer::serialize(scope, args.get(0), args.get(1)) {
1112+
crate::worker_threads::queue_outgoing(message);
1113+
}
1114+
}
1115+
10831116
/// Runs on a worker's creating thread when the worker has queued messages: hands them to that
10841117
/// `Worker` object (`__nsWorkerDeliver`, installed by the Worker shim).
10851118
pub(crate) fn deliver_worker_events(worker_id: u64) {
@@ -5788,6 +5821,8 @@ pub(crate) fn init_async_helpers(
57885821
register!("__nsDescribeWinRTType", handle_describe_winrt_type);
57895822
register!("__nsWorkerCreateThreaded", handle_worker_create_threaded);
57905823
register!("__nsWorkerPostMessage", handle_worker_post_message);
5824+
register!("__nsWorkerQueueMessage", handle_worker_queue_message);
5825+
register!("__nsWorkerPushesMessages", handle_worker_pushes_messages);
57915826
register!("__nsWorkerPollMessages", handle_worker_poll_messages);
57925827
register!("__nsWorkerTerminate", handle_worker_terminate);
57935828
register!(
@@ -5879,6 +5914,7 @@ pub(crate) fn init_async_helpers(
58795914
}
58805915

58815916
crate::message_port::install_message_port_runtime(scope);
5917+
crate::transfer::install_transfer_runtime(scope);
58825918
crate::worker_support::install_worker_runtime(scope);
58835919
crate::hmr_support::install_hmr_support(scope);
58845920
crate::livesync::install_livesync_support(scope);

0 commit comments

Comments
 (0)