use ::std::cell::RefCell;
use ::std::ffi::CString;
use ::std::ptr::NonNull;
use ::std::sync::atomic::{AtomicU32, Ordering};
use ::std::sync::mpsc::{self, Receiver, Sender};
use ::std::sync::OnceLock;
use dashmap::DashMap;
use mozjs::conversions::unsafe_jsstr_to_string;
use mozjs::glue::{
CopyJSStructuredCloneData, GetLengthOfJSStructuredCloneData, WriteBytesToJSStructuredCloneData,
};
use mozjs::jsapi::*;
use mozjs::jsval::{BooleanValue, Int32Value, JSVal, ObjectValue, StringValue, UndefinedValue};
use mozjs::realm::AutoRealm;
use mozjs::rooted;
use mozjs::rust::wrappers2 as w2;
use mozjs::rust::JSAutoStructuredCloneBufferWrapper;
use crate::require::cache_builtin;
static NEXT_THREAD_ID: AtomicU32 = AtomicU32::new(1);
static WORKER_REGISTRY: OnceLock<DashMap<u32, WorkerHandle>> = OnceLock::new();
fn worker_registry() -> &'static DashMap<u32, WorkerHandle> {
WORKER_REGISTRY.get_or_init(DashMap::new)
}
struct WorkerHandle {
sender: Sender<WorkerMessage>,
thread: Option<::std::thread::JoinHandle<()>>,
main_rx: Option<::std::sync::Mutex<Receiver<WorkerToMainMessage>>>,
}
enum WorkerMessage {
Data(Vec<u8>),
Terminate,
}
enum WorkerToMainMessage {
Data(Vec<u8>),
Error(String),
}
fn sc_clone_policy() -> CloneDataPolicy {
CloneDataPolicy {
allowIntraClusterClonableSharedObjects_: false,
allowSharedMemoryObjects_: false,
}
}
const SC_BOUNDARY_HELPERS_SRC: &str = r#"(function(g){
function _baoSCWalk(v, seen, visit) {
if (v === null || typeof v !== 'object') return false;
if (seen.has(v)) return false;
seen.add(v);
if (visit(v)) return true;
var descs = Object.getOwnPropertyDescriptors(v);
var keys = Object.keys(descs);
for (var i = 0; i < keys.length; i++) {
var d = descs[keys[i]];
if (d && 'value' in d && _baoSCWalk(d.value, seen, visit)) return true;
}
if (v instanceof Map) {
var entries = [];
var it = v.entries();
for (;;) { var e = it.next(); if (e.done) break; entries.push(e.value); }
for (var j = 0; j < entries.length; j++) {
if (_baoSCWalk(entries[j][0], seen, visit)) return true;
if (_baoSCWalk(entries[j][1], seen, visit)) return true;
}
} else if (v instanceof Set) {
var vals = [];
var it2 = v.values();
for (;;) { var e2 = it2.next(); if (e2.done) break; vals.push(e2.value); }
for (var k = 0; k < vals.length; k++) {
if (_baoSCWalk(vals[k], seen, visit)) return true;
}
}
return false;
}
Object.defineProperty(g, '__baoSCHasAggregateError', {
configurable: true, enumerable: false, writable: true,
value: function(v) {
return _baoSCWalk(v, new WeakSet(), function(o) {
return o instanceof AggregateError;
});
}
});
Object.defineProperty(g, '__baoSCDemoteAggregateErrors', {
configurable: true, enumerable: false, writable: true,
value: function(v) {
_baoSCWalk(v, new WeakSet(), function(o) {
if (o instanceof AggregateError) {
Object.setPrototypeOf(o, Error.prototype);
try { delete o.errors; } catch (e) {}
}
return false;
});
}
});
})(globalThis)"#;
#[allow(unsafe_op_in_unsafe_fn)]
pub(crate) unsafe fn install_sc_boundary_helpers(
cx: &mut mozjs::context::JSContext,
global: mozjs::rust::Handle<*mut JSObject>,
) {
let mut source_text = mozjs::rust::transform_str_to_source_text(SC_BOUNDARY_HELPERS_SRC);
let opts = mozjs::glue::NewCompileOptions(cx.raw_cx(), c"<sc-boundary-helpers>".as_ptr(), 1);
if opts.is_null() {
return;
}
let mut rval = UndefinedValue();
let rval_handle = MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut rval,
};
let _ = JS::Evaluate2(cx.raw_cx(), opts, &mut source_text, rval_handle);
libc::free(opts as *mut _);
rooted!(&in(cx) let _global_guard = global.get());
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe fn sc_boundary_helper_call(
raw_cx: *mut JSContext,
name: &str,
value: JSVal,
) -> ::std::result::Result<bool, ()> {
let mut wrapped = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(raw_cx));
let cx = &mut wrapped;
let global = CurrentGlobalOrNull(raw_cx);
if global.is_null() {
return Err(());
}
rooted!(&in(cx) let global_root = global);
let c_name = CString::new(name).map_err(|_| ())?;
let mut fn_val = UndefinedValue();
JS_GetProperty(
raw_cx,
global_root.handle().into(),
c_name.as_ptr(),
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut fn_val,
},
);
if !fn_val.is_object() {
if w2::JS_IsExceptionPending(cx) {
JS_ClearPendingException(raw_cx);
}
return Err(());
}
rooted!(&in(cx) let fn_obj = fn_val.to_object());
rooted!(&in(cx) let arg = value);
let args = HandleValueArray {
length_: 1,
elements_: &arg.get() as *const Value,
};
rooted!(&in(cx) let fn_handle = ObjectValue(fn_obj.get()));
let mut rval = UndefinedValue();
if !JS_CallFunctionValue(
raw_cx,
global_root.handle().into(),
fn_handle.handle().into(),
&args,
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut rval,
},
) {
if w2::JS_IsExceptionPending(cx) {
JS_ClearPendingException(raw_cx);
}
return Err(());
}
Ok(rval.is_boolean() && rval.to_boolean())
}
#[allow(unsafe_op_in_unsafe_fn)]
pub(crate) unsafe fn sc_serialize(raw_cx: *mut JSContext, value: JSVal) -> ::std::result::Result<Vec<u8>, ()> {
sc_serialize_with_transfer(raw_cx, value, None)
}
#[allow(unsafe_op_in_unsafe_fn)]
pub(crate) unsafe fn sc_serialize_with_transfer(
raw_cx: *mut JSContext,
value: JSVal,
transferable: ::std::option::Option<mozjs::rust::Handle<'_, Value>>,
) -> ::std::result::Result<Vec<u8>, ()> {
let bytes = sc_write_bytes(raw_cx, value, transferable)?;
if !sc_boundary_helper_call(raw_cx, "__baoSCHasAggregateError", value)? {
return Ok(bytes);
}
let mut wrapped = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(raw_cx));
let cx = &mut wrapped;
rooted!(&in(cx) let mut mirror = UndefinedValue());
if !sc_deserialize(raw_cx, &bytes, mirror.handle_mut()) {
return Err(());
}
rooted!(&in(cx) let mirror_val = mirror.get());
sc_boundary_helper_call(raw_cx, "__baoSCDemoteAggregateErrors", mirror_val.get())?;
sc_write_bytes(raw_cx, mirror_val.get(), None)
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe fn sc_write_bytes(
raw_cx: *mut JSContext,
value: JSVal,
transferable: ::std::option::Option<mozjs::rust::Handle<'_, Value>>,
) -> ::std::result::Result<Vec<u8>, ()> {
let mut wrapped_cx = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(raw_cx));
let cx = &mut wrapped_cx;
rooted!(&in(cx) let val = value);
rooted!(&in(cx) let no_transfer = UndefinedValue());
let transfer_handle = match transferable {
Some(h) => h,
None => no_transfer.handle(),
};
let scbuf = unsafe {
JSAutoStructuredCloneBufferWrapper::new(
StructuredCloneScope::DifferentProcess,
::std::ptr::null(),
)
};
let scdata = unsafe { &mut ((*scbuf.as_raw_ptr()).data_) };
let ok = unsafe {
w2::JS_WriteStructuredClone(
cx,
val.handle(),
scdata,
StructuredCloneScope::DifferentProcess,
&sc_clone_policy(),
::std::ptr::null(),
::std::ptr::null_mut(),
transfer_handle,
)
};
if !ok {
return Err(());
}
let nbytes = unsafe { GetLengthOfJSStructuredCloneData(scdata) };
let mut bytes = Vec::with_capacity(nbytes);
unsafe {
CopyJSStructuredCloneData(scdata, bytes.as_mut_ptr());
bytes.set_len(nbytes);
}
Ok(bytes)
}
#[allow(unsafe_op_in_unsafe_fn)]
pub(crate) unsafe fn sc_deserialize(
raw_cx: *mut JSContext,
bytes: &[u8],
rval: mozjs::gc::MutableHandleValue<'_>,
) -> bool {
let mut wrapped_cx = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(raw_cx));
let cx = &mut wrapped_cx;
let scbuf = unsafe {
JSAutoStructuredCloneBufferWrapper::new(
StructuredCloneScope::DifferentProcess,
::std::ptr::null(),
)
};
let scdata = unsafe { &mut ((*scbuf.as_raw_ptr()).data_) };
if !bytes.is_empty()
&& !unsafe { WriteBytesToJSStructuredCloneData(bytes.as_ptr(), bytes.len(), scdata) }
{
return false;
}
unsafe {
w2::JS_ReadStructuredClone(
cx,
scdata,
JS_STRUCTURED_CLONE_VERSION,
StructuredCloneScope::DifferentProcess,
rval,
&sc_clone_policy(),
::std::ptr::null(),
::std::ptr::null_mut(),
)
}
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe fn report_data_clone_error(raw_cx: *mut JSContext) {
let wrapped_cx = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(raw_cx));
if w2::JS_IsExceptionPending(&wrapped_cx) {
JS_ClearPendingException(raw_cx);
}
JS_ReportErrorUTF8(
raw_cx,
c"DataCloneError: The object could not be cloned.".as_ptr(),
);
}
fn worker_entry(
filename: String,
thread_id: u32,
receiver: Receiver<WorkerMessage>,
main_sender: Sender<WorkerToMainMessage>,
worker_data_bytes: Option<Vec<u8>>,
) {
let engine_handle = match bao_engine::context::ensure_engine_handle() {
Ok(h) => h,
Err(_) => return,
};
let _runtime = mozjs::rust::Runtime::new(engine_handle);
let mut ctx = match unsafe { bao_engine::context::JsContext::from_servo_runtime() } {
Ok(c) => c,
Err(_) => return,
};
ctx.set_global_setup(worker_global_setup);
let mut cx = ctx.cx();
let raw_cx = ctx.raw_cx();
if !bao_engine::job_queue::JobQueue::init(&cx) {
return;
}
bao_engine::module_loader::ModuleLoader::init_thread_local(&cx);
bao_engine::module_loader::set_job_queue_drain(bao_engine::job_queue::JobQueue::drain);
let source = match ::std::fs::read_to_string(&filename) {
Ok(s) => s,
Err(e) => {
let msg = format!("Worker: failed to read '{}': {}", filename, e);
let _ = main_sender.send(WorkerToMainMessage::Error(msg));
return;
}
};
let bootstrap = format!(
r#"(function() {{
// workerData — deserialized from structured-clone bytes by the host before
// this script runs; non-enumerable raw value is deleted after capture.
var workerData = (typeof __baoWorkerDataRaw === 'undefined') ? null : __baoWorkerDataRaw;
delete self.__baoWorkerDataRaw;
var __pendingMessages = [];
// self.postMessage — structured clone via native host fn; uncloneable
// values (functions, ...) throw DataCloneError, matching Node.
self.postMessage = __baoPostToMain;
// Queue messages until onmessage handler is set. `data` arrives already
// deserialized from structured-clone bytes by the host.
self.__baoDeliverMessage = function(data) {{
if (typeof self.onmessage === 'function') {{
self.onmessage({{ data: data }});
}} else {{
__pendingMessages.push(data);
}}
}};
// When onmessage is set, deliver any queued messages
var __origOnMessage = null;
Object.defineProperty(self, 'onmessage', {{
configurable: true,
enumerable: true,
get: function() {{ return __origOnMessage; }},
set: function(fn) {{
__origOnMessage = fn;
// Deliver queued messages
while (__pendingMessages.length > 0 && typeof fn === 'function') {{
var data = __pendingMessages.shift();
fn({{ data: data }});
}}
}}
}});
// parentPort stub (worker_threads compat)
var parentPort = {{
postMessage: self.postMessage,
on: function() {{}},
once: function() {{}},
removeListener: function() {{}},
}};
// isMainThread is false inside workers
self.isMainThread = false;
self.threadId = {thread_id};
self.parentPort = parentPort;
// Execute the worker script
{source}
}})();"#,
thread_id = thread_id,
source = source,
);
let global_ptr = match ctx.ensure_realm_global(&mut cx, Some(worker_global_setup)) {
Ok(g) if !g.is_null() => g,
Ok(_) => {
let _ = main_sender.send(WorkerToMainMessage::Error(
"Worker realm global null after ensure_realm_global".into(),
));
return;
}
Err(e) => {
let _ = main_sender.send(WorkerToMainMessage::Error(format!(
"Worker realm init failed: {}",
e.message
)));
return;
}
};
rooted!(&in(cx) let global = global_ptr);
if let Some(wd_bytes) = worker_data_bytes.as_ref() {
let mut realm = AutoRealm::new_from_handle(&mut cx, global.handle());
let realm_cx: &mut mozjs::context::JSContext = &mut realm;
rooted!(&in(realm_cx) let mut wd_val = UndefinedValue());
let wd_ok = unsafe { sc_deserialize(realm_cx.raw_cx(), wd_bytes, wd_val.handle_mut()) };
if !wd_ok {
let _ = main_sender.send(WorkerToMainMessage::Error(
"Worker: workerData structured-clone deserialization failed".into(),
));
return;
}
unsafe {
JS_DefineProperty(
realm_cx.raw_cx(),
global.handle().into(),
c"__baoWorkerDataRaw".as_ptr(),
wd_val.handle().into(),
0u32, );
}
}
let eval_result = bao_engine::module_loader::ModuleLoader::eval_module_in_realm(
&mut cx,
&bootstrap,
&filename,
None,
global.handle(),
);
if let Err(e) = eval_result {
let msg = format!(
"Worker script error: {} ({}:{})",
e.message, e.filename, e.line
);
let _ = main_sender.send(WorkerToMainMessage::Error(msg));
return;
}
bao_engine::job_queue::JobQueue::drain(&mut cx);
WORKER_MAIN_SENDER.with(|s| {
*s.borrow_mut() = Some(main_sender);
});
loop {
match receiver.recv() {
Ok(WorkerMessage::Data(sc_bytes)) => {
deliver_message_to_worker(raw_cx, &sc_bytes);
bao_engine::job_queue::JobQueue::drain(&mut cx);
}
Ok(WorkerMessage::Terminate) | Err(_) => {
break;
}
}
}
}
unsafe fn worker_global_setup(
cx: &mut mozjs::context::JSContext,
global: mozjs::rust::Handle<*mut JSObject>,
) {
w2::JS_DefineFunction(
cx,
global,
c"__baoPostToMain".as_ptr(),
Some(worker_post_to_main),
1,
JSPROP_ENUMERATE as u32,
);
install_sc_boundary_helpers(cx, global);
rooted!(&in(cx) let global_val = ObjectValue(global.get()));
JS_DefineProperty(
cx.raw_cx(),
global.into(),
c"self".as_ptr(),
global_val.handle().into(),
(JSPROP_ENUMERATE | JSPROP_READONLY | JSPROP_PERMANENT) as u32,
);
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe extern "C" fn worker_post_to_main(cx: *mut JSContext, argc: u32, vp: *mut JSVal) -> bool {
let args = CallArgs::from_vp(vp, argc);
if argc == 0 {
JS_ReportErrorUTF8(cx, c"__baoPostToMain requires a value argument".as_ptr());
return false;
}
let data_val = *args.get(0).ptr;
let sc_bytes = match unsafe { sc_serialize(cx, data_val) } {
Ok(bytes) => bytes,
Err(()) => {
unsafe { report_data_clone_error(cx) };
return false;
}
};
WORKER_MAIN_SENDER.with(|sender| {
if let Some(tx) = sender.borrow().as_ref() {
let _ = tx.send(WorkerToMainMessage::Data(sc_bytes));
}
});
args.rval().set(UndefinedValue());
true
}
thread_local! {
static WORKER_MAIN_SENDER: RefCell<Option<Sender<WorkerToMainMessage>>> =
RefCell::new(None);
}
fn deliver_message_to_worker(raw_cx: *mut JSContext, sc_bytes: &[u8]) {
unsafe {
let global = match bao_engine::context::thread_realm_global() {
Some(g) if !g.is_null() => g,
_ => return,
};
let mut wrapped_cx = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(raw_cx));
let cx = &mut wrapped_cx;
rooted!(&in(cx) let global_root = global);
let mut realm = AutoRealm::new_from_handle(cx, global_root.handle());
let cx: &mut mozjs::context::JSContext = &mut realm;
rooted!(&in(cx) let mut data_val = UndefinedValue());
if !sc_deserialize(raw_cx, sc_bytes, data_val.handle_mut()) {
WORKER_MAIN_SENDER.with(|sender| {
if let Some(tx) = sender.borrow().as_ref() {
let _ = tx.send(WorkerToMainMessage::Error(
"Worker: message structured-clone deserialization failed".into(),
));
}
});
return;
}
rooted!(&in(cx) let mut fn_val = UndefinedValue());
JS_GetProperty(
raw_cx,
global_root.handle().into(),
c"__baoDeliverMessage".as_ptr(),
fn_val.handle_mut().into(),
);
if !fn_val.is_object() {
return;
}
rooted!(&in(cx) let fn_obj = fn_val.to_object());
let call_args_elements = [data_val.get()];
let call_args = HandleValueArray {
length_: 1,
elements_: call_args_elements.as_ptr() as *const Value,
};
rooted!(&in(cx) let fn_obj_val = ObjectValue(fn_obj.get()));
rooted!(&in(cx) let mut rval = UndefinedValue());
JS_CallFunctionValue(
raw_cx,
global_root.handle().into(),
fn_obj_val.handle().into(),
&call_args,
rval.handle_mut().into(),
);
}
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe extern "C" fn worker_constructor(cx: *mut JSContext, argc: u32, vp: *mut JSVal) -> bool {
let args = CallArgs::from_vp(vp, argc);
if argc == 0 {
JS_ReportErrorUTF8(cx, c"Worker requires a filename argument".as_ptr());
return false;
}
let filename_val = *args.get(0).ptr;
if !filename_val.is_string() {
JS_ReportErrorUTF8(
cx,
c"Worker first argument must be a string filename".as_ptr(),
);
return false;
}
let mut wrapped_cx = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(cx));
let filename = unsafe_jsstr_to_string(
wrapped_cx.raw_cx(),
NonNull::new_unchecked(filename_val.to_string()),
);
if filename.is_empty() {
JS_ReportErrorUTF8(cx, c"Worker: entry file path must not be empty".as_ptr());
return false;
}
let abs_filename = if ::std::path::Path::new(&filename).is_absolute() {
filename.clone()
} else {
match ::std::env::current_dir() {
Ok(cwd) => cwd.join(&filename).to_string_lossy().to_string(),
Err(_) => filename.clone(),
}
};
if !::std::path::Path::new(&abs_filename).exists() {
let msg = format!("Worker: entry file not found: {}", abs_filename);
let c_msg = ::std::ffi::CString::new(msg).unwrap_or_default();
JS_ReportErrorUTF8(cx, c_msg.as_ptr());
return false;
}
let mut worker_data_bytes: Option<Vec<u8>> = None;
if argc > 1 {
let opts_val = *args.get(1).ptr;
if opts_val.is_object() {
let opts_obj = opts_val.to_object();
let cx_ref = &mut wrapped_cx;
rooted!(&in(cx_ref) let opts_root = opts_obj);
let mut wd_val = UndefinedValue();
JS_GetProperty(
cx,
opts_root.handle().into(),
c"workerData".as_ptr(),
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut wd_val,
},
);
if !wd_val.is_undefined() {
match unsafe { sc_serialize(cx, wd_val) } {
Ok(bytes) => worker_data_bytes = Some(bytes),
Err(()) => {
unsafe { report_data_clone_error(cx) };
return false;
}
}
}
}
}
let thread_id = NEXT_THREAD_ID.fetch_add(1, Ordering::Relaxed);
let (main_to_worker_tx, main_to_worker_rx): (Sender<WorkerMessage>, Receiver<WorkerMessage>) =
mpsc::channel();
let (worker_to_main_tx, worker_to_main_rx): (
Sender<WorkerToMainMessage>,
Receiver<WorkerToMainMessage>,
) = mpsc::channel();
let worker_filename = abs_filename.clone();
let join_handle = ::std::thread::Builder::new()
.name(format!("bao-worker-{}", thread_id))
.spawn(move || {
worker_entry(
worker_filename,
thread_id,
main_to_worker_rx,
worker_to_main_tx,
worker_data_bytes,
);
});
let join_handle = match join_handle {
Ok(h) => h,
Err(e) => {
let msg = format!("Worker: failed to spawn thread: {}", e);
let c_msg = CString::new(msg).unwrap_or_default();
JS_ReportErrorUTF8(cx, c_msg.as_ptr());
return false;
}
};
worker_registry().insert(
thread_id,
WorkerHandle {
sender: main_to_worker_tx,
thread: Some(join_handle),
main_rx: Some(::std::sync::Mutex::new(worker_to_main_rx)),
},
);
let cx_ref = &mut wrapped_cx;
rooted!(&in(cx_ref) let worker_obj = w2::JS_NewPlainObject(cx_ref));
if worker_obj.get().is_null() {
args.rval().set(UndefinedValue());
return true;
}
rooted!(&in(cx_ref) let tid_val = Int32Value(thread_id as i32));
JS_DefineProperty(
cx,
worker_obj.handle().into(),
c"__threadId".as_ptr(),
tid_val.handle().into(),
0u32, );
w2::JS_DefineFunction(
cx_ref,
worker_obj.handle(),
c"postMessage".as_ptr(),
Some(worker_post_message),
1,
JSPROP_ENUMERATE as u32,
);
w2::JS_DefineFunction(
cx_ref,
worker_obj.handle(),
c"terminate".as_ptr(),
Some(worker_terminate),
0,
JSPROP_ENUMERATE as u32,
);
w2::JS_DefineFunction(
cx_ref,
worker_obj.handle(),
c"ref".as_ptr(),
Some(worker_noop),
0,
JSPROP_ENUMERATE as u32,
);
w2::JS_DefineFunction(
cx_ref,
worker_obj.handle(),
c"unref".as_ptr(),
Some(worker_noop),
0,
JSPROP_ENUMERATE as u32,
);
rooted!(&in(cx_ref) let tid_enum = Int32Value(thread_id as i32));
JS_DefineProperty(
cx,
worker_obj.handle().into(),
c"threadId".as_ptr(),
tid_enum.handle().into(),
(JSPROP_ENUMERATE | JSPROP_READONLY) as u32,
);
args.rval().set(ObjectValue(worker_obj.get()));
true
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe extern "C" fn worker_post_message(cx: *mut JSContext, argc: u32, vp: *mut JSVal) -> bool {
let args = CallArgs::from_vp(vp, argc);
let this_val = args.thisv();
if !this_val.is_object() {
JS_ReportErrorUTF8(
cx,
c"Worker.prototype.postMessage called on non-object".as_ptr(),
);
return false;
}
let this_obj = this_val.to_object();
let mut wrapped_cx = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(cx));
let cx_ref = &mut wrapped_cx;
rooted!(&in(cx_ref) let this_root = this_obj);
let mut tid_val = UndefinedValue();
JS_GetProperty(
cx,
this_root.handle().into(),
c"__threadId".as_ptr(),
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut tid_val,
},
);
if !tid_val.is_int32() {
JS_ReportErrorUTF8(cx, c"Worker: invalid threadId".as_ptr());
return false;
}
let thread_id = tid_val.to_int32() as u32;
if argc > 1 {
let transfer_val = *args.get(1).ptr;
if transfer_val.is_object() {
let transfer_obj = transfer_val.to_object();
rooted!(&in(cx_ref) let transfer_root = transfer_obj);
let mut len_val = UndefinedValue();
JS_GetProperty(
cx,
transfer_root.handle().into(),
c"length".as_ptr(),
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut len_val,
},
);
if len_val.is_int32() && len_val.to_int32() > 0 {
JS_ReportErrorUTF8(
cx,
c"DataCloneError: postMessage transfer list is not supported in Bao".as_ptr(),
);
return false;
}
}
}
let sc_bytes = if argc > 0 {
let data_val = *args.get(0).ptr;
match unsafe { sc_serialize(cx, data_val) } {
Ok(bytes) => bytes,
Err(()) => {
unsafe { report_data_clone_error(cx) };
return false;
}
}
} else {
match unsafe { sc_serialize(cx, UndefinedValue()) } {
Ok(bytes) => bytes,
Err(()) => {
unsafe { report_data_clone_error(cx) };
return false;
}
}
};
if let Some(handle) = worker_registry().get_mut(&thread_id) {
let _ = handle.sender.send(WorkerMessage::Data(sc_bytes));
}
args.rval().set(UndefinedValue());
true
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe extern "C" fn worker_terminate(cx: *mut JSContext, _argc: u32, vp: *mut JSVal) -> bool {
let args = CallArgs::from_vp(vp, _argc);
let this_val = args.thisv();
if !this_val.is_object() {
args.rval().set(UndefinedValue());
return true;
}
let this_obj = this_val.to_object();
let mut wrapped_cx = mozjs::context::JSContext::from_ptr(NonNull::new_unchecked(cx));
let cx_ref = &mut wrapped_cx;
rooted!(&in(cx_ref) let this_root = this_obj);
let mut tid_val = UndefinedValue();
JS_GetProperty(
cx,
this_root.handle().into(),
c"__threadId".as_ptr(),
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut tid_val,
},
);
if !tid_val.is_int32() {
args.rval().set(UndefinedValue());
return true;
}
let thread_id = tid_val.to_int32() as u32;
if let Some((_, mut handle)) = worker_registry().remove(&thread_id) {
let _ = handle.sender.send(WorkerMessage::Terminate);
if let Some(join) = handle.thread.take() {
let _ = join.join();
}
}
args.rval().set(UndefinedValue());
true
}
#[allow(unsafe_op_in_unsafe_fn)]
unsafe extern "C" fn worker_noop(_cx: *mut JSContext, _argc: u32, vp: *mut JSVal) -> bool {
let args = CallArgs::from_vp(vp, _argc);
args.rval().set(UndefinedValue());
true
}
#[derive(Debug, PartialEq)]
pub enum WorkerIncoming {
Data,
Error(String),
Empty,
}
pub fn worker_try_recv(
cx: &mut mozjs::context::JSContext,
thread_id: u32,
rval: mozjs::gc::MutableHandleValue<'_>,
) -> WorkerIncoming {
let msg = {
let handle = match worker_registry().get_mut(&thread_id) {
Some(h) => h,
None => return WorkerIncoming::Empty,
};
match handle.main_rx.as_ref() {
Some(rx) => rx.lock().ok().and_then(|rx| rx.try_recv().ok()),
None => return WorkerIncoming::Empty,
}
};
match msg {
Some(WorkerToMainMessage::Data(bytes)) => unsafe {
let global = match bao_engine::context::thread_realm_global() {
Some(g) if !g.is_null() => g,
_ => {
return WorkerIncoming::Error(
"worker_try_recv: main thread realm not initialized".into(),
)
}
};
rooted!(&in(cx) let global_root = global);
let mut realm = AutoRealm::new_from_handle(cx, global_root.handle());
let realm_cx: &mut mozjs::context::JSContext = &mut realm;
if sc_deserialize(realm_cx.raw_cx(), &bytes, rval) {
WorkerIncoming::Data
} else {
WorkerIncoming::Error(
"worker_try_recv: structured-clone deserialization failed".into(),
)
}
},
Some(WorkerToMainMessage::Error(msg)) => WorkerIncoming::Error(msg),
None => WorkerIncoming::Empty,
}
}
pub fn install(cx: &mut mozjs::context::JSContext) {
let raw_cx = unsafe { cx.raw_cx() };
rooted!(&in(cx) let exports = unsafe { w2::JS_NewPlainObject(cx) });
if exports.get().is_null() {
return;
}
unsafe {
let worker_fn = JS_NewFunction(
raw_cx,
Some(worker_constructor),
1, 0x400, c"Worker".as_ptr(),
);
if !worker_fn.is_null() {
let fn_obj = JS_GetFunctionObject(worker_fn);
rooted!(&in(cx) let fn_root = fn_obj);
rooted!(&in(cx) let proto = w2::JS_NewPlainObject(cx));
if !proto.get().is_null() {
w2::JS_DefineFunction(
cx,
proto.handle(),
c"postMessage".as_ptr(),
Some(worker_post_message),
1,
JSPROP_ENUMERATE as u32,
);
w2::JS_DefineFunction(
cx,
proto.handle(),
c"terminate".as_ptr(),
Some(worker_terminate),
0,
JSPROP_ENUMERATE as u32,
);
w2::JS_DefineFunction(
cx,
proto.handle(),
c"ref".as_ptr(),
Some(worker_noop),
0,
JSPROP_ENUMERATE as u32,
);
w2::JS_DefineFunction(
cx,
proto.handle(),
c"unref".as_ptr(),
Some(worker_noop),
0,
JSPROP_ENUMERATE as u32,
);
rooted!(&in(cx) let proto_val = ObjectValue(proto.get()));
JS_DefineProperty(
raw_cx,
fn_root.handle().into(),
c"prototype".as_ptr(),
proto_val.handle().into(),
0u32,
);
}
rooted!(&in(cx) let fn_val = ObjectValue(fn_root.get()));
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"Worker".as_ptr(),
fn_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
}
let source = r#"(function() {
var MC = (typeof globalThis.MessageChannel === 'function')
? globalThis.MessageChannel
: function MessageChannel() {
var queue1 = [];
var queue2 = [];
var onmsg1 = null;
var onmsg2 = null;
this.port1 = {
postMessage: function(data) {
if (typeof onmsg2 === 'function') {
onmsg2({ data: data });
} else {
queue2.push(data);
}
},
get onmessage() { return onmsg1; },
set onmessage(fn) {
onmsg1 = fn;
while (queue1.length > 0 && typeof fn === 'function') {
fn({ data: queue1.shift() });
}
},
close: function() {},
start: function() {},
addEventListener: function() {},
removeEventListener: function() {},
};
this.port2 = {
postMessage: function(data) {
if (typeof onmsg1 === 'function') {
onmsg1({ data: data });
} else {
queue1.push(data);
}
},
get onmessage() { return onmsg2; },
set onmessage(fn) {
onmsg2 = fn;
while (queue2.length > 0 && typeof fn === 'function') {
fn({ data: queue2.shift() });
}
},
close: function() {},
start: function() {},
addEventListener: function() {},
removeEventListener: function() {},
};
};
return MC;
})()"#;
let mut source_text = mozjs::rust::transform_str_to_source_text(source);
let mut rval = UndefinedValue();
let rval_handle = MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut rval,
};
let opts =
mozjs::glue::NewCompileOptions(raw_cx, c"<worker_threads:MessageChannel>".as_ptr(), 1);
if !opts.is_null() {
let ok = mozjs_sys::jsapi::JS::Evaluate2(raw_cx, opts, &mut source_text, rval_handle);
libc::free(opts as *mut _);
if ok && rval.is_object() {
rooted!(&in(cx) let mc_val = ObjectValue(rval.to_object()));
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"MessageChannel".as_ptr(),
mc_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
}
}
let mp_source = r#"(typeof globalThis.MessagePort === 'function'
? globalThis.MessagePort
: function MessagePort() {
throw new TypeError("worker_threads.MessagePort must be obtained from MessageChannel (port1/port2) or Worker — bare construction is not supported and would return an inert fake port.");
})"#;
let mut mp_text = mozjs::rust::transform_str_to_source_text(mp_source);
let mut mp_val = UndefinedValue();
let mp_opts =
mozjs::glue::NewCompileOptions(raw_cx, c"<worker_threads:MessagePort>".as_ptr(), 1);
if !mp_opts.is_null() {
let mp_ok = mozjs_sys::jsapi::JS::Evaluate2(
raw_cx,
mp_opts,
&mut mp_text,
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut mp_val,
},
);
libc::free(mp_opts as *mut _);
if mp_ok && mp_val.is_object() {
rooted!(&in(cx) let mp_obj = ObjectValue(mp_val.to_object()));
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"MessagePort".as_ptr(),
mp_obj.handle().into(),
JSPROP_ENUMERATE as u32,
);
}
}
let bc_source = r#"(typeof globalThis.BroadcastChannel === 'function'
? globalThis.BroadcastChannel
: (function() {
var registry = globalThis.__baoBroadcastRegistry || (globalThis.__baoBroadcastRegistry = {});
function BroadcastChannel(name) {
if (!(this instanceof BroadcastChannel)) return new BroadcastChannel(name);
if (typeof name !== 'string' || name === '') throw new TypeError('BroadcastChannel: name must be a non-empty string');
this.name = name;
this.onmessage = null;
this.onmessageerror = null;
this._closed = false;
this._listeners = [];
(registry[name] || (registry[name] = [])).push(this);
}
function _deliver(port, ev) {
var firstErr = null;
if (typeof port.onmessage === 'function') {
try { port.onmessage(ev); } catch (e) { if (!firstErr) firstErr = e; }
}
for (var i = 0; i < port._listeners.length; i++) {
try { port._listeners[i](ev); } catch (e) { if (!firstErr) firstErr = e; }
}
return firstErr;
}
BroadcastChannel.prototype.postMessage = function(message) {
if (this._closed) throw new Error('BroadcastChannel "' + this.name + '" is closed');
var peers = registry[this.name] || [];
var firstErr = null;
for (var i = 0; i < peers.length; i++) {
var peer = peers[i];
if (peer === this || peer._closed) continue;
var err = _deliver(peer, { data: message });
if (err && !firstErr) firstErr = err;
}
// Delivery completes for every peer before a throwing handler's
// error surfaces — no peer is starved, no error is swallowed.
if (firstErr) throw firstErr;
};
BroadcastChannel.prototype.close = function() {
if (this._closed) return;
this._closed = true;
var peers = registry[this.name] || [];
var idx = peers.indexOf(this);
if (idx >= 0) peers.splice(idx, 1);
if (peers.length === 0) delete registry[this.name];
};
BroadcastChannel.prototype.addEventListener = function(type, fn) {
if (typeof fn !== 'function') throw new TypeError('BroadcastChannel.addEventListener: listener must be a function');
// 'messageerror' can never fire in-process (no structured-clone
// failures without serialization); anything else is refused.
if (type !== 'message' && type !== 'messageerror') throw new TypeError('BroadcastChannel.addEventListener: unsupported event type "' + type + '"');
this._listeners.push(fn);
};
BroadcastChannel.prototype.removeEventListener = function(type, fn) {
if (type !== 'message' && type !== 'messageerror') return;
var idx = this._listeners.indexOf(fn);
if (idx >= 0) this._listeners.splice(idx, 1);
};
Object.defineProperty(BroadcastChannel.prototype, 'closed', { get: function() { return this._closed; }, configurable: true });
return BroadcastChannel;
})())"#;
let mut bc_text = mozjs::rust::transform_str_to_source_text(bc_source);
let mut bc_val = UndefinedValue();
let bc_opts = mozjs::glue::NewCompileOptions(
raw_cx,
c"<worker_threads:BroadcastChannel>".as_ptr(),
1,
);
if !bc_opts.is_null() {
let bc_ok = mozjs_sys::jsapi::JS::Evaluate2(
raw_cx,
bc_opts,
&mut bc_text,
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut bc_val,
},
);
libc::free(bc_opts as *mut _);
if bc_ok && bc_val.is_object() {
rooted!(&in(cx) let bc_obj = ObjectValue(bc_val.to_object()));
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"BroadcastChannel".as_ptr(),
bc_obj.handle().into(),
JSPROP_ENUMERATE as u32,
);
}
}
rooted!(&in(cx) let true_val = BooleanValue(true));
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"isMainThread".as_ptr(),
true_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
rooted!(&in(cx) let zero_val = Int32Value(0));
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"threadId".as_ptr(),
zero_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
rooted!(&in(cx) let undef_val = UndefinedValue());
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"workerData".as_ptr(),
undef_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"parentPort".as_ptr(),
undef_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
rooted!(&in(cx) let empty_obj = w2::JS_NewPlainObject(cx));
rooted!(&in(cx) let empty_obj_val = ObjectValue(empty_obj.get()));
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"resourceLimits".as_ptr(),
empty_obj_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
let share_env_source = r#"Symbol('nodejs.worker_threads.SHARE_ENV')"#;
let mut se_text = mozjs::rust::transform_str_to_source_text(share_env_source);
rooted!(&in(cx) let mut se_val = UndefinedValue());
let se_opts =
mozjs::glue::NewCompileOptions(raw_cx, c"<worker_threads:SHARE_ENV>".as_ptr(), 1);
if !se_opts.is_null() {
let se_ok = mozjs_sys::jsapi::JS::Evaluate2(
raw_cx,
se_opts,
&mut se_text,
se_val.handle_mut().into(),
);
libc::free(se_opts as *mut _);
if se_ok {
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c"SHARE_ENV".as_ptr(),
se_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
}
}
let utils_source = r#"({
getEnvironmentData: function() {},
setEnvironmentData: function() {},
getHeapSnapshot: function() { return {}; },
markAsUntransferable: function() { throw new Error('markAsUntransferable is not implemented in Bao'); },
moveMessagePortToContext: function() { throw new Error('moveMessagePortToContext is not implemented in Bao'); },
receiveMessageOnPort: function() { return undefined; },
})"#;
let mut ut_text = mozjs::rust::transform_str_to_source_text(utils_source);
let mut ut_val = UndefinedValue();
let ut_opts = mozjs::glue::NewCompileOptions(raw_cx, c"<worker_threads:utils>".as_ptr(), 1);
if !ut_opts.is_null() {
let ut_ok = mozjs_sys::jsapi::JS::Evaluate2(
raw_cx,
ut_opts,
&mut ut_text,
MutableHandle::<Value> {
_phantom_0: ::std::marker::PhantomData,
ptr: &mut ut_val,
},
);
libc::free(ut_opts as *mut _);
if ut_ok && ut_val.is_object() {
let utils_obj = ut_val.to_object();
rooted!(&in(cx) let utils_root = utils_obj);
for name in &[
"getEnvironmentData",
"setEnvironmentData",
"getHeapSnapshot",
"markAsUntransferable",
"moveMessagePortToContext",
"receiveMessageOnPort",
] {
let c_name = CString::new(*name).unwrap_or_default();
rooted!(&in(cx) let mut prop_val = UndefinedValue());
JS_GetProperty(
raw_cx,
utils_root.handle().into(),
c_name.as_ptr(),
prop_val.handle_mut().into(),
);
if !prop_val.is_undefined() {
JS_DefineProperty(
raw_cx,
exports.handle().into(),
c_name.as_ptr(),
prop_val.handle().into(),
JSPROP_ENUMERATE as u32,
);
}
}
}
}
}
cache_builtin(cx, "worker_threads", exports.get());
}