From 9d1385fe44648f2433521ecd5d2f233047d1b361 Mon Sep 17 00:00:00 2001 From: Pierre Le Fevre Date: Sat, 16 May 2026 14:21:34 +0200 Subject: [PATCH] Implement SharedArrayBuffer, Atomics, Int32Array, and importScripts (Phase 19) - shared_memory.rs: global registry of Arc (Mutex+Condvar) keyed by u64 ID - SharedArrayBuffer constructor: allocates a shared buffer, stores ID as __sab_id__ - Int32Array constructor: view over SharedArrayBuffer or standalone (also backed by SAB) - Atomics object: load, store, add, sub, and, or, xor, exchange, compareExchange, wait, notify all operations go through the shared buffer registry for cross-thread visibility - structured_clone: TAG_SHARED_ARRAY_BUFFER transfers by ID only (Arc-clone semantics) - importScripts: fetches scripts synchronously, queues them in thread-local IMPORT_QUEUE, drained by run_worker_script after each vm.execute() call - fetch.rs: fetch_text_sync() helper for importScripts synchronous fetch - 13 new Atomics/SAB unit tests + 4 shared_memory tests, all passing Co-Authored-By: Claude Sonnet 4.6 --- crates/js/src/builtins.rs | 741 ++++++++++++++++++++++++++++++ crates/js/src/fetch.rs | 12 + crates/js/src/lib.rs | 1 + crates/js/src/shared_memory.rs | 160 +++++++ crates/js/src/structured_clone.rs | 47 ++ crates/js/src/worker_global.rs | 61 +++ 6 files changed, 1022 insertions(+) create mode 100644 crates/js/src/shared_memory.rs diff --git a/crates/js/src/builtins.rs b/crates/js/src/builtins.rs index 666e1e5..ec834ed 100644 --- a/crates/js/src/builtins.rs +++ b/crates/js/src/builtins.rs @@ -265,6 +265,15 @@ pub fn init_builtins(vm: &mut Vm) { // Register global utility functions. init_global_functions(vm); + // Register SharedArrayBuffer constructor. + init_shared_array_buffer_builtin(vm); + + // Register TypedArray constructors (Int32Array). + init_typed_arrays_builtin(vm); + + // Register Atomics built-in object. + init_atomics_builtin(vm); + // Execute JS preamble for callback-based methods. init_js_preamble(vm); } @@ -7181,3 +7190,735 @@ Promise.any = __Promise_static_any; }; let _ = vm.execute(&func); } + +// ── SharedArrayBuffer ────────────────────────────────────────── + +fn shared_array_buffer_constructor( + args: &[Value], + ctx: &mut NativeContext, +) -> Result { + let byte_len = args.first().map(|v| v.to_number() as usize).unwrap_or(0); + let (id, _) = crate::shared_memory::alloc_shared_buffer(byte_len); + + let mut obj = ObjectData::new(); + obj.insert_property( + "__sab_id__".to_string(), + Property::builtin(Value::Number(id as f64)), + ctx.shapes, + ); + obj.insert_property( + "byteLength".to_string(), + Property::builtin(Value::Number(byte_len as f64)), + ctx.shapes, + ); + let r = ctx.gc.alloc(HeapObject::Object(obj)); + Ok(Value::Object(r)) +} + +pub fn init_shared_array_buffer_builtin(vm: &mut Vm) { + let ctor = make_native( + &mut vm.gc, + "SharedArrayBuffer", + shared_array_buffer_constructor, + ); + vm.set_global("SharedArrayBuffer", Value::Function(ctor)); +} + +// ── Int32Array (minimal TypedArray for Atomics) ─────────────── + +fn make_int32_array_object( + gc: &mut Gc, + shapes: &mut ShapeTable, + sab_id: u64, + byte_offset: usize, + byte_len: usize, + buffer_val: Value, +) -> GcRef { + let length = byte_len / 4; + let mut obj = ObjectData::new(); + obj.insert_property( + "__typed_array_kind__".to_string(), + Property::builtin(Value::String("Int32".to_string())), + shapes, + ); + obj.insert_property( + "__sab_id__".to_string(), + Property::builtin(Value::Number(sab_id as f64)), + shapes, + ); + obj.insert_property( + "__byte_offset__".to_string(), + Property::builtin(Value::Number(byte_offset as f64)), + shapes, + ); + obj.insert_property( + "length".to_string(), + Property::builtin(Value::Number(length as f64)), + shapes, + ); + obj.insert_property( + "byteLength".to_string(), + Property::builtin(Value::Number(byte_len as f64)), + shapes, + ); + obj.insert_property("buffer".to_string(), Property::builtin(buffer_val), shapes); + gc.alloc(HeapObject::Object(obj)) +} + +fn int32_array_constructor(args: &[Value], ctx: &mut NativeContext) -> Result { + match args.first() { + Some(Value::Object(sab_ref)) => { + // View over SharedArrayBuffer. + let sab_ref = *sab_ref; + let (sab_id, byte_len) = { + match ctx.gc.get(sab_ref) { + Some(HeapObject::Object(data)) => { + let id = data + .get_property("__sab_id__", ctx.shapes) + .map(|p| p.value.to_number() as u64) + .filter(|&id| id != 0) + .ok_or_else(|| { + RuntimeError::type_error( + "Int32Array: argument must be a SharedArrayBuffer", + ) + })?; + let bl = data + .get_property("byteLength", ctx.shapes) + .map(|p| p.value.to_number() as usize) + .unwrap_or(0); + (id, bl) + } + _ => { + return Err(RuntimeError::type_error( + "Int32Array: argument must be a SharedArrayBuffer", + )) + } + } + }; + let r = make_int32_array_object( + ctx.gc, + ctx.shapes, + sab_id, + 0, + byte_len, + Value::Object(sab_ref), + ); + Ok(Value::Object(r)) + } + Some(Value::Number(n)) => { + // New standalone buffer. + let length = (*n).max(0.0) as usize; + let byte_len = length * 4; + let (sab_id, _) = crate::shared_memory::alloc_shared_buffer(byte_len); + + // Build a SAB wrapper object first. + let mut sab_obj = ObjectData::new(); + sab_obj.insert_property( + "__sab_id__".to_string(), + Property::builtin(Value::Number(sab_id as f64)), + ctx.shapes, + ); + sab_obj.insert_property( + "byteLength".to_string(), + Property::builtin(Value::Number(byte_len as f64)), + ctx.shapes, + ); + let sab_ref = ctx.gc.alloc(HeapObject::Object(sab_obj)); + + let r = make_int32_array_object( + ctx.gc, + ctx.shapes, + sab_id, + 0, + byte_len, + Value::Object(sab_ref), + ); + Ok(Value::Object(r)) + } + _ => Err(RuntimeError::type_error( + "Int32Array: argument must be a SharedArrayBuffer or length", + )), + } +} + +pub fn init_typed_arrays_builtin(vm: &mut Vm) { + let ctor = make_native(&mut vm.gc, "Int32Array", int32_array_constructor); + vm.set_global("Int32Array", Value::Function(ctor)); +} + +// ── Atomics ─────────────────────────────────────────────────── + +/// Extract (sab_id, byte_offset) from a TypedArray argument. +fn typed_array_sab( + val: &Value, + gc: &Gc, + shapes: &ShapeTable, +) -> Result<(u64, usize), RuntimeError> { + if let Value::Object(r) = val { + if let Some(HeapObject::Object(data)) = gc.get(*r) { + let sab_id = data + .get_property("__sab_id__", shapes) + .map(|p| p.value.to_number() as u64) + .filter(|&id| id != 0) + .ok_or_else(|| { + RuntimeError::type_error( + "Atomics: argument must be an integer TypedArray backed by SharedArrayBuffer", + ) + })?; + let byte_offset = data + .get_property("__byte_offset__", shapes) + .map(|p| p.value.to_number() as usize) + .unwrap_or(0); + return Ok((sab_id, byte_offset)); + } + } + Err(RuntimeError::type_error( + "Atomics: first argument must be an Int32Array", + )) +} + +fn atomics_load(args: &[Value], ctx: &mut NativeContext) -> Result { + let ta = args + .first() + .ok_or_else(|| RuntimeError::type_error("Atomics.load: missing typed array"))?; + let index = args.get(1).map(|v| v.to_number() as usize).unwrap_or(0); + let (sab_id, byte_offset) = typed_array_sab(ta, ctx.gc, ctx.shapes)?; + let buf = crate::shared_memory::get_shared_buffer(sab_id) + .ok_or_else(|| RuntimeError::type_error("Atomics.load: invalid SharedArrayBuffer"))?; + let data = buf + .data + .lock() + .map_err(|_| RuntimeError::type_error("Atomics.load: buffer poisoned"))?; + let byte_idx = byte_offset + index * 4; + let val = crate::shared_memory::read_i32(&data, byte_idx) + .ok_or_else(|| RuntimeError::range_error("Atomics.load: index out of bounds"))?; + Ok(Value::Number(val as f64)) +} + +fn atomics_store(args: &[Value], ctx: &mut NativeContext) -> Result { + let ta = args + .first() + .ok_or_else(|| RuntimeError::type_error("Atomics.store: missing typed array"))?; + let index = args.get(1).map(|v| v.to_number() as usize).unwrap_or(0); + let new_val = args.get(2).map(|v| v.to_number() as i32).unwrap_or(0); + let (sab_id, byte_offset) = typed_array_sab(ta, ctx.gc, ctx.shapes)?; + let buf = crate::shared_memory::get_shared_buffer(sab_id) + .ok_or_else(|| RuntimeError::type_error("Atomics.store: invalid SharedArrayBuffer"))?; + let mut data = buf + .data + .lock() + .map_err(|_| RuntimeError::type_error("Atomics.store: buffer poisoned"))?; + let byte_idx = byte_offset + index * 4; + if !crate::shared_memory::write_i32(&mut data, byte_idx, new_val) { + return Err(RuntimeError::range_error( + "Atomics.store: index out of bounds", + )); + } + Ok(Value::Number(new_val as f64)) +} + +/// Generic read-modify-write operation for Atomics. +fn atomics_rmw( + args: &[Value], + ctx: &mut NativeContext, + op: fn(i32, i32) -> i32, +) -> Result { + let ta = args + .first() + .ok_or_else(|| RuntimeError::type_error("Atomics: missing typed array"))?; + let index = args.get(1).map(|v| v.to_number() as usize).unwrap_or(0); + let operand = args.get(2).map(|v| v.to_number() as i32).unwrap_or(0); + let (sab_id, byte_offset) = typed_array_sab(ta, ctx.gc, ctx.shapes)?; + let buf = crate::shared_memory::get_shared_buffer(sab_id) + .ok_or_else(|| RuntimeError::type_error("Atomics RMW: invalid SharedArrayBuffer"))?; + let mut data = buf + .data + .lock() + .map_err(|_| RuntimeError::type_error("Atomics RMW: buffer poisoned"))?; + let byte_idx = byte_offset + index * 4; + let old_val = crate::shared_memory::read_i32(&data, byte_idx) + .ok_or_else(|| RuntimeError::range_error("Atomics RMW: index out of bounds"))?; + let new_val = op(old_val, operand); + crate::shared_memory::write_i32(&mut data, byte_idx, new_val); + Ok(Value::Number(old_val as f64)) +} + +fn atomics_add(args: &[Value], ctx: &mut NativeContext) -> Result { + atomics_rmw(args, ctx, |a, b| a.wrapping_add(b)) +} + +fn atomics_sub(args: &[Value], ctx: &mut NativeContext) -> Result { + atomics_rmw(args, ctx, |a, b| a.wrapping_sub(b)) +} + +fn atomics_and(args: &[Value], ctx: &mut NativeContext) -> Result { + atomics_rmw(args, ctx, |a, b| a & b) +} + +fn atomics_or(args: &[Value], ctx: &mut NativeContext) -> Result { + atomics_rmw(args, ctx, |a, b| a | b) +} + +fn atomics_xor(args: &[Value], ctx: &mut NativeContext) -> Result { + atomics_rmw(args, ctx, |a, b| a ^ b) +} + +fn atomics_exchange(args: &[Value], ctx: &mut NativeContext) -> Result { + atomics_rmw(args, ctx, |_old, new| new) +} + +fn atomics_compare_exchange( + args: &[Value], + ctx: &mut NativeContext, +) -> Result { + let ta = args + .first() + .ok_or_else(|| RuntimeError::type_error("Atomics.compareExchange: missing typed array"))?; + let index = args.get(1).map(|v| v.to_number() as usize).unwrap_or(0); + let expected = args.get(2).map(|v| v.to_number() as i32).unwrap_or(0); + let replacement = args.get(3).map(|v| v.to_number() as i32).unwrap_or(0); + let (sab_id, byte_offset) = typed_array_sab(ta, ctx.gc, ctx.shapes)?; + let buf = crate::shared_memory::get_shared_buffer(sab_id) + .ok_or_else(|| RuntimeError::type_error("Atomics.compareExchange: invalid SAB"))?; + let mut data = buf + .data + .lock() + .map_err(|_| RuntimeError::type_error("Atomics.compareExchange: buffer poisoned"))?; + let byte_idx = byte_offset + index * 4; + let old_val = crate::shared_memory::read_i32(&data, byte_idx) + .ok_or_else(|| RuntimeError::range_error("Atomics.compareExchange: index out of bounds"))?; + if old_val == expected { + crate::shared_memory::write_i32(&mut data, byte_idx, replacement); + } + Ok(Value::Number(old_val as f64)) +} + +fn atomics_wait(args: &[Value], ctx: &mut NativeContext) -> Result { + let ta = args + .first() + .ok_or_else(|| RuntimeError::type_error("Atomics.wait: missing typed array"))?; + let index = args.get(1).map(|v| v.to_number() as usize).unwrap_or(0); + let expected = args.get(2).map(|v| v.to_number() as i32).unwrap_or(0); + let timeout_ms = args.get(3).map(|v| v.to_number()).unwrap_or(f64::INFINITY); + let (sab_id, byte_offset) = typed_array_sab(ta, ctx.gc, ctx.shapes)?; + let buf = crate::shared_memory::get_shared_buffer(sab_id) + .ok_or_else(|| RuntimeError::type_error("Atomics.wait: invalid SharedArrayBuffer"))?; + let byte_idx = byte_offset + index * 4; + + let guard = buf + .data + .lock() + .map_err(|_| RuntimeError::type_error("Atomics.wait: buffer poisoned"))?; + let current = crate::shared_memory::read_i32(&guard, byte_idx) + .ok_or_else(|| RuntimeError::range_error("Atomics.wait: index out of bounds"))?; + if current != expected { + return Ok(Value::String("not-equal".to_string())); + } + + if timeout_ms <= 0.0 { + return Ok(Value::String("timed-out".to_string())); + } + + if timeout_ms.is_infinite() || timeout_ms.is_nan() { + let _guard = buf + .condvar + .wait(guard) + .map_err(|_| RuntimeError::type_error("Atomics.wait: condvar poisoned"))?; + Ok(Value::String("ok".to_string())) + } else { + let duration = std::time::Duration::from_millis(timeout_ms as u64); + let (_guard, timed_out) = buf + .condvar + .wait_timeout(guard, duration) + .map_err(|_| RuntimeError::type_error("Atomics.wait: condvar poisoned"))?; + if timed_out.timed_out() { + Ok(Value::String("timed-out".to_string())) + } else { + Ok(Value::String("ok".to_string())) + } + } +} + +fn atomics_notify(args: &[Value], ctx: &mut NativeContext) -> Result { + let ta = args + .first() + .ok_or_else(|| RuntimeError::type_error("Atomics.notify: missing typed array"))?; + let _index = args.get(1).map(|v| v.to_number() as usize).unwrap_or(0); + let count_arg = args.get(2).map(|v| v.to_number()); + let (sab_id, _byte_offset) = typed_array_sab(ta, ctx.gc, ctx.shapes)?; + let buf = crate::shared_memory::get_shared_buffer(sab_id) + .ok_or_else(|| RuntimeError::type_error("Atomics.notify: invalid SharedArrayBuffer"))?; + + match count_arg { + Some(n) if !n.is_infinite() && !n.is_nan() && n >= 0.0 => { + let count = n as u32; + for _ in 0..count { + buf.condvar.notify_one(); + } + Ok(Value::Number(count as f64)) + } + _ => { + buf.condvar.notify_all(); + Ok(Value::Number(f64::INFINITY)) + } + } +} + +pub fn init_atomics_builtin(vm: &mut Vm) { + let methods: &[NativeMethod] = &[ + ("load", atomics_load), + ("store", atomics_store), + ("add", atomics_add), + ("sub", atomics_sub), + ("and", atomics_and), + ("or", atomics_or), + ("xor", atomics_xor), + ("exchange", atomics_exchange), + ("compareExchange", atomics_compare_exchange), + ("wait", atomics_wait), + ("notify", atomics_notify), + ]; + + let mut data = ObjectData::new(); + for &(name, cb) in methods { + let f = make_native(&mut vm.gc, name, cb); + data.insert_property( + name.to_string(), + Property::builtin(Value::Function(f)), + &mut vm.shapes, + ); + } + let atomics_ref = vm.gc.alloc(HeapObject::Object(data)); + vm.set_global("Atomics", Value::Object(atomics_ref)); +} + +// ── Tests: SharedArrayBuffer, Int32Array, Atomics ───────────── + +#[cfg(test)] +mod sab_atomics_tests { + use super::*; + use crate::compiler; + use crate::parser::Parser; + + fn eval(src: &str) -> Result { + let mut vm = Vm::new(); + let prog = Parser::parse(src).map_err(|e| e.to_string())?; + let func = compiler::compile(&prog).map_err(|e| format!("{e:?}"))?; + let val = vm.execute(&func).map_err(|e| e.to_string())?; + Ok(val.to_js_string(&vm.gc)) + } + + #[test] + fn shared_array_buffer_constructor() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(16); + sab.byteLength; + "#, + ) + .unwrap(); + assert_eq!(result, "16"); + } + + #[test] + fn int32_array_from_sab() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(16); + var ta = new Int32Array(sab); + ta.length; + "#, + ) + .unwrap(); + assert_eq!(result, "4"); + } + + #[test] + fn int32_array_from_length() { + let result = eval( + r#" + var ta = new Int32Array(8); + ta.length; + "#, + ) + .unwrap(); + assert_eq!(result, "8"); + } + + #[test] + fn atomics_store_load() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(16); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 42); + Atomics.load(ta, 0); + "#, + ) + .unwrap(); + assert_eq!(result, "42"); + } + + #[test] + fn atomics_add_returns_old() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(16); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 10); + Atomics.add(ta, 0, 5); + "#, + ) + .unwrap(); + assert_eq!(result, "10"); + } + + #[test] + fn atomics_add_new_value() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(16); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 10); + Atomics.add(ta, 0, 5); + Atomics.load(ta, 0); + "#, + ) + .unwrap(); + assert_eq!(result, "15"); + } + + #[test] + fn atomics_sub() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(8); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 100); + Atomics.sub(ta, 0, 30); + Atomics.load(ta, 0); + "#, + ) + .unwrap(); + assert_eq!(result, "70"); + } + + #[test] + fn atomics_compare_exchange_success() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(8); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 5); + Atomics.compareExchange(ta, 0, 5, 99); + Atomics.load(ta, 0); + "#, + ) + .unwrap(); + assert_eq!(result, "99"); + } + + #[test] + fn atomics_compare_exchange_fail() { + let result = eval( + r#" + var sab = new SharedArrayBuffer(8); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 5); + Atomics.compareExchange(ta, 0, 999, 42); + Atomics.load(ta, 0); + "#, + ) + .unwrap(); + assert_eq!(result, "5"); + } + + #[test] + fn atomics_cross_thread_counter() { + use crate::structured_clone::{self, SerializedData}; + use crate::worker::WorkerMessage; + use crate::worker_global::run_worker_script; + use std::sync::mpsc; + use std::time::Duration; + + // Allocate a SAB via the JS API on the "main thread". + let sab_bytes = { + let mut vm = Vm::new(); + let prog = Parser::parse("new SharedArrayBuffer(4)").unwrap(); + let func = compiler::compile(&prog).unwrap(); + let sab_val = vm.execute(&func).unwrap(); + + // Extract the sab_id from the JS object so we can check it later. + if let Value::Object(r) = &sab_val { + if let Some(HeapObject::Object(data)) = vm.gc.get(*r) { + let id = data + .get_property("__sab_id__", &vm.shapes) + .map(|p| p.value.to_number() as u64) + .unwrap_or(0); + assert_ne!(id, 0, "SAB ID should be non-zero"); + } + } + + structured_clone::serialize(&sab_val, &vm.gc, &vm.shapes) + .unwrap() + .into_bytes() + }; + + // Recover the SAB ID from the serialized bytes so we can check the buffer later. + let sab_id = { + let mut vm = Vm::new(); + let sab_val = structured_clone::deserialize( + &SerializedData::from_bytes(sab_bytes.clone()), + &mut vm.gc, + &mut vm.shapes, + ); + if let Value::Object(r) = sab_val { + if let Some(HeapObject::Object(data)) = vm.gc.get(r) { + data.get_property("__sab_id__", &vm.shapes) + .map(|p| p.value.to_number() as u64) + .unwrap_or(0) + } else { + 0 + } + } else { + 0 + } + }; + assert_ne!(sab_id, 0); + + // Worker script: receive SAB, increment counter, postMessage 'done'. + let worker_script = r#" + self.onmessage = function(e) { + var sab = e.data; + var ta = new Int32Array(sab); + Atomics.add(ta, 0, 1); + self.postMessage('done'); + }; + "#; + + let (inbound_tx, inbound_rx) = mpsc::channel::(); + let (outbound_tx, outbound_rx) = mpsc::channel::(); + let (error_tx, error_rx) = mpsc::channel::(); + + let thread = std::thread::spawn(move || { + run_worker_script( + worker_script, + "http://localhost/worker.js", + inbound_rx, + outbound_tx, + error_tx, + ); + }); + + // Send the SAB to the worker. + inbound_tx.send(WorkerMessage { data: sab_bytes }).unwrap(); + + // Wait for 'done'. + let reply = outbound_rx + .recv_timeout(Duration::from_millis(2000)) + .expect("worker did not respond"); + let msg = { + let mut vm = Vm::new(); + let val = structured_clone::deserialize( + &SerializedData::from_bytes(reply.data), + &mut vm.gc, + &mut vm.shapes, + ); + val.to_js_string(&vm.gc) + }; + assert_eq!(msg, "done"); + + drop(inbound_tx); + thread.join().unwrap(); + assert!(error_rx.try_recv().is_err()); + + // The counter in the SAB should now be 1. + let buf = crate::shared_memory::get_shared_buffer(sab_id).unwrap(); + let data = buf.data.lock().unwrap(); + let counter = crate::shared_memory::read_i32(&data, 0).unwrap(); + assert_eq!(counter, 1); + } + + #[test] + fn atomics_wait_not_equal() { + // Atomics.wait returns "not-equal" when the value at index differs. + let result = eval( + r#" + var sab = new SharedArrayBuffer(4); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 7); + Atomics.wait(ta, 0, 99, 0); + "#, + ) + .unwrap(); + assert_eq!(result, "not-equal"); + } + + #[test] + fn atomics_wait_timeout() { + // Atomics.wait times out when the value matches but no notify comes. + let result = eval( + r#" + var sab = new SharedArrayBuffer(4); + var ta = new Int32Array(sab); + Atomics.store(ta, 0, 0); + Atomics.wait(ta, 0, 0, 1); + "#, + ) + .unwrap(); + assert_eq!(result, "timed-out"); + } + + #[test] + fn atomics_wait_notify_unblocks() { + use crate::shared_memory::{alloc_shared_buffer, get_shared_buffer, write_i32}; + use std::thread; + use std::time::Duration; + + // Allocate a SAB and build an Int32Array backed by it. + let (sab_id, _) = alloc_shared_buffer(4); + + // Notifier thread: wait a bit, write 1, then notify. + let notifier_id = sab_id; + let notifier = thread::spawn(move || { + thread::sleep(Duration::from_millis(20)); + let buf = get_shared_buffer(notifier_id).unwrap(); + let mut data = buf.data.lock().unwrap(); + write_i32(&mut data, 0, 1); + drop(data); + buf.condvar.notify_all(); + }); + + // Main: build a VM with the SAB, call Atomics.wait. + let result = { + let mut vm = Vm::new(); + + // Build the SAB JS object directly with our known sab_id. + let mut sab_obj = crate::vm::ObjectData::new(); + sab_obj.insert_property( + "__sab_id__".to_string(), + crate::vm::Property::builtin(Value::Number(sab_id as f64)), + &mut vm.shapes, + ); + sab_obj.insert_property( + "byteLength".to_string(), + crate::vm::Property::builtin(Value::Number(4.0)), + &mut vm.shapes, + ); + let sab_ref = vm.gc.alloc(HeapObject::Object(sab_obj)); + vm.set_global("__sab__", Value::Object(sab_ref)); + + let src = r#" + var ta = new Int32Array(__sab__); + Atomics.wait(ta, 0, 0, 500); + "#; + let prog = Parser::parse(src).unwrap(); + let func = compiler::compile(&prog).unwrap(); + let val = vm.execute(&func).unwrap(); + val.to_js_string(&vm.gc) + }; + + notifier.join().unwrap(); + assert_eq!(result, "ok"); + } +} diff --git a/crates/js/src/fetch.rs b/crates/js/src/fetch.rs index 3bb58b2..f438404 100644 --- a/crates/js/src/fetch.rs +++ b/crates/js/src/fetch.rs @@ -693,6 +693,18 @@ fn write_header_entries( // ── Registration ──────────────────────────────────────────────── +/// Synchronously fetch a URL and return the response body as a UTF-8 string. +/// +/// Used by `importScripts()` to load worker scripts without going through the +/// async Promise-based `fetch()` API. Blocks the calling thread. +pub fn fetch_text_sync(url: &str) -> Result { + let result = do_fetch(url, "GET", &[], None, None, "no-cors", "omit")?; + if !(200..300).contains(&result.status) { + return Err(format!("HTTP {}: {}", result.status, result.status_text)); + } + Ok(String::from_utf8_lossy(&result.body).into_owned()) +} + /// Register the `fetch` global function in the VM. pub fn init_fetch_api(vm: &mut Vm) { let fetch_fn = make_native(&mut vm.gc, "fetch", fetch_native); diff --git a/crates/js/src/lib.rs b/crates/js/src/lib.rs index 0ceb409..ebee1f0 100644 --- a/crates/js/src/lib.rs +++ b/crates/js/src/lib.rs @@ -16,6 +16,7 @@ pub mod location; pub mod parser; pub mod regex; pub mod shape; +pub mod shared_memory; pub mod storage; pub mod structured_clone; pub mod timers; diff --git a/crates/js/src/shared_memory.rs b/crates/js/src/shared_memory.rs new file mode 100644 index 0000000..0ffd890 --- /dev/null +++ b/crates/js/src/shared_memory.rs @@ -0,0 +1,160 @@ +//! Shared memory backing for SharedArrayBuffer. +//! +//! Maintains a process-wide registry of live shared buffers, identified by +//! a monotonically-increasing `u64` ID. A JS SharedArrayBuffer object stores +//! its ID as `__sab_id__` and both the worker and main thread look up the +//! same `Arc` via that ID. + +use std::collections::HashMap; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Condvar, Mutex, OnceLock}; + +// ── Registry ───────────────────────────────────────────────────────────────── + +static NEXT_ID: AtomicU64 = AtomicU64::new(1); + +fn registry() -> &'static Mutex>> { + static REGISTRY: OnceLock>>> = OnceLock::new(); + REGISTRY.get_or_init(|| Mutex::new(HashMap::new())) +} + +// ── SharedBuffer ───────────────────────────────────────────────────────────── + +/// The backing store for a `SharedArrayBuffer`. +/// +/// Wraps a `Mutex>` so that cross-thread reads and writes are safe. +/// The `Condvar` is used by `Atomics.wait` / `Atomics.notify`. +pub struct SharedBuffer { + pub data: Mutex>, + pub condvar: Condvar, +} + +impl SharedBuffer { + fn new(byte_len: usize) -> Arc { + Arc::new(Self { + data: Mutex::new(vec![0u8; byte_len]), + condvar: Condvar::new(), + }) + } +} + +// ── Public API ──────────────────────────────────────────────────────────────── + +/// Allocate a zeroed shared buffer of `byte_len` bytes. +/// +/// Returns the unique ID for the buffer and an `Arc` to it. The `Arc` is also +/// stored in the global registry so that any thread can look it up by ID. +pub fn alloc_shared_buffer(byte_len: usize) -> (u64, Arc) { + let id = NEXT_ID.fetch_add(1, Ordering::Relaxed); + let buf = SharedBuffer::new(byte_len); + registry() + .lock() + .expect("shared buffer registry poisoned") + .insert(id, buf.clone()); + (id, buf) +} + +/// Look up a shared buffer by its ID. +/// +/// Returns `None` if no buffer with this ID exists (e.g. it has been GC-ed +/// and its registry entry cleaned up, which we don't implement yet). +pub fn get_shared_buffer(id: u64) -> Option> { + registry() + .lock() + .expect("shared buffer registry poisoned") + .get(&id) + .cloned() +} + +// ── Atomic helpers ──────────────────────────────────────────────────────────── + +/// Read a little-endian `i32` at `byte_offset` inside a locked byte slice. +pub fn read_i32(data: &[u8], byte_offset: usize) -> Option { + if byte_offset + 4 > data.len() { + return None; + } + let bytes: [u8; 4] = data[byte_offset..byte_offset + 4].try_into().ok()?; + Some(i32::from_le_bytes(bytes)) +} + +/// Write a little-endian `i32` at `byte_offset` inside a mutable byte slice. +pub fn write_i32(data: &mut [u8], byte_offset: usize, val: i32) -> bool { + if byte_offset + 4 > data.len() { + return false; + } + data[byte_offset..byte_offset + 4].copy_from_slice(&val.to_le_bytes()); + true +} + +#[cfg(test)] +mod tests { + use super::*; + use std::thread; + use std::time::Duration; + + #[test] + fn alloc_and_get() { + let (id, _buf) = alloc_shared_buffer(16); + let got = get_shared_buffer(id); + assert!(got.is_some()); + } + + #[test] + fn read_write_i32() { + let mut data = vec![0u8; 8]; + assert!(write_i32(&mut data, 0, 42)); + assert_eq!(read_i32(&data, 0), Some(42)); + assert!(write_i32(&mut data, 4, -1)); + assert_eq!(read_i32(&data, 4), Some(-1)); + } + + #[test] + fn cross_thread_visibility() { + let (id, _) = alloc_shared_buffer(4); + + let writer = thread::spawn(move || { + let buf = get_shared_buffer(id).unwrap(); + let mut data = buf.data.lock().unwrap(); + write_i32(&mut data, 0, 99); + }); + writer.join().unwrap(); + + let buf = get_shared_buffer(id).unwrap(); + let data = buf.data.lock().unwrap(); + assert_eq!(read_i32(&data, 0), Some(99)); + } + + #[test] + fn wait_notify_synchronize() { + let (id, _) = alloc_shared_buffer(4); + + // Writer thread: write 1, then notify. + let writer_id = id; + let writer = thread::spawn(move || { + thread::sleep(Duration::from_millis(10)); + let buf = get_shared_buffer(writer_id).unwrap(); + let mut data = buf.data.lock().unwrap(); + write_i32(&mut data, 0, 1); + drop(data); + buf.condvar.notify_all(); + }); + + // Main thread: wait until value == 0 changes (i.e., wait while it's 0). + let buf = get_shared_buffer(id).unwrap(); + let guard = buf.data.lock().unwrap(); + let initial = read_i32(&guard, 0).unwrap(); + assert_eq!(initial, 0); + + let (woken_guard, _timed_out) = buf + .condvar + .wait_timeout(guard, Duration::from_millis(500)) + .unwrap(); + + // Read the value while still holding the guard. + let counter = read_i32(&woken_guard, 0).unwrap(); + drop(woken_guard); + + writer.join().unwrap(); + assert_eq!(counter, 1); + } +} diff --git a/crates/js/src/structured_clone.rs b/crates/js/src/structured_clone.rs index d791de3..b94f061 100644 --- a/crates/js/src/structured_clone.rs +++ b/crates/js/src/structured_clone.rs @@ -38,6 +38,7 @@ const TAG_DATE: u8 = 0x08; const TAG_MAP: u8 = 0x09; const TAG_SET: u8 = 0x0A; const TAG_ARRAY_BUFFER: u8 = 0x0B; +const TAG_SHARED_ARRAY_BUFFER: u8 = 0x0C; // ── Public types ──────────────────────────────────────────────────────────── @@ -166,6 +167,7 @@ enum ObjectKind { Map, Set, ArrayBuffer, + SharedArrayBuffer(u64, usize), PlainObject, } @@ -183,6 +185,18 @@ fn classify_object(gc: &Gc, shapes: &ShapeTable, r: GcRef) -> Object return ObjectKind::Date(ms); } + // SharedArrayBuffer — identified by __sab_id__ (must come before ArrayBuffer check). + if let Some(id_prop) = data.get_property("__sab_id__", shapes) { + let sab_id = id_prop.value.to_number() as u64; + if sab_id != 0 { + let byte_len = data + .get_property("byteLength", shapes) + .map(|p| p.value.to_number() as usize) + .unwrap_or(0); + return ObjectKind::SharedArrayBuffer(sab_id, byte_len); + } + } + if data.contains_key("__array_buffer_bytes__", shapes) { return ObjectKind::ArrayBuffer; } @@ -217,6 +231,13 @@ fn write_object( buf.extend_from_slice(&ms.to_le_bytes()); } + ObjectKind::SharedArrayBuffer(sab_id, byte_len) => { + // Transfer by ID — both threads see the same Arc. + buf.push(TAG_SHARED_ARRAY_BUFFER); + buf.extend_from_slice(&sab_id.to_le_bytes()); + buf.extend_from_slice(&(byte_len as u64).to_le_bytes()); + } + ObjectKind::ArrayBuffer => { if let Some(HeapObject::Object(data)) = gc.get(r) { let byte_len = data @@ -551,6 +572,32 @@ fn read_value( Some(Value::Object(set_ref)) } + TAG_SHARED_ARRAY_BUFFER => { + // Re-wrap the existing SAB by ID — no new allocation. + if *pos + 16 > buf.len() { + return None; + } + let id_bytes: [u8; 8] = buf[*pos..*pos + 8].try_into().ok()?; + *pos += 8; + let len_bytes: [u8; 8] = buf[*pos..*pos + 8].try_into().ok()?; + *pos += 8; + let sab_id = u64::from_le_bytes(id_bytes); + let byte_len = u64::from_le_bytes(len_bytes) as usize; + + let mut obj = ObjectData::new(); + obj.insert_property( + "__sab_id__".to_string(), + Property::builtin(Value::Number(sab_id as f64)), + shapes, + ); + obj.insert_property( + "byteLength".to_string(), + Property::builtin(Value::Number(byte_len as f64)), + shapes, + ); + Some(Value::Object(gc.alloc(HeapObject::Object(obj)))) + } + TAG_ARRAY_BUFFER => { let byte_len = read_u32(buf, pos)?; if *pos + byte_len as usize > buf.len() { diff --git a/crates/js/src/worker_global.rs b/crates/js/src/worker_global.rs index 9849b27..4d528ac 100644 --- a/crates/js/src/worker_global.rs +++ b/crates/js/src/worker_global.rs @@ -25,6 +25,8 @@ thread_local! { static CLOSE_REQUESTED: Cell = const { Cell::new(false) }; /// GcRef of the `self` global object, for dispatching onmessage. static SELF_GCREF: Cell> = const { Cell::new(None) }; + /// Script sources fetched by importScripts(), executed after the current vm.execute() returns. + static IMPORT_QUEUE: RefCell> = const { RefCell::new(Vec::new()) }; } // ── Base64 helpers (atob / btoa) ────────────────────────────────────────────── @@ -141,6 +143,17 @@ fn worker_atob(args: &[Value], ctx: &mut NativeContext) -> Result Result { + for arg in args { + let url = arg.to_js_string(ctx.gc); + let src = crate::fetch::fetch_text_sync(&url).map_err(|e| { + RuntimeError::type_error(format!("importScripts: failed to fetch '{url}': {e}")) + })?; + IMPORT_QUEUE.with(|q| q.borrow_mut().push(src)); + } + Ok(Value::Undefined) +} + fn worker_queue_microtask(args: &[Value], _ctx: &mut NativeContext) -> Result { match args.first() { Some(Value::Function(r)) => { @@ -183,6 +196,9 @@ pub fn init_dedicated_worker_globals( vm.define_native("postMessage", worker_post_message); vm.define_native("close", worker_close); + // Register importScripts. + vm.define_native("importScripts", worker_import_scripts); + // ── Build the WorkerLocation stub ──────────────────────────────────────── let location_obj = { let mut data = ObjectData::new(); @@ -338,6 +354,45 @@ fn split_host_path(rest: &str, protocol: &str) -> (String, String, String) { } } +// ── Import queue drain ──────────────────────────────────────────────────────── + +/// Drain scripts queued by `importScripts()` and execute them in `vm`. +/// +/// Returns `false` if execution of any script failed (error forwarded via `error_tx`). +fn drain_import_queue(vm: &mut Vm, error_tx: &mpsc::Sender) -> bool { + loop { + let src = IMPORT_QUEUE.with(|q| q.borrow_mut().pop()); + let src = match src { + Some(s) => s, + None => return true, + }; + let program = match parser::Parser::parse(&src) { + Ok(p) => p, + Err(e) => { + let _ = error_tx.send(ErrorRecord { + message: format!("importScripts SyntaxError: {e}"), + }); + return false; + } + }; + let func = match compiler::compile(&program) { + Ok(f) => f, + Err(e) => { + let _ = error_tx.send(ErrorRecord { + message: format!("importScripts CompileError: {e:?}"), + }); + return false; + } + }; + if let Err(e) = vm.execute(&func) { + let _ = error_tx.send(ErrorRecord { + message: format!("importScripts RuntimeError: {e}"), + }); + return false; + } + } +} + // ── Worker run-loop ─────────────────────────────────────────────────────────── /// Execute a worker script on the current thread. @@ -387,6 +442,11 @@ pub fn run_worker_script( return; } + // Drain any scripts queued by importScripts() during the initial execution. + if !drain_import_queue(&mut vm, &error_tx) { + return; + } + // Run the worker event loop. loop { if CLOSE_REQUESTED.with(|c| c.get()) { @@ -446,6 +506,7 @@ pub fn run_worker_script( OUTBOUND_TX.with(|cell| cell.borrow_mut().take()); CLOSE_REQUESTED.with(|c| c.set(false)); SELF_GCREF.with(|c| c.set(None)); + IMPORT_QUEUE.with(|q| q.borrow_mut().clear()); } /// Deserialize a `WorkerMessage` and dispatch it to the worker's `onmessage`. -- 2.51.2