From 3bb14896bff25a05a704eb7a742e1482090b7f65 Mon Sep 17 00:00:00 2001 From: Osei Fortune Date: Sun, 4 Oct 2026 11:25:43 -0400 Subject: [PATCH 1/2] feat: transfer lists in worker postMessage, worker-owned Node-API envs - postMessage(data, transfer) on Worker and in workers: the value is cloned and the list taken when postMessage is called. ArrayBuffers are copied, then detached; other objects need a type registered in each isolate with __nsRegisterTransferable(name, { test, detach, attach, release }), which turns them into a token and back (@nativescript/canvas's OffscreenCanvas). - Node-API threadsafe-function calls and async-work completions run on the thread whose env made them; a worker's ran on the UI thread. - A terminating worker runs its addons' cleanup hooks and frees their envs; external ArrayBuffers still alive are finalized during env teardown, as Node does. An external ArrayBuffer without a finalizer had a null deleter, which crashed when a worker's isolate was disposed. - requestAnimationFrame in workers, following the UI thread's frames. - new Worker('~/chunk.js') (NativeScript's webpack worker loader) resolves against the app directory. - Worker.postMessage no longer waits up to 500 ms for a reply under a XAML host. --- napi-v8-shim/csrc/env_ext.cpp | 7 +- packages/windows-v8/vendor/shim/v8-api.cpp | 68 +++- runtime/src/animation_frames.rs | 48 ++- runtime/src/global_fns.rs | 68 +++- runtime/src/interop_test.rs | 75 +++++ runtime/src/lib.rs | 125 +------ runtime/src/node_api.rs | 101 ++++-- runtime/src/transfer.rs | 358 +++++++++++++++++++++ runtime/src/worker_support.rs | 12 +- runtime/src/worker_threads.rs | 63 ++-- 10 files changed, 732 insertions(+), 193 deletions(-) create mode 100644 runtime/src/transfer.rs diff --git a/napi-v8-shim/csrc/env_ext.cpp b/napi-v8-shim/csrc/env_ext.cpp index 32b15b5..c08aa8b 100644 --- a/napi-v8-shim/csrc/env_ext.cpp +++ b/napi-v8-shim/csrc/env_ext.cpp @@ -14,6 +14,9 @@ static_assert(sizeof(v8::Local) == sizeof(void*), "v8::Local must b // From node_api_types.h, which the vendored headers don't carry. typedef napi_value (*napi_addon_register_func)(napi_env env, napi_value exports); +// v8-api.cpp. +void ns_napi_finalize_external_arraybuffers(napi_env env); + extern "C" { // `context` is a `v8::Local` (as its underlying pointer) of the isolate current on @@ -53,7 +56,8 @@ bool ns_napi_has_pending_finalizers(napi_env env) { return !env->pending_finalizers.empty(); } -// Tears the env down: finalizes remaining references (running their finalizers) and frees it. +// Tears the env down: finalizes remaining references and external ArrayBuffers (running their +// finalizers) and frees it. void ns_napi_env_teardown(napi_env env) { v8::Isolate* isolate = env->isolate; v8::HandleScope handle_scope(isolate); @@ -61,6 +65,7 @@ void ns_napi_env_teardown(napi_env env) { // frees before this scope exits. v8::Local context = v8::Local::New(isolate, env->context()); v8::Context::Scope context_scope(context); + ns_napi_finalize_external_arraybuffers(env); env->DeleteMe(); } diff --git a/packages/windows-v8/vendor/shim/v8-api.cpp b/packages/windows-v8/vendor/shim/v8-api.cpp index 14b258b..64bca43 100644 --- a/packages/windows-v8/vendor/shim/v8-api.cpp +++ b/packages/windows-v8/vendor/shim/v8-api.cpp @@ -3,6 +3,10 @@ #include #include // string_view, u16string_view #include +#include +#include +#include +#include #define NAPI_EXPERIMENTAL @@ -2995,6 +2999,41 @@ const v8::ArrayBuffer *v8__ArrayBuffer__New__with_backing_store( void std__shared_ptr__v8__BackingStore__reset(void *shared_ptr_ref); } +namespace { +// An external ArrayBuffer's finalizer. Its backing store can outlive the env (it goes with the +// isolate's heap), so the env's teardown runs it instead and clears `cb`, as Node does. +struct AbFinalize { + napi_env env; + napi_finalize cb; + void *data; + void *hint; +}; + +std::mutex ab_finalizers_mutex; +std::unordered_map> ab_finalizers; +} + +// [windows port] Called by the env's teardown (napi-v8-shim/csrc/env_ext.cpp). +void ns_napi_finalize_external_arraybuffers(napi_env env) { + // Copies: once `cb` is cleared, a deleter on another thread may free the record. + std::vector pending; + { + std::lock_guard lock(ab_finalizers_mutex); + auto found = ab_finalizers.find(env); + if (found == ab_finalizers.end()) { + return; + } + for (AbFinalize *fd : found->second) { + pending.push_back(*fd); + fd->cb = nullptr; + } + ab_finalizers.erase(found); + } + for (const AbFinalize &fd : pending) { + fd.cb(env, fd.data, fd.hint); + } +} + napi_status NAPI_CDECL napi_create_external_arraybuffer(napi_env env, void *external_data, @@ -3010,19 +3049,30 @@ napi_create_external_arraybuffer(napi_env env, // [windows port] TRUE zero-copy: build a BackingStore that aliases external_data with a deleter // that invokes the napi finalizer when the ArrayBuffer is GC'd. Goes through rusty_v8's C // bindings (see above) to avoid the libc++ std::unique_ptr ABI boundary. - struct AbFinalize { - napi_env env; - napi_finalize cb; - void *hint; - }; void *deleter_data = nullptr; - void (*deleter)(void *, size_t, void *) = nullptr; + // V8 calls the deleter unconditionally when the backing store goes (at the latest when the + // isolate is disposed, as a terminated worker's is), so it can't be null. + void (*deleter)(void *, size_t, void *) = [](void *, size_t, void *) {}; if (finalize_cb != nullptr) { - deleter_data = new AbFinalize{env, finalize_cb, finalize_hint}; + auto *fd = new AbFinalize{env, finalize_cb, external_data, finalize_hint}; + { + std::lock_guard lock(ab_finalizers_mutex); + ab_finalizers[env].insert(fd); + } + deleter_data = fd; deleter = [](void *data, size_t, void *dd) { auto *fd = static_cast(dd); - if (fd->cb) { - fd->cb(fd->env, data, fd->hint); + napi_finalize cb; + { + std::lock_guard lock(ab_finalizers_mutex); + cb = fd->cb; + auto found = ab_finalizers.find(fd->env); + if (cb && found != ab_finalizers.end()) { + found->second.erase(fd); + } + } + if (cb) { + cb(fd->env, data, fd->hint); } delete fd; }; diff --git a/runtime/src/animation_frames.rs b/runtime/src/animation_frames.rs index 8cd31a1..3c1db4f 100644 --- a/runtime/src/animation_frames.rs +++ b/runtime/src/animation_frames.rs @@ -5,16 +5,35 @@ //! once per compositor frame from `CompositionTarget.Rendering`, outside the render walk) then //! runs the queued callbacks once and drains microtasks, which is where rendering work such as //! canvas presents happens. Nothing waits for vsync on the UI thread, and a continuous rAF loop -//! gives the dispatcher back between frames. +//! gives the dispatcher back between frames. A worker's frames follow the UI thread's: each pump +//! there hands a frame to every worker that asked for one. The compositor stops raising frames +//! while nothing on screen changes, so a worker not handed one in time takes its own. use std::cell::Cell; +use std::time::{Duration, Instant}; use crate::DELEGATE_ISOLATE_PTR; thread_local! { static REQUESTED: Cell = const { Cell::new(false) }; + static LAST_FRAME: Cell> = const { Cell::new(None) }; } +#[cfg(feature = "classic")] +const WORKER_FRAME_INTERVAL: Duration = Duration::from_micros(16_667); + +#[cfg(feature = "classic")] +pub(crate) fn worker_frame_wait() -> Option { + if !REQUESTED.with(|r| r.get()) { + return None; + } + let last = LAST_FRAME.with(|l| l.get()); + Some(last.map_or(Duration::ZERO, |last| WORKER_FRAME_INTERVAL.saturating_sub(last.elapsed()))) +} + +#[cfg(feature = "classic")] +static WORKERS_REQUESTED: std::sync::Mutex> = std::sync::Mutex::new(Vec::new()); + /// `__nsRequestFrame()`: run animation callbacks at the next pump. pub(crate) fn handle_request_frame( _scope: &mut v8::PinScope<'_, '_>, @@ -22,6 +41,30 @@ pub(crate) fn handle_request_frame( _retval: v8::ReturnValue, ) { REQUESTED.with(|r| r.set(true)); + #[cfg(feature = "classic")] + { + let isolate = DELEGATE_ISOLATE_PTR.with(|c| c.get()) as usize; + if crate::worker_threads::is_worker_isolate(isolate) { + let mut workers = WORKERS_REQUESTED.lock().unwrap_or_else(|e| e.into_inner()); + if !workers.contains(&isolate) { + workers.push(isolate); + } + } + } +} + +#[cfg(feature = "classic")] +fn hand_frames_to_workers() { + let isolate = DELEGATE_ISOLATE_PTR.with(|c| c.get()) as usize; + if crate::worker_threads::is_worker_isolate(isolate) { + return; + } + let workers = std::mem::take(&mut *WORKERS_REQUESTED.lock().unwrap_or_else(|e| e.into_inner())); + for worker in workers { + crate::worker_threads::post_to_worker(worker, || { + pump(); + }); + } } /// Drops a pending request (the runtime on this thread is going away). @@ -40,9 +83,12 @@ fn now_ms() -> f64 { /// Runs this thread's pending animation-frame callbacks, if a frame was requested, then drains /// microtasks. Returns whether callbacks ran. pub fn pump() -> bool { + #[cfg(feature = "classic")] + hand_frames_to_workers(); if !REQUESTED.with(|r| r.replace(false)) { return false; } + LAST_FRAME.with(|l| l.set(Some(Instant::now()))); let isolate_ptr = DELEGATE_ISOLATE_PTR.with(|c| c.get()); if isolate_ptr.is_null() { return false; diff --git a/runtime/src/global_fns.rs b/runtime/src/global_fns.rs index 279bb91..4cbfa7d 100644 --- a/runtime/src/global_fns.rs +++ b/runtime/src/global_fns.rs @@ -14,7 +14,7 @@ use windows::Win32::UI::WindowsAndMessaging::{ use crate::type_description::build_runtime_type_descriptor; use crate::dotnet::{bin_write_str16, bin_write_str32}; -use crate::{normalize_js_path, proxy_manifests, throw_js_error, try_resolve_with_known_extensions, Runtime, ASYNC_PUMP_HOOK}; +use crate::{normalize_js_path, proxy_manifests, throw_js_error, try_resolve_with_known_extensions, ASYNC_PUMP_HOOK}; use std::cell::RefCell; use std::ffi::c_void; @@ -209,23 +209,29 @@ fn default_auto_capture_path() -> PathBuf { PathBuf::from("sbg_output").join("sbg_metadata.json") } +/// Delivered as a `messageerror`. +fn worker_error<'s>(scope: &mut v8::PinScope<'s, '_>, error: &str) -> Option> { + let obj = v8::Object::new(scope); + if let Some(key) = v8::String::new(scope, "__workerError") { + if let Some(val) = v8::String::new(scope, error) { + obj.set(scope, key.into(), val.into()); + } + } + Some(obj.into()) +} + fn polled_event_to_v8<'s>( scope: &mut v8::PinScope<'s, '_>, event: crate::worker_threads::PolledWorkerEvent, ) -> Option> { match event { - crate::worker_threads::PolledWorkerEvent::Message(bytes) => { - Runtime::deserialize_value(scope, &bytes) - } - crate::worker_threads::PolledWorkerEvent::Error(error) => { - let obj = v8::Object::new(scope); - if let Some(key) = v8::String::new(scope, "__workerError") { - if let Some(val) = v8::String::new(scope, error.as_str()) { - obj.set(scope, key.into(), val.into()); - } + crate::worker_threads::PolledWorkerEvent::Message(message) => { + match crate::transfer::deserialize(scope, message) { + Ok(value) => Some(value), + Err(error) => worker_error(scope, &error), } - Some(obj.into()) } + crate::worker_threads::PolledWorkerEvent::Error(error) => worker_error(scope, &error), crate::worker_threads::PolledWorkerEvent::Exited => { let obj = v8::Object::new(scope); if let Some(key) = v8::String::new(scope, "__workerExit") { @@ -919,6 +925,11 @@ pub(crate) fn handle_resolve_module_path( throw_js_error(scope, "__nsResolveModulePath: module specifier is empty"); return; } + // The app directory, as NativeScript's webpack worker loader names worker chunks. + let specifier = match specifier.strip_prefix("~/") { + Some(rest) => rest.to_string(), + None => specifier, + }; let parent_path = if args.length() >= 2 { value_to_string(scope, args.get(1)) } else { @@ -1061,7 +1072,7 @@ pub(crate) fn handle_worker_post_message( if args.length() < 2 { throw_js_error( scope, - "__nsWorkerPostMessage(workerId, value) expects 2 arguments", + "__nsWorkerPostMessage(workerId, value, transfer?) expects 2 arguments", ); return; } @@ -1070,16 +1081,38 @@ pub(crate) fn handle_worker_post_message( throw_js_error(scope, "Invalid worker id"); return; } - let value = args.get(1); - let Some(bytes) = Runtime::serialize_value(scope, value) else { - throw_js_error(scope, "DataCloneError: value could not be cloned."); + // On failure the exception is pending. + let Some(message) = crate::transfer::serialize(scope, args.get(1), args.get(2)) else { return; }; - if let Err(err) = crate::worker_threads::post_message(worker_id as u64, bytes) { + if let Err(err) = crate::worker_threads::post_message(worker_id as u64, message) { throw_js_error(scope, err.as_str()); } } +/// A dispatcher the runtime made itself only runs when its host pumps messages: there a worker's +/// replies only arrive by polling. +pub(crate) fn handle_worker_pushes_messages( + _scope: &mut v8::PinScope<'_, '_>, + _args: v8::FunctionCallbackArguments, + mut retval: v8::ReturnValue, +) { + let isolate = crate::DELEGATE_ISOLATE_PTR.with(|c| c.get()) as usize; + let xaml_host = crate::ui_dispatcher::is_initialized() && !crate::ui_dispatcher::needs_win32_pump(); + retval.set_bool(crate::worker_threads::is_worker_isolate(isolate) || xaml_host); +} + +/// Cloned and transferred now, as on the web; sent once the current job is done. +pub(crate) fn handle_worker_queue_message( + scope: &mut v8::PinScope<'_, '_>, + args: v8::FunctionCallbackArguments, + _retval: v8::ReturnValue, +) { + if let Some(message) = crate::transfer::serialize(scope, args.get(0), args.get(1)) { + crate::worker_threads::queue_outgoing(message); + } +} + /// Runs on a worker's creating thread when the worker has queued messages: hands them to that /// `Worker` object (`__nsWorkerDeliver`, installed by the Worker shim). pub(crate) fn deliver_worker_events(worker_id: u64) { @@ -5788,6 +5821,8 @@ pub(crate) fn init_async_helpers( register!("__nsDescribeWinRTType", handle_describe_winrt_type); register!("__nsWorkerCreateThreaded", handle_worker_create_threaded); register!("__nsWorkerPostMessage", handle_worker_post_message); + register!("__nsWorkerQueueMessage", handle_worker_queue_message); + register!("__nsWorkerPushesMessages", handle_worker_pushes_messages); register!("__nsWorkerPollMessages", handle_worker_poll_messages); register!("__nsWorkerTerminate", handle_worker_terminate); register!( @@ -5879,6 +5914,7 @@ pub(crate) fn init_async_helpers( } crate::message_port::install_message_port_runtime(scope); + crate::transfer::install_transfer_runtime(scope); crate::worker_support::install_worker_runtime(scope); crate::hmr_support::install_hmr_support(scope); crate::livesync::install_livesync_support(scope); diff --git a/runtime/src/interop_test.rs b/runtime/src/interop_test.rs index 3648dbd..1e318fd 100644 --- a/runtime/src/interop_test.rs +++ b/runtime/src/interop_test.rs @@ -1231,6 +1231,81 @@ fn worker_multiple_sequential_roundtrips() { ); } +#[test] +fn worker_transfer_list_moves_buffers_and_registered_objects() { + run_js_assert( + "worker_transfer_list_moves_buffers_and_registered_objects", + r#" + return new Promise((resolve, reject) => { + const worker = new Worker(` + class Box { constructor(value) { this.value = value; } } + __nsRegisterTransferable('Box', { test: (v) => v instanceof Box, detach: (box) => box.value, attach: (value) => new Box(value) }); + self.onmessage = (event) => { + const { box, buffer } = event.data; + const back = new Box('from the worker'); + self.postMessage({ isBox: box instanceof Box, value: box.value, bytes: new Uint8Array(buffer)[3], back }, [back]); + }; + `, { eval: true }); + + class Box { constructor(value) { this.value = value; this.sent = false; } } + __nsRegisterTransferable('Box', { + test: (v) => v instanceof Box, + detach: (box) => { box.sent = true; return box.value; }, + attach: (value) => new Box(value), + }); + const box = new Box(7); + const buffer = new Uint8Array([0, 1, 2, 3]).buffer; + + worker.onmessage = (event) => { + worker.terminate(); + const { isBox, value, bytes, back } = event.data; + if (!isBox || value !== 7 || bytes !== 3) { + reject(new Error(`received ${JSON.stringify({ isBox, value, bytes })}`)); + } else if (!box.sent || buffer.byteLength !== 0) { + reject(new Error('the sender kept what it transferred')); + } else if (!(back instanceof Box) || back.value !== 'from the worker') { + reject(new Error('the worker\'s Box did not come back as one')); + } else { + resolve(); + } + }; + worker.postMessage({ box, buffer }, [box, buffer]); + }); + "#, + ); +} + +#[test] +fn worker_transfer_list_rejects_what_it_cannot_take() { + run_js_assert( + "worker_transfer_list_rejects_what_it_cannot_take", + r#" + const worker = new Worker("self.onmessage = function () {};", { eval: true }); + const nameOf = (fn) => { + try { + fn(); + } catch (e) { + return e.name; + } + return 'no error'; + }; + const buffer = new ArrayBuffer(4); + const results = [ + nameOf(() => worker.postMessage({}, [{}])), + nameOf(() => worker.postMessage(buffer, [buffer, buffer])), + nameOf(() => worker.postMessage(1, 'not a list')), + ]; + worker.terminate(); + if (results[0] !== 'DataCloneError' || results[1] !== 'DataCloneError' || results[2] !== 'TypeError') { + throw new Error(`got ${results.join(', ')}`); + } + if (buffer.byteLength !== 4) { + throw new Error('a failed transfer took the buffer'); + } + "#, + ); +} + #[test] fn worker_complex_object_roundtrip() { run_js_assert( diff --git a/runtime/src/lib.rs b/runtime/src/lib.rs index 211ff20..87a6cbd 100644 --- a/runtime/src/lib.rs +++ b/runtime/src/lib.rs @@ -61,6 +61,8 @@ pub(crate) mod win32_known_fns; mod websocket; mod winhttp; #[cfg(feature = "classic")] +mod transfer; +#[cfg(feature = "classic")] mod worker_support; #[cfg(feature = "classic")] mod worker_threads; @@ -9473,130 +9475,19 @@ impl Drop for Runtime { } } -#[cfg(feature = "classic")] -struct WorkerValueSerializer; - -#[cfg(feature = "classic")] -impl v8::ValueSerializerImpl for WorkerValueSerializer { - fn throw_data_clone_error<'s>( - &self, - scope: &mut v8::PinScope<'s, '_>, - message: v8::Local<'s, v8::String>, - ) { - let error = v8::Exception::error(scope, message); - scope.throw_exception(error); - } -} - -#[cfg(feature = "classic")] -struct WorkerValueDeserializer; - -#[cfg(feature = "classic")] -impl v8::ValueDeserializerImpl for WorkerValueDeserializer {} - #[cfg(feature = "classic")] impl Runtime { - /// Serialize a single V8 value to structured-clone bytes using V8's own - /// `ValueSerializer`. Returns `None` if the value is not cloneable (e.g. - /// a function or a circular object); in that case an exception has already - /// been thrown into `scope`. - #[cfg(feature = "classic")] - pub fn serialize_value<'s, 'v>( - scope: &mut v8::PinScope<'s, '_>, - value: v8::Local<'v, v8::Value>, - ) -> Option> { - use v8::ValueSerializerHelper; - let context = scope.get_current_context(); - let ser = v8::ValueSerializer::new(scope, Box::new(WorkerValueSerializer)); - ser.write_header(); - if ser.write_value(context, value).unwrap_or(false) { - Some(ser.release()) - } else { - None - } - } - - /// Deserialize structured-clone bytes produced by `serialize_value` back - /// into a V8 value in the current context. - #[cfg(feature = "classic")] - pub fn deserialize_value<'s>( - scope: &mut v8::PinScope<'s, '_>, - bytes: &[u8], - ) -> Option> { - use v8::ValueDeserializerHelper; - let context = scope.get_current_context(); - let de = v8::ValueDeserializer::new(scope, Box::new(WorkerValueDeserializer), bytes); - if !de.read_header(context).unwrap_or(false) { - return None; - } - de.read_value(context) - } - - /// Drain `globalThis.__nsWorkerOutbox`, serialize every item with V8's - /// structured-clone algorithm, and return the resulting byte blobs. - #[cfg(feature = "classic")] - pub fn drain_outbox_bytes(&mut self) -> Vec, String>> { - v8::scope!(scope, &mut self.isolate); - let context = v8::Local::new(scope, &self.global_context); - let scope = &mut v8::ContextScope::new(scope, context); - v8::tc_scope!(tc, scope); - - let script_src = - "(function(){var o=globalThis.__nsWorkerOutbox||[];return o.splice(0);})()"; - let Some(src) = v8::String::new(tc, script_src) else { - return Vec::new(); - }; - let Some(script) = v8::Script::compile(tc, src, None) else { - return Vec::new(); - }; - let Some(result) = script.run(tc) else { - return Vec::new(); - }; - let Ok(array) = v8::Local::::try_from(result) else { - return Vec::new(); - }; - - let len = array.length(); - let mut out = Vec::with_capacity(len as usize); - - for i in 0..len { - let Some(item) = array.get_index(tc, i) else { - continue; - }; - match Self::serialize_value(tc, item) { - Some(bytes) => out.push(Ok(bytes)), - None => { - let msg = if tc.has_caught() { - let s = tc - .message() - .and_then(|m| Some(m.get(tc).to_rust_string_lossy(tc))) - .unwrap_or_else(|| "DataCloneError".to_string()); - tc.reset(); - s - } else { - "DataCloneError: value could not be cloned".to_string() - }; - out.push(Err(msg)); - } - } - } - - out - } - - /// Deserialize `payload_bytes` and deliver them to the worker's - /// `__nsDispatchToWorker` JS function. - #[cfg(feature = "classic")] - pub fn dispatch_to_worker(&mut self, payload_bytes: &[u8]) { + pub fn dispatch_to_worker(&mut self, message: crate::transfer::Message) { v8::scope!(scope, &mut self.isolate); let context = v8::Local::new(scope, &self.global_context); let scope = &mut v8::ContextScope::new(scope, context); v8::tc_scope!(tc, scope); - let data_value = match Self::deserialize_value(tc, payload_bytes) { - Some(v) => v, - None => { - eprintln!("[NativeScript] Worker dispatch: failed to deserialize message"); + let data_value = match crate::transfer::deserialize(tc, message) { + Ok(value) => value, + Err(error) => { + tc.reset(); + eprintln!("[NativeScript] Worker dispatch: {error}"); return; } }; diff --git a/runtime/src/node_api.rs b/runtime/src/node_api.rs index 0897142..ca3e014 100644 --- a/runtime/src/node_api.rs +++ b/runtime/src/node_api.rs @@ -6,10 +6,12 @@ //! and exported from `nativescript.dll`, so an addon built for Node (napi-rs, node-addon-api, …) //! resolves every symbol it imports. //! -//! Anything that must run on the JS thread goes through one job queue. Other threads push jobs and -//! wake the UI thread with a `DispatcherQueue` work item; `runtime_pump_timers` drains it too, so a -//! host without a dispatcher still makes progress. Each job runs in its own V8 scope and ends with -//! a microtask checkpoint, like a timer task. +//! Anything that must run on a JS thread goes through that thread's job queue: an env belongs to the +//! thread (the UI thread's or a worker's runtime) that loaded the addon. Other threads push jobs and +//! wake the UI thread with a `DispatcherQueue` work item, or a worker through its command channel; +//! `runtime_pump_timers` drains the UI thread's too, so a host without a dispatcher still makes +//! progress. Each job runs in its own V8 scope and ends with a microtask checkpoint, like a timer +//! task. use std::cell::{Cell, RefCell}; use std::collections::{HashMap, VecDeque}; @@ -128,42 +130,77 @@ enum Job { // Jobs only carry `Arc` (Send+Sync below) and pointers used on the JS thread. unsafe impl Send for Job {} -static QUEUE: Mutex> = Mutex::new(VecDeque::new()); -static WAKE_QUEUED: AtomicBool = AtomicBool::new(false); +#[derive(Default)] +struct Queue { + jobs: VecDeque, + wake_queued: bool, +} + +static QUEUES: Mutex>> = Mutex::new(None); /// Set by [`teardown`]: the runtime is going away, so late work (a cleanup hook releasing a /// threadsafe function, a dispatcher item that runs during shutdown) must not touch V8. static SHUT_DOWN: AtomicBool = AtomicBool::new(false); -fn push(job: Job) { +fn with_queues(f: impl FnOnce(&mut HashMap) -> R) -> R { + let mut guard = QUEUES.lock().unwrap(); + f(guard.get_or_insert_with(HashMap::new)) +} + +/// The isolate of this thread's runtime: what an env made here belongs to. +fn current_owner() -> usize { + DELEGATE_ISOLATE_PTR.with(|c| c.get()) as usize +} + +fn push(owner: usize, job: Job) { if SHUT_DOWN.load(Ordering::Acquire) { return; } - QUEUE.lock().unwrap().push_back(job); - wake(); + let wake = with_queues(|queues| { + let queue = queues.entry(owner).or_default(); + queue.jobs.push_back(job); + !std::mem::replace(&mut queue.wake_queued, true) + }); + if wake { + wake_owner(owner); + } } -fn wake() { - if SHUT_DOWN.load(Ordering::Acquire) || WAKE_QUEUED.swap(true, Ordering::AcqRel) { +fn wake_owner(owner: usize) { + let drain = || { + let _ = std::panic::catch_unwind(drain); + }; + if crate::worker_threads::is_worker_isolate(owner) { + if !crate::worker_threads::post_to_worker(owner, drain) { + // The worker is gone, and its env with it. + with_queues(|queues| queues.remove(&owner)); + } return; } // A dispatcher work item runs between frames, never inside a XAML callout. Hosts without one // (console apps, tests) drain from their pump loop instead. - if !crate::ui_dispatcher::enqueue_on_ui_thread(|| { - let _ = std::panic::catch_unwind(drain); - }) { - WAKE_QUEUED.store(false, Ordering::Release); + if !crate::ui_dispatcher::enqueue_on_ui_thread(drain) { + with_queues(|queues| { + if let Some(queue) = queues.get_mut(&owner) { + queue.wake_queued = false; + } + }); } } -/// Runs every queued JS-thread job and every env's deferred finalizers. Called from the dispatcher -/// wake-up and from `runtime_pump_timers`; a no-op when there is nothing to do. +/// Runs this thread's queued jobs and its envs' deferred finalizers. Called from the dispatcher +/// (or worker) wake-up and from `runtime_pump_timers`; a no-op when there is nothing to do. pub fn drain() { - WAKE_QUEUED.store(false, Ordering::Release); - if SHUT_DOWN.load(Ordering::Acquire) { + let owner = current_owner(); + if SHUT_DOWN.load(Ordering::Acquire) || owner == 0 { return; } + with_queues(|queues| { + if let Some(queue) = queues.get_mut(&owner) { + queue.wake_queued = false; + } + }); loop { - let job = QUEUE.lock().unwrap().pop_front(); + let job = with_queues(|queues| queues.get_mut(&owner).and_then(|queue| queue.jobs.pop_front())); let Some(job) = job else { break }; match job { Job::Tsfn(tsfn) => in_js_scope(|| unsafe { tsfn.dispatch() }), @@ -204,7 +241,18 @@ pub fn teardown() { // Deliver what is already queued while everything is still alive, then stop accepting work. drain(); SHUT_DOWN.store(true, Ordering::Release); - QUEUE.lock().unwrap().clear(); + with_queues(|queues| queues.clear()); + for &env in envs.iter().rev() { + run_cleanup_hooks(env); + in_js_scope(|| unsafe { napi_v8_shim::ns_napi_env_teardown(env) }); + } +} + +/// A worker's [`teardown`], before its isolate goes: the other runtimes keep theirs. +pub(crate) fn teardown_thread() { + let envs: Vec = ENVS.with(|e| std::mem::take(&mut *e.borrow_mut())); + drain(); + with_queues(|queues| queues.remove(¤t_owner())); for &env in envs.iter().rev() { run_cleanup_hooks(env); in_js_scope(|| unsafe { napi_v8_shim::ns_napi_env_teardown(env) }); @@ -330,6 +378,7 @@ struct TsfnState { struct Tsfn { env: napi_env, + owner: usize, func: napi_ref, context: *mut c_void, call_js: napi_threadsafe_function_call_js, @@ -435,7 +484,7 @@ impl Tsfn { fn queue_finalize(self: &Arc, state: &mut TsfnState) { if !state.finalize_queued { state.finalize_queued = true; - push(Job::Tsfn(self.clone())); + push(self.owner, Job::Tsfn(self.clone())); } } } @@ -463,6 +512,7 @@ pub unsafe extern "C" fn napi_create_threadsafe_function( } let tsfn = Arc::new(Tsfn { env, + owner: current_owner(), func: reference, context, call_js: call_js_cb, @@ -501,7 +551,7 @@ pub unsafe extern "C" fn napi_call_threadsafe_function(func: *mut c_void, data: } state.queue.push_back(data); drop(state); - push(Job::Tsfn(Tsfn::arc(func))); + push(tsfn.owner, Job::Tsfn(Tsfn::arc(func))); NAPI_OK } @@ -570,6 +620,7 @@ const WORK_CANCELLED: u8 = 3; struct AsyncWork { env: napi_env, + owner: usize, execute: napi_async_execute_callback, complete: napi_async_complete_callback, data: *mut c_void, @@ -612,6 +663,7 @@ pub unsafe extern "C" fn napi_create_async_work( } *result = Box::into_raw(Box::new(AsyncWork { env, + owner: current_owner(), execute, complete, data, @@ -643,6 +695,7 @@ pub unsafe extern "C" fn napi_queue_async_work(_env: napi_env, work: *mut c_void return NAPI_GENERIC_FAILURE; } let address = work as usize; + let owner = entry.owner; let task: Task = Box::new(move || { let work = unsafe { &*(address as *const AsyncWork) }; if work @@ -654,7 +707,7 @@ pub unsafe extern "C" fn napi_queue_async_work(_env: napi_env, work: *mut c_void unsafe { execute(work.env, work.data) }; } } - push(Job::AsyncComplete(address)); + push(owner, Job::AsyncComplete(address)); }); match pool().lock() { Ok(tx) if tx.send(task).is_ok() => NAPI_OK, diff --git a/runtime/src/transfer.rs b/runtime/src/transfer.rs new file mode 100644 index 0000000..ada08fb --- /dev/null +++ b/runtime/src/transfer.rs @@ -0,0 +1,358 @@ +use v8::{ValueDeserializerHelper, ValueSerializerHelper}; + +enum Token { + Number(f64), + String(String), +} + +struct Transferred { + kind: String, + token: Token, +} + +pub struct Message { + bytes: Vec, + transferred: Vec, +} + +/// `__nsRegisterTransferable(name, { test, detach, attach, release? })`: `detach` returns a number +/// or string token; the receiving isolate's handler of the same name `attach`es it, or `release`s +/// it when the message can't be received. ArrayBuffers in a transfer list are copied, then detached. +pub fn install_transfer_runtime(scope: &mut v8::ContextScope) { + let source = r#" + (function () { + if (typeof globalThis.__nsRegisterTransferable === 'function') { + return; + } + var handlers = Object.create(null); + + function hidden(name, value) { + Object.defineProperty(globalThis, name, { value: value, configurable: true, writable: true }); + } + + function release(kinds, tokens, from) { + for (var i = from; i < kinds.length; i++) { + var handler = handlers[kinds[i]]; + if (handler && typeof handler.release === 'function') { + try { handler.release(tokens[i]); } catch (_) {} + } + } + } + + hidden('__nsRegisterTransferable', function (name, handler) { + if (typeof name !== 'string' || !handler || typeof handler.test !== 'function' || + typeof handler.detach !== 'function' || typeof handler.attach !== 'function') { + throw new TypeError('__nsRegisterTransferable(name, { test, detach, attach, release? })'); + } + handlers[name] = handler; + }); + + hidden('__nsTransferHooks', { + typeOf: function (value) { + for (var name in handlers) { + if (handlers[name].test(value)) { + return name; + } + } + return undefined; + }, + // Sender: all or none. A throw releases what was taken. + detachAll: function (kinds, values) { + var tokens = []; + try { + for (var i = 0; i < values.length; i++) { + var token = handlers[kinds[i]].detach(values[i]); + if (typeof token !== 'number' && typeof token !== 'string') { + tokens.push(token); + throw new TypeError("A '" + kinds[i] + "' transfer handler returned neither a number nor a string."); + } + tokens.push(token); + } + } catch (e) { + release(kinds, tokens, 0); + throw e; + } + return tokens; + }, + // Receiver: never throws. On an error the rest are released. + attachAll: function (kinds, tokens) { + var objects = []; + for (var i = 0; i < kinds.length; i++) { + try { + var handler = handlers[kinds[i]]; + if (!handler) { + throw new Error("DataCloneError: nothing here receives a transferred '" + kinds[i] + "'."); + } + var object = handler.attach(tokens[i]); + if (object === null || (typeof object !== 'object' && typeof object !== 'function')) { + throw new TypeError("A '" + kinds[i] + "' transfer handler returned no object."); + } + objects.push(object); + } catch (e) { + release(kinds, tokens, i + 1); + return { error: String((e && e.message) || e) }; + } + } + return { objects: objects }; + } + }); + })(); + "#; + let Some(source) = v8::String::new(scope, source) else { + return; + }; + if let Some(script) = v8::Script::compile(scope, source, None) { + script.run(scope); + } +} + +fn throw_data_clone_error(scope: &mut v8::PinScope<'_, '_>, message: &str) { + let Some(text) = v8::String::new(scope, &format!("DataCloneError: {message}")) else { + return; + }; + let error = v8::Exception::error(scope, text); + if let (Ok(object), Some(key), Some(name)) = ( + v8::Local::::try_from(error), + v8::String::new(scope, "name"), + v8::String::new(scope, "DataCloneError"), + ) { + object.set(scope, key.into(), name.into()); + } + scope.throw_exception(error); +} + +fn throw_type_error(scope: &mut v8::PinScope<'_, '_>, message: &str) { + if let Some(text) = v8::String::new(scope, message) { + let error = v8::Exception::type_error(scope, text); + scope.throw_exception(error); + } +} + +fn hooks<'s>(scope: &mut v8::PinScope<'s, '_>) -> Option> { + let context = scope.get_current_context(); + let global = context.global(scope); + let key = v8::String::new(scope, "__nsTransferHooks")?; + global.get(scope, key.into())?.try_into().ok() +} + +/// `None`: an exception is pending. +fn call_hook<'s>( + scope: &mut v8::PinScope<'s, '_>, + hooks: v8::Local<'s, v8::Object>, + name: &str, + args: &[v8::Local<'s, v8::Value>], +) -> Option> { + let key = v8::String::new(scope, name)?; + let function: v8::Local = hooks.get(scope, key.into())?.try_into().ok()?; + function.call(scope, hooks.into(), args) +} + +/// The transfer list: `postMessage(data, [..])` or `postMessage(data, { transfer: [..] })`. +fn transfer_list<'s>(scope: &mut v8::PinScope<'s, '_>, transfer: v8::Local<'s, v8::Value>) -> Option>> { + if transfer.is_null_or_undefined() { + return Some(Vec::new()); + } + let list = if transfer.is_array() { + transfer + } else if let Ok(options) = v8::Local::::try_from(transfer) { + let key = v8::String::new(scope, "transfer")?; + let list = options.get(scope, key.into())?; + if list.is_null_or_undefined() { + return Some(Vec::new()); + } + list + } else { + throw_type_error(scope, "postMessage's transfer must be an array"); + return None; + }; + let Ok(array) = v8::Local::::try_from(list) else { + throw_type_error(scope, "postMessage's transfer must be an array"); + return None; + }; + (0..array.length()).map(|index| array.get_index(scope, index)).collect() +} + +struct Serializer { + objects: Vec>, +} + +impl Serializer { + fn index_of(&self, scope: &mut v8::PinScope<'_, '_>, object: v8::Local) -> Option { + self.objects + .iter() + .position(|known| v8::Local::new(scope, known).strict_equals(object.into())) + .map(|index| index as u32) + } +} + +impl v8::ValueSerializerImpl for Serializer { + fn throw_data_clone_error<'s>(&self, scope: &mut v8::PinScope<'s, '_>, message: v8::Local<'s, v8::String>) { + let error = v8::Exception::error(scope, message); + scope.throw_exception(error); + } + + fn has_custom_host_object(&self, _isolate: &v8::Isolate) -> bool { + !self.objects.is_empty() + } + + fn is_host_object<'s>(&self, scope: &mut v8::PinScope<'s, '_>, object: v8::Local<'s, v8::Object>) -> Option { + Some(self.index_of(scope, object).is_some()) + } + + fn write_host_object<'s>( + &self, + scope: &mut v8::PinScope<'s, '_>, + object: v8::Local<'s, v8::Object>, + serializer: &dyn ValueSerializerHelper, + ) -> Option { + serializer.write_uint32(self.index_of(scope, object)?); + Some(true) + } +} + +struct Deserializer { + objects: Vec>, +} + +impl v8::ValueDeserializerImpl for Deserializer { + fn read_host_object<'s>( + &self, + scope: &mut v8::PinScope<'s, '_>, + deserializer: &dyn ValueDeserializerHelper, + ) -> Option> { + let mut index = 0; + let object = deserializer.read_uint32(&mut index).then(|| self.objects.get(index as usize)).flatten(); + match object { + Some(object) => Some(v8::Local::new(scope, object)), + None => { + throw_data_clone_error(scope, "the message names a transferred object it doesn't carry"); + None + } + } + } +} + +/// Clones `value` and takes what `transfer` lists. `None`: an exception is pending, and nothing +/// was taken. +pub fn serialize<'s>( + scope: &mut v8::PinScope<'s, '_>, + value: v8::Local<'_, v8::Value>, + transfer: v8::Local<'_, v8::Value>, +) -> Option { + let value = v8::Local::new(scope, value); + let transfer = v8::Local::new(scope, transfer); + let items = transfer_list(scope, transfer)?; + let mut buffers = Vec::new(); + let mut objects = Vec::new(); + let mut kinds = Vec::new(); + for (index, item) in items.iter().enumerate() { + if items[..index].iter().any(|seen| seen.strict_equals(*item)) { + throw_data_clone_error(scope, &format!("value at index {index} is listed twice in the transfer list")); + return None; + } + if let Ok(buffer) = v8::Local::::try_from(*item) { + if buffer.was_detached() || !buffer.is_detachable() { + throw_data_clone_error(scope, &format!("the ArrayBuffer at index {index} can't be transferred")); + return None; + } + buffers.push(buffer); + continue; + } + let kind = match (v8::Local::::try_from(*item), hooks(scope)) { + (Ok(object), Some(hooks)) => { + let kind = call_hook(scope, hooks, "typeOf", &[object.into()])?; + kind.is_string().then(|| (object, kind)) + } + _ => None, + }; + let Some((object, kind)) = kind else { + throw_data_clone_error(scope, &format!("value at index {index} is not transferable")); + return None; + }; + objects.push(object); + kinds.push(kind); + } + + let serializer = Serializer { + objects: objects.iter().map(|object| v8::Global::new(scope, *object)).collect(), + }; + let context = scope.get_current_context(); + let bytes = { + let serializer = v8::ValueSerializer::new(scope, Box::new(serializer)); + serializer.write_header(); + if !serializer.write_value(context, value).unwrap_or(false) { + return None; + } + serializer.release() + }; + + let mut transferred = Vec::with_capacity(objects.len()); + if !objects.is_empty() { + let hooks = hooks(scope)?; + let kinds_array = v8::Array::new_with_elements(scope, &kinds); + let values: Vec> = objects.iter().map(|object| (*object).into()).collect(); + let values_array = v8::Array::new_with_elements(scope, &values); + let tokens = call_hook(scope, hooks, "detachAll", &[kinds_array.into(), values_array.into()])?; + let tokens: v8::Local = tokens.try_into().ok()?; + for (index, kind) in kinds.iter().enumerate() { + let token = tokens.get_index(scope, index as u32)?; + let token = if token.is_number() { + Token::Number(token.number_value(scope)?) + } else { + Token::String(token.to_rust_string_lossy(scope)) + }; + transferred.push(Transferred { kind: kind.to_rust_string_lossy(scope), token }); + } + } + for buffer in buffers { + buffer.detach(None); + } + Some(Message { bytes, transferred }) +} + +/// The message's value in this isolate. `Err`: why it can't be received (what it carried is +/// released). +pub fn deserialize<'s>(scope: &mut v8::PinScope<'s, '_>, message: Message) -> Result, String> { + let mut objects = Vec::with_capacity(message.transferred.len()); + if !message.transferred.is_empty() { + let hooks = hooks(scope).ok_or("DataCloneError: the transfer runtime is missing")?; + let mut kinds = Vec::with_capacity(message.transferred.len()); + let mut tokens = Vec::with_capacity(message.transferred.len()); + for transferred in &message.transferred { + let kind = v8::String::new(scope, &transferred.kind).ok_or("DataCloneError")?; + kinds.push(kind.into()); + tokens.push(match &transferred.token { + Token::Number(number) => v8::Number::new(scope, *number).into(), + Token::String(string) => v8::String::new(scope, string).ok_or("DataCloneError")?.into(), + }); + } + let kinds = v8::Array::new_with_elements(scope, &kinds); + let tokens = v8::Array::new_with_elements(scope, &tokens); + let result = call_hook(scope, hooks, "attachAll", &[kinds.into(), tokens.into()]) + .and_then(|result| v8::Local::::try_from(result).ok()) + .ok_or("DataCloneError: the transferred objects could not be received")?; + let key = v8::String::new(scope, "objects").ok_or("DataCloneError")?; + let attached = result.get(scope, key.into()).and_then(|value| v8::Local::::try_from(value).ok()); + let Some(attached) = attached else { + let key = v8::String::new(scope, "error").ok_or("DataCloneError")?; + let error = result.get(scope, key.into()).map(|error| error.to_rust_string_lossy(scope)); + return Err(error.unwrap_or_else(|| "DataCloneError".into())); + }; + for index in 0..attached.length() { + let object = attached + .get_index(scope, index) + .and_then(|value| v8::Local::::try_from(value).ok()) + .ok_or("DataCloneError")?; + objects.push(v8::Global::new(scope, object)); + } + } + + let context = scope.get_current_context(); + let deserializer = v8::ValueDeserializer::new(scope, Box::new(Deserializer { objects }), &message.bytes); + if !deserializer.read_header(context).unwrap_or(false) { + return Err("DataCloneError: the message could not be read".into()); + } + deserializer + .read_value(context) + .ok_or_else(|| "DataCloneError: the message could not be read".into()) +} diff --git a/runtime/src/worker_support.rs b/runtime/src/worker_support.rs index 3211345..da678d1 100644 --- a/runtime/src/worker_support.rs +++ b/runtime/src/worker_support.rs @@ -160,11 +160,19 @@ pub fn install_worker_runtime(scope: &mut ContextScope) { return delivered; } - this.postMessage = function (data) { + // Replies arrive through __nsWorkerDeliver where this thread has a dispatcher; without one + // (console hosts, tests) postMessage waits a little for them. + var pushesMessages = typeof globalThis.__nsWorkerPushesMessages === 'function' && + globalThis.__nsWorkerPushesMessages(); + + this.postMessage = function (data, transfer) { if (terminated) { return; } if (canUseThreadedHost && workerId >= 0) { - globalThis.__nsWorkerPostMessage(workerId, data); + globalThis.__nsWorkerPostMessage(workerId, data, transfer); + if (pushesMessages) { + return; + } for (var attempt = 0; attempt < 100; attempt++) { var blockingMessages = globalThis.__nsWorkerPollMessagesBlocking(workerId, 5); diff --git a/runtime/src/worker_threads.rs b/runtime/src/worker_threads.rs index 0bf2bcf..8b536bc 100644 --- a/runtime/src/worker_threads.rs +++ b/runtime/src/worker_threads.rs @@ -1,5 +1,7 @@ +use crate::transfer::Message; use crate::Runtime; use parking_lot::{Mutex, RwLock}; +use std::cell::RefCell; use std::collections::HashMap; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::mpsc::{self, Receiver, Sender, TryRecvError}; @@ -11,7 +13,7 @@ use windows::Win32::System::Threading::{CreateEventW, SetEvent, INFINITE}; use windows::Win32::UI::WindowsAndMessaging::{MsgWaitForMultipleObjectsEx, MWMO_INPUTAVAILABLE, QS_ALLINPUT}; enum WorkerCommand { - PostMessage(Vec), + PostMessage(Message), /// Work handed to the worker's thread from another thread: a callback created in the worker /// (see [`run_on_worker_sync`]) or the release of one. Run(Box), @@ -74,9 +76,8 @@ pub(crate) fn post_to_worker(isolate: usize, f: impl FnOnce() + Send + 'static) sent } -#[derive(Debug)] enum WorkerEvent { - Message(Vec), + Message(Message), Error(String), Exited, } @@ -91,13 +92,24 @@ struct WorkerHandle { join: thread::JoinHandle<()>, } -#[derive(Debug)] pub enum PolledWorkerEvent { - Message(Vec), + Message(Message), Error(String), Exited, } +thread_local! { + static OUTGOING: RefCell> = const { RefCell::new(Vec::new()) }; +} + +pub(crate) fn queue_outgoing(message: Message) { + OUTGOING.with(|outgoing| outgoing.borrow_mut().push(message)); +} + +fn take_outgoing() -> Vec { + OUTGOING.with(|outgoing| std::mem::take(&mut *outgoing.borrow_mut())) +} + static NEXT_WORKER_ID: AtomicU64 = AtomicU64::new(1); static WORKERS: OnceLock>> = OnceLock::new(); @@ -118,10 +130,10 @@ fn worker_bootstrap_script(source: &str, filename: &str) -> Result Resu } } - let forward_outbox = |runtime: &mut Runtime| { + let forward_outbox = || { let mut sent = false; - for result in runtime.drain_outbox_bytes() { + for message in take_outgoing() { sent = true; - let _ = evt_tx.send(match result { - Ok(b) => WorkerEvent::Message(b), - Err(e) => WorkerEvent::Error(e), - }); + let _ = evt_tx.send(WorkerEvent::Message(message)); } if sent { notify_creator(); } }; - forward_outbox(&mut runtime); + forward_outbox(); // Commands, Windows messages (WinRT calls marshaled to this STA thread, async // completions) and the worker's timers, until terminated. 'run: loop { loop { match cmd_rx.try_recv() { - Ok(WorkerCommand::PostMessage(bytes)) => runtime.dispatch_to_worker(&bytes), + Ok(WorkerCommand::PostMessage(message)) => runtime.dispatch_to_worker(message), Ok(WorkerCommand::Run(job)) => job(), Ok(WorkerCommand::Terminate) | Err(TryRecvError::Disconnected) => break 'run, Err(TryRecvError::Empty) => break, } - forward_outbox(&mut runtime); + forward_outbox(); } crate::pump_messages(); crate::timers::pump(); - forward_outbox(&mut runtime); - let timeout = if crate::timers::has_pending() { 10 } else { INFINITE }; + if crate::animation_frames::worker_frame_wait() == Some(Duration::ZERO) { + crate::animation_frames::pump(); + } + forward_outbox(); + let mut timeout = if crate::timers::has_pending() { 10 } else { INFINITE }; + if let Some(wait) = crate::animation_frames::worker_frame_wait() { + timeout = timeout.min(wait.as_millis().max(1) as u32); + } unsafe { MsgWaitForMultipleObjectsEx(Some(&[wake.handle()]), timeout, QS_ALLINPUT, MWMO_INPUTAVAILABLE); } } + // Its addons' cleanup hooks, while its isolate is alive. + crate::node_api::teardown_thread(); worker_executors().lock().remove(&isolate); let _ = evt_tx.send(WorkerEvent::Exited); notify_creator(); @@ -274,7 +291,7 @@ pub fn create_worker(app_root: String, source: String, filename: String) -> Resu Ok(worker_id) } -pub fn post_message(worker_id: u64, payload_bytes: Vec) -> Result<(), String> { +pub fn post_message(worker_id: u64, message: Message) -> Result<(), String> { let workers = workers().read(); let Some(worker) = workers.get(&worker_id) else { return Err(format!("Unknown worker id: {worker_id}")); @@ -282,7 +299,7 @@ pub fn post_message(worker_id: u64, payload_bytes: Vec) -> Result<(), String let sent = worker .tx - .send(WorkerCommand::PostMessage(payload_bytes)) + .send(WorkerCommand::PostMessage(message)) .map_err(|e| format!("Failed to send worker message: {e}")); worker.wake.signal(); sent @@ -292,7 +309,7 @@ fn collect_events(rx: &Receiver) -> Vec { let mut events = Vec::new(); loop { match rx.try_recv() { - Ok(WorkerEvent::Message(bytes)) => events.push(PolledWorkerEvent::Message(bytes)), + Ok(WorkerEvent::Message(message)) => events.push(PolledWorkerEvent::Message(message)), Ok(WorkerEvent::Error(err)) => events.push(PolledWorkerEvent::Error(err)), Ok(WorkerEvent::Exited) => { events.push(PolledWorkerEvent::Exited); @@ -340,7 +357,7 @@ pub fn poll_events_blocking( let mut events = Vec::new(); match rx.recv_timeout(Duration::from_millis(timeout_ms)) { - Ok(WorkerEvent::Message(bytes)) => events.push(PolledWorkerEvent::Message(bytes)), + Ok(WorkerEvent::Message(message)) => events.push(PolledWorkerEvent::Message(message)), Ok(WorkerEvent::Error(err)) => events.push(PolledWorkerEvent::Error(err)), Ok(WorkerEvent::Exited) => events.push(PolledWorkerEvent::Exited), Err(_) => return Ok(events), From d12793c90ed8969f45960b59e71b1d3cb5c0d8e0 Mon Sep 17 00:00:00 2001 From: Osei Fortune Date: Sun, 4 Oct 2026 11:32:37 -0400 Subject: [PATCH 2/2] chore: 1.0.0-beta.8 --- template/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/template/package.json b/template/package.json index 27eafaf..9caed01 100644 --- a/template/package.json +++ b/template/package.json @@ -1,6 +1,6 @@ { "name": "@nativescript/windows", - "version": "0.1.0-beta.4", + "version": "1.0.0-beta.8", "description": "NativeScript Windows runtime with a WinUI 3 app template", "repository": { "type": "git",