From 19f4d5be5c06392a3f8d51837c0cf16cc5b6e294 Mon Sep 17 00:00:00 2001 From: Tomasz Andrzejak Date: Tue, 15 Sep 2026 10:13:15 +0200 Subject: [PATCH] fix(runtime): correct resource ownership and async lifecycles - track imported ownership, borrow lifetimes, and deferred cleanup, - settle cancellations without prematurely releasing buffers, - keep future handles non-thenable and trace cached read promises - upgrade rquickjs to 0.13 and consolidate typed-array helpers - fix disposal symbols, empty-buffer alignment, and error handling --- Cargo.lock | 31 ++- README.md | 35 ++- crates/core/src/codegen.rs | 14 ++ crates/runtime/Cargo.toml | 2 +- crates/runtime/src/abi.rs | 30 ++- crates/runtime/src/bindings.rs | 110 ++++++---- crates/runtime/src/buffer.rs | 4 +- crates/runtime/src/call.rs | 183 ++++++++-------- crates/runtime/src/futures.rs | 146 +++++++++---- crates/runtime/src/interpreter.rs | 41 ++-- crates/runtime/src/lib.rs | 45 ++-- crates/runtime/src/module/mod.rs | 43 +++- crates/runtime/src/resources.rs | 341 ++++++++++++++++++++++++++++-- crates/runtime/src/result.rs | 4 +- crates/runtime/src/streams.rs | 123 ++++------- crates/runtime/src/task.rs | 258 ++++++++++++++++++---- crates/runtime/src/trivia.rs | 18 +- crates/runtime/src/typed_array.rs | 95 +++++++++ docs/runtime-intrinsics.md | 87 +++++++- tests/common/mod.rs | 2 +- tests/wasi.rs | 40 ++++ tests/wit/test.wit | 12 ++ tests/wit_types.rs | 49 ++++- 23 files changed, 1315 insertions(+), 398 deletions(-) create mode 100644 crates/runtime/src/typed_array.rs diff --git a/Cargo.lock b/Cargo.lock index fbfcac7..f75a692 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -483,6 +483,15 @@ dependencies = [ "unicode-segmentation", ] +[[package]] +name = "convert_case" +version = "0.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1af709f1f33454bf52eadfc8c78b3b9ef9cb26fb54d16dc9cd9a7299f899fd1b" +dependencies = [ + "unicode-segmentation", +] + [[package]] name = "cow-utils" version = "0.1.3" @@ -1647,7 +1656,7 @@ version = "3.6.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0fa55ea69990c90b888e9e77044410e304ce7f35de599dc6d0b5c1923d2e59af" dependencies = [ - "convert_case", + "convert_case 0.11.0", "ctor", "napi-derive-backend", "proc-macro2", @@ -1661,7 +1670,7 @@ version = "6.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "df4056ac7c18e4438ccf0edaed4340ca0d269278c8ec19284f7b23cb039fd0ae" dependencies = [ - "convert_case", + "convert_case 0.11.0", "proc-macro2", "quote", "semver", @@ -2539,9 +2548,9 @@ dependencies = [ [[package]] name = "rquickjs" -version = "0.12.2" +version = "0.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4e04e4eedfb060b503b5f0a2644abb890b0b3620d3fb674f9455f230014964e4" +checksum = "b7d96fb23e8ff51c8d4772ea84e44b5238543b7d81e9cc23c0219511a8f16482" dependencies = [ "rquickjs-core", "rquickjs-macro", @@ -2549,9 +2558,9 @@ dependencies = [ [[package]] name = "rquickjs-core" -version = "0.12.2" +version = "0.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "16e4f499ac5b943d97ee6dbc44f23c2c10426f420f7d2f1793d6318911b6608c" +checksum = "7c6dbfedfbf458dc119c21ccd032c7527d185210b11a1ae796cd04106e6db280" dependencies = [ "hashbrown 0.17.1", "relative-path", @@ -2560,11 +2569,11 @@ dependencies = [ [[package]] name = "rquickjs-macro" -version = "0.12.2" +version = "0.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3cbcc8219b70ee2faa08d5339f47f15741cba2ce0cb9640e8495f2ab51293f50" +checksum = "4446670c6ae57191ac325a4c9675fee5b31172533ad68753f7afa83759f082fd" dependencies = [ - "convert_case", + "convert_case 0.12.0", "fnv", "ident_case", "indexmap", @@ -2577,9 +2586,9 @@ dependencies = [ [[package]] name = "rquickjs-sys" -version = "0.12.2" +version = "0.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a13ac243b86a74120814ef7e9e30ad5a2c1199b7b9963b1cf7c84e4cdc1cad99" +checksum = "53d0aaff245bed1c6f3c39e477fb6b98d710d9d298bfec7c8dfc589bed0e5cef" dependencies = [ "bindgen", "cc", diff --git a/README.md b/README.md index 5cacf82..5727bb9 100644 --- a/README.md +++ b/README.md @@ -227,6 +227,15 @@ output.blockingWriteAndFlush(chunk); `[static]` methods are exposed on the resource class and `[constructor]` makes the class callable with `new`. +Owned imported resources support `[Symbol.dispose]()` for deterministic +cleanup. Passing one to a WIT `own` parameter transfers ownership and +invalidates the original wrapper. Borrowed wrappers expire when their component +call completes and cannot be transferred or disposed as owners. Resources lent +to an in-flight import remain alive and cannot be disposed or transferred until +that import completes. Ownership is restored for unconsumed stream/future writes +and for imported calls cancelled before they start. Await imports using an +incoming borrowed resource before returning from its export. + ### Async Exports Async exports are declared with the `async` keyword in WIT and implemented @@ -405,7 +414,7 @@ world async-value { ``` ```js -async function compute() { +export async function compute() { const { readable, writable } = wit.Future(); // Write the value (fire-and-forget; completes when reader reads) @@ -418,6 +427,16 @@ async function compute() { **Future type constants** follow the same pattern: `wit.Future.U32`, `wit.Future.STRING`, etc. +Future handles are deliberately not thenable: use `await readable.read()` to +obtain the payload. This keeps `return readable` in an async export from +implicitly consuming the future, and preserves nested future handles. +Repeated `read()` calls share the same Promise; a cancelled read rejects that +Promise and allows a subsequent read to retry. + +`wit.Future.from(value, type)` adapts a value or Promise and returns +`{ readable, completion }`. Return its `readable` from an async export rather +than returning the payload Promise directly, which JavaScript would await. + **FutureReadable methods:** | Method | Returns | Description | @@ -436,17 +455,25 @@ async function compute() { ### Resource Cleanup -Stream and future handles support +Owned imported resources, stream endpoints, and future endpoints support [Explicit Resource Management](https://github.com/tc39/proposal-explicit-resource-management) via `Symbol.dispose`. In environments that support `using`: ```js { - using stream = wit.Stream(); - // stream.writable and stream.readable are auto-dropped when leaving scope + const stream = wit.Stream(); + using writable = stream.writable; + using readable = stream.readable; + // Each endpoint is disposed when leaving scope. } ``` +The factory's `{ readable, writable }` pair is not itself disposable. +Complete an endpoint's pending operation (or cancel it and await its +settlement) before leaving its `using` scope or calling `.drop()`. +Host-requested task cancellation retains pending ABI buffers until the host +acknowledges each operation's completion or cancellation. + Otherwise, call `.drop()` explicitly to release handles. ## Node.js API diff --git a/crates/core/src/codegen.rs b/crates/core/src/codegen.rs index 371dedc..3c6fe95 100644 --- a/crates/core/src/codegen.rs +++ b/crates/core/src/codegen.rs @@ -137,6 +137,20 @@ impl<'a> EmitContext<'a> { return { readable, completion }; };"#, ); + } else { + self.multiline( + r#"wit.Future.from = function(value, type) { + const { readable, writable } = wit.Future(type); + const completion = (async () => { + try { + await writable.write(await value); + } finally { + writable.drop(); + } + })(); + return { readable, completion }; + };"#, + ); } } diff --git a/crates/runtime/Cargo.toml b/crates/runtime/Cargo.toml index 9bf8cd8..2273bfd 100644 --- a/crates/runtime/Cargo.toml +++ b/crates/runtime/Cargo.toml @@ -12,7 +12,7 @@ publish = false crate-type = ["cdylib"] [dependencies] -rquickjs = { version = "0.12", default-features = false, features = ["bindgen", "disable-assertions", "loader", "std", "macro"] } +rquickjs = { version = "0.13", default-features = false, features = ["bindgen", "disable-assertions", "loader", "std", "macro"] } wit-bindgen = "0.61" wit-dylib-ffi = { version = "0.1.0", git = "https://github.com/bytecodealliance/wasm-tools", tag = "v1.258.0", default-features = false, features = ["async-raw"] } heck = "0.5" diff --git a/crates/runtime/src/abi.rs b/crates/runtime/src/abi.rs index 4da84d3..1e79f7c 100644 --- a/crates/runtime/src/abi.rs +++ b/crates/runtime/src/abi.rs @@ -184,10 +184,10 @@ impl CopyEnd { /// Validate that the end is idle and ready for a new read/write. /// Returns `(handle, type_index)` on success. pub(crate) fn begin_op(&self) -> rquickjs::Result<(u32, u32)> { - if self.state.copying() { + if self.state != CopyState::Idle { return Err(rquickjs::Error::new_from_js( self.kind.label(), - "operation while copy in progress", + "operation requires an idle endpoint", )); } @@ -198,6 +198,30 @@ impl CopyEnd { Ok((h, self.type_index)) } + pub(crate) fn begin_transfer(&mut self, type_index: u32) -> rquickjs::Result { + if self.type_index != type_index { + return Err(rquickjs::Error::new_from_js( + self.kind.label(), + "matching WIT type", + )); + } + let (handle, _) = self.begin_op()?; + self.handle = None; + self.state = CopyState::Done; + Ok(handle) + } + + pub(crate) fn begin_drop(&mut self) -> rquickjs::Result> { + if self.state.copying() { + return Err(rquickjs::Error::new_from_js( + self.kind.label(), + "cancel and await the active operation before dropping", + )); + } + self.state = CopyState::Done; + Ok(self.handle.take()) + } + /// Validate that the end has an active async copy that can be cancelled. /// Returns `(handle, type_index)` on success. pub(crate) fn begin_cancel(&self) -> rquickjs::Result<(u32, u32)> { @@ -308,7 +332,6 @@ mod async_builtins { pub(crate) fn subtask_drop(task: u32); #[link_name = "[subtask-cancel]"] - #[allow(dead_code)] pub(crate) fn subtask_cancel(task: u32) -> u32; #[link_name = "[context-get-0]"] @@ -325,7 +348,6 @@ mod async_builtins { #[link(wasm_import_module = "[export]$root")] unsafe extern "C" { #[link_name = "[task-cancel]"] - #[allow(dead_code)] pub(crate) fn task_cancel(); #[link_name = "[backpressure-set]"] diff --git a/crates/runtime/src/bindings.rs b/crates/runtime/src/bindings.rs index fbfb575..825a59e 100644 --- a/crates/runtime/src/bindings.rs +++ b/crates/runtime/src/bindings.rs @@ -5,16 +5,17 @@ use rquickjs::function; use rquickjs::function::{Constructor, Rest, This}; use rquickjs::{Ctx, Function, Object, Value}; use smallvec::SmallVec; -use wit_dylib_ffi::{Resource, Wit}; +use wit_dylib_ffi::{Resource, Type, Wit}; use crate::CtxExt; use crate::futures::{make_future, register_future_classes}; +use crate::resources::{imported_resource_prototype, validate_imported_resource}; use crate::result::ResultBoundary; use crate::streams::{make_stream, register_stream_classes}; use crate::task::Pending; use crate::trivia::{fn_lookup, iface_lookup}; use crate::wit_imports::{FuncKind, WitInterface, classify, find_resource}; -use crate::{DetHashSet, DetIndexMap, QjsCallContext, coerce_fn}; +use crate::{DetHashMap, DetHashSet, DetIndexMap, QjsCallContext, coerce_fn}; /// Register all wit bindings on the js global scope. pub(crate) fn register(ctx: &rquickjs::Ctx<'_>, wit_def: Wit) -> rquickjs::Result<()> { @@ -73,12 +74,15 @@ fn register_resource_classes<'js>(ctx: &Ctx<'js>, wit: Wit) -> rquickjs::Result< let mut built: Vec<( usize, - Persistent>, - Persistent>, + Persistent>, + Persistent>, )> = Vec::new(); + let native_proto = imported_resource_prototype(ctx)?; for (index, group) in groups { - let prototype = Object::new(ctx.clone())?; + let proto = Object::new(ctx.clone())?; + proto.set_prototype(Some(&native_proto))?; + for (method, func_index) in group.methods { let js_func = Function::new( ctx.clone(), @@ -90,13 +94,13 @@ fn register_resource_classes<'js>(ctx: &Ctx<'js>, wit: Wit) -> rquickjs::Result< call_import(ctx, func_index, call_args) }, )?; - prototype.set(method.to_lower_camel_case(), js_func)?; + proto.set(method.to_lower_camel_case(), js_func)?; } let class: Constructor = match group.ctor { Some(func_index) => Constructor::new_prototype( ctx, - prototype.clone(), + proto.clone(), move |ctx: Ctx<'js>, args: Rest>| { call_import(ctx, func_index, SmallVec::from_vec(args.0)) }, @@ -105,7 +109,7 @@ fn register_resource_classes<'js>(ctx: &Ctx<'js>, wit: Wit) -> rquickjs::Result< let resource_name = group.resource.name(); Constructor::new_prototype( ctx, - prototype.clone(), + proto.clone(), move |ctx: Ctx<'js>, _args: Rest>| -> rquickjs::Result> { Err(rquickjs::Exception::throw_type( &ctx, @@ -126,8 +130,8 @@ fn register_resource_classes<'js>(ctx: &Ctx<'js>, wit: Wit) -> rquickjs::Result< built.push(( index, - Persistent::save(ctx, class.into_value()), - Persistent::save(ctx, prototype.into_value()), + Persistent::save(ctx, class), + Persistent::save(ctx, proto), )); } @@ -202,6 +206,38 @@ fn call_import<'js>( let wit_def = ctx.wit(); let func = wit_def.import_func(func_index); + let param_count = func.params().len(); + if args.len() < param_count { + return Err(rquickjs::Exception::throw_type( + &ctx, + &format!("{} requires {param_count} arguments", func.name()), + )); + } + + let mut resources = DetHashMap::default(); + for (mut ty, arg) in func.params().zip(&args) { + while let Type::Alias(alias) = ty { + ty = alias.ty(); + } + let (resource, owned) = match ty { + Type::Own(resource) => (resource, true), + Type::Borrow(resource) => (resource, false), + _ => continue, + }; + if resource.new().is_some() { + continue; + } + let handle = validate_imported_resource(resource, arg, owned)?; + if let Some(was_owned) = resources.insert((resource.index(), handle), owned) + && (was_owned || owned) + { + return Err(rquickjs::Exception::throw_type( + &ctx, + "a transferred resource cannot be used by another argument in the same call", + )); + } + } + let boundary = ResultBoundary::new(func.result()); let mut call = QjsCallContext::default(); for arg in args.into_iter().rev() { @@ -209,14 +245,15 @@ fn call_import<'js>( } if func.is_async() { + ctx.task().ensure_active(&ctx)?; let (promise, resolve, reject) = ctx.promise()?; if let Some(pending) = unsafe { func.call_import_async(&mut call) } { let handle = pending.subtask; let buffer = pending.buffer; - let resolve = Persistent::save(&ctx, resolve.into_value()); - let reject = Persistent::save(&ctx, reject.into_value()); + let resolve = Persistent::save(&ctx, resolve); + let reject = Persistent::save(&ctx, reject); let pending = Pending::ImportCall { func_index, call, @@ -226,15 +263,16 @@ fn call_import<'js>( }; ctx.task().register(handle, pending); } else { + call.complete_transfers(1); boundary .lift(&ctx, call.maybe_pop_value(&ctx)?)? - .settle(&resolve, &reject) - .expect("Failed to settle async import"); + .settle(&resolve, &reject)?; } Ok(promise.into_value()) } else { func.call_import_sync(&mut call); + call.complete_transfers(1); boundary .lift(&ctx, call.maybe_pop_value(&ctx)?)? .into_result(&ctx) @@ -332,6 +370,9 @@ fn build_async_exports<'js>( let then_cb = Function::new( ctx.clone(), coerce_fn(move |ctx: Ctx<'_>, args: Rest>| { + if ctx.task().is_cancelling() { + return Ok(Value::new_undefined(ctx)); + } let value = args .0 .into_iter() @@ -342,14 +383,16 @@ fn build_async_exports<'js>( let boundary = ResultBoundary::new(func.result()); let mut call = QjsCallContext::default(); - let value = boundary.lower_value(&ctx, value).unwrap_or_else(|e| { - panic!("Call failed '{}': {:?}", "async export", e) - }); + let value = boundary + .lower_value(&ctx, value) + .expect("Call failed 'async export'"); if let Some(value) = value { call.push_value(&ctx, value); } + ctx.task().finish_export(); func.call_task_return(&mut call); + call.complete_transfers(1); Ok(Value::new_undefined(ctx)) }), )?; @@ -357,6 +400,9 @@ fn build_async_exports<'js>( let catch_cb = Function::new( ctx.clone(), coerce_fn(move |ctx: Ctx<'_>, args: Rest>| { + if ctx.task().is_cancelling() { + return Ok(Value::new_undefined(ctx)); + } let reason = args .0 .into_iter() @@ -365,15 +411,17 @@ fn build_async_exports<'js>( let func = ctx.wit().export_func(func_index); let boundary = ResultBoundary::new(func.result()); let mut call = QjsCallContext::default(); - let value = boundary.lower_throw(&ctx, reason).unwrap_or_else(|e| { - panic!("Call failed '{}': {:?}", "async export", e) - }); + let value = boundary + .lower_throw(&ctx, reason) + .expect("Call failed 'async export'"); if let Some(value) = value { call.push_value(&ctx, value); } + ctx.task().finish_export(); func.call_task_return(&mut call); + call.complete_transfers(1); Ok(Value::new_undefined(ctx)) }), )?; @@ -414,21 +462,9 @@ fn register_cqjs_namespace(ctx: &rquickjs::Ctx<'_>, wit_def: Wit) -> rquickjs::R let ns = rquickjs::Object::new(ctx.clone())?; // Stream/future factories - ns.set( - "makeStream", - Function::new( - ctx.clone(), - coerce_fn(move |ctx: Ctx<'_>, args: Rest>| make_stream(ctx, args)), - )?, - )?; + ns.set("makeStream", Function::new(ctx.clone(), make_stream)?)?; - ns.set( - "makeFuture", - Function::new( - ctx.clone(), - coerce_fn(move |ctx: Ctx<'_>, args: Rest>| make_future(ctx, args)), - )?, - )?; + ns.set("makeFuture", Function::new(ctx.clone(), make_future)?)?; // Memory introspection ns.set( @@ -466,10 +502,8 @@ fn register_cqjs_namespace(ctx: &rquickjs::Ctx<'_>, wit_def: Wit) -> rquickjs::R ctx.clone(), coerce_fn( move |ctx: Ctx<'_>, _args: Rest>| -> rquickjs::Result> { - unsafe { - let rt = rquickjs::qjs::JS_GetRuntime(ctx.as_raw().as_ptr()); - rquickjs::qjs::JS_RunGC(rt); - } + ctx.run_gc(); + crate::resources::drain_resource_drops(&ctx); Ok(Value::new_undefined(ctx)) }, ), diff --git a/crates/runtime/src/buffer.rs b/crates/runtime/src/buffer.rs index fc29b33..74cd3cc 100644 --- a/crates/runtime/src/buffer.rs +++ b/crates/runtime/src/buffer.rs @@ -30,7 +30,7 @@ impl BufferGuard { fn allocate(size: usize, align: usize, zeroed: bool) -> Self { let layout = Layout::from_size_align(size, align).expect("invalid layout"); let ptr = if size == 0 { - std::ptr::NonNull::::dangling().as_ptr() + std::ptr::without_provenance_mut(layout.align()) } else if zeroed { unsafe { std::alloc::alloc_zeroed(layout) } } else { @@ -61,7 +61,7 @@ impl BufferGuard { /// /// # Safety /// The pointer must have been allocated with the given layout or be - /// dangling if `layout.size() == 0`. + /// non-null and aligned to the layout if `layout.size() == 0`. #[allow(dead_code)] pub(crate) unsafe fn from_raw(ptr: *mut u8, layout: Layout) -> Self { Self { ptr, layout } diff --git a/crates/runtime/src/call.rs b/crates/runtime/src/call.rs index c0a9a49..efa8551 100644 --- a/crates/runtime/src/call.rs +++ b/crates/runtime/src/call.rs @@ -1,16 +1,19 @@ //! `Call` trait implementation for quickjs to/from wit type conversions. use crate::CtxExt; -use crate::buffer::BufferGuard; use crate::futures::{FutureReadable, FutureWritable}; -use crate::resources::{exported_resource_to_handle, imported_resource_to_handle}; +use crate::resources::{ + borrow_imported_resource, exported_resource_to_handle, make_imported_borrowed, + make_imported_owned, transfer_imported_resource, +}; use crate::streams::{StreamReadable, StreamWritable}; use crate::tagged::decode_tagged; use crate::trivia::fn_lookup; -use crate::{BorrowedResource, QjsCallContext, with_ctx}; +use crate::typed_array::try_typed_array_copy; +use crate::{QjsCallContext, with_ctx}; use rquickjs::class::Class; use rquickjs::function::This; -use rquickjs::{Coerced, Constructor, Function, IntoJs, Persistent, Symbol, Value}; +use rquickjs::{CatchResultExt, Coerced, Constructor, Function, IntoJs, Persistent, Symbol, Value}; use smallvec::SmallVec; use wit_dylib_ffi::{ Call, Enum, Flags, Future, List, Map, Record, Resource, Stream, Tuple, Type, Variant, @@ -48,49 +51,6 @@ fn push_with(cx: &mut QjsCallContext, f: impl for<'js> FnOnce(&rquickjs::Ctx<'js }); } -/// Assign the imported resource prototype -fn set_imported_prototype<'js>( - ctx: &rquickjs::Ctx<'js>, - obj: &rquickjs::Object<'js>, - ty: Resource, -) { - let Some(proto) = ctx.resource_classes().prototype(ty.index()) else { - return; - }; - let Ok(proto_val) = proto.restore(ctx) else { - return; - }; - if let Some(proto_obj) = proto_val.into_object() { - let _ = obj.set_prototype(Some(&proto_obj)); - } -} - -fn copy_typed_array(slice: &[T]) -> (*const u8, usize, Layout) { - let count = slice.len(); - let byte_len = count - .checked_mul(std::mem::size_of::()) - .expect("typed array byte length overflow"); - - let buffer = BufferGuard::new_uninit(byte_len, std::mem::align_of::()); - if byte_len > 0 { - unsafe { - std::ptr::copy_nonoverlapping(slice.as_ptr().cast::(), buffer.ptr(), byte_len) - }; - } - let (ptr, layout) = buffer.into_raw(); - (ptr.cast_const(), count, layout) -} - -/// Extract a TypedArray and copy its bytes into a new buffer. -/// This has to be a macro because TypedArrayItem is not public. -macro_rules! try_typed_array_copy { - ($val:expr, $t:ty) => { - $val.as_object() - .and_then(|object| object.as_typed_array::<$t>()) - .map(|array| copy_typed_array::<$t>(array.as_ref())) - }; -} - impl Call for QjsCallContext { unsafe fn defer_deallocate(&mut self, ptr: *mut u8, layout: Layout) { self.deferred_deallocs.push((ptr, layout)); @@ -168,7 +128,7 @@ impl Call for QjsCallContext { let result = with_ctx(|ctx| { let val = persistent.restore(ctx).unwrap(); - match ty.ty() { + let result = match ty.ty() { Type::U8 => try_typed_array_copy!(val, u8), Type::S8 => try_typed_array_copy!(val, i8), Type::U16 => try_typed_array_copy!(val, u16), @@ -179,16 +139,18 @@ impl Call for QjsCallContext { Type::S64 => try_typed_array_copy!(val, i64), Type::F32 => try_typed_array_copy!(val, f32), Type::F64 => try_typed_array_copy!(val, f64), - _ => None, - } + _ => Ok(None), + }; + result.catch(ctx).expect("Failed to copy typed array") }); - result.map(|(ptr, count, layout)| { + result.map(|(buffer, count)| { + let (ptr, layout) = buffer.into_raw(); if layout.size() > 0 { - self.deferred_deallocs.push((ptr as *mut u8, layout)); + self.deferred_deallocs.push((ptr, layout)); } self.stack.pop(); - (ptr, count) + (ptr.cast_const(), count) }) } @@ -274,7 +236,8 @@ impl Call for QjsCallContext { // Nested option: { tag: "some", val } | { tag: "none" }. let (discriminant, payload) = decode_tagged(ctx, val, "nested option", [("none", false), ("some", true)]) - .unwrap_or_else(|err| panic!("invalid nested option: {err}")); + .catch(ctx) + .expect("invalid nested option"); if let Some(inner) = payload { self.stack.push(Persistent::save(ctx, inner)); } @@ -298,7 +261,8 @@ impl Call for QjsCallContext { "result", [("ok", ty.ok().is_some()), ("err", ty.err().is_some())], ) - .unwrap_or_else(|err| panic!("invalid result: {err}")); + .catch(ctx) + .expect("invalid result"); if let Some(inner) = payload { self.stack.push(Persistent::save(ctx, inner)); } @@ -315,7 +279,8 @@ impl Call for QjsCallContext { .map(|(name, payload_ty)| (name, payload_ty.is_some())); let (discriminant, payload) = decode_tagged(ctx, val, "variant", cases) - .unwrap_or_else(|err| panic!("invalid variant: {err}")); + .catch(ctx) + .expect("invalid variant"); if let Some(inner) = payload { self.stack.push(Persistent::save(ctx, inner)); @@ -347,8 +312,9 @@ impl Call for QjsCallContext { for (i, name) in ty.names().enumerate() { let set = obj .get::<_, Coerced>(fn_lookup(ctx, name)) - .map(|c| c.0) - .unwrap_or(false); + .catch(ctx) + .unwrap_or_else(|err| panic!("Failed to read flag '{name}': {err}")) + .0; if set { bits |= 1 << i; } @@ -364,7 +330,11 @@ impl Call for QjsCallContext { if ty.new().is_some() { exported_resource_to_handle(ctx, ty, &val) } else { - imported_resource_to_handle(&val) + let (handle, loan) = borrow_imported_resource(ty, &val) + .catch(ctx) + .expect("Failed to borrow imported resource"); + self.loans.push(loan); + handle } }) } @@ -376,7 +346,11 @@ impl Call for QjsCallContext { if ty.new().is_some() { exported_resource_to_handle(ctx, ty, &val) } else { - imported_resource_to_handle(&val) + let (handle, transfer) = transfer_imported_resource(ty, &val) + .catch(ctx) + .expect("Failed to transfer imported resource"); + self.transfers.push((self.transfer_group, transfer)); + handle } }) } @@ -405,23 +379,33 @@ impl Call for QjsCallContext { }); } - fn pop_future(&mut self, _ty: Future) -> u32 { + fn pop_future(&mut self, ty: Future) -> u32 { pop_with(self, |v| { if let Ok(class) = Class::::from_value(&v) { return class .borrow_mut() .end - .handle - .take() - .expect("future already transferred"); + .begin_transfer(ty.index() as u32) + .expect("future cannot be transferred"); } + if let Ok(class) = Class::::from_value(&v) { return class .borrow_mut() .end - .handle - .take() - .expect("future already transferred"); + .begin_transfer(ty.index() as u32) + .expect("future cannot be transferred"); + } + + let ctx = v.ctx().clone(); + + if is_thenable(&v) + .catch(&ctx) + .expect("Failed to inspect future value") + { + return crate::futures::lower_thenable(&ctx, ty.index() as u32, v) + .catch(&ctx) + .expect("Failed to lower WIT future"); } v.get().expect("expected future handle") }) @@ -436,22 +420,24 @@ impl Call for QjsCallContext { return class .borrow_mut() .end - .handle - .take() - .expect("stream already transferred"); + .begin_transfer(ty.index() as u32) + .expect("stream cannot be transferred"); } if let Ok(class) = Class::::from_value(&v) { return class .borrow_mut() .end - .handle - .take() - .expect("stream already transferred"); + .begin_transfer(ty.index() as u32) + .expect("stream cannot be transferred"); } - if is_iterable(ctx, &v) { + if is_iterable(ctx, &v) + .catch(ctx) + .expect("Failed to inspect stream value") + { return crate::streams::lower_iterable(ctx, ty.index() as u32, v) - .expect("lower async iterable to WIT stream"); + .catch(ctx) + .expect("Failed to lower WIT stream"); } v.get().expect("expected stream handle") @@ -660,15 +646,11 @@ impl Call for QjsCallContext { let rep = handle as usize; ctx.resources().get(rep).restore(ctx).unwrap() } else { - self.borrows.push(BorrowedResource { - handle, - drop_fn: ty.drop(), - }); - - let obj = rquickjs::Object::new(ctx.clone()).unwrap(); - obj.set("__cqjs_handle", handle).unwrap(); - set_imported_prototype(ctx, &obj, ty); - obj.into_value() + let (val, borrow) = make_imported_borrowed(ctx, ty, handle) + .catch(ctx) + .expect("Failed to wrap borrowed imported resource"); + self.borrows.push(borrow); + val }; self.push_value(ctx, val); }); @@ -684,10 +666,9 @@ impl Call for QjsCallContext { } val } else { - let obj = rquickjs::Object::new(ctx.clone()).unwrap(); - obj.set("__cqjs_handle", handle).unwrap(); - set_imported_prototype(ctx, &obj, ty); - obj.into_value() + make_imported_owned(ctx, ty, handle) + .catch(ctx) + .expect("Failed to wrap owned imported resource") }; self.push_value(ctx, val); }); @@ -739,19 +720,25 @@ impl Call for QjsCallContext { } } -fn is_iterable<'js>(ctx: &rquickjs::Ctx<'js>, value: &Value<'js>) -> bool { +fn is_iterable<'js>(ctx: &rquickjs::Ctx<'js>, value: &Value<'js>) -> rquickjs::Result { let Some(object) = value.as_object() else { - return false; + return Ok(false); }; - [ + for symbol in [ Symbol::async_iterator(ctx.clone()), Symbol::iterator(ctx.clone()), - ] - .into_iter() - .any(|symbol| { - object - .get::<_, Value>(symbol.as_atom()) - .is_ok_and(|value| value.is_function()) - }) + ] { + if object.get::<_, Value>(symbol.as_atom())?.is_function() { + return Ok(true); + } + } + Ok(false) +} + +fn is_thenable(value: &Value<'_>) -> rquickjs::Result { + let Some(object) = value.as_object() else { + return Ok(false); + }; + Ok(object.get::<_, Value>("then")?.is_function()) } diff --git a/crates/runtime/src/futures.rs b/crates/runtime/src/futures.rs index 7d7c0c0..3527c1c 100644 --- a/crates/runtime/src/futures.rs +++ b/crates/runtime/src/futures.rs @@ -1,7 +1,7 @@ //! component-model future operations using rquickjs classes. use rquickjs::class::{Class, JsClass, Trace}; use rquickjs::function::This; -use rquickjs::{Ctx, Function, Object, Persistent, Value}; +use rquickjs::{CatchResultExt, Ctx, Function, Object, Persistent, Promise, Value}; use rquickjs::{JsLifetime, function}; use crate::CtxExt; @@ -11,20 +11,22 @@ use crate::task::Pending; use crate::{QjsCallContext, resolve_promise, symbol_dispose, with_ctx}; #[derive(Trace, JsLifetime)] -pub(crate) struct FutureReadable { +pub(crate) struct FutureReadable<'js> { #[qjs(skip_trace)] pub(crate) end: CopyEnd, + promise: Option>, } -impl FutureReadable { +impl FutureReadable<'_> { fn new(type_index: u32, handle: u32) -> Self { Self { end: CopyEnd::new_future(type_index, handle), + promise: None, } } } -impl<'js> JsClass<'js> for FutureReadable { +impl<'js> JsClass<'js> for FutureReadable<'js> { const NAME: &'static str = "FutureReadable"; type Mutable = rquickjs::class::Writable; @@ -105,12 +107,12 @@ pub(crate) fn make_future_readable<'js>( Ok(instance.into_inner()) } -pub(crate) fn make_future<'js>( - ctx: Ctx<'js>, - args: function::Rest>, -) -> rquickjs::Result> { - let type_index: u32 = args.0[0].get()?; - let ty = ctx.wit().future(type_index as usize); +pub(crate) fn make_future<'js>(ctx: Ctx<'js>, type_index: u32) -> rquickjs::Result> { + let ty = ctx + .wit() + .iter_futures() + .nth(type_index as usize) + .ok_or_else(|| rquickjs::Exception::throw_range(&ctx, "unknown WIT future type"))?; let handles = unsafe { ty.new()() }; let tx_handle = (handles >> 32) as u32; @@ -126,13 +128,51 @@ pub(crate) fn make_future<'js>( Ok(result.into_value()) } +pub(crate) fn lower_thenable<'js>( + ctx: &Ctx<'js>, + type_index: u32, + thenable: Value<'js>, +) -> rquickjs::Result { + if !ctx.task().is_active() { + return Err(rquickjs::Error::new_from_js( + "thenable", + "WIT future lowering requires an active async call", + )); + } + + let wit: Object = ctx.globals().get("wit")?; + let future: Function = wit.get("Future")?; + let from: Function = future.get("from")?; + let pair: Object = from.call((thenable, type_index))?; + let readable: Value = pair.get("readable")?; + let readable = Class::::from_value(&readable)?; + let handle = readable.borrow_mut().end.begin_transfer(type_index)?; + + Ok(handle) +} + fn future_read<'js>( - this: This>, + this: This>>, ctx: Ctx<'js>, ) -> rquickjs::Result> { + { + let readable = this.0.borrow(); + if readable.end.handle.is_none() { + return Err(rquickjs::Exception::throw_type( + &ctx, + "future already dropped or transferred", + )); + } + if let Some(promise) = &readable.promise { + return Ok(promise.clone().into_value()); + } + } + + ctx.task().ensure_active(&ctx)?; let (handle, type_index) = this.0.borrow().end.begin_op()?; let (promise, resolve, reject) = ctx.promise()?; + this.0.borrow_mut().promise = Some(promise.clone()); let ty = ctx.wit().future(type_index as usize); let buffer = BufferGuard::new_zeroed(ty.abi_payload_size(), ty.abi_payload_align()); @@ -145,33 +185,33 @@ fn future_read<'js>( let pending = Pending::FutureRead { call, buffer, - resolve: Persistent::save(&ctx, resolve.into_value()), - reject: Persistent::save(&ctx, reject.into_value()), + resolve: Persistent::save(&ctx, resolve), + reject: Persistent::save(&ctx, reject), wrapper: Persistent::save(&ctx, this.0.into_inner().into_value()), }; ctx.task().register(handle, pending); } else { let result_code = CopyResult::try_from(code & 0xF).expect("unknown copy result"); - this.0.borrow_mut().end.mark_completed(result_code); + finish_read_state(&this.0, result_code); match result_code { CopyResult::Completed => { unsafe { ty.lift(&mut call, buffer.ptr()) }; - let result = call.pop_value(&ctx); - resolve - .call::<_, Value>((result,)) - .expect("resolve future read"); + let result = call + .maybe_pop_value(&ctx)? + .unwrap_or_else(|| Value::new_undefined(ctx.clone())); + resolve.call::<_, Value>((result,))?; } CopyResult::Dropped => { let msg = rquickjs::String::from_str(ctx.clone(), "future writer dropped")?.into_value(); - reject.call::<_, Value>((msg,)).ok(); + reject.call::<_, Value>((msg,))?; } CopyResult::Cancelled => { let msg = rquickjs::String::from_str(ctx.clone(), "future read cancelled")?.into_value(); - reject.call::<_, Value>((msg,)).ok(); + reject.call::<_, Value>((msg,))?; } } drop(buffer); @@ -180,33 +220,48 @@ fn future_read<'js>( Ok(promise.into_value()) } -fn future_cancel_read<'js>( - this: This>, +fn finish_read_state<'js>(class: &Class<'js, FutureReadable<'js>>, result: CopyResult) { + let mut readable = class.borrow_mut(); + readable.end.mark_completed(result); + if result == CopyResult::Cancelled { + readable.promise = None; + } +} + +pub(crate) fn future_cancel_read<'js>( + this: This>>, ctx: Ctx<'js>, ) -> rquickjs::Result> { let (handle, type_index) = this.0.borrow().end.begin_cancel()?; let ty = ctx.wit().future(type_index as usize); + ctx.task().unjoin(handle); let code = unsafe { ty.cancel_read()(handle) }; match unpack_copy_result(code) { None => { + ctx.task().rejoin(handle); this.0.borrow_mut().end.mark_cancel_blocked(); Ok(Value::new_undefined(ctx)) } Some((_progress, result)) => { - this.0.borrow_mut().end.mark_completed(result); + handle_read_event(handle, code); Ok(Value::new_number(ctx, result as u32 as f64)) } } } fn future_drop_readable<'js>( - this: This>, + this: This>>, ctx: Ctx<'js>, ) -> rquickjs::Result<()> { - let mut w = this.0.borrow_mut(); - if let Some(handle) = w.end.handle.take() { - let ty = ctx.wit().future(w.end.type_index as usize); + let (handle, type_index) = { + let mut readable = this.0.borrow_mut(); + let handle = readable.end.begin_drop()?; + readable.promise = None; + (handle, readable.end.type_index) + }; + if let Some(handle) = handle { + let ty = ctx.wit().future(type_index as usize); unsafe { ty.drop_readable()(handle) }; } Ok(()) @@ -217,6 +272,7 @@ fn future_write<'js>( ctx: Ctx<'js>, value: Value<'js>, ) -> rquickjs::Result> { + ctx.task().ensure_active(&ctx)?; let (handle, type_index) = this.0.borrow().end.begin_op()?; let (promise, resolve, _reject) = ctx.promise()?; @@ -235,7 +291,7 @@ fn future_write<'js>( let pending = Pending::FutureWrite { call, buffer, - resolve: Persistent::save(&ctx, resolve.into_value()), + resolve: Persistent::save(&ctx, resolve), wrapper: Persistent::save(&ctx, this.0.into_inner().into_value()), }; ctx.task().register(handle, pending); @@ -243,6 +299,7 @@ fn future_write<'js>( drop(buffer); let result_code = CopyResult::try_from(code & 0xF).expect("unknown copy result"); let success = result_code == CopyResult::Completed; + call.complete_transfers(usize::from(success)); this.0.borrow_mut().end.mark_completed(result_code); let result = Value::new_bool(ctx.clone(), success); @@ -255,21 +312,23 @@ fn future_write<'js>( Ok(promise.into_value()) } -fn future_cancel_write<'js>( +pub(crate) fn future_cancel_write<'js>( this: This>, ctx: Ctx<'js>, ) -> rquickjs::Result> { let (handle, type_index) = this.0.borrow().end.begin_cancel()?; let ty = ctx.wit().future(type_index as usize); + ctx.task().unjoin(handle); let code = unsafe { ty.cancel_write()(handle) }; match unpack_copy_result(code) { None => { + ctx.task().rejoin(handle); this.0.borrow_mut().end.mark_cancel_blocked(); Ok(Value::new_undefined(ctx)) } Some((_progress, result)) => { - this.0.borrow_mut().end.mark_completed(result); + handle_write_event(handle, code); Ok(Value::new_number(ctx, result as u32 as f64)) } } @@ -279,9 +338,13 @@ fn future_drop_writable<'js>( this: This>, ctx: Ctx<'js>, ) -> rquickjs::Result<()> { - let mut w = this.0.borrow_mut(); - if let Some(handle) = w.end.handle.take() { - let ty = ctx.wit().future(w.end.type_index as usize); + let (handle, type_index) = { + let mut writable = this.0.borrow_mut(); + (writable.end.begin_drop()?, writable.end.type_index) + }; + + if let Some(handle) = handle { + let ty = ctx.wit().future(type_index as usize); unsafe { ty.drop_writable()(handle) }; } Ok(()) @@ -292,7 +355,7 @@ pub(crate) fn handle_write_event(handle: u32, result: u32) { let pending = with_ctx(|ctx| ctx.task().take(handle)); let Pending::FutureWrite { - call: _call, + mut call, resolve, wrapper, .. @@ -303,6 +366,7 @@ pub(crate) fn handle_write_event(handle: u32, result: u32) { let copy_result = CopyResult::try_from(result & 0xF).expect("unknown copy result"); let success = copy_result == CopyResult::Completed; + call.complete_transfers(usize::from(success)); let result = with_ctx(|ctx| { let w = wrapper.restore(ctx).unwrap(); @@ -339,13 +403,13 @@ pub(crate) fn handle_read_event(handle: u32, result: u32) { let result = with_ctx(|ctx| { let w = wrapper.restore(ctx).unwrap(); let class = Class::::from_value(&w).unwrap(); - class.borrow_mut().end.mark_completed(copy_result); + finish_read_state(&class, copy_result); let type_index = class.borrow().end.type_index; let ty = ctx.wit().future(type_index as usize); unsafe { ty.lift(&mut call, buffer.ptr()) }; - Some(call.pop_persistent()) + call.maybe_pop_persistent() }); drop(buffer); @@ -356,10 +420,9 @@ pub(crate) fn handle_read_event(handle: u32, result: u32) { with_ctx(|ctx| { let w = wrapper.restore(ctx).unwrap(); let class = Class::::from_value(&w).unwrap(); - class.borrow_mut().end.mark_completed(copy_result); + finish_read_state(&class, copy_result); - let reject_fn = reject.restore(ctx).unwrap(); - let reject_fn: rquickjs::Function = reject_fn.get().unwrap(); + let reject_fn = reject.restore(ctx).expect("restore future rejecter"); let msg = if copy_result == CopyResult::Dropped { "future writer dropped" @@ -369,7 +432,10 @@ pub(crate) fn handle_read_event(handle: u32, result: u32) { let msg_val = rquickjs::String::from_str(ctx.clone(), msg) .unwrap() .into_value(); - reject_fn.call::<_, Value>((msg_val,)).ok(); + reject_fn + .call::<_, Value>((msg_val,)) + .catch(ctx) + .expect("Failed to reject future read"); }); } } diff --git a/crates/runtime/src/interpreter.rs b/crates/runtime/src/interpreter.rs index 7a9d3a3..32eac66 100644 --- a/crates/runtime/src/interpreter.rs +++ b/crates/runtime/src/interpreter.rs @@ -1,8 +1,8 @@ //! `Interpreter` trait implementation for quickjs. use crate::CtxExt; -use crate::abi::{CallbackCode, Event}; +use crate::abi::Event; use crate::bindings::register; -use crate::resources::{ResourceClasses, ResourceTable}; +use crate::resources::{ResourceTable, drain_resource_drops}; use crate::result::ResultBoundary; use crate::task::TaskState; use crate::trivia::{fn_lookup, iface_lookup}; @@ -12,7 +12,7 @@ use crate::{abi, futures, streams}; use heck::ToUpperCamelCase; use rquickjs::function::{Args, Constructor}; -use rquickjs::{Ctx, Function, JsLifetime, Object, Value}; +use rquickjs::{CatchResultExt, Ctx, Function, JsLifetime, Object, Value}; use wit_dylib_ffi::{ExportFunction, Interpreter, Resource, Wit}; /// Newtype wrapper for `Wit` so it can be stored as rquickjs userdata. @@ -64,8 +64,6 @@ impl Interpreter for QjsInterpreter { .expect("Failed to store WIT userdata"); ctx.store_userdata(ResourceTable::default()) .expect("Failed to store ResourceTable userdata"); - ctx.store_userdata(ResourceClasses::default()) - .expect("Failed to store ResourceClasses userdata"); ctx.store_userdata(TaskState::new()) .expect("Failed to store TaskState userdata"); ctx.store_userdata(WitImportRegistry::new(wit)) @@ -75,6 +73,7 @@ impl Interpreter for QjsInterpreter { } fn export_start<'a>(_wit: Wit, _func: ExportFunction) -> Box> { + with_ctx(drain_resource_drops); Box::new(QjsCallContext::default()) } @@ -98,7 +97,7 @@ impl Interpreter for QjsInterpreter { let self_val = cx.shift_value(ctx); let self_obj = self_val .as_object() - .unwrap_or_else(|| panic!("method receiver is not an object")); + .expect("method receiver is not an object"); let method: Function = self_obj .get(method_name) .unwrap_or_else(|err| panic!("method '{method_name}' not found: {err:?}")); @@ -135,6 +134,7 @@ impl Interpreter for QjsInterpreter { call_export(ctx, func, func_name, js_func, args, cx); } }); + with_ctx(drain_resource_drops); } fn export_async_start( @@ -143,7 +143,9 @@ impl Interpreter for QjsInterpreter { mut cx: Box>, ) -> u32 { with_ctx(|ctx| { - ctx.task().init(); + drain_resource_drops(ctx); + let args = cx.stack_into_args(ctx); + ctx.task().init(cx); let globals = ctx.globals(); @@ -164,11 +166,10 @@ impl Interpreter for QjsInterpreter { .get(func_name) .unwrap_or_else(|e| panic!("Failed to get async export '{}': {:?}", func_name, e)); - let args = cx.stack_into_args(ctx); - let _result = js_func .call_arg::(args) - .unwrap_or_else(|e| panic!("Failed to call async '{}': {:?}", func.name(), e)); + .catch(ctx) + .unwrap_or_else(|e| panic!("Failed to call async '{}': {e}", func.name())); }); with_ctx(|ctx| ctx.task().poll()) @@ -191,19 +192,27 @@ impl Interpreter for QjsInterpreter { Event::StreamRead { handle, result } => streams::handle_read_event(handle, result), Event::FutureWrite { handle, result } => futures::handle_write_event(handle, result), Event::FutureRead { handle, result } => futures::handle_read_event(handle, result), - Event::TaskCancelled => with_ctx(|ctx| ctx.task().cancel()), + Event::TaskCancelled => with_ctx(|ctx| { + ctx.task() + .cancel(ctx) + .catch(ctx) + .expect("Failed to cancel async task"); + }), } - if matches!(evt, Event::TaskCancelled) { - CallbackCode::Exit.encode(0) - } else { - with_ctx(|ctx| ctx.task().poll()) - } + with_ctx(|ctx| ctx.task().poll()) + } + + fn export_finish(mut cx: Box>, _func: ExportFunction) { + // Post-return cannot invoke host destructors; queued drops wait for the next entry. + cx.complete_transfers(1); + drop(cx); } fn resource_dtor(_ty: Resource, handle: usize) { with_ctx(|ctx| { ctx.resources().remove(handle); + drain_resource_drops(ctx); }); } } diff --git a/crates/runtime/src/lib.rs b/crates/runtime/src/lib.rs index d7b8466..b8ebe09 100644 --- a/crates/runtime/src/lib.rs +++ b/crates/runtime/src/lib.rs @@ -11,6 +11,7 @@ mod streams; mod tagged; mod task; mod trivia; +mod typed_array; mod wit_imports; use std::cell::{Cell, OnceCell, RefCell}; @@ -23,9 +24,7 @@ use task::TaskState; use wit_dylib_ffi::Wit; use crate::interpreter::WitData; -use crate::resources::BorrowedResource; -use crate::resources::ResourceClasses; -use crate::resources::ResourceTable; +use crate::resources::*; use crate::trivia::*; /// Deterministic, fixed-seed hash map/set used everywhere in the runtime so the @@ -157,6 +156,9 @@ impl JsState { context.with(|ctx| { ctx.store_userdata(FnNameCache::default()) .expect("Failed to store function name cache"); + // Startup cleanup also runs for empty worlds without WIT initialization. + ctx.store_userdata(ResourceClasses::default()) + .expect("Failed to store ResourceClasses userdata"); module::init_state(&ctx); }); @@ -214,11 +216,25 @@ pub struct QjsCallContext { temp_strings: SmallVec<[String; 4]>, /// Raw allocations to free when this context is dropped deferred_deallocs: SmallVec<[(*mut u8, std::alloc::Layout); 4]>, + /// Keeps imported resources alive while an outgoing call borrows them + loans: SmallVec<[ImportedResourceLoan; 4]>, + /// Own transfers grouped by stream element, or group zero for ordinary calls + transfers: SmallVec<[(usize, ImportedResourceTransfer); 4]>, + /// Number of transfer groups consumed by the callee, + transfer_group: usize, /// Imported resource borrows to drop when this context is dropped borrows: SmallVec<[BorrowedResource; 4]>, } impl QjsCallContext { + pub(crate) fn complete_transfers(&mut self, consumed_groups: usize) { + for (group, transfer) in self.transfers.drain(..) { + if group < consumed_groups { + transfer.commit(); + } + } + } + pub(crate) fn push_value<'js>(&mut self, ctx: &rquickjs::Ctx<'js>, val: Value<'js>) { self.stack.push(Persistent::save(ctx, val)); } @@ -261,16 +277,16 @@ impl QjsCallContext { impl Drop for QjsCallContext { fn drop(&mut self) { + self.transfers.clear(); for (ptr, layout) in self.deferred_deallocs.drain(..) { - unsafe { - std::alloc::dealloc(ptr, layout); - } - } - for borrow in self.borrows.drain(..) { - unsafe { - (borrow.drop_fn)(borrow.handle); + if layout.size() > 0 { + unsafe { + std::alloc::dealloc(ptr, layout); + } } } + self.loans.clear(); + self.borrows.clear(); } } @@ -294,15 +310,14 @@ fn init_js( } if disable_gc { - state.with_ctx(|ctx| unsafe { - let rt = rquickjs::qjs::JS_GetRuntime(ctx.as_raw().as_ptr()); - rquickjs::qjs::JS_SetGCThreshold(rt, usize::MAX as _); - }); + state.context.runtime().set_gc_threshold(usize::MAX); } state.with_ctx(|ctx| { module::evaluate_shim(ctx, shim)?; - module::evaluate_user(ctx, js_source, entry_path) + let result = module::evaluate_user(ctx, js_source, entry_path); + resources::drain_resource_drops(ctx); + result })?; unsafe { diff --git a/crates/runtime/src/module/mod.rs b/crates/runtime/src/module/mod.rs index 486aa79..e73e750 100644 --- a/crates/runtime/src/module/mod.rs +++ b/crates/runtime/src/module/mod.rs @@ -4,7 +4,7 @@ mod wit; use std::cell::RefCell; -use rquickjs::{CaughtError, JsLifetime, Module, Persistent, Runtime}; +use rquickjs::{CaughtError, CaughtResult, Ctx, JsLifetime, Module, Persistent, Runtime}; use crate::CtxExt; @@ -88,9 +88,46 @@ fn evaluate<'js>( let (module, promise) = CaughtError::catch(ctx, module.eval()) .map_err(|e| format!("Failed to evaluate JavaScript module: {e}"))?; - CaughtError::catch(ctx, promise.finish::<()>()) - .map_err(|e| format!("Failed to finish JavaScript module evaluation: {e}"))?; + loop { + if let Some(result) = promise.result::<()>() { + CaughtError::catch(ctx, result) + .map_err(|e| format!("Failed to finish JavaScript module evaluation: {e}"))?; + break; + } + + if !execute_pending_job(ctx).map_err(|e| format!("JavaScript module job failed: {e}"))? { + return Err(format!( + "Failed to finish JavaScript module evaluation: {}", + rquickjs::Error::WouldBlock + )); + } + } CaughtError::catch(ctx, module.namespace()) .map_err(|e| format!("Failed to read JavaScript module namespace: {e}")) } + +/// Execute a job without reacquiring the runtime lock or losing its exception. +pub(crate) fn execute_pending_job<'js>(ctx: &Ctx<'js>) -> CaughtResult<'js, bool> { + let mut job_ctx = std::ptr::null_mut(); + + // The caller holds the runtime lock through Context::with. QuickJS returns + // the context that owns an exception, which need not be the caller's context. + let result = unsafe { + let runtime = rquickjs::qjs::JS_GetRuntime(ctx.as_raw().as_ptr()); + rquickjs::qjs::JS_ExecutePendingJob(runtime, &mut job_ctx) + }; + + if result < 0 { + let ptr = std::ptr::NonNull::new(job_ctx).expect("job exception without a context"); + // SAFETY: this context belongs to the locked runtime and cannot escape 'js. + let job_ctx = unsafe { Ctx::from_raw(ptr) }; + + Err(CaughtError::from_error( + &job_ctx, + rquickjs::Error::Exception, + )) + } else { + Ok(result > 0) + } +} diff --git a/crates/runtime/src/resources.rs b/crates/runtime/src/resources.rs index 9ba1f8c..7636f99 100644 --- a/crates/runtime/src/resources.rs +++ b/crates/runtime/src/resources.rs @@ -1,21 +1,270 @@ -//! Exported resource table: maps `rep` indices to JS objects. +//! Native ownership of imported resources and the exported JS object table. //! -//! When JS defines a resource type and exports it, the component model needs -//! a table mapping internal representation indices (`rep`) to the js objects -//! that back them. +//! Imported resource finalizers only enqueue owned handles. The runtime must +//! drain those handles outside QuickJS evaluations and jobs. -use std::cell::RefCell; +use std::cell::{Cell, RefCell}; +use std::collections::VecDeque; +use std::rc::Rc; -use rquickjs::{JsLifetime, Persistent, Value}; +use rquickjs::class::{Class, JsClass, Readable, Trace}; +use rquickjs::function::{Constructor, This}; +use rquickjs::{Ctx, Exception, Function, JsLifetime, Object, Persistent, Value}; use wit_dylib_ffi::Resource; -use crate::CtxExt; -use crate::DetHashMap; +use crate::{CtxExt, DetHashMap, symbol_dispose}; -/// A borrowed imported resource handle that must be dropped when the call ends. +#[derive(Default)] +struct ResourceDropQueue { + pending: RefCell>, + draining: Cell, +} + +#[derive(Clone, Copy)] +enum Ownership { + Owned(u32), + Borrowed(u32), + Transferring, + Transferred, + Disposed, + Expired, +} + +struct ImportedResourceState { + ty: Resource, + ownership: Cell, + loans: Cell, + drops: Rc, +} + +impl ImportedResourceState { + fn take_owned(&self) -> Option<(Resource, u32)> { + let Ownership::Owned(handle) = self.ownership.get() else { + return None; + }; + self.ownership.set(Ownership::Disposed); + Some((self.ty, handle)) + } + + fn validate(&self, ty: Resource, owned: bool) -> Result { + if self.ty != ty { + return Err("imported resource has the wrong WIT resource type"); + } + match self.ownership.get() { + Ownership::Owned(_) if owned && self.loans.get() != 0 => { + Err("cannot transfer ownership of a resource while it is borrowed") + } + Ownership::Owned(handle) => Ok(handle), + Ownership::Borrowed(handle) if !owned => Ok(handle), + Ownership::Borrowed(_) => Err("cannot transfer ownership of a borrowed resource"), + Ownership::Transferring => Err("resource ownership transfer is in progress"), + Ownership::Transferred => Err("resource ownership has already been transferred"), + Ownership::Disposed => Err("resource has been disposed"), + Ownership::Expired => Err("borrowed resource has expired"), + } + } +} + +impl Drop for ImportedResourceState { + fn drop(&mut self) { + // This can run inside QuickJS's GC: no JS values or host calls here. + if let Some(resource) = self.take_owned() { + self.drops.pending.borrow_mut().push_back(resource); + } + } +} + +/// Opaque, type-checked state for an imported resource's JS wrapper. +#[derive(Trace, JsLifetime)] +pub(crate) struct ImportedResource { + #[qjs(skip_trace)] + state: Rc, +} + +impl<'js> JsClass<'js> for ImportedResource { + const NAME: &'static str = "ImportedResource"; + type Mutable = Readable; + + fn prototype(ctx: &Ctx<'js>) -> rquickjs::Result>> { + let prototype = Object::new(ctx.clone())?; + let dispose = Function::new(ctx.clone(), dispose_imported_resource)?; + prototype.set("drop", dispose.clone())?; + prototype.set(symbol_dispose(ctx)?, dispose)?; + Ok(Some(prototype)) + } + + fn constructor(_ctx: &Ctx<'js>) -> rquickjs::Result>> { + Ok(None) + } +} + +/// Keeps a canonical imported borrow alive until its call ends. +/// +/// Dropping the guard expires any retained JS wrapper before releasing the +/// canonical borrowed handle. That release does not destroy the borrowed owner. +#[must_use = "keep this guard alive until the borrowed resource's call ends"] pub(crate) struct BorrowedResource { - pub(crate) handle: u32, - pub(crate) drop_fn: unsafe extern "C" fn(u32), + state: Rc, + handle: u32, +} + +impl Drop for BorrowedResource { + fn drop(&mut self) { + assert_eq!( + self.state.loans.get(), + 0, + "borrowed resource still used by an outstanding import at the end of its call" + ); + self.state.ownership.set(Ownership::Expired); + unsafe { self.state.ty.drop()(self.handle) }; + } +} + +/// Retains native owner state while an outbound call borrows its handle. +/// +/// Live loans prevent disposal and ownership transfer, even across suspension. +/// This does not extend an incoming canonical borrow's lifetime: its original +/// `BorrowedResource` guard must also outlive every call that uses it. +#[must_use = "keep this loan alive until the borrowing import completes"] +pub(crate) struct ImportedResourceLoan { + state: Rc, +} + +impl Drop for ImportedResourceLoan { + fn drop(&mut self) { + self.state.loans.set(self.state.loans.get() - 1); + } +} + +/// Restores an unconsumed handle after a partial write or cancelled import. +#[must_use = "commit consumed transfers before releasing their call context"] +pub(crate) struct ImportedResourceTransfer { + state: Rc, + handle: u32, +} + +impl ImportedResourceTransfer { + pub(crate) fn commit(self) { + self.state.ownership.set(Ownership::Transferred); + } +} + +impl Drop for ImportedResourceTransfer { + fn drop(&mut self) { + if matches!(self.state.ownership.get(), Ownership::Transferring) { + self.state.ownership.set(Ownership::Owned(self.handle)); + } + } +} + +/// Native disposal prototype inherited by each imported WIT resource prototype. +pub(crate) fn imported_resource_prototype<'js>(ctx: &Ctx<'js>) -> rquickjs::Result> { + Class::::prototype(ctx)?.ok_or_else(|| { + Exception::throw_type(ctx, "native imported resource prototype is unavailable") + }) +} + +fn make_imported<'js>(ctx: &Ctx<'js>, resource: ImportedResource) -> rquickjs::Result> { + let prototype = ctx.resource_classes().prototype(resource.state.ty.index()); + let prototype = match prototype { + Some(prototype) => prototype.restore(ctx)?, + None => imported_resource_prototype(ctx)?, + }; + Class::instance_proto(resource, prototype).map(Class::into_value) +} + +/// Take ownership of an imported handle, including on wrapper creation failure. +/// +/// Failed construction queues the handle for the next safe-point drain. +pub(crate) fn make_imported_owned<'js>( + ctx: &Ctx<'js>, + ty: Resource, + handle: u32, +) -> rquickjs::Result> { + let state = Rc::new(ImportedResourceState { + ty, + ownership: Cell::new(Ownership::Owned(handle)), + loans: Cell::new(0), + drops: Rc::clone(&ctx.resource_classes().drops), + }); + make_imported(ctx, ImportedResource { state }) +} + +/// Wrap a canonical imported borrow and return its call-lifetime guard. +/// +/// Failed construction releases the borrowed handle immediately. +pub(crate) fn make_imported_borrowed<'js>( + ctx: &Ctx<'js>, + ty: Resource, + handle: u32, +) -> rquickjs::Result<(Value<'js>, BorrowedResource)> { + let state = Rc::new(ImportedResourceState { + ty, + ownership: Cell::new(Ownership::Borrowed(handle)), + loans: Cell::new(0), + drops: Rc::clone(&ctx.resource_classes().drops), + }); + let guard = BorrowedResource { + state: Rc::clone(&state), + handle, + }; + let value = make_imported(ctx, ImportedResource { state })?; + Ok((value, guard)) +} + +fn dispose_imported_resource<'js>( + this: This>, + ctx: Ctx<'js>, +) -> rquickjs::Result<()> { + let resource = { + let resource = this.0.try_borrow()?; + match resource.state.ownership.get() { + Ownership::Borrowed(_) => Err("cannot dispose a borrowed resource"), + Ownership::Transferring => Err("cannot dispose a resource during ownership transfer"), + Ownership::Owned(_) if resource.state.loans.get() != 0 => { + Err("cannot dispose a resource while it is borrowed") + } + _ => Ok(resource.state.take_owned()), + } + } + .map_err(|message| Exception::throw_type(&ctx, message))?; + + // The handle is invalidated and the class borrow is released before re-entry. + if let Some((ty, handle)) = resource { + unsafe { ty.drop()(handle) }; + } + Ok(()) +} + +/// Run deferred owned-resource destructors at a safe point outside JS execution. +/// +/// Call after evaluations/jobs, loan completion and explicit GC, while the runtime +/// and its WIT metadata are still alive. Never call from canonical post-return, +/// where imports and intrinsics are forbidden. Teardown must drain before +/// destroying runtime state; neither finalizers nor dropping the queue call hosts. +pub(crate) fn drain_resource_drops(ctx: &Ctx<'_>) { + let drops = Rc::clone(&ctx.resource_classes().drops); + if drops.draining.replace(true) { + return; + } + + struct DrainGuard<'a>(&'a Cell); + + impl Drop for DrainGuard<'_> { + fn drop(&mut self) { + self.0.set(false); + } + } + + let _guard = DrainGuard(&drops.draining); + loop { + let resource = drops.pending.borrow_mut().pop_front(); + let Some((ty, handle)) = resource else { + break; + }; + // No userdata or RefCell guard may survive this potentially re-entrant call. + unsafe { ty.drop()(handle) }; + } } /// Table mapping `rep` indices to JS objects for exported resources. @@ -64,6 +313,7 @@ impl ResourceTable { #[derive(Default, JsLifetime)] pub(crate) struct ResourceClasses { inner: RefCell, + drops: Rc, } #[derive(Default)] @@ -72,8 +322,8 @@ struct ClassInner { } struct ResourceClass { - class: Persistent>, - prototype: Persistent>, + class: Persistent>, + prototype: Persistent>, } impl ResourceClasses { @@ -81,8 +331,8 @@ impl ResourceClasses { pub(crate) fn insert( &self, index: usize, - class: Persistent>, - prototype: Persistent>, + class: Persistent>, + prototype: Persistent>, ) { self.inner .borrow_mut() @@ -91,12 +341,12 @@ impl ResourceClasses { } /// Get the class (constructor) for a resource, if any. - pub(crate) fn class(&self, index: usize) -> Option>> { + pub(crate) fn class(&self, index: usize) -> Option>> { self.inner.borrow().map.get(&index).map(|c| c.class.clone()) } /// Get the prototype object for a resource, if any. - pub(crate) fn prototype(&self, index: usize) -> Option>> { + pub(crate) fn prototype(&self, index: usize) -> Option>> { self.inner .borrow() .map @@ -105,11 +355,58 @@ impl ResourceClasses { } } -/// Extract the canonical handle from an imported resource wrapper object. -pub(crate) fn imported_resource_to_handle(val: &Value<'_>) -> u32 { - val.as_object() - .and_then(|obj| obj.get::<_, u32>("__cqjs_handle").ok()) - .expect("expected resource wrapper with __cqjs_handle") +fn imported_resource_state(val: &Value<'_>) -> rquickjs::Result> { + let resource = Class::::from_value(val) + .map_err(|_| Exception::throw_type(val.ctx(), "expected an imported resource wrapper"))?; + let state = Rc::clone(&resource.try_borrow()?.state); + Ok(state) +} + +/// Check type, lifetime and ownership without consuming or borrowing the handle. +/// +/// Preflight each argument before entering the infallible `Call` ABI. This does +/// not reserve ownership or detect duplicate/mixed uses among several arguments; +/// the caller must check such aliasing before starting any transfers. +pub(crate) fn validate_imported_resource( + ty: Resource, + val: &Value<'_>, + owned: bool, +) -> rquickjs::Result { + imported_resource_state(val)? + .validate(ty, owned) + .map_err(|message| Exception::throw_type(val.ctx(), message)) +} + +/// Borrow a handle and retain its native state until the returned loan is dropped. +/// +/// No JS root is needed to keep an owned handle alive. If its wrapper is collected, +/// releasing the final loan queues the owned handle for a safe-point drop. +pub(crate) fn borrow_imported_resource( + ty: Resource, + val: &Value<'_>, +) -> rquickjs::Result<(u32, ImportedResourceLoan)> { + let state = imported_resource_state(val)?; + let handle = state + .validate(ty, false) + .map_err(|message| Exception::throw_type(val.ctx(), message))?; + let loans = state.loans.get().checked_add(1).ok_or_else(|| { + Exception::throw_type(val.ctx(), "too many outstanding imported resource borrows") + })?; + state.loans.set(loans); + Ok((handle, ImportedResourceLoan { state })) +} + +/// Reserve ownership until the canonical ABI reports whether it consumed a value. +pub(crate) fn transfer_imported_resource( + ty: Resource, + val: &Value<'_>, +) -> rquickjs::Result<(u32, ImportedResourceTransfer)> { + let state = imported_resource_state(val)?; + let handle = state + .validate(ty, true) + .map_err(|message| Exception::throw_type(val.ctx(), message))?; + state.ownership.set(Ownership::Transferring); + Ok((handle, ImportedResourceTransfer { state, handle })) } /// Convert a js object to a canonical handle for an exported resource. diff --git a/crates/runtime/src/result.rs b/crates/runtime/src/result.rs index c490e34..274acdb 100644 --- a/crates/runtime/src/result.rs +++ b/crates/runtime/src/result.rs @@ -56,8 +56,8 @@ impl<'js> JsCompletion<'js> { pub(crate) fn settle_persistent( self, ctx: &Ctx<'js>, - resolve: Persistent>, - reject: Persistent>, + resolve: Persistent>, + reject: Persistent>, ) { match self { JsCompletion::Return(result) => { diff --git a/crates/runtime/src/streams.rs b/crates/runtime/src/streams.rs index a55ee88..8582758 100644 --- a/crates/runtime/src/streams.rs +++ b/crates/runtime/src/streams.rs @@ -8,54 +8,18 @@ use crate::CtxExt; use crate::abi::{CopyEnd, CopyResult, CopyState, is_blocked_raw, unpack_copy_result}; use crate::buffer::BufferGuard; use crate::task::Pending; +use crate::typed_array::{copy_typed_array_as, typed_array_len_as}; use crate::{QjsCallContext, resolve_promise, symbol_dispose, with_ctx}; use rquickjs::JsLifetime; use rquickjs::class::{Class, JsClass, Trace}; -use rquickjs::function::{self, Rest, This}; +use rquickjs::function::{self, Opt, Rest, This}; use rquickjs::{Ctx, Function, Object, Persistent, Symbol, Value}; use std::cell::Cell; const BYTE_ITERATOR_CHUNK_SIZE: usize = 64 * 1024; -macro_rules! copy_typed_array_as { - ($obj:expr, $ty:expr, $t:ty) => {{ - let Some(ta) = $obj.as_typed_array::<$t>() else { - return Ok(None); - }; - - let slice: &[$t] = ta.as_ref(); - let count = slice.len(); - - assert_eq!($ty.abi_payload_size(), std::mem::size_of::<$t>()); - assert!($ty.abi_payload_align() >= std::mem::align_of::<$t>()); - - let byte_len = count - .checked_mul(std::mem::size_of::<$t>()) - .ok_or_else(|| rquickjs::Error::new_from_js("number", "buffer size overflow"))?; - - let buf = BufferGuard::new_zeroed(byte_len, $ty.abi_payload_align()); - if byte_len > 0 { - unsafe { - let src = slice.as_ptr() as *const u8; - let dst = buf.ptr(); - std::ptr::copy_nonoverlapping(src, dst, byte_len); - } - } - Some((buf, count)) - }}; -} - -macro_rules! typed_array_len_as { - ($obj:expr, $t:ty) => { - $obj.as_typed_array::<$t>().map(|array| { - let slice: &[$t] = array.as_ref(); - slice.len() - }) - }; -} - /// Rust side state for the readable end of a component-model stream. #[derive(Trace, JsLifetime)] pub(crate) struct StreamReadable { @@ -169,12 +133,12 @@ pub(crate) fn make_stream_readable<'js>( } /// Create a `[StreamWritable, StreamReadable]` pair. -pub(crate) fn make_stream<'js>( - ctx: Ctx<'js>, - args: Rest>, -) -> rquickjs::Result> { - let type_index: u32 = args.0[0].get()?; - let ty = ctx.wit().stream(type_index as usize); +pub(crate) fn make_stream<'js>(ctx: Ctx<'js>, type_index: u32) -> rquickjs::Result> { + let ty = ctx + .wit() + .iter_streams() + .nth(type_index as usize) + .ok_or_else(|| rquickjs::Exception::throw_range(&ctx, "unknown WIT stream type"))?; let handles = unsafe { ty.new()() }; let tx_handle = (handles >> 32) as u32; @@ -208,19 +172,19 @@ pub(crate) fn lower_iterable<'js>( let pair: Object = from.call((iterable, type_index))?; let readable: Value = pair.get("readable")?; let readable = Class::::from_value(&readable)?; - let handle = readable - .borrow_mut() - .end - .handle - .take() - .ok_or_else(|| rquickjs::Error::new_from_js("stream", "already transferred"))?; + let handle = readable.borrow_mut().end.begin_transfer(type_index)?; Ok(handle) } -fn typed_array_batch_len<'js>(data: &Value<'js>, ty: &wit_dylib_ffi::Stream) -> Option { - let obj = data.as_object()?; - match ty.ty()? { +fn typed_array_batch_len<'js>( + data: &Value<'js>, + ty: &wit_dylib_ffi::Stream, +) -> rquickjs::Result> { + let Some((obj, elem_ty)) = data.as_object().zip(ty.ty()) else { + return Ok(None); + }; + match elem_ty { wit_dylib_ffi::Type::U8 => typed_array_len_as!(obj, u8), wit_dylib_ffi::Type::S8 => typed_array_len_as!(obj, i8), wit_dylib_ffi::Type::U16 => typed_array_len_as!(obj, u16), @@ -231,7 +195,7 @@ fn typed_array_batch_len<'js>(data: &Value<'js>, ty: &wit_dylib_ffi::Stream) -> wit_dylib_ffi::Type::S64 => typed_array_len_as!(obj, i64), wit_dylib_ffi::Type::F32 => typed_array_len_as!(obj, f32), wit_dylib_ffi::Type::F64 => typed_array_len_as!(obj, f64), - _ => None, + _ => Ok(None), } } @@ -248,7 +212,7 @@ fn try_typed_array_to_buffer<'js>( return Ok(None); }; - let pair = match elem_ty { + match elem_ty { wit_dylib_ffi::Type::U8 => copy_typed_array_as!(obj, ty, u8), wit_dylib_ffi::Type::S8 => copy_typed_array_as!(obj, ty, i8), wit_dylib_ffi::Type::U16 => copy_typed_array_as!(obj, ty, u16), @@ -259,19 +223,16 @@ fn try_typed_array_to_buffer<'js>( wit_dylib_ffi::Type::S64 => copy_typed_array_as!(obj, ty, i64), wit_dylib_ffi::Type::F32 => copy_typed_array_as!(obj, ty, f32), wit_dylib_ffi::Type::F64 => copy_typed_array_as!(obj, ty, f64), - _ => return Ok(None), - }; - - Ok(pair) + _ => Ok(None), + } } fn stream_read<'js>( this: This>, ctx: Ctx<'js>, - args: Rest>, + count: Opt, ) -> rquickjs::Result> { - let count: usize = args.0.first().and_then(|v| v.get().ok()).unwrap_or(1); - stream_read_impl(this, ctx, count, false) + stream_read_impl(this, ctx, count.0.unwrap_or(1), false) } fn stream_next<'js>( @@ -308,6 +269,7 @@ fn stream_read_impl<'js>( count: usize, iterator: bool, ) -> rquickjs::Result> { + ctx.task().ensure_active(&ctx)?; if count == 0 { return Err(rquickjs::Error::new_from_js( "number", @@ -336,7 +298,7 @@ fn stream_read_impl<'js>( buffer, iterator, iterator_return: None, - resolve: Persistent::save(&ctx, resolve.into_value()), + resolve: Persistent::save(&ctx, resolve), wrapper: Persistent::save(&ctx, this.0.into_inner().into_value()), }; ctx.task().register(handle, pending); @@ -478,7 +440,7 @@ fn stream_iterator_return<'js>( ctx.task().unjoin(handle); let code = unsafe { ty.cancel_read()(handle) }; ctx.task() - .set_stream_iterator_return(handle, Persistent::save(&ctx, resolve.into_value())); + .set_stream_iterator_return(handle, Persistent::save(&ctx, resolve)); if is_blocked_raw(code) { ctx.task().rejoin(handle); @@ -490,7 +452,7 @@ fn stream_iterator_return<'js>( Ok(promise.into_value()) } -fn stream_cancel_read<'js>( +pub(crate) fn stream_cancel_read<'js>( this: This>, ctx: Ctx<'js>, ) -> rquickjs::Result> { @@ -519,10 +481,12 @@ fn stream_drop_readable<'js>( this: This>, ctx: Ctx<'js>, ) -> rquickjs::Result<()> { - let mut w = this.0.borrow_mut(); - - if let Some(handle) = w.end.handle.take() { - let ty = ctx.wit().stream(w.end.type_index as usize); + let (handle, type_index) = { + let mut readable = this.0.borrow_mut(); + (readable.end.begin_drop()?, readable.end.type_index) + }; + if let Some(handle) = handle { + let ty = ctx.wit().stream(type_index as usize); unsafe { ty.drop_readable()(handle) }; } @@ -569,7 +533,7 @@ fn stream_write_iterable_item<'js>( } let ty = ctx.wit().stream(type_index as usize); - let batch_len = typed_array_batch_len(&data, &ty); + let batch_len = typed_array_batch_len(&data, &ty)?; let expected = batch_len.unwrap_or(1); let result = if batch_len.is_some() { @@ -619,6 +583,7 @@ fn stream_write_impl<'js>( data: Value<'js>, mode: StreamWriteMode, ) -> rquickjs::Result> { + ctx.task().ensure_active(&ctx)?; let (handle, type_index) = this.0.borrow().end.begin_op()?; let (promise, resolve, _reject) = ctx.promise()?; @@ -646,6 +611,7 @@ fn stream_write_impl<'js>( for i in 0..count { let elem: Value = arr.get(i)?; + call.transfer_group = i; call.push_value(&ctx, elem); unsafe { ty.lower(&mut call, buf.ptr().add(ty.abi_payload_size() * i)) }; } @@ -663,7 +629,7 @@ fn stream_write_impl<'js>( this.0.borrow_mut().end.mark_blocked(); let pending = Pending::StreamWrite { call, - resolve: Persistent::save(&ctx, resolve.into_value()), + resolve: Persistent::save(&ctx, resolve), wrapper: Persistent::save(&ctx, this.0.into_inner().into_value()), buffer, }; @@ -671,6 +637,7 @@ fn stream_write_impl<'js>( } else { drop(buffer); let (progress, copy_result) = unpack_copy_result(code).expect("non-blocked"); + call.complete_transfers(progress as usize); let dropped_handle = { let mut writable = this.0.borrow_mut(); writable.end.mark_completed(copy_result); @@ -811,7 +778,7 @@ fn write_all_step<'js>( then_fn.call_arg(then_args) } -fn stream_cancel_write<'js>( +pub(crate) fn stream_cancel_write<'js>( this: This>, ctx: Ctx<'js>, ) -> rquickjs::Result> { @@ -840,9 +807,12 @@ fn stream_drop_writable<'js>( this: This>, ctx: Ctx<'js>, ) -> rquickjs::Result<()> { - let mut w = this.0.borrow_mut(); - if let Some(handle) = w.end.handle.take() { - let ty = ctx.wit().stream(w.end.type_index as usize); + let (handle, type_index) = { + let mut writable = this.0.borrow_mut(); + (writable.end.begin_drop()?, writable.end.type_index) + }; + if let Some(handle) = handle { + let ty = ctx.wit().stream(type_index as usize); unsafe { ty.drop_writable()(handle) }; } Ok(()) @@ -853,7 +823,7 @@ pub(crate) fn handle_write_event(handle: u32, result: u32) { let pending = with_ctx(|ctx| ctx.task().take(handle)); let Pending::StreamWrite { - call: _call, + mut call, resolve, wrapper, .. @@ -864,6 +834,7 @@ pub(crate) fn handle_write_event(handle: u32, result: u32) { let (progress, copy_result) = unpack_copy_result(result).expect("StreamWrite callback should not be BLOCKED"); + call.complete_transfers(progress as usize); let result = with_ctx(|ctx| { let w = wrapper.restore(ctx).unwrap(); diff --git a/crates/runtime/src/task.rs b/crates/runtime/src/task.rs index 2de5c2a..a53dfa1 100644 --- a/crates/runtime/src/task.rs +++ b/crates/runtime/src/task.rs @@ -3,14 +3,20 @@ use std::cell::RefCell; -use rquickjs::{JsLifetime, Persistent, Value}; +use rquickjs::class::Class; +use rquickjs::function::This; +use rquickjs::{CatchResultExt, Ctx, Function, JsLifetime, Persistent, Value}; use crate::CtxExt; use crate::DetHashMap; use crate::abi::*; use crate::buffer::BufferGuard; +use crate::futures::{self, FutureReadable, FutureWritable}; +use crate::module::execute_pending_job; +use crate::resources::drain_resource_drops; use crate::result::ResultBoundary; -use crate::{QjsCallContext, resolve_promise, with_ctx}; +use crate::streams::{self, StreamReadable, StreamWritable}; +use crate::{QjsCallContext, reject_promise, with_ctx}; /// A pending async operation awaiting a callback event. /// @@ -24,13 +30,13 @@ pub(crate) enum Pending { call: QjsCallContext, func_index: usize, buffer: *mut u8, - resolve: Persistent>, - reject: Persistent>, + resolve: Persistent>, + reject: Persistent>, }, /// A stream write that blocked. StreamWrite { call: QjsCallContext, - resolve: Persistent>, + resolve: Persistent>, wrapper: Persistent>, buffer: BufferGuard, }, @@ -39,14 +45,14 @@ pub(crate) enum Pending { call: QjsCallContext, buffer: BufferGuard, iterator: bool, - iterator_return: Option>>, - resolve: Persistent>, + iterator_return: Option>>, + resolve: Persistent>, wrapper: Persistent>, }, /// A future write that blocked. FutureWrite { call: QjsCallContext, - resolve: Persistent>, + resolve: Persistent>, wrapper: Persistent>, buffer: BufferGuard, }, @@ -54,12 +60,75 @@ pub(crate) enum Pending { FutureRead { call: QjsCallContext, buffer: BufferGuard, - resolve: Persistent>, - reject: Persistent>, + resolve: Persistent>, + reject: Persistent>, wrapper: Persistent>, }, } +struct PendingEntry { + op: Pending, + cancel: bool, +} + +enum CancelRequest { + Import, + StreamRead(Persistent>), + StreamWrite(Persistent>), + FutureRead(Persistent>), + FutureWrite(Persistent>), +} + +impl Pending { + fn cancel_request(&self) -> CancelRequest { + match self { + Self::ImportCall { .. } => CancelRequest::Import, + Self::StreamRead { wrapper, .. } => CancelRequest::StreamRead(wrapper.clone()), + Self::StreamWrite { wrapper, .. } => CancelRequest::StreamWrite(wrapper.clone()), + Self::FutureRead { wrapper, .. } => CancelRequest::FutureRead(wrapper.clone()), + Self::FutureWrite { wrapper, .. } => CancelRequest::FutureWrite(wrapper.clone()), + } + } +} + +impl CancelRequest { + fn cancel(self, ctx: &Ctx<'_>, handle: u32) -> rquickjs::Result<()> { + match self { + Self::Import => { + ctx.task().unjoin(handle); + let code = unsafe { subtask_cancel(handle) }; + if is_blocked_raw(code) { + ctx.task().rejoin(handle); + } else { + let state = SubtaskState::try_from(code).expect("invalid cancellation state"); + handle_subtask(handle, state); + } + } + Self::StreamRead(wrapper) => { + let value = wrapper.restore(ctx)?; + let class = Class::::from_value(&value)?; + streams::stream_cancel_read(This(class), ctx.clone())?; + } + Self::StreamWrite(wrapper) => { + let value = wrapper.restore(ctx)?; + let class = Class::::from_value(&value)?; + streams::stream_cancel_write(This(class), ctx.clone())?; + } + Self::FutureRead(wrapper) => { + let value = wrapper.restore(ctx)?; + let class = Class::::from_value(&value)?; + futures::future_cancel_read(This(class), ctx.clone())?; + } + Self::FutureWrite(wrapper) => { + let value = wrapper.restore(ctx)?; + let class = Class::::from_value(&value)?; + futures::future_cancel_write(This(class), ctx.clone())?; + } + } + Ok(()) + } +} + /// Inflight async operations for a single export call. /// /// Every entry is joined to `waitable_set` while the host may wake it. @@ -67,19 +136,35 @@ pub(crate) enum Pending { /// remain synchronized. #[derive(Default)] struct TaskInner { - pending: DetHashMap, + pending: DetHashMap, waitable_set: Option, + export_call: Option>, + cancel: bool, + returned: bool, } impl TaskInner { /// Store an operation and join its handle to the lazily-created waitable set. fn register(&mut self, handle: u32, pending: Pending) { + assert!( + !self.pending.contains_key(&handle), + "operation already pending" + ); + if self.waitable_set.is_none() { self.waitable_set = Some(unsafe { waitable_set_new() }); } + let set = self.waitable_set.unwrap(); unsafe { waitable_join(handle, set) }; - self.pending.insert(handle, pending); + + self.pending.insert( + handle, + PendingEntry { + op: pending, + cancel: false, + }, + ); } /// Unjoin a completed handle and take ownership of its pending state. @@ -88,11 +173,17 @@ impl TaskInner { self.pending .remove(&handle) .expect("no pending entry for handle") + .op } /// Temporarily remove a handle while issuing a cancellation request. fn unjoin(&mut self, handle: u32) { - assert!(self.pending.contains_key(&handle)); + let pending = self + .pending + .get_mut(&handle) + .expect("no pending entry for handle"); + assert!(!pending.cancel, "operation cancellation already requested"); + pending.cancel = true; unsafe { waitable_join(handle, 0) }; } @@ -101,16 +192,6 @@ impl TaskInner { assert!(self.pending.contains_key(&handle)); unsafe { waitable_join(handle, self.waitable_set.unwrap()) }; } - - fn cancel(&mut self) { - for &handle in self.pending.keys() { - unsafe { waitable_join(handle, 0) }; - } - self.pending.clear(); - if let Some(set) = self.waitable_set.take() { - unsafe { waitable_set_drop(set) } - } - } } /// Task state for the active async export. @@ -127,7 +208,26 @@ impl TaskState { } pub(crate) fn is_active(&self) -> bool { - self.0.borrow().is_some() + self.0.borrow().as_ref().is_some_and(|inner| !inner.cancel) + } + + pub(crate) fn ensure_active(&self, ctx: &Ctx<'_>) -> rquickjs::Result<()> { + let inner = self.0.borrow(); + match inner.as_ref() { + Some(inner) if !inner.cancel => Ok(()), + Some(_) => Err(rquickjs::Exception::throw_message( + ctx, + "component task cancelled", + )), + None => Err(rquickjs::Exception::throw_type( + ctx, + "operation requires an active async call", + )), + } + } + + pub(crate) fn is_cancelling(&self) -> bool { + self.with(|inner| inner.cancel) } fn with(&self, f: impl FnOnce(&mut TaskInner) -> R) -> R { @@ -136,21 +236,60 @@ impl TaskState { } /// Initialize a fresh task state for a new async export call. - pub(crate) fn init(&self) { - *self.0.borrow_mut() = Some(TaskInner::default()); + pub(crate) fn init(&self, export_call: Box) { + let mut inner = self.0.borrow_mut(); + assert!(inner.is_none(), "task state already active"); + + *inner = Some(TaskInner { + export_call: Some(export_call), + ..TaskInner::default() + }); + } + + pub(crate) fn finish_export(&self) { + let call = self.with(|inner| { + assert!(!inner.cancel && !inner.returned, "task already completed"); + inner.returned = true; + inner.export_call.take() + }); + drop(call); } /// Restore task state previously transferred to the host by [`Self::poll`]. pub(crate) fn restore(&self, ptr: usize) { + assert_ne!(ptr, 0, "missing suspended task state"); // `poll` created this allocation with `Box::into_raw`; the host returns // the same pointer exactly once on the next callback. let inner = unsafe { *Box::from_raw(ptr as *mut TaskInner) }; - *self.0.borrow_mut() = Some(inner); + let mut state = self.0.borrow_mut(); + assert!(state.is_none(), "task state already active"); + *state = Some(inner); } - /// Cancel and clean up the current task state. - pub(crate) fn cancel(&self) { - self.with(|inner| inner.cancel()); + /// Request cancellation without releasing buffers still owned by the host. + pub(crate) fn cancel(&self, ctx: &Ctx<'_>) -> rquickjs::Result<()> { + let requests: Vec<_> = self.with(|inner| { + inner.cancel = true; + inner + .pending + .iter() + .filter(|(_, entry)| !entry.cancel) + .map(|(&handle, entry)| (handle, entry.op.cancel_request())) + .collect() + }); + + for (handle, request) in requests { + let pending = self.with(|inner| { + inner + .pending + .get(&handle) + .is_some_and(|entry| !entry.cancel) + }); + if pending { + request.cancel(ctx, handle)?; + } + } + Ok(()) } /// Register a pending operation, joining it to the waitable set. @@ -178,15 +317,19 @@ impl TaskState { pub(crate) fn set_stream_iterator_return( &self, handle: u32, - resolve: Persistent>, + resolve: Persistent>, ) { self.with(|inner| { - let Some(Pending::StreamRead { - iterator_return, .. + let Some(PendingEntry { + op: Pending::StreamRead { + iterator_return, .. + }, + .. }) = inner.pending.get_mut(&handle) else { panic!("no pending stream read for handle"); }; + assert!( iterator_return.replace(resolve).is_none(), "stream iterator return already pending" @@ -199,11 +342,23 @@ impl TaskState { /// Suspending transfers `TaskInner` into a raw host-context pointer. No Rust /// owner remains until the host supplies that pointer to [`Self::restore`]. pub(crate) fn poll(&self) -> u32 { - with_ctx(|ctx| while ctx.execute_pending_job() {}); + with_ctx(|ctx| { + while execute_pending_job(ctx).expect("QuickJS job failed") { + drain_resource_drops(ctx); + } + drain_resource_drops(ctx); + }); let mut inner = self.0.borrow_mut().take().expect("no active task state"); if inner.pending.is_empty() { + drop(inner.export_call.take()); + with_ctx(drain_resource_drops); + + if inner.cancel && !inner.returned { + unsafe { task_cancel() }; + } + if let Some(set) = inner.waitable_set.take() { unsafe { waitable_set_drop(set) } } @@ -221,11 +376,19 @@ impl TaskState { /// Reconcile a host subtask event with its pending JavaScript import promise. /// /// Returned subtasks lift their canonical result before settling the promise; -/// cancellation drops the subtask and resolves with `undefined`. +/// cancellation drops the subtask and rejects its promise. pub(crate) fn handle_subtask(handle: u32, state: SubtaskState) { match state { SubtaskState::Starting => unreachable!("Starting should not reach callback"), - SubtaskState::Started => { /* subtask started, nothing to do yet */ } + SubtaskState::Started => with_ctx(|ctx| { + ctx.task().with(|inner| { + let entry = inner.pending.get_mut(&handle).expect("no pending import"); + let Pending::ImportCall { call, .. } = &mut entry.op else { + unreachable!("expected ImportCall pending for started subtask"); + }; + call.complete_transfers(1); + }); + }), SubtaskState::Returned => { let pending = with_ctx(|ctx| ctx.task().take(handle)); unsafe { subtask_drop(handle) }; @@ -240,6 +403,7 @@ pub(crate) fn handle_subtask(handle: u32, state: SubtaskState) { else { unreachable!("expected ImportCall pending"); }; + call.complete_transfers(1); let func = with_ctx(|ctx| ctx.wit()).import_func(func_index); unsafe { func.lift_import_async_result(&mut call, buffer) }; @@ -252,13 +416,31 @@ pub(crate) fn handle_subtask(handle: u32, state: SubtaskState) { }); } SubtaskState::CancelledBeforeStarted | SubtaskState::CancelledBeforeReturned => { - let Pending::ImportCall { resolve, .. } = with_ctx(|ctx| ctx.task().take(handle)) + let pending = with_ctx(|ctx| ctx.task().take(handle)); + let Pending::ImportCall { + reject, mut call, .. + } = pending else { unreachable!("expected ImportCall pending for cancelled subtask"); }; unsafe { subtask_drop(handle) }; - resolve_promise(resolve, None); + + if state == SubtaskState::CancelledBeforeReturned { + call.complete_transfers(1); + } + + drop(call); + + let reason = with_ctx(|ctx| { + let error = + rquickjs::Exception::from_message(ctx.clone(), "async import cancelled") + .catch(ctx) + .expect("Failed to create cancellation error"); + Persistent::save(ctx, error.into_value()) + }); + + reject_promise(reject, reason); } } } diff --git a/crates/runtime/src/trivia.rs b/crates/runtime/src/trivia.rs index f94af5e..d14cd93 100644 --- a/crates/runtime/src/trivia.rs +++ b/crates/runtime/src/trivia.rs @@ -2,7 +2,7 @@ use crate::CtxExt; use crate::with_ctx; use heck::ToLowerCamelCase; -use rquickjs::{Atom, Function, Persistent, Symbol}; +use rquickjs::{Atom, Function, Object, Persistent, Symbol}; use rquickjs::{Result, Value, function::Rest}; /// Coerce closure lifetimes so the returned `Value<'js>` gets the same @@ -18,12 +18,11 @@ where /// Resolve a JS Promise by calling the stored resolve function with the given value. pub(crate) fn resolve_promise( - resolve: Persistent>, + resolve: Persistent>, result: Option>>, ) { with_ctx(|ctx| { - let resolve_val = resolve.restore(ctx).unwrap(); - let resolve_fn = resolve_val.get::().unwrap(); + let resolve_fn = resolve.restore(ctx).expect("restore promise resolver"); let result_val = result.map_or(Value::new_undefined(ctx.clone()), |res| { res.restore(ctx).unwrap() }); @@ -35,12 +34,11 @@ pub(crate) fn resolve_promise( } pub(crate) fn reject_promise( - reject: Persistent>, + reject: Persistent>, reason: Persistent>, ) { with_ctx(|ctx| { - let reject_val = reject.restore(ctx).unwrap(); - let reject_fn = reject_val.get::().unwrap(); + let reject_fn = reject.restore(ctx).expect("restore promise rejecter"); let reason_val = reason.restore(ctx).unwrap(); reject_fn @@ -49,9 +47,11 @@ pub(crate) fn reject_promise( }); } -/// Get `Symbol.for("dispose")` via the rquickjs API. +/// Get the well-known `Symbol.dispose`. pub(crate) fn symbol_dispose<'js>(ctx: &rquickjs::Ctx<'js>) -> Result> { - Ok(Symbol::new_global(ctx.clone(), "dispose")?.as_atom()) + let symbol: Object = ctx.globals().get("Symbol")?; + let dispose: Symbol = symbol.get("dispose")?; + Ok(dispose.as_atom()) } /// Convert a WIT function name to lower camel case, caching the result. diff --git a/crates/runtime/src/typed_array.rs b/crates/runtime/src/typed_array.rs new file mode 100644 index 0000000..9cac2ee --- /dev/null +++ b/crates/runtime/src/typed_array.rs @@ -0,0 +1,95 @@ +//! Checked typed-array access and copying into owned ABI buffers. +#![allow(unsafe_code)] + +use std::ptr::NonNull; + +use rquickjs::{Error, Exception, Result, TypedArray}; + +use crate::buffer::BufferGuard; + +pub(crate) trait TypedArrayExt { + /// Read the native element count without consulting the JS `length` property. + fn checked_len(&self) -> Result; + + fn copy_to_buffer(&self, align: usize) -> Result<(BufferGuard, usize)>; +} + +impl TypedArrayExt for TypedArray<'_, T> { + fn checked_len(&self) -> Result { + checked_view(self).map(|(_, count)| count) + } + + fn copy_to_buffer(&self, align: usize) -> Result<(BufferGuard, usize)> { + let (bytes, count) = checked_view(self)?; + let byte_len = bytes.len(); + let buffer = BufferGuard::new_uninit(byte_len, align); + + if byte_len > 0 { + // SAFETY: the runtime is single-threaded and this Rust allocation + // cannot execute JS. The source stays live until the disjoint copy + // finishes; no JS-backed pointer or slice escapes this function. + unsafe { + std::ptr::copy_nonoverlapping(bytes.cast::().as_ptr(), buffer.ptr(), byte_len); + } + } + Ok((buffer, count)) + } +} + +fn checked_view(array: &TypedArray<'_, T>) -> Result<(NonNull<[u8]>, usize)> { + let bytes = array.as_raw().ok_or_else(|| { + if array.ctx().has_exception() { + Error::Exception + } else { + Exception::throw_type(array.ctx(), "typed array backing store is unavailable") + } + })?; + + let width = std::mem::size_of::(); + if width == 0 || !bytes.len().is_multiple_of(width) { + return Err(Exception::throw_type( + array.ctx(), + "invalid typed array element layout", + )); + } + + Ok((bytes, bytes.len() / width)) +} + +// Type-specific casts need macros because rquickjs does not export TypedArrayItem. +macro_rules! try_typed_array_copy { + ($val:expr, $t:ty) => { + $val.as_object() + .and_then(|object| object.as_typed_array::<$t>()) + .map(|array| { + $crate::typed_array::TypedArrayExt::copy_to_buffer( + array, + std::mem::align_of::<$t>(), + ) + }) + .transpose() + }; +} + +macro_rules! copy_typed_array_as { + ($obj:expr, $ty:expr, $t:ty) => { + $obj.as_typed_array::<$t>() + .map(|array| { + assert_eq!($ty.abi_payload_size(), std::mem::size_of::<$t>()); + assert!($ty.abi_payload_align() >= std::mem::align_of::<$t>()); + + $crate::typed_array::TypedArrayExt::copy_to_buffer(array, $ty.abi_payload_align()) + }) + .transpose() + }; +} + +macro_rules! typed_array_len_as { + ($obj:expr, $t:ty) => { + $obj.as_typed_array::<$t>() + .map($crate::typed_array::TypedArrayExt::checked_len) + .transpose() + }; +} + +pub(crate) use {copy_typed_array_as, try_typed_array_copy, typed_array_len_as}; diff --git a/docs/runtime-intrinsics.md b/docs/runtime-intrinsics.md index 1dee064..af07a05 100644 --- a/docs/runtime-intrinsics.md +++ b/docs/runtime-intrinsics.md @@ -120,6 +120,13 @@ const { readable, writable } = wit.Future(wit.Future.STRING); If only one future type exists in the WIT world, `type` may be omitted. +### `wit.Future.from(value, type)` + +Adapt a value or Promise to a WIT future. Returns `{ readable, completion }`. +Use the readable handle when returning a future from an async export; returning +the payload Promise directly would cause JavaScript to await it before the WIT +boundary is reached. + ### Type Constants Type constants are generated for each stream/future element type found in @@ -195,6 +202,9 @@ Writable endpoint of a component-model stream. ### `FutureReadable` Readable endpoint of a component-model future. +This object is not a thenable; `await readable.read()` obtains its payload. +The native class traces its cached read Promise so repeated reads share one +operation without hiding JavaScript references from the garbage collector. | Method | Description | |---|---| @@ -203,6 +213,10 @@ Readable endpoint of a component-model future. | `drop()` | Drop the readable end. | | `[Symbol.dispose]()` | Alias for `drop()`. | +An immediately completed cancellation settles the original operation just like +a callback completion. A cancelled read clears the cached Promise, allowing a +new `read()` to retry. A completed future cannot be written a second time. + ### `FutureWritable` Writable endpoint of a component-model future. @@ -214,21 +228,78 @@ Writable endpoint of a component-model future. | `drop()` | Drop the writable end. | | `[Symbol.dispose]()` | Alias for `drop()`. | +All endpoint disposal methods use the well-known `Symbol.dispose`, not +`Symbol.for("dispose")`. Disposal and ownership transfer reject endpoints with +an in-progress operation; cancellation must settle first. + +### Imported Resources + +Imported resource wrappers store their resource type, handle, and ownership +state in native class data, not JavaScript properties. Their WIT-specific +prototypes inherit native disposal methods. Owned handles are invalidated on +transfer or explicit disposal; borrowed wrappers expire at the end of their +component call. + +Outbound loans keep native owner state alive and block disposal or transfer until +the import completes. Own transfers remain reserved until consumption is known: +partial writes restore unconsumed resource wrappers, and imports cancelled before +starting restore their owned arguments. Incoming borrow guards reject scope exit +while a borrowing import remains outstanding. + +The imported-resource registry and deferred-drop queue are initialized with the +JavaScript context, including for empty WIT worlds with no bindings to initialize. + +Garbage-collection finalizers enqueue abandoned owned handles. Host destructors +run from safe runtime boundaries after QuickJS has left its finalizer, with no +Rust queue borrow held across a host call. Explicit disposal is preferred when +prompt cleanup is required. + +### Task Cancellation + +Cancellation unjoins each pending waitable before invoking its cancellation +intrinsic. An immediate completion uses the normal completion handler; a +blocked cancellation rejoins the waitable and retains its conversion context, +callbacks, and ABI buffer until completion. + +Cancelled async imports reject their JavaScript Promises. During task +cancellation, new asynchronous operations are rejected and export completion +callbacks do not call `task.return`. The runtime acknowledges `task.cancel` +only after pending operations and their borrowed arguments have been released. + --- +## rquickjs Integration + +The runtime uses rquickjs 0.13. Numeric typed arrays are copied into aligned, +Rust-owned ABI buffers through `typed_array::TypedArrayExt`. The same module +contains native-view validation and shared type-specific casting macros. +`BufferGuard` only handles allocation and ownership; WIT type dispatch remains +in the call and stream conversion code. + +The adapter uses `TypedArray::as_raw()`. The raw view supplies the +actual view offset and byte length; element counts do not read an overridable +JavaScript `length` property. No JavaScript executes between obtaining the raw +view and copying its bytes, and no JS-backed pointer survives the copy. Attached +empty views are valid; detached or out-of-bounds views report an error. + +The custom job helper remains necessary: `Ctx::execute_pending_job()` still +loses exception status, while `Runtime::execute_pending_job()` would reacquire +the lock already held by `Context::with()`. `Promise::finish()` uses the same +lossy context method, so module evaluation retains its explicit job loop. +Likewise, rquickjs does not yet expose a typed `Symbol.dispose` constructor; +the runtime retrieves the well-known symbol from the global `Symbol` object. + ## Hidden Object Properties ### `__cqjs_handle` -A numeric property set on JS objects that wrap imported or exported WIT -resources. Stores the canonical component-model resource handle (`u32`). +A numeric property used on exported, JS-backed WIT resources. Stores the +canonical component-model resource handle (`u32`). Imported resources instead +use native opaque state and never trust this property. -- Set on: Resource wrapper objects during `push_borrow`, `push_own`, and - `exported_resource_to_handle` calls. -- Read by: `imported_resource_to_handle` and `exported_resource_to_handle` - to retrieve the canonical handle. -- Removed: When an owned resource is lifted back to JS via `push_own`, the - property is removed since the handle is no longer valid. +- Set/read by: `exported_resource_to_handle`. +- Removed: When an exported owned resource is lifted back to JS via + `push_own`, since the handle is no longer valid. ## WIT Import/Export Naming diff --git a/tests/common/mod.rs b/tests/common/mod.rs index b792951..5fe5f2d 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -256,7 +256,7 @@ impl ComponentInstance { let mut results = vec![Val::Bool(false); result_count]; func.call(&mut self.store, params, &mut results) - .unwrap_or_else(|e| panic!("calling `{name}` failed: {e}")); + .unwrap_or_else(|e| panic!("calling `{name}` failed: {e:#}")); results } diff --git a/tests/wasi.rs b/tests/wasi.rs index 04fb944..1dbde03 100644 --- a/tests/wasi.rs +++ b/tests/wasi.rs @@ -185,6 +185,46 @@ fn test_wasi_stdio() { let result = inst.call1("echo-stdin-to-stdout", &[]); assert_eq!(result, Val::Result(Ok(None))); assert_eq!(inst.stdout_bytes(), b"hello from stdin"); + assert!( + inst.parts().1.data().table.is_empty(), + "stdio resources should be released before post-return" + ); +} + +#[test] +fn test_wasi_resources_finalized_during_result_lowering() { + let mut inst = TestCase::new() + .wit_dir(wasi_wit_dir()) + .world("wasi-resource-cleanup") + .script( + r#" + import stdin from "wasi:cli/stdin@0.2.12"; + + export function makeResult() { + return { value: 7, extra: stdin.getStdin() }; + } + + export function flush() {} + "#, + ) + .build() + .expect("should build wasi-resource-cleanup component"); + + // Each next call must release the discarded resource before opening another stream. + inst.parts().1.data_mut().table.set_max_capacity(1); + + for _ in 0..2 { + assert_eq!( + inst.call1("make-result", &[]), + Val::Record(vec![("value".into(), Val::U32(7))]) + ); + } + + inst.call("flush", &[], 0); + assert!( + inst.parts().1.data().table.is_empty(), + "resources finalized during result lowering should be released on the next entry" + ); } #[tokio::test] diff --git a/tests/wit/test.wit b/tests/wit/test.wit index 86a3a78..720e5c9 100644 --- a/tests/wit/test.wit +++ b/tests/wit/test.wit @@ -28,6 +28,18 @@ world wasi-stdio { export echo-stdin-to-stdout: func() -> result; } +world wasi-resource-cleanup { + import wasi:cli/stdin@0.2.12; + import wasi:io/streams@0.2.12; + + record marker { + value: u32, + } + + export make-result: func() -> marker; + export flush: func(); +} + world wasi-import-types { import wasi:filesystem/types@0.2.12; diff --git a/tests/wit_types.rs b/tests/wit_types.rs index af41067..e6b1b78 100644 --- a/tests/wit_types.rs +++ b/tests/wit_types.rs @@ -886,16 +886,45 @@ fn test_list_of_strings() { #[test] fn test_empty_world() { - TestCase::new() - .wit( - r#" - package test:empty; - world empty {} - "#, - ) - .script("// empty module\n") - .build() - .unwrap(); + for script in ["// empty module\n", "await Promise.resolve();"] { + TestCase::new() + .wit( + r#" + package test:empty; + world empty {} + "#, + ) + .script(script) + .build() + .unwrap(); + } +} + +#[test] +fn test_empty_world_initialization_error() { + for script in [ + "throw new Error('empty-world initialization failed');", + "await Promise.reject(new Error('empty-world initialization failed'));", + ] { + let Err(err) = TestCase::new() + .wit( + r#" + package test:empty; + world empty {} + "#, + ) + .script(script) + .build() + else { + panic!("module initialization should fail"); + }; + + let message = format!("{err:#}"); + assert!( + message.contains("empty-world initialization failed"), + "{message}" + ); + } } #[test]