From ead50af50d23fd4bcc34b6fe94bcbde3d5a8f40d Mon Sep 17 00:00:00 2001 From: Tomasz Andrzejak Date: Wed, 30 Sep 2026 17:29:51 +0200 Subject: [PATCH] refactor: simplify runtime internals and artifact builds - share stream/future completion and endpoint lifecycle handling, - unify export resolution and remove async JS dispatch wrappers, - extract js adapters but keep hot-path stream helpers in rust, - share runtime builds between xargo and xtask --- .gitattributes | 3 + .../actions/setup-runtime-cache/action.yml | 11 +- .github/scripts/prepare-runtime-artifacts.sh | 9 +- .github/workflows/ci.yml | 11 +- Cargo.lock | 13 + Cargo.toml | 6 +- crates/core/Cargo.toml | 9 +- crates/core/build.rs | 515 +++-------------- crates/core/build/runtime.rs | 527 ++++++++++++++++++ crates/core/src/codegen.rs | 30 +- crates/core/src/js/future-from.js | 11 + crates/core/src/js/stream-from.js | 13 + crates/runtime/Cargo.toml | 3 + crates/runtime/src/abi.rs | 173 ------ crates/runtime/src/bindings.rs | 182 +----- crates/runtime/src/endpoint.rs | 251 +++++++++ crates/runtime/src/exports.rs | 177 ++++++ crates/runtime/src/futures.rs | 178 +++--- crates/runtime/src/interpreter.rs | 131 +---- crates/runtime/src/lib.rs | 31 +- crates/runtime/src/streams.rs | 447 ++++----------- crates/runtime/src/streams/helpers.rs | 244 ++++++++ crates/runtime/src/trivia.rs | 16 + crates/runtime/wit/init.wit | 11 - docs/runtime-intrinsics.md | 23 +- tests/async_types.rs | 388 +++++++++++-- tests/js/stream-helpers.js | 128 +++++ xtask/Cargo.toml | 14 + xtask/src/main.rs | 52 ++ 29 files changed, 2148 insertions(+), 1459 deletions(-) create mode 100644 crates/core/build/runtime.rs create mode 100644 crates/core/src/js/future-from.js create mode 100644 crates/core/src/js/stream-from.js create mode 100644 crates/runtime/src/endpoint.rs create mode 100644 crates/runtime/src/exports.rs create mode 100644 crates/runtime/src/streams/helpers.rs delete mode 100644 crates/runtime/wit/init.wit create mode 100644 tests/js/stream-helpers.js create mode 100644 xtask/Cargo.toml create mode 100644 xtask/src/main.rs diff --git a/.gitattributes b/.gitattributes index 9a33b88..84f2943 100644 --- a/.gitattributes +++ b/.gitattributes @@ -2,4 +2,7 @@ /Cargo.toml text eol=lf /Cargo.lock text eol=lf /crates/core/build.rs text eol=lf +/crates/core/build/** text=auto eol=lf +/crates/core/wit/** text=auto eol=lf /crates/runtime/** text=auto eol=lf +/xtask/** text=auto eol=lf diff --git a/.github/actions/setup-runtime-cache/action.yml b/.github/actions/setup-runtime-cache/action.yml index c428b63..50736a7 100644 --- a/.github/actions/setup-runtime-cache/action.yml +++ b/.github/actions/setup-runtime-cache/action.yml @@ -9,17 +9,10 @@ runs: uses: actions/cache@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 # v6.1.0 with: path: crates/core/prebuilt - key: runtime-wasm-v3-${{ hashFiles('Cargo.toml', 'Cargo.lock', 'crates/runtime/**', 'crates/core/build.rs') }} + key: runtime-wasm-v4-${{ hashFiles('Cargo.toml', 'Cargo.lock', 'crates/runtime/**', 'crates/core/build.rs', 'crates/core/build/**', 'crates/core/wit/**', 'xtask/**') }} enableCrossOsArchive: true - name: Build runtime wasm if: steps.cache.outputs.cache-hit != 'true' shell: bash - run: | - cargo build --locked -p componentize-qjs - mkdir -p crates/core/prebuilt - for f in runtime.wasm runtime-opt-size.wasm runtime-sync.wasm runtime-opt-size-sync.wasm; do - src=$(find target -path "*/out/$f" -type f | sort | tail -n 1) - test -n "$src" || { echo "ERROR: built runtime $f not found"; exit 1; } - cp "$src" "crates/core/prebuilt/$f" - done + run: cargo run --locked -p xtask -- --output crates/core/prebuilt diff --git a/.github/scripts/prepare-runtime-artifacts.sh b/.github/scripts/prepare-runtime-artifacts.sh index 6e3cc55..bf8bcf0 100644 --- a/.github/scripts/prepare-runtime-artifacts.sh +++ b/.github/scripts/prepare-runtime-artifacts.sh @@ -5,17 +5,18 @@ target_dir="$1" destination_dir="$2" shift 2 -rm -rf crates/core/prebuilt "$destination_dir" "$target_dir" -COMPONENTIZE_QJS_RUNTIME_AUDITABLE=1 cargo build --release -p componentize-qjs --target-dir "$target_dir" +artifact_dir="$target_dir/runtime" +cargo run --locked -p xtask -- --release --auditable \ + --output "$artifact_dir" --build-dir "$target_dir/build" mkdir -p "$destination_dir" for mapping in "$@"; do source_name="${mapping%%=*}" destination_name="${mapping#*=}" - source_path=$(find "$target_dir" -path "*/out/$source_name" -type f | sort | tail -n 1) + source_path="$artifact_dir/$source_name" destination_path="$destination_dir/$destination_name" - test -n "$source_path" || { echo "ERROR: $source_name not found"; exit 1; } + test -f "$source_path" || { echo "ERROR: $source_name not found"; exit 1; } cp "$source_path" "$destination_path" test -f "$destination_path" || { echo "ERROR: $destination_name not created"; exit 1; } diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4dbdf3a..6784809 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -40,10 +40,13 @@ jobs: uses: ./.github/actions/setup-runtime-cache - name: Check formatting - run: cargo fmt --check + run: cargo fmt --all --check - name: Clippy - run: cargo clippy --all-features -- -D warnings + run: cargo clippy --workspace --all-features -- -D warnings + + - name: Runtime and build-tool unit tests + run: cargo test --locked -p componentize-qjs-runtime -p xtask build-test: name: Build and test (${{ matrix.os }}) @@ -78,6 +81,10 @@ jobs: - name: Setup runtime wasm cache uses: ./.github/actions/setup-runtime-cache + - name: Build Windows runtime from source + if: runner.os == 'Windows' + run: cargo run --locked -p xtask -- --sync-only --output target/runtime-source-check + - name: Run tests run: cargo test -p componentize-qjs-cli diff --git a/Cargo.lock b/Cargo.lock index f5b53a2..4dbed5d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -468,6 +468,7 @@ dependencies = [ "heck", "indexmap", "num_enum", + "quickcheck", "rquickjs", "smallvec", "wit-bindgen 0.61.1", @@ -4167,6 +4168,18 @@ dependencies = [ "rustix", ] +[[package]] +name = "xtask" +version = "0.0.0" +dependencies = [ + "anyhow", + "clap", + "flate2", + "glob", + "tar", + "ureq", +] + [[package]] name = "yoke" version = "0.8.3" diff --git a/Cargo.toml b/Cargo.toml index fe24f7b..5738b84 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = ["crates/core", "crates/runtime", "napi"] +members = ["crates/core", "crates/runtime", "napi", "xtask"] default-members = [".", "napi"] resolver = "2" @@ -16,6 +16,10 @@ componentize-qjs = { path = "crates/core", version = "0.4.5", default-features = componentize-qjs-cli = { path = ".", version = "0.4.5", default-features = false } anyhow = "1.0" clap = { version = "4", features = ["derive"] } +flate2 = "1.1" +glob = "0.3" +tar = "0.4" +ureq = "3" tokio = { version = "1", default-features = false, features = ["rt", "macros"] } wasmtime = { version = "47", features = ["component-model", "async"] } wasmtime-wasi = { version = "47", default-features = false, features = ["p2", "p3"] } diff --git a/crates/core/Cargo.toml b/crates/core/Cargo.toml index 88bf260..1c7faac 100644 --- a/crates/core/Cargo.toml +++ b/crates/core/Cargo.toml @@ -12,6 +12,7 @@ include = [ "src/**/*", "wit/**/*", "build.rs", + "build/**/*", "../../README.md", "prebuilt/runtime.wasm", "prebuilt/runtime.wasm.cdx.json", @@ -49,10 +50,10 @@ tempfile = "3" [build-dependencies] anyhow.workspace = true -flate2 = "1.1" -glob = "0.3" -tar = "0.4" -ureq = "3" +flate2.workspace = true +glob.workspace = true +tar.workspace = true +ureq.workspace = true [features] default = ["component-model-async"] diff --git a/crates/core/build.rs b/crates/core/build.rs index 27379de..d73d9b2 100644 --- a/crates/core/build.rs +++ b/crates/core/build.rs @@ -1,484 +1,99 @@ -use std::io::Read; +//! Select prepared runtimes, falling back to the shared source-build recipe. + +#[path = "build/runtime.rs"] +pub mod runtime; + use std::path::{Path, PathBuf}; -use std::process::Command; use std::{env, fs}; use anyhow::{Context, Result, bail}; -use flate2::read::GzDecoder; - -const WASI_SDK_VERSION: &str = "33"; -const WASI_SKD_DL_URL: &str = "https://github.com/WebAssembly/wasi-sdk/releases/download"; - -const BINARYEN_VERSION: &str = "130"; -const BINARYEN_DL_URL: &str = "https://github.com/WebAssembly/binaryen/releases/download"; -const RUNTIME_AUDITABLE_ENV: &str = "COMPONENTIZE_QJS_RUNTIME_AUDITABLE"; -const MAX_ARCHIVE_BYTES: u64 = 1_000_000_000; - -/// Description of one embedded runtime artifact. -/// -/// `fallback_const` is used for async constants when the async Cargo feature is -/// disabled; in that configuration they alias the corresponding sync runtime. -#[derive(Clone, Copy)] -struct RuntimeBuild { - name: &'static str, - filename: &'static str, - const_name: &'static str, - optimize_size: bool, - async_support: bool, - fallback_const: Option<&'static str>, -} - -/// All runtime artifacts in generated-constant order. -const RUNTIME_BUILDS: [RuntimeBuild; 4] = [ - RuntimeBuild { - name: "default-sync", - filename: "runtime-sync.wasm", - const_name: "DEFAULT_SYNC_RUNTIME_WASM", - optimize_size: false, - async_support: false, - fallback_const: None, - }, - RuntimeBuild { - name: "opt-size-sync", - filename: "runtime-opt-size-sync.wasm", - const_name: "OPT_SIZE_SYNC_RUNTIME_WASM", - optimize_size: true, - async_support: false, - fallback_const: None, - }, - RuntimeBuild { - name: "default", - filename: "runtime.wasm", - const_name: "DEFAULT_RUNTIME_WASM", - optimize_size: false, - async_support: true, - fallback_const: Some("DEFAULT_SYNC_RUNTIME_WASM"), - }, - RuntimeBuild { - name: "opt-size", - filename: "runtime-opt-size.wasm", - const_name: "OPT_SIZE_RUNTIME_WASM", - optimize_size: true, - async_support: true, - fallback_const: Some("OPT_SIZE_SYNC_RUNTIME_WASM"), - }, -]; - -struct CargoProfile { - name: String, - release: bool, -} - -impl CargoProfile { - fn current() -> Self { - let name = env::var("PROFILE").unwrap_or_else(|_| "debug".to_string()); - let release = name == "release"; - Self { name, release } - } - - fn runtime_rustflags(&self, optimize_size: bool) -> String { - let flags = "-Clink-arg=-shared -Clink-arg=-Wl,--no-entry -Clink-arg=-Wl,--allow-undefined"; - match (self.release, optimize_size) { - (true, true) => format!("{flags} -Clto=fat -Copt-level=z"), - (true, false) => format!("{flags} -Clto=fat -Copt-level=3"), - (false, _) => flags.to_string(), - } - } - - fn runtime_cflags(&self, optimize_size: bool) -> String { - let flags = "-fPIC"; - match (self.release, optimize_size) { - (true, true) => format!("{flags} -Oz"), - (true, false) => format!("{flags} -O3"), - (false, _) => flags.to_string(), - } - } - - fn configure_nested_build(&self, cargo: &mut Command) { - if !self.release { - set_env_if_unset(cargo, "CARGO_PROFILE_DEV_DEBUG", "0"); - } - // Runtime target directories are disposable, so incremental state - // cannot be reused after this build script finishes. - set_env_if_unset(cargo, "CARGO_INCREMENTAL", "0"); - } -} - -type RuntimePaths = [Option; RUNTIME_BUILDS.len()]; - -/// Shared nested Cargo targets for speed- and size-optimized variants. -/// -/// Sync and async builds with the same optimization flags share dependencies. -/// Both directories remain alive until every variant has finished. -struct RuntimeTargetDirs { - default: PathBuf, - opt_size: PathBuf, -} - -impl RuntimeTargetDirs { - fn new(out_dir: &Path) -> Self { - Self { - default: out_dir.join("runtime-default"), - opt_size: out_dir.join("runtime-opt-size"), - } - } - - fn get(&self, optimize_size: bool) -> &Path { - if optimize_size { - &self.opt_size - } else { - &self.default - } - } -} - -impl Drop for RuntimeTargetDirs { - fn drop(&mut self) { - cleanup_runtime_target_dir(&self.default); - cleanup_runtime_target_dir(&self.opt_size); - } -} +use runtime::{RUNTIME_AUDITABLE_ENV, RUNTIME_BUILDS, RuntimeOptions}; +/// Track build inputs, obtain the selected artifacts, and emit embedding constants. fn main() -> Result<()> { let manifest_dir = PathBuf::from(env::var("CARGO_MANIFEST_DIR").context("CARGO_MANIFEST_DIR not set")?); - let runtime_dir = manifest_dir.join("../runtime"); - - println!("cargo:rerun-if-changed={}/src", runtime_dir.display()); - println!( - "cargo:rerun-if-changed={}/Cargo.toml", - runtime_dir.display() - ); - println!("cargo:rerun-if-changed=prebuilt/runtime.wasm"); - println!("cargo:rerun-if-changed=prebuilt/runtime-opt-size.wasm"); - println!("cargo:rerun-if-changed=prebuilt/runtime-sync.wasm"); - println!("cargo:rerun-if-changed=prebuilt/runtime-opt-size-sync.wasm"); - println!("cargo:rerun-if-env-changed={RUNTIME_AUDITABLE_ENV}"); - let out_dir = PathBuf::from(env::var("OUT_DIR").context("OUT_DIR not set")?); - let async_on = component_model_async_enabled(); - // Check for pre-built runtimes (used when installing from crates.io) - let prebuilt_dir = manifest_dir.join("prebuilt"); - let prebuilt_sync = prebuilt_dir.join("runtime-sync.wasm"); + for input in [ + "../../Cargo.toml", + "../../Cargo.lock", + "../runtime/src", + "../runtime/Cargo.toml", + "build", + "wit", + ] { + let path = manifest_dir.join(input); + + if path.exists() { + println!("cargo:rerun-if-changed={}", path.display()); + } + } - if prebuilt_sync.exists() { - return emit_from_prebuilt(&prebuilt_dir, async_on, &out_dir); + for build in RUNTIME_BUILDS { + println!("cargo:rerun-if-changed=prebuilt/{}", build.filename); } - // Check that runtime source is available (won't be when installed from crates.io - // without a pre-built runtime) - let runtime_src_dir = runtime_dir.join("src"); - if !runtime_src_dir.exists() { - bail!( - "Runtime source not found at {} and no pre-built runtime at {}. \ - If installing from crates.io, this is a packaging bug.", - runtime_src_dir.display(), - prebuilt_sync.display(), - ); + for variable in [RUNTIME_AUDITABLE_ENV, "WASI_SDK_PATH", "WASM_OPT"] { + println!("cargo:rerun-if-env-changed={variable}"); } - let profile = CargoProfile::current(); - let target_dirs = RuntimeTargetDirs::new(&out_dir); - let mut paths = RuntimePaths::default(); + let async_support = env::var_os("CARGO_FEATURE_COMPONENT_MODEL_ASYNC").is_some(); + let prebuilt_dir = manifest_dir.join("prebuilt"); - for (index, build) in RUNTIME_BUILDS.iter().copied().enumerate() { - if build.async_support && !async_on { - continue; - } - paths[index] = Some(build_runtime(&out_dir, &target_dirs, build, &profile)?); - } + let artifact_dir = if prebuilt_dir.join("runtime-sync.wasm").exists() { + eprintln!("Using prebuilt runtimes from: {}", prebuilt_dir.display()); + prebuilt_dir + } else { + runtime::prepare( + &manifest_dir.join("../runtime"), + &out_dir, + &out_dir, + RuntimeOptions { + release: env::var("PROFILE").is_ok_and(|profile| profile == "release"), + async_support, + auditable: env::var_os(RUNTIME_AUDITABLE_ENV).is_some(), + }, + )?; + out_dir.clone() + }; - emit_runtime_wasms(&paths, &out_dir) + let source = embedding_source(&artifact_dir, async_support)?; + fs::write(out_dir.join("output.rs"), source).context("Failed to write output.rs") } -/// Emit runtime constants from the pre-built runtimes packaged with the crate. -fn emit_from_prebuilt(prebuilt_dir: &Path, async_on: bool, out_dir: &Path) -> Result<()> { - let mut paths = RuntimePaths::default(); - for (index, build) in RUNTIME_BUILDS.iter().copied().enumerate() { - if build.async_support && !async_on { +/// Embed every selected artifact, aliasing disabled async variants to sync ones. +fn embedding_source(artifact_dir: &Path, async_support: bool) -> Result { + let mut source = String::new(); + + for build in RUNTIME_BUILDS { + if build.async_support && !async_support { + let fallback = build + .fallback_const + .expect("async variants have sync fallbacks"); + source.push_str(&format!( + "const {}: &[u8] = {fallback};\n", + build.const_name + )); continue; } - let path = prebuilt_dir.join(build.filename); - if !path.exists() { + let path = artifact_dir.join(build.filename); + + if !path.is_file() { bail!( - "Pre-built {} runtime is missing at {}. If installing from crates.io, \ + "Prepared {} runtime is missing at {}. If installing from crates.io, \ this is a packaging bug.", build.name, path.display(), ); } - paths[index] = Some(path); - } - - eprintln!("Using prebuilt runtimes from: {}", prebuilt_dir.display()); - - emit_runtime_wasms(&paths, out_dir) -} - -fn emit_runtime_wasms(paths: &RuntimePaths, out_dir: &Path) -> Result<()> { - let mut output = String::new(); - for (build, path) in RUNTIME_BUILDS.iter().zip(paths) { - if let Some(path) = path { - output.push_str(&const_line(build.const_name, path)); - continue; - } - let Some(fallback_const) = build.fallback_const else { - bail!("missing {} runtime artifact", build.name); - }; - - output.push_str(&format!( - "const {}: &[u8] = {fallback_const};\n", + source.push_str(&format!( + "const {}: &[u8] = include_bytes!({path:?});\n", build.const_name )); } - fs::write(out_dir.join("output.rs"), output).context("Failed to write output.rs")?; - - Ok(()) -} - -fn const_line(name: &str, path: &Path) -> String { - format!("const {name}: &[u8] = include_bytes!({path:?});\n") -} - -fn set_env_if_unset(cargo: &mut Command, key: &str, value: &str) { - if env::var_os(key).is_none() { - cargo.env(key, value); - } -} - -fn build_runtime( - out_dir: &Path, - target_dirs: &RuntimeTargetDirs, - build: RuntimeBuild, - profile: &CargoProfile, -) -> Result { - let target = "wasm32-wasip2"; - let upcase = target.to_uppercase().replace('-', "_"); - - // Get wasi-sdk - from env, cached, or download - let wasi_sdk = get_wasi_sdk(out_dir)?; - eprintln!("Using wasi-sdk at: {}", wasi_sdk.display()); - - let optimize_size = build.optimize_size; - let rustflags = profile.runtime_rustflags(optimize_size); - let cflags = profile.runtime_cflags(optimize_size); - - let clang = executable(&wasi_sdk, "bin/clang"); - let target_dir = target_dirs.get(optimize_size); - let mut cargo = Command::new("cargo"); - if env::var_os(RUNTIME_AUDITABLE_ENV).is_some() { - cargo.arg("auditable"); - } - cargo - .arg("build") - .arg("--target") - .arg(target) - .arg("--package=componentize-qjs-runtime") - .arg("--no-default-features") - .env("CARGO_TARGET_DIR", target_dir) - .env(format!("CARGO_TARGET_{upcase}_RUSTFLAGS"), rustflags) - .env(format!("CARGO_TARGET_{upcase}_LINKER"), &clang) - .env(format!("CFLAGS_{}", target.replace('-', "_")), cflags) - .env(format!("CC_{}", target.replace('-', "_")), &clang) - .env("WASI_SDK_PATH", &wasi_sdk) - .env("WASI_SDK", &wasi_sdk) - .env_remove("CARGO_ENCODED_RUSTFLAGS"); - - profile.configure_nested_build(&mut cargo); - - if profile.release { - cargo.arg("--release"); - } - - if build.async_support { - cargo.arg("--features").arg("component-model-async"); - } - - eprintln!("Building {} runtime: {cargo:?}", build.name); - let status = cargo.status().context("Failed to run cargo build")?; - if !status.success() { - bail!("Failed to build {} runtime", build.name); - } - - let runtime_src = target_dir - .join(target) - .join(&profile.name) - .join("componentize_qjs_runtime.wasm"); - - let runtime_dst = out_dir.join(build.filename); - - fs::copy(&runtime_src, &runtime_dst) - .with_context(|| format!("Failed to copy {}", runtime_src.display()))?; - - if profile.release { - let wasm_opt = get_wasm_opt(out_dir)?; - let opt_level = if optimize_size { "-Oz" } else { "-O3" }; - - let status = Command::new(&wasm_opt) - .arg(opt_level) - .arg("--all-features") - .arg("--disable-gc") - .arg("--disable-reference-types") - .arg("--strip-debug") - .arg("--strip-producers") - .arg(&runtime_dst) - .arg("-o") - .arg(&runtime_dst) - .status() - .context("Failed to run wasm-opt")?; - - if !status.success() { - bail!("wasm-opt failed"); - } - } - - Ok(runtime_dst) -} - -fn cleanup_runtime_target_dir(target_dir: &Path) { - match fs::remove_dir_all(target_dir) { - Ok(()) => {} - Err(err) if err.kind() == std::io::ErrorKind::NotFound => {} - Err(err) => { - eprintln!( - "warning: failed to clean nested runtime target dir {}: {err}", - target_dir.display() - ); - } - } -} - -fn component_model_async_enabled() -> bool { - env::var_os("CARGO_FEATURE_COMPONENT_MODEL_ASYNC").is_some() -} - -fn get_wasi_sdk(out_dir: &Path) -> Result { - // Check environment first - if let Ok(path) = env::var("WASI_SDK_PATH") { - let p = PathBuf::from(path); - if executable(&p, "bin/clang").exists() { - return Ok(p); - } - } - - // Check cached location - let stable = out_dir.join("wasi-sdk"); - if executable(&stable, "bin/clang").exists() { - return Ok(stable); - } - - // Download wasi-sdk - let (arch, os) = system()?; - let filename = format!("wasi-sdk-{WASI_SDK_VERSION}.0-{arch}-{os}.tar.gz"); - let url = format!("{WASI_SKD_DL_URL}/wasi-sdk-{WASI_SDK_VERSION}/{filename}"); - - http_archive(&url, out_dir)?; - - // Rename extracted directory to stable location - let extracted = find_wasi_sdk(out_dir).context("Could not find extracted wasi-sdk")?; - fs::rename(&extracted, &stable).context("Failed to rename wasi-sdk directory")?; - - Ok(stable) -} - -fn find_wasi_sdk(target_dir: &Path) -> Option { - let pattern = target_dir.join("wasi-sdk*"); - glob::glob(pattern.to_str()?) - .ok()? - .filter_map(Result::ok) - .find(|entry| entry.is_dir() && executable(entry, "bin/clang").exists()) -} - -fn get_wasm_opt(out_dir: &Path) -> Result { - // Check WASM_OPT environment variable first - if let Ok(path) = env::var("WASM_OPT") { - let p = PathBuf::from(path); - if p.exists() { - return Ok(p); - } - } - - // Check cached location - let stable = out_dir.join("binaryen"); - let wasm_opt = executable(&stable, "bin/wasm-opt"); - if wasm_opt.exists() { - return Ok(wasm_opt); - } - - // Download binaryen - let (arch, os) = system()?; - let tag = format!("version_{BINARYEN_VERSION}"); - let filename = format!("binaryen-{tag}-{arch}-{os}.tar.gz"); - let url = format!("{BINARYEN_DL_URL}/{tag}/{filename}"); - - http_archive(&url, out_dir)?; - - // Rename extracted directory to stable location - let extracted = find_binaryen(out_dir).context("Could not find extracted binaryen")?; - fs::rename(&extracted, &stable).context("Failed to rename binaryen directory")?; - - Ok(executable(&stable, "bin/wasm-opt")) -} - -fn find_binaryen(target_dir: &Path) -> Option { - let pattern = target_dir.join("binaryen*"); - glob::glob(pattern.to_str()?) - .ok()? - .filter_map(Result::ok) - .find(|entry| entry.is_dir() && executable(entry, "bin/wasm-opt").exists()) -} - -fn executable(root: &Path, relative: &str) -> PathBuf { - let mut path = root.join(relative); - if !env::consts::EXE_SUFFIX.is_empty() { - path.set_extension(&env::consts::EXE_SUFFIX[1..]); - } - path -} - -fn system() -> Result<(&'static str, &'static str)> { - let (arch, os) = match (env::consts::ARCH, env::consts::OS) { - ("x86_64", "linux") => ("x86_64", "linux"), - ("aarch64", "linux") => ("arm64", "linux"), - ("x86_64", "macos") => ("x86_64", "macos"), - ("aarch64", "macos") => ("arm64", "macos"), - ("x86_64", "windows") => ("x86_64", "windows"), - ("aarch64", "windows") => ("arm64", "windows"), - (arch, os) => bail!("Unsupported platform: {arch}-{os}"), - }; - - Ok((arch, os)) -} - -fn http_archive(url: &str, out_dir: &Path) -> Result<()> { - eprintln!("Downloading archive from {url}..."); - - let response = ureq::get(url) - .call() - .context("Failed to download wasi-sdk")?; - - let mut bytes = Vec::new(); - response - .into_body() - .into_reader() - .take(MAX_ARCHIVE_BYTES + 1) - .read_to_end(&mut bytes) - .context("Failed to download archive")?; - - if bytes.len() as u64 > MAX_ARCHIVE_BYTES { - bail!("Archive exceeds maximum download size of {MAX_ARCHIVE_BYTES} bytes"); - } - - let decoder = GzDecoder::new(bytes.as_slice()); - - let mut archive = tar::Archive::new(decoder); - archive - .unpack(out_dir) - .context("Failed to extract archive")?; - - Ok(()) + Ok(source) } diff --git a/crates/core/build/runtime.rs b/crates/core/build/runtime.rs new file mode 100644 index 0000000..2fd746e --- /dev/null +++ b/crates/core/build/runtime.rs @@ -0,0 +1,527 @@ +//! Portable runtime-artifact preparation shared by xtask and the Cargo fallback. + +use std::io::Read; +use std::path::{Path, PathBuf, absolute}; +use std::process::Command; +use std::{env, fs}; + +use anyhow::{Context, Result, bail}; +use flate2::read::GzDecoder; + +const WASI_SDK_VERSION: &str = "33"; +const WASI_SKD_DL_URL: &str = "https://github.com/WebAssembly/wasi-sdk/releases/download"; + +const BINARYEN_VERSION: &str = "130"; +const BINARYEN_DL_URL: &str = "https://github.com/WebAssembly/binaryen/releases/download"; +/// Compatibility switch for dependency-auditable source builds. +pub const RUNTIME_AUDITABLE_ENV: &str = "COMPONENTIZE_QJS_RUNTIME_AUDITABLE"; +const MAX_ARCHIVE_BYTES: u64 = 1_000_000_000; + +/// Description of one embedded runtime artifact. +/// +/// `fallback_const` is used for async constants when the async Cargo feature is +/// disabled; in that configuration they alias the corresponding sync runtime. +#[derive(Clone, Copy)] +pub struct RuntimeBuild { + /// Human-readable variant name. + pub name: &'static str, + /// Stable artifact filename, independent of Cargo's OUT_DIR layout. + pub filename: &'static str, + /// Constant emitted when the artifact is embedded. + pub const_name: &'static str, + /// Whether this artifact is optimized for size rather than speed. + pub optimize_size: bool, + /// Whether this artifact requires component-model async support. + pub async_support: bool, + /// Sync constant used when async embedding is disabled. + pub fallback_const: Option<&'static str>, +} + +/// All runtime artifacts in generated-constant order. +pub const RUNTIME_BUILDS: [RuntimeBuild; 4] = [ + RuntimeBuild { + name: "default-sync", + filename: "runtime-sync.wasm", + const_name: "DEFAULT_SYNC_RUNTIME_WASM", + optimize_size: false, + async_support: false, + fallback_const: None, + }, + RuntimeBuild { + name: "opt-size-sync", + filename: "runtime-opt-size-sync.wasm", + const_name: "OPT_SIZE_SYNC_RUNTIME_WASM", + optimize_size: true, + async_support: false, + fallback_const: None, + }, + RuntimeBuild { + name: "default", + filename: "runtime.wasm", + const_name: "DEFAULT_RUNTIME_WASM", + optimize_size: false, + async_support: true, + fallback_const: Some("DEFAULT_SYNC_RUNTIME_WASM"), + }, + RuntimeBuild { + name: "opt-size", + filename: "runtime-opt-size.wasm", + const_name: "OPT_SIZE_RUNTIME_WASM", + optimize_size: true, + async_support: true, + fallback_const: Some("OPT_SIZE_SYNC_RUNTIME_WASM"), + }, +]; + +/// Options common to explicit preparation and automatic source builds. +#[derive(Clone, Copy, Debug)] +pub struct RuntimeOptions { + /// Optimize the Rust/C runtime and run wasm-opt. + pub release: bool, + /// Include the two async artifacts in addition to the sync variants. + pub async_support: bool, + /// Build through cargo-auditable to retain dependency metadata. + pub auditable: bool, +} + +/// Compiler settings for one Cargo profile. +struct CargoProfile { + name: &'static str, + release: bool, +} + +impl CargoProfile { + /// Select the nested Cargo profile independently of the caller's profile. + fn new(release: bool) -> Self { + let name = if release { "release" } else { "debug" }; + Self { name, release } + } + + /// Link a relocatable runtime with the requested optimization objective. + fn runtime_rustflags(&self, optimize_size: bool) -> String { + let flags = "-Clink-arg=-shared -Clink-arg=-Wl,--no-entry -Clink-arg=-Wl,--allow-undefined"; + + match (self.release, optimize_size) { + (true, true) => format!("{flags} -Clto=fat -Copt-level=z"), + (true, false) => format!("{flags} -Clto=fat -Copt-level=3"), + (false, _) => flags.to_string(), + } + } + + /// Compile QuickJS C sources as position-independent code. + fn runtime_cflags(&self, optimize_size: bool) -> String { + let flags = "-fPIC"; + + match (self.release, optimize_size) { + (true, true) => format!("{flags} -Oz"), + (true, false) => format!("{flags} -O3"), + (false, _) => flags.to_string(), + } + } + + /// Avoid disposable debug/incremental data while respecting caller overrides. + fn configure_nested_build(&self, cargo: &mut Command) { + if !self.release { + set_env_if_unset(cargo, "CARGO_PROFILE_DEV_DEBUG", "0"); + } + // Runtime target directories are disposable, so incremental state + // cannot be reused after this build script finishes. + set_env_if_unset(cargo, "CARGO_INCREMENTAL", "0"); + } +} + +/// Shared nested Cargo targets for speed- and size-optimized variants. +/// +/// Sync and async builds with the same optimization flags share dependencies. +/// Both directories remain alive until every variant has finished. +struct RuntimeTargetDirs { + default: PathBuf, + opt_size: PathBuf, +} + +impl RuntimeTargetDirs { + /// Allocate separate target paths for incompatible optimization flags. + fn new(out_dir: &Path) -> Self { + Self { + default: out_dir.join("runtime-default"), + opt_size: out_dir.join("runtime-opt-size"), + } + } + + /// Share one target directory between sync/async variants with matching flags. + fn get(&self, optimize_size: bool) -> &Path { + if optimize_size { + &self.opt_size + } else { + &self.default + } + } +} + +impl Drop for RuntimeTargetDirs { + fn drop(&mut self) { + cleanup_runtime_target_dir(&self.default); + cleanup_runtime_target_dir(&self.opt_size); + } +} + +/// Prepare all selected runtime variants at their stable artifact paths. +/// +/// Validate source/output locations, share nested dependency builds between +/// compatible variants, then copy and optionally optimize each resulting Wasm. +pub fn prepare( + runtime_dir: &Path, + out_dir: &Path, + build_dir: &Path, + options: RuntimeOptions, +) -> Result<()> { + if !runtime_dir.join("src").is_dir() { + bail!( + "Runtime source not found at {}. If installing from crates.io, \ + missing pre-built runtimes are a packaging bug.", + runtime_dir.display(), + ); + } + + fs::create_dir_all(out_dir).context("Failed to create runtime artifact directory")?; + fs::create_dir_all(build_dir).context("Failed to create runtime build directory")?; + + // Canonicalization introduces Windows verbatim paths that break nested + // dependencies using include!(concat!(env!("OUT_DIR"), "/bindings.rs")). + let out_dir = absolute(out_dir).context("Failed to resolve runtime artifact directory")?; + let build_dir = absolute(build_dir).context("Failed to resolve runtime build directory")?; + let runtime_dir = + absolute(runtime_dir).context("Failed to resolve runtime source directory")?; + let profile = CargoProfile::new(options.release); + let target_dirs = RuntimeTargetDirs::new(&build_dir); + + for build in RUNTIME_BUILDS { + if build.async_support && !options.async_support { + continue; + } + + build_runtime( + &runtime_dir, + &out_dir, + &build_dir, + &target_dirs, + build, + &profile, + options.auditable, + )?; + } + + Ok(()) +} + +/// Supply a nested-build default only when the caller did not set it. +fn set_env_if_unset(cargo: &mut Command, key: &str, value: &str) { + if env::var_os(key).is_none() { + cargo.env(key, value); + } +} + +/// Build, stage, and optionally optimize one runtime artifact. +fn build_runtime( + runtime_dir: &Path, + out_dir: &Path, + build_dir: &Path, + target_dirs: &RuntimeTargetDirs, + build: RuntimeBuild, + profile: &CargoProfile, + auditable: bool, +) -> Result<()> { + let target = "wasm32-wasip2"; + let upcase = target.to_uppercase().replace('-', "_"); + + // Get wasi-sdk - from env, cached, or download + let wasi_sdk = get_wasi_sdk(build_dir)?; + eprintln!("Using wasi-sdk at: {}", wasi_sdk.display()); + + let optimize_size = build.optimize_size; + let rustflags = profile.runtime_rustflags(optimize_size); + let cflags = profile.runtime_cflags(optimize_size); + + let clang = executable(&wasi_sdk, "bin/clang"); + let target_dir = target_dirs.get(optimize_size); + let mut cargo = Command::new("cargo"); + + if auditable { + cargo.arg("auditable"); + } + + cargo + .current_dir(runtime_dir) + .arg("build") + .arg("--target") + .arg(target) + .arg("--package=componentize-qjs-runtime") + .arg("--no-default-features") + .env("CARGO_TARGET_DIR", target_dir) + .env(format!("CARGO_TARGET_{upcase}_RUSTFLAGS"), rustflags) + .env(format!("CARGO_TARGET_{upcase}_LINKER"), &clang) + .env(format!("CFLAGS_{}", target.replace('-', "_")), cflags) + .env(format!("CC_{}", target.replace('-', "_")), &clang) + .env("WASI_SDK_PATH", &wasi_sdk) + .env("WASI_SDK", &wasi_sdk) + .env_remove("CARGO_ENCODED_RUSTFLAGS"); + + profile.configure_nested_build(&mut cargo); + + if profile.release { + cargo.arg("--release"); + } + + if build.async_support { + cargo.arg("--features").arg("component-model-async"); + } + + eprintln!("Building {} runtime: {cargo:?}", build.name); + let status = cargo.status().context("Failed to run cargo build")?; + if !status.success() { + bail!("Failed to build {} runtime", build.name); + } + + let runtime_src = target_dir + .join(target) + .join(profile.name) + .join("componentize_qjs_runtime.wasm"); + + let runtime_dst = out_dir.join(build.filename); + + fs::copy(&runtime_src, &runtime_dst) + .with_context(|| format!("Failed to copy {}", runtime_src.display()))?; + + if profile.release { + let wasm_opt = get_wasm_opt(build_dir)?; + let opt_level = if optimize_size { "-Oz" } else { "-O3" }; + + let status = Command::new(&wasm_opt) + .arg(opt_level) + .arg("--all-features") + .arg("--disable-gc") + .arg("--disable-reference-types") + .arg("--strip-debug") + .arg("--strip-producers") + .arg(&runtime_dst) + .arg("-o") + .arg(&runtime_dst) + .status() + .context("Failed to run wasm-opt")?; + + if !status.success() { + bail!("wasm-opt failed"); + } + } + + Ok(()) +} + +/// Remove only the nested target directory owned by this preparation run. +fn cleanup_runtime_target_dir(target_dir: &Path) { + match fs::remove_dir_all(target_dir) { + Ok(()) => {} + Err(err) if err.kind() == std::io::ErrorKind::NotFound => {} + Err(err) => { + eprintln!( + "warning: failed to clean nested runtime target dir {}: {err}", + target_dir.display() + ); + } + } +} + +/// Resolve the WASI SDK from an override, the output cache, or its pinned release. +fn get_wasi_sdk(out_dir: &Path) -> Result { + // Check environment first + if let Ok(path) = env::var("WASI_SDK_PATH") { + let p = PathBuf::from(path); + if executable(&p, "bin/clang").exists() { + return absolute(p).context("Failed to resolve WASI_SDK_PATH"); + } + } + + // Check cached location + let stable = out_dir.join("wasi-sdk"); + if executable(&stable, "bin/clang").exists() { + return Ok(stable); + } + + // Download wasi-sdk + let (arch, os) = system()?; + let filename = format!("wasi-sdk-{WASI_SDK_VERSION}.0-{arch}-{os}.tar.gz"); + let url = format!("{WASI_SKD_DL_URL}/wasi-sdk-{WASI_SDK_VERSION}/{filename}"); + + http_archive(&url, out_dir)?; + + // Rename extracted directory to stable location + let extracted = find_wasi_sdk(out_dir).context("Could not find extracted wasi-sdk")?; + fs::rename(&extracted, &stable).context("Failed to rename wasi-sdk directory")?; + + Ok(stable) +} + +/// Locate the extracted SDK directory. +fn find_wasi_sdk(target_dir: &Path) -> Option { + let pattern = target_dir.join("wasi-sdk*"); + glob::glob(pattern.to_str()?) + .ok()? + .filter_map(Result::ok) + .find(|entry| entry.is_dir() && executable(entry, "bin/clang").exists()) +} + +/// Resolve wasm-opt from an override, the output cache, or its pinned release. +fn get_wasm_opt(out_dir: &Path) -> Result { + // Check WASM_OPT environment variable first + if let Ok(path) = env::var("WASM_OPT") { + let p = PathBuf::from(path); + if p.exists() { + return absolute(p).context("Failed to resolve WASM_OPT"); + } + } + + // Check cached location + let stable = out_dir.join("binaryen"); + let wasm_opt = executable(&stable, "bin/wasm-opt"); + if wasm_opt.exists() { + return Ok(wasm_opt); + } + + // Download binaryen + let (arch, os) = system()?; + let tag = format!("version_{BINARYEN_VERSION}"); + let filename = format!("binaryen-{tag}-{arch}-{os}.tar.gz"); + let url = format!("{BINARYEN_DL_URL}/{tag}/{filename}"); + + http_archive(&url, out_dir)?; + + // Rename extracted directory to stable location + let extracted = find_binaryen(out_dir).context("Could not find extracted binaryen")?; + fs::rename(&extracted, &stable).context("Failed to rename binaryen directory")?; + + Ok(executable(&stable, "bin/wasm-opt")) +} + +/// Locate the extracted Binaryen directory. +fn find_binaryen(target_dir: &Path) -> Option { + let pattern = target_dir.join("binaryen*"); + glob::glob(pattern.to_str()?) + .ok()? + .filter_map(Result::ok) + .find(|entry| entry.is_dir() && executable(entry, "bin/wasm-opt").exists()) +} + +/// Apply the host executable suffix to a tool path. +fn executable(root: &Path, relative: &str) -> PathBuf { + let mut path = root.join(relative); + if !env::consts::EXE_SUFFIX.is_empty() { + path.set_extension(&env::consts::EXE_SUFFIX[1..]); + } + path +} + +/// Select the pinned tool archives for the build host, not the Wasm target. +fn system() -> Result<(&'static str, &'static str)> { + let (arch, os) = match (env::consts::ARCH, env::consts::OS) { + ("x86_64", "linux") => ("x86_64", "linux"), + ("aarch64", "linux") => ("arm64", "linux"), + ("x86_64", "macos") => ("x86_64", "macos"), + ("aarch64", "macos") => ("arm64", "macos"), + ("x86_64", "windows") => ("x86_64", "windows"), + ("aarch64", "windows") => ("arm64", "windows"), + (arch, os) => bail!("Unsupported platform: {arch}-{os}"), + }; + + Ok((arch, os)) +} + +/// Download a size-bounded tool archive and extract it into the artifact cache. +fn http_archive(url: &str, out_dir: &Path) -> Result<()> { + eprintln!("Downloading archive from {url}..."); + + let response = ureq::get(url) + .call() + .context("Failed to download wasi-sdk")?; + + let mut bytes = Vec::new(); + response + .into_body() + .into_reader() + .take(MAX_ARCHIVE_BYTES + 1) + .read_to_end(&mut bytes) + .context("Failed to download archive")?; + + if bytes.len() as u64 > MAX_ARCHIVE_BYTES { + bail!("Archive exceeds maximum download size of {MAX_ARCHIVE_BYTES} bytes"); + } + + let decoder = GzDecoder::new(bytes.as_slice()); + + let mut archive = tar::Archive::new(decoder); + archive + .unpack(out_dir) + .context("Failed to extract archive")?; + + Ok(()) +} + +#[cfg(test)] +mod tests { + use std::collections::HashSet; + + use super::*; + + /// Each embedding constant and artifact has one definition and a valid fallback. + #[test] + fn artifact_manifest_is_consistent() { + let mut filenames = HashSet::new(); + let mut constants = HashSet::new(); + let mut variants = HashSet::new(); + + for build in RUNTIME_BUILDS { + assert!(filenames.insert(build.filename)); + assert!(variants.insert((build.async_support, build.optimize_size))); + + if let Some(fallback) = build.fallback_const { + assert!(build.async_support); + assert!( + constants.contains(fallback), + "fallback must precede its alias" + ); + } else { + assert!(!build.async_support); + } + + assert!(constants.insert(build.const_name)); + } + + assert_eq!(variants.len(), 4); + } + + /// Explicit preparation and the Cargo fallback use identical profile flags. + #[test] + fn profile_flags_preserve_runtime_variants() { + let debug = CargoProfile::new(false); + assert_eq!(debug.name, "debug"); + assert_eq!( + debug.runtime_rustflags(false), + debug.runtime_rustflags(true) + ); + assert_eq!(debug.runtime_cflags(true), "-fPIC"); + + let release = CargoProfile::new(true); + assert_eq!(release.name, "release"); + assert!( + release + .runtime_rustflags(false) + .ends_with("-Clto=fat -Copt-level=3") + ); + assert!( + release + .runtime_rustflags(true) + .ends_with("-Clto=fat -Copt-level=z") + ); + assert_eq!(release.runtime_cflags(false), "-fPIC -O3"); + assert_eq!(release.runtime_cflags(true), "-fPIC -Oz"); + } +} diff --git a/crates/core/src/codegen.rs b/crates/core/src/codegen.rs index 3c6fe95..ee121b9 100644 --- a/crates/core/src/codegen.rs +++ b/crates/core/src/codegen.rs @@ -122,35 +122,9 @@ impl<'a> EmitContext<'a> { } if name == "Stream" { - self.multiline( - r#"wit.Stream.from = function(iterable, type) { - const { readable, writable } = wit.Stream(type); - const completion = (async () => { - try { - for await (const item of iterable) { - if (!await writable.writeIterableItem(item)) break; - } - } finally { - writable.drop(); - } - })(); - return { readable, completion }; - };"#, - ); + self.multiline(include_str!("js/stream-from.js")); } 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 }; - };"#, - ); + self.multiline(include_str!("js/future-from.js")); } } diff --git a/crates/core/src/js/future-from.js b/crates/core/src/js/future-from.js new file mode 100644 index 0000000..4bcb927 --- /dev/null +++ b/crates/core/src/js/future-from.js @@ -0,0 +1,11 @@ +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/core/src/js/stream-from.js b/crates/core/src/js/stream-from.js new file mode 100644 index 0000000..d8c5d41 --- /dev/null +++ b/crates/core/src/js/stream-from.js @@ -0,0 +1,13 @@ +wit.Stream.from = function(iterable, type) { + const { readable, writable } = wit.Stream(type); + const completion = (async () => { + try { + for await (const item of iterable) { + if (!await writable.writeIterableItem(item)) break; + } + } finally { + writable.drop(); + } + })(); + return { readable, completion }; +}; diff --git a/crates/runtime/Cargo.toml b/crates/runtime/Cargo.toml index 2273bfd..fde5d85 100644 --- a/crates/runtime/Cargo.toml +++ b/crates/runtime/Cargo.toml @@ -20,6 +20,9 @@ num_enum = { version = "0.7", default-features = false } smallvec = "1" indexmap = { version = "2", default-features = false } +[dev-dependencies] +quickcheck = "1" + [features] default = ["component-model-async"] component-model-async = [] diff --git a/crates/runtime/src/abi.rs b/crates/runtime/src/abi.rs index 1e79f7c..33e4b26 100644 --- a/crates/runtime/src/abi.rs +++ b/crates/runtime/src/abi.rs @@ -88,179 +88,6 @@ pub(crate) fn unpack_copy_result(packed: u32) -> Option<(u32, CopyResult)> { Some((progress, result)) } -/// Tracks the lifecycle of a copy operation on a stream/future end. -#[derive(Debug, Clone, Copy, PartialEq, Eq, TryFromPrimitive)] -#[repr(u32)] -#[allow(dead_code)] -pub(crate) enum CopyState { - /// Ready to start an operation. - Idle = 1, - /// Executing an operation that has not returned to JavaScript yet. - SyncCopying = 2, - /// Waiting for the host to deliver a completion callback. - AsyncCopying = 3, - /// Waiting for a blocked cancellation request to complete. - CancellingCopy = 4, - /// Closed, dropped, or consumed. - Done = 5, -} - -impl CopyState { - /// Whether a copy operation is currently in progress. - pub(crate) fn copying(self) -> bool { - !matches!(self, Self::Idle | Self::Done) - } - - /// Whether an active copy operation can be cancelled. - pub(crate) fn cancellable(self) -> bool { - matches!(self, CopyState::AsyncCopying) - } -} - -/// Whether a copy end allows re-use after a completed operation. -/// -/// Streams are reusable (can do more reads/writes after completion), -/// futures are one-shot (always transition to Done). -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -enum CopyKind { - Future, - Stream, -} - -impl CopyKind { - fn completion(self) -> CopyState { - match self { - // Futures are one-shot, they can't be used again after completion. - CopyKind::Future => CopyState::Done, - // Streams are reusable, they can be used again after completion. - CopyKind::Stream => CopyState::Idle, - } - } - - fn label(self) -> &'static str { - match self { - CopyKind::Future => "future", - CopyKind::Stream => "stream", - } - } -} - -/// Common state for a readable or writable stream/future end. -/// -/// A present `handle` is owned by this end. Streams return to [`CopyState::Idle`] -/// after successful copies, while futures transition to [`CopyState::Done`] -/// because they are one-shot. -pub(crate) struct CopyEnd { - kind: CopyKind, - /// Index into the WIT stream/future metadata table. - pub(crate) type_index: u32, - /// Canonical ABI handle, removed when ownership is transferred or dropped. - pub(crate) handle: Option, - /// Current copy/cancellation state. - pub(crate) state: CopyState, -} - -impl CopyEnd { - /// Create a new stream end in the Idle state. - pub(crate) fn new_stream(type_index: u32, handle: u32) -> Self { - Self { - kind: CopyKind::Stream, - type_index, - handle: Some(handle), - state: CopyState::Idle, - } - } - - /// Create a new future end in the Idle state. - pub(crate) fn new_future(type_index: u32, handle: u32) -> Self { - Self { - kind: CopyKind::Future, - type_index, - handle: Some(handle), - state: CopyState::Idle, - } - } - - /// 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 != CopyState::Idle { - return Err(rquickjs::Error::new_from_js( - self.kind.label(), - "operation requires an idle endpoint", - )); - } - - let h = self - .handle - .ok_or_else(|| rquickjs::Error::new_from_js("object", "already dropped"))?; - - 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)> { - if !self.state.cancellable() { - return Err(rquickjs::Error::new_from_js( - self.kind.label(), - "cancel without active operation", - )); - } - - let h = self - .handle - .ok_or_else(|| rquickjs::Error::new_from_js("object", "already dropped"))?; - - Ok((h, self.type_index)) - } - - /// Transition to AsyncCopying when the operation blocks. - pub(crate) fn mark_blocked(&mut self) { - debug_assert_eq!(self.state, CopyState::Idle); - self.state = CopyState::AsyncCopying; - } - - /// Transition after a sync or async completion. - pub(crate) fn mark_completed(&mut self, result: CopyResult) { - self.state = match result { - CopyResult::Dropped => CopyState::Done, - CopyResult::Cancelled => CopyState::Idle, - CopyResult::Completed => self.kind.completion(), - }; - } - - /// Transition to CancellingCopy when a cancel itself blocks. - pub(crate) fn mark_cancel_blocked(&mut self) { - debug_assert_eq!(self.state, CopyState::AsyncCopying); - self.state = CopyState::CancellingCopy; - } -} - /// A decoded async callback event. pub(crate) enum Event { None, diff --git a/crates/runtime/src/bindings.rs b/crates/runtime/src/bindings.rs index 825a59e..672ba28 100644 --- a/crates/runtime/src/bindings.rs +++ b/crates/runtime/src/bindings.rs @@ -1,7 +1,6 @@ //! WIT to/from JS binding registration. use heck::{ToLowerCamelCase, ToUpperCamelCase}; use rquickjs::Persistent; -use rquickjs::function; use rquickjs::function::{Constructor, Rest, This}; use rquickjs::{Ctx, Function, Object, Value}; use smallvec::SmallVec; @@ -13,7 +12,6 @@ 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::{DetHashMap, DetHashSet, DetIndexMap, QjsCallContext, coerce_fn}; @@ -23,7 +21,7 @@ pub(crate) fn register(ctx: &rquickjs::Ctx<'_>, wit_def: Wit) -> rquickjs::Resul register_future_classes(ctx)?; register_resource_classes(ctx, wit_def)?; register_root_imports(ctx)?; - register_cqjs_namespace(ctx, wit_def)?; + register_cqjs_namespace(ctx)?; Ok(()) } @@ -279,177 +277,6 @@ fn call_import<'js>( } } -/// Build the `asyncExports` object for the `__cqjs` namespace. -/// -/// Each wrapper calls the user's export function or resource member, then -/// chains `.then()` to signal `task_return` back to the host. -fn build_async_exports<'js>( - ctx: &rquickjs::Ctx<'js>, - wit_def: Wit, -) -> rquickjs::Result> { - let exports = rquickjs::Object::new(ctx.clone())?; - // Insertion-ordered so the resulting object's property order is deterministic - // (and follows WIT declaration order) for a reproducible Wizer snapshot. - let mut iface_objs: DetIndexMap> = DetIndexMap::default(); - - for (func_index, func) in wit_def.iter_export_funcs().enumerate() { - let wrapper_name = func.name().to_lower_camel_case(); - let func_name = func.name(); - let kind = classify(func_name); - let iface_name = func - .interface() - .map(|interface| iface_lookup(ctx, interface).to_string()); - - let iface = iface_name.clone(); - - let wrapper = Function::new( - ctx.clone(), - coerce_fn(move |ctx: Ctx<'_>, args: Rest>| { - let exports = ctx.user_module().exports(&ctx)?; - let export_scope = || { - if let Some(ref iface) = iface { - exports.get(iface.as_str()) - } else { - Ok(exports.clone()) - } - }; - - let mut values = args.0; - let (user_fn, this): (Function<'_>, Option>) = match kind { - FuncKind::Freestanding => { - let scope: Object = export_scope()?; - (scope.get(fn_lookup(&ctx, func_name))?, None) - } - FuncKind::Constructor { .. } => { - return Err(rquickjs::Exception::throw_type( - &ctx, - "resource constructors cannot be async", - )); - } - FuncKind::Method { method, .. } => { - if values.is_empty() { - return Err(rquickjs::Exception::throw_type( - &ctx, - "async resource method receiver is missing", - )); - } - let receiver = values.remove(0); - let receiver_obj = receiver.as_object().ok_or_else(|| { - rquickjs::Exception::throw_type( - &ctx, - "async resource method receiver is not an object", - ) - })?; - let method: Function = receiver_obj.get(fn_lookup(&ctx, method))?; - (method, Some(receiver)) - } - FuncKind::Static { resource, method } => { - let scope: Object = export_scope()?; - let class_name = resource.to_upper_camel_case(); - let class: Object = scope.get(class_name.as_str())?; - let method: Function = class.get(fn_lookup(&ctx, method))?; - (method, Some(class.into_value())) - } - }; - - let mut js_args = function::Args::new(ctx.clone(), values.len()); - for arg in values { - js_args.push_arg(arg)?; - } - if let Some(this) = this { - js_args.this(this)?; - } - let result = user_fn.call_arg::(js_args)?; - - let promise_obj = result - .as_object() - .ok_or_else(|| rquickjs::Error::new_from_js("value", "promise"))?; - - let then_fn: Function = promise_obj.get("then")?; - - 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() - .next() - .unwrap_or_else(|| Value::new_undefined(ctx.clone())); - - let func = ctx.wit().export_func(func_index); - let boundary = ResultBoundary::new(func.result()); - let mut call = QjsCallContext::default(); - - 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)) - }), - )?; - - 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() - .next() - .unwrap_or_else(|| Value::new_undefined(ctx.clone())); - 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) - .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)) - }), - )?; - - let mut call_args = function::Args::new(ctx.clone(), 2); - call_args.this(result)?; - call_args.push_arg(then_cb)?; - call_args.push_arg(catch_cb)?; - then_fn.call_arg(call_args) - }), - )?; - - let target = match &iface_name { - Some(iface) => iface_objs - .entry(iface.clone()) - .or_insert_with(|| rquickjs::Object::new(ctx.clone()).unwrap()), - None => &exports, - }; - target.set(wrapper_name.as_str(), wrapper)?; - } - - for (name, obj) in iface_objs { - exports.set(name.as_str(), obj)?; - } - - Ok(exports) -} - /// Register the `__cqjs` namespace object on globalThis. /// /// Consolidates all internal bridge globals into a single frozen object: @@ -457,8 +284,7 @@ fn build_async_exports<'js>( /// - `makeFuture(typeIndex)` — create a future pair /// - `getMemoryUsage()` — return QuickJS memory statistics /// - `runGc()` — trigger QuickJS garbage collection -/// - `asyncExports` — object containing async export wrappers -fn register_cqjs_namespace(ctx: &rquickjs::Ctx<'_>, wit_def: Wit) -> rquickjs::Result<()> { +fn register_cqjs_namespace(ctx: &rquickjs::Ctx<'_>) -> rquickjs::Result<()> { let ns = rquickjs::Object::new(ctx.clone())?; // Stream/future factories @@ -510,10 +336,6 @@ fn register_cqjs_namespace(ctx: &rquickjs::Ctx<'_>, wit_def: Wit) -> rquickjs::R )?, )?; - // Async export wrappers - let async_exports = build_async_exports(ctx, wit_def)?; - ns.set("asyncExports", async_exports)?; - // Freeze and install on globalThis let object_ctor: rquickjs::Object = ctx.globals().get("Object")?; let freeze_fn: Function = object_ctor.get("freeze")?; diff --git a/crates/runtime/src/endpoint.rs b/crates/runtime/src/endpoint.rs new file mode 100644 index 0000000..894a649 --- /dev/null +++ b/crates/runtime/src/endpoint.rs @@ -0,0 +1,251 @@ +//! Stream/future endpoint ownership and copy-state transitions. + +use crate::abi::CopyResult; + +/// A rejected endpoint operation, translated to a JS error at the binding boundary. +#[derive(Debug, PartialEq, Eq)] +pub(crate) enum EndpointError { + /// Another operation has already used or reserved this endpoint. + NotIdle(&'static str), + /// The endpoint belongs to a different WIT stream/future type. + WrongType(&'static str), + /// An outstanding operation must settle before disposal. + InFlight(&'static str), + /// Only a pending, not-yet-cancelled operation can be cancelled. + NotCancellable(&'static str), + /// Ownership has already been transferred or released. + Dropped, +} + +/// The lifecycle of a copy operation, independent of JS and host handles. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum CopyState { + Idle, + Copying, + Cancelling, + Done, +} + +/// Streams are reusable; futures are consumed by their first successful copy. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum CopyKind { + Future, + Stream, +} + +impl CopyKind { + /// Describe the endpoint in conversion errors. + fn label(self) -> &'static str { + match self { + Self::Future => "future", + Self::Stream => "stream", + } + } +} + +/// Owns one canonical endpoint handle and validates its lifecycle transitions. +/// +/// State transitions never call the host. Disposal returns the owned handle so +/// callers can release their JS class borrow before invoking its destructor. +pub(crate) struct CopyEnd { + kind: CopyKind, + type_index: u32, + handle: Option, + state: CopyState, +} + +impl CopyEnd { + /// Create an idle stream endpoint. + pub(crate) fn new_stream(type_index: u32, handle: u32) -> Self { + Self { + kind: CopyKind::Stream, + type_index, + handle: Some(handle), + state: CopyState::Idle, + } + } + + /// Create an idle future endpoint. + pub(crate) fn new_future(type_index: u32, handle: u32) -> Self { + Self { + kind: CopyKind::Future, + type_index, + handle: Some(handle), + state: CopyState::Idle, + } + } + + /// Return the index into the endpoint's WIT metadata table. + pub(crate) fn type_index(&self) -> u32 { + self.type_index + } + + /// Whether the endpoint still owns a handle, even if its operation is done. + pub(crate) fn has_handle(&self) -> bool { + self.handle.is_some() + } + + /// Whether no further values can be copied through this endpoint. + pub(crate) fn is_closed(&self) -> bool { + self.state == CopyState::Done + } + + /// Whether an operation or its cancellation still owns the endpoint. + pub(crate) fn is_pending(&self) -> bool { + matches!(self.state, CopyState::Copying | CopyState::Cancelling) + } + + /// Whether an outstanding operation can receive its first cancel request. + pub(crate) fn can_cancel(&self) -> bool { + self.state == CopyState::Copying + } + + /// Validate an idle endpoint and return its handle and WIT type index. + pub(crate) fn begin_op(&self) -> Result<(u32, u32), EndpointError> { + if self.state != CopyState::Idle { + return Err(EndpointError::NotIdle(self.kind.label())); + } + + let handle = self.handle.ok_or(EndpointError::Dropped)?; + Ok((handle, self.type_index)) + } + + /// Transfer an idle endpoint to a matching canonical ABI parameter. + pub(crate) fn begin_transfer(&mut self, type_index: u32) -> Result { + if self.type_index != type_index { + return Err(EndpointError::WrongType(self.kind.label())); + } + + let (handle, _) = self.begin_op()?; + self.handle = None; + self.state = CopyState::Done; + Ok(handle) + } + + /// Close a settled endpoint and return its handle exactly once. + pub(crate) fn begin_drop(&mut self) -> Result, EndpointError> { + if self.is_pending() { + return Err(EndpointError::InFlight(self.kind.label())); + } + + self.state = CopyState::Done; + Ok(self.handle.take()) + } + + /// Validate a pending copy before issuing a cancellation request. + pub(crate) fn begin_cancel(&self) -> Result<(u32, u32), EndpointError> { + if !self.can_cancel() { + return Err(EndpointError::NotCancellable(self.kind.label())); + } + + let handle = self.handle.ok_or(EndpointError::Dropped)?; + Ok((handle, self.type_index)) + } + + /// Retain ownership while a copy waits for a host callback. + pub(crate) fn mark_blocked(&mut self) { + debug_assert_eq!(self.state, CopyState::Idle); + self.state = CopyState::Copying; + } + + /// Finish a copy, permitting retries after cancellation but not after closure. + pub(crate) fn mark_completed(&mut self, result: CopyResult) { + debug_assert!(!self.is_closed()); + + self.state = match (result, self.kind) { + (CopyResult::Dropped, _) | (CopyResult::Completed, CopyKind::Future) => CopyState::Done, + (CopyResult::Cancelled, _) | (CopyResult::Completed, CopyKind::Stream) => { + CopyState::Idle + } + }; + } + + /// Retain ownership when the cancellation request itself blocks. + pub(crate) fn mark_cancel_blocked(&mut self) { + debug_assert_eq!(self.state, CopyState::Copying); + self.state = CopyState::Cancelling; + } +} + +#[cfg(test)] +mod tests { + use quickcheck::{Arbitrary, Gen, TestResult, quickcheck}; + + use super::*; + + impl Arbitrary for CopyResult { + /// Generate one of the three settled copy outcomes. + fn arbitrary(generator: &mut Gen) -> Self { + *generator + .choose(&[Self::Completed, Self::Dropped, Self::Cancelled]) + .unwrap() + } + } + + quickcheck! { + /// Exercise immediate, callback, and cancellation completion for both kinds. + fn completion_transitions( + future: bool, + blocked: bool, + cancelled: bool, + result: CopyResult + ) -> TestResult { + if cancelled && !blocked { + return TestResult::discard(); + } + + let mut end = if future { + CopyEnd::new_future(3, 7) + } else { + CopyEnd::new_stream(3, 7) + }; + + assert_eq!(end.begin_op(), Ok((7, 3))); + + if blocked { + end.mark_blocked(); + assert!(end.begin_op().is_err()); + assert!(end.begin_drop().is_err()); + assert!(end.begin_transfer(3).is_err()); + assert_eq!(end.begin_cancel(), Ok((7, 3))); + } + + if cancelled { + end.mark_cancel_blocked(); + assert!(end.begin_cancel().is_err()); + assert!(end.begin_drop().is_err()); + } + + end.mark_completed(result); + let closed = result == CopyResult::Dropped + || (future && result == CopyResult::Completed); + + assert_eq!(end.is_closed(), closed); + assert_eq!(end.begin_op().is_err(), closed); + + assert!(!end.is_pending()); + assert_eq!(end.begin_drop(), Ok(Some(7))); + assert_eq!(end.begin_drop(), Ok(None)); + + TestResult::passed() + } + } + + /// Failed transfers preserve ownership; successful transfers consume it. + #[test] + fn transfer_ownership() { + let mut end = CopyEnd::new_stream(3, 7); + assert_eq!( + end.begin_transfer(4), + Err(EndpointError::WrongType("stream")) + ); + + assert_eq!(end.begin_op(), Ok((7, 3))); + assert_eq!(end.begin_transfer(3), Ok(7)); + + assert!(!end.has_handle()); + assert!(end.is_closed()); + assert!(end.begin_op().is_err()); + assert_eq!(end.begin_drop(), Ok(None)); + } +} diff --git a/crates/runtime/src/exports.rs b/crates/runtime/src/exports.rs new file mode 100644 index 0000000..6b5b516 --- /dev/null +++ b/crates/runtime/src/exports.rs @@ -0,0 +1,177 @@ +//! Shared export lookup, argument binding, invocation, and async completion. + +use heck::ToUpperCamelCase; +use rquickjs::function::{Args, Constructor, Rest}; +use rquickjs::{Ctx, Exception, Function, Object, Persistent, Value}; +use wit_dylib_ffi::ExportFunction; + +use crate::result::{JsCompletion, ResultBoundary}; +use crate::trivia::{fn_lookup, iface_lookup}; +use crate::wit_imports::{FuncKind, classify}; +use crate::{CtxExt, QjsCallContext, coerce_fn}; + +/// A resolved export with its arguments and JavaScript receiver already bound. +pub(crate) enum ExportCall<'js> { + /// A freestanding function, resource method, or static method. + Function(Function<'js>, Args<'js>), + /// A resource constructor, invoked with JavaScript construction semantics. + Constructor(Constructor<'js>, Args<'js>), +} + +impl<'js> ExportCall<'js> { + /// Invoke the resolved function or resource constructor. + pub(crate) fn invoke(self) -> rquickjs::Result> { + match self { + Self::Function(function, args) => function.call_arg(args), + Self::Constructor(constructor, args) => constructor.construct_args(args), + } + } + + /// Invoke an async export and attach its canonical completion callbacks. + /// + /// Keep thenable semantics: successful values and rejection reasons pass + /// through the WIT result boundary when the returned promise settles. + pub(crate) fn invoke_async( + self, + ctx: &Ctx<'js>, + func: ExportFunction, + ) -> rquickjs::Result> { + let result = self.invoke()?; + let promise = result + .as_object() + .ok_or_else(|| rquickjs::Error::new_from_js("value", "promise"))?; + + let then: Function = promise.get("then")?; + + let resolve = Function::new( + ctx.clone(), + coerce_fn(move |ctx: Ctx<'_>, args: Rest>| { + let value = args + .0 + .into_iter() + .next() + .unwrap_or_else(|| Value::new_undefined(ctx.clone())); + + finish_async_export(ctx, func, JsCompletion::Return(value)) + }), + )?; + + let reject = Function::new( + ctx.clone(), + coerce_fn(move |ctx: Ctx<'_>, args: Rest>| { + let reason = args + .0 + .into_iter() + .next() + .unwrap_or_else(|| Value::new_undefined(ctx.clone())); + + finish_async_export(ctx, func, JsCompletion::Throw(reason)) + }), + )?; + + let mut args = Args::new(ctx.clone(), 2); + args.this(result)?; + args.push_arg(resolve)?; + args.push_arg(reject)?; + then.call_arg(args) + } +} + +/// Lower one promise settlement, release argument guards, and return to the host. +/// +/// Cancellation owns task completion once requested. Otherwise, keep the result +/// conversion context alive until `task.return` has consumed its transfers. +fn finish_async_export<'js>( + ctx: Ctx<'js>, + func: ExportFunction, + compl: JsCompletion<'js>, +) -> rquickjs::Result> { + if ctx.task().is_cancelling() { + return Ok(Value::new_undefined(ctx)); + } + + let boundary = ResultBoundary::new(func.result()); + let maybe_val = match compl { + JsCompletion::Return(value) => boundary.lower_value(&ctx, value), + JsCompletion::Throw(reason) => boundary.lower_throw(&ctx, reason), + }; + + let val = maybe_val.expect("Failed to complete"); + let mut call = QjsCallContext::default(); + + if let Some(val) = val { + call.push_value(&ctx, val); + } + + ctx.task().finish_export(); + func.call_task_return(&mut call); + call.complete_transfers(1); + Ok(Value::new_undefined(ctx)) +} + +/// Resolve an export's namespace/member and bind its receiver and arguments. +/// +/// Resource methods consume the first canonical argument as `this`; static +/// methods bind their class instead. Other exports preserve every argument. +pub(crate) fn prepare_export<'js>( + ctx: &Ctx<'js>, + func: ExportFunction, + mut values: impl ExactSizeIterator>>, +) -> rquickjs::Result> { + let kind = classify(func.name()); + + let receiver = if matches!(kind, FuncKind::Method { .. }) { + let recv = values + .next() + .ok_or_else(|| Exception::throw_type(ctx, "no resource method receiver"))? + .restore(ctx)?; + Some(recv) + } else { + None + }; + + let mut args = Args::new(ctx.clone(), values.len()); + + for val in values { + args.push_arg(val.restore(ctx)?)?; + } + + match kind { + FuncKind::Constructor { resource } => { + let scope = export_scope(ctx, func.interface())?; + let ctor = scope.get(resource.to_upper_camel_case())?; + Ok(ExportCall::Constructor(ctor, args)) + } + FuncKind::Method { method, .. } => { + let recv = receiver.expect("method receiver was consumed above"); + let obj = recv.as_object().ok_or_else(|| { + Exception::throw_type(ctx, "resource method receiver is not an object") + })?; + let func = obj.get(fn_lookup(ctx, method))?; + args.this(recv)?; + Ok(ExportCall::Function(func, args)) + } + FuncKind::Static { resource, method } => { + let scope = export_scope(ctx, func.interface())?; + let class: Object = scope.get(resource.to_upper_camel_case())?; + let function = class.get(fn_lookup(ctx, method))?; + args.this(class)?; + Ok(ExportCall::Function(function, args)) + } + FuncKind::Freestanding => { + let scope = export_scope(ctx, func.interface())?; + let function = scope.get(fn_lookup(ctx, func.name()))?; + Ok(ExportCall::Function(function, args)) + } + } +} + +/// Select the root or interface-scoped user module namespace. +fn export_scope<'js>(ctx: &Ctx<'js>, ifc: Option<&'static str>) -> rquickjs::Result> { + let exports = ctx.user_module().exports(ctx)?; + + match ifc { + Some(ifc) => exports.get(iface_lookup(ctx, ifc)), + None => Ok(exports), + } +} diff --git a/crates/runtime/src/futures.rs b/crates/runtime/src/futures.rs index 3527c1c..e656183 100644 --- a/crates/runtime/src/futures.rs +++ b/crates/runtime/src/futures.rs @@ -5,8 +5,10 @@ use rquickjs::{CatchResultExt, Ctx, Function, Object, Persistent, Promise, Value use rquickjs::{JsLifetime, function}; use crate::CtxExt; -use crate::abi::{CopyEnd, CopyResult, is_blocked_raw, unpack_copy_result}; +use crate::abi::{CopyResult, is_blocked_raw, unpack_copy_result}; use crate::buffer::BufferGuard; +use crate::endpoint::CopyEnd; +use crate::result::JsCompletion; use crate::task::Pending; use crate::{QjsCallContext, resolve_promise, symbol_dispose, with_ctx}; @@ -157,7 +159,7 @@ fn future_read<'js>( ) -> rquickjs::Result> { { let readable = this.0.borrow(); - if readable.end.handle.is_none() { + if !readable.end.has_handle() { return Err(rquickjs::Exception::throw_type( &ctx, "future already dropped or transferred", @@ -192,40 +194,76 @@ fn future_read<'js>( ctx.task().register(handle, pending); } else { - let result_code = CopyResult::try_from(code & 0xF).expect("unknown copy result"); - finish_read_state(&this.0, result_code); - - match result_code { - CopyResult::Completed => { - unsafe { ty.lift(&mut call, buffer.ptr()) }; - 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,))?; - } - CopyResult::Cancelled => { - let msg = - rquickjs::String::from_str(ctx.clone(), "future read cancelled")?.into_value(); - reject.call::<_, Value>((msg,))?; - } - } - drop(buffer); + finish_read(&ctx, &this.0, &mut call, buffer, code)?.settle(&resolve, &reject)?; } Ok(promise.into_value()) } -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; - } +/// Finish any future read: update its cached promise, then lift or reject. +/// +/// The ABI buffer is released after lifting. The caller retains the conversion +/// context through promise settlement so its resource guards remain alive. +fn finish_read<'js>( + ctx: &Ctx<'js>, + class: &Class<'js, FutureReadable<'js>>, + call: &mut QjsCallContext, + buffer: BufferGuard, + code: u32, +) -> rquickjs::Result> { + let (_, result) = unpack_copy_result(code).expect("future read completion must not block"); + let type_index = { + let mut readable = class.borrow_mut(); + readable.end.mark_completed(result); + + if result == CopyResult::Cancelled { + readable.promise = None; + } + + readable.end.type_index() + }; + + let completion = match result { + CopyResult::Completed => { + let ty = ctx.wit().future(type_index as usize); + unsafe { ty.lift(call, buffer.ptr()) }; + + let val = call + .maybe_pop_value(ctx)? + .unwrap_or_else(|| Value::new_undefined(ctx.clone())); + JsCompletion::Return(val) + } + CopyResult::Dropped | CopyResult::Cancelled => { + let message = if result == CopyResult::Dropped { + "future writer dropped" + } else { + "future read cancelled" + }; + + let reason = rquickjs::String::from_str(ctx.clone(), message)?.into_value(); + JsCompletion::Throw(reason) + } + }; + + drop(buffer); + Ok(completion) +} + +/// Finish any future write, committing ownership only if its value was consumed. +fn finish_write<'js>( + ctx: &Ctx<'js>, + class: &Class<'js, FutureWritable>, + call: &mut QjsCallContext, + buffer: BufferGuard, + code: u32, +) -> Value<'js> { + let (_, result) = unpack_copy_result(code).expect("future write completion must not block"); + let success = result == CopyResult::Completed; + drop(buffer); + + call.complete_transfers(usize::from(success)); + class.borrow_mut().end.mark_completed(result); + Value::new_bool(ctx.clone(), success) } pub(crate) fn future_cancel_read<'js>( @@ -234,6 +272,7 @@ pub(crate) fn future_cancel_read<'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) }; @@ -258,12 +297,14 @@ fn future_drop_readable<'js>( let mut readable = this.0.borrow_mut(); let handle = readable.end.begin_drop()?; readable.promise = None; - (handle, readable.end.type_index) + (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(()) } @@ -296,13 +337,7 @@ fn future_write<'js>( }; ctx.task().register(handle, pending); } else { - 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); + let result = finish_write(&ctx, &this.0, &mut call, buffer, code); resolve .call::<_, Value>((result,)) @@ -340,7 +375,7 @@ fn future_drop_writable<'js>( ) -> rquickjs::Result<()> { let (handle, type_index) = { let mut writable = this.0.borrow_mut(); - (writable.end.begin_drop()?, writable.end.type_index) + (writable.end.begin_drop()?, writable.end.type_index()) }; if let Some(handle) = handle { @@ -356,26 +391,19 @@ pub(crate) fn handle_write_event(handle: u32, result: u32) { let Pending::FutureWrite { mut call, + buffer, resolve, wrapper, - .. } = pending else { unreachable!("expected FutureWrite pending"); }; - 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(); let class = Class::::from_value(&w).unwrap(); - class.borrow_mut().end.mark_completed(copy_result); - - let val = Value::new_bool(ctx.clone(), success); - let res = Persistent::save(ctx, val); - Some(res) + let value = finish_write(ctx, &class, &mut call, buffer, result); + Some(Persistent::save(ctx, value)) }); resolve_promise(resolve, result); @@ -396,47 +424,13 @@ pub(crate) fn handle_read_event(handle: u32, result: u32) { unreachable!("expected FutureRead pending"); }; - let copy_result = CopyResult::try_from(result & 0xF).expect("unknown copy result"); + with_ctx(|ctx| { + let value = wrapper.restore(ctx).unwrap(); + let class = Class::::from_value(&value).unwrap(); - match copy_result { - CopyResult::Completed => { - let result = with_ctx(|ctx| { - let w = wrapper.restore(ctx).unwrap(); - let class = Class::::from_value(&w).unwrap(); - 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()) }; - - call.maybe_pop_persistent() - }); - - drop(buffer); - resolve_promise(resolve, result); - } - CopyResult::Dropped | CopyResult::Cancelled => { - drop(buffer); - with_ctx(|ctx| { - let w = wrapper.restore(ctx).unwrap(); - let class = Class::::from_value(&w).unwrap(); - finish_read_state(&class, copy_result); - - let reject_fn = reject.restore(ctx).expect("restore future rejecter"); - - let msg = if copy_result == CopyResult::Dropped { - "future writer dropped" - } else { - "future read cancelled" - }; - let msg_val = rquickjs::String::from_str(ctx.clone(), msg) - .unwrap() - .into_value(); - reject_fn - .call::<_, Value>((msg_val,)) - .catch(ctx) - .expect("Failed to reject future read"); - }); - } - } + finish_read(ctx, &class, &mut call, buffer, result) + .catch(ctx) + .expect("Failed to finish future read") + .settle_persistent(ctx, resolve, reject); + }); } diff --git a/crates/runtime/src/interpreter.rs b/crates/runtime/src/interpreter.rs index 32eac66..c4164ec 100644 --- a/crates/runtime/src/interpreter.rs +++ b/crates/runtime/src/interpreter.rs @@ -2,17 +2,15 @@ use crate::CtxExt; use crate::abi::Event; use crate::bindings::register; +use crate::exports::prepare_export; use crate::resources::{ResourceTable, drain_resource_drops}; use crate::result::ResultBoundary; use crate::task::TaskState; -use crate::trivia::{fn_lookup, iface_lookup}; -use crate::wit_imports::{FuncKind, WitImportRegistry, classify}; +use crate::wit_imports::WitImportRegistry; use crate::{QjsCallContext, with_ctx}; use crate::{abi, futures, streams}; -use heck::ToUpperCamelCase; -use rquickjs::function::{Args, Constructor}; -use rquickjs::{CatchResultExt, Ctx, Function, JsLifetime, Object, Value}; +use rquickjs::{CatchResultExt, JsLifetime}; use wit_dylib_ffi::{ExportFunction, Interpreter, Resource, Wit}; /// Newtype wrapper for `Wit` so it can be stored as rquickjs userdata. @@ -22,39 +20,6 @@ pub(crate) struct WitData(pub(crate) Wit); /// quickjs interpreter implementation of the `Interpreter` trait. pub struct QjsInterpreter; -/// Select the root or interface-scoped namespace for an exported function. -fn export_scope<'js>(ctx: &Ctx<'js>, interface: Option<&'static str>) -> Object<'js> { - let exports = ctx - .user_module() - .exports(ctx) - .expect("user module exports not found"); - - match interface { - Some(interface) => exports - .get(iface_lookup(ctx, interface)) - .unwrap_or_else(|err| panic!("interface '{interface}' not found: {err:?}")), - None => exports, - } -} - -/// Call a JavaScript export and lower its return/throw through the WIT boundary. -fn call_export<'js>( - ctx: &Ctx<'js>, - func: ExportFunction, - name: &str, - js_func: Function<'js>, - args: Args<'js>, - cx: &mut QjsCallContext, -) { - let value = ResultBoundary::new(func.result()) - .lower_call(ctx, js_func.call_arg::(args)) - .unwrap_or_else(|err| panic!("Failed to call '{name}': {err:?}")); - - if let Some(value) = value { - cx.push_value(ctx, value); - } -} - impl Interpreter for QjsInterpreter { type CallCx<'a> = QjsCallContext; @@ -78,60 +43,17 @@ impl Interpreter for QjsInterpreter { } fn export_call(_wit: Wit, func: ExportFunction, cx: &mut Self::CallCx<'_>) { - with_ctx(|ctx| match classify(func.name()) { - FuncKind::Constructor { resource } => { - let class_name = resource.to_upper_camel_case(); - let scope = export_scope(ctx, func.interface()); - let ctor: Constructor = scope - .get(class_name.as_str()) - .unwrap_or_else(|err| panic!("class '{class_name}' not found: {err:?}")); - let args = cx.stack_into_args(ctx); - let instance: Value = ctor - .construct_args(args) - .unwrap_or_else(|err| panic!("Failed to construct '{class_name}': {err:?}")); - cx.push_value(ctx, instance); - } - FuncKind::Method { method, .. } => { - assert!(!method.is_empty(), "invalid method name: {}", func.name()); - let method_name = fn_lookup(ctx, method); - let self_val = cx.shift_value(ctx); - let self_obj = self_val - .as_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:?}")); - let mut args = cx.stack_into_args(ctx); - args.this(self_val).expect("failed to set this"); - call_export(ctx, func, method_name, method, args, cx); - } - FuncKind::Static { resource, method } => { - assert!( - !method.is_empty(), - "invalid static method name: {}", - func.name() - ); - let method_name = fn_lookup(ctx, method); - let class_name = resource.to_upper_camel_case(); - let scope = export_scope(ctx, func.interface()); - let class: Object = scope - .get(class_name.as_str()) - .unwrap_or_else(|err| panic!("class '{class_name}' not found: {err:?}")); - let js_func: Function = class.get(method_name).unwrap_or_else(|err| { - panic!("static method '{method_name}' not found: {err:?}") - }); - let mut args = cx.stack_into_args(ctx); - args.this(class).expect("failed to set static this"); - call_export(ctx, func, method_name, js_func, args, cx); - } - FuncKind::Freestanding => { - let func_name = fn_lookup(ctx, func.name()); - let scope = export_scope(ctx, func.interface()); - let js_func: Function = scope - .get(func_name) - .unwrap_or_else(|err| panic!("function '{func_name}' not found: {err:?}")); - let args = cx.stack_into_args(ctx); - call_export(ctx, func, func_name, js_func, args, cx); + with_ctx(|ctx| { + let call = prepare_export(ctx, func, cx.drain_values()) + .catch(ctx) + .unwrap_or_else(|err| panic!("Failed to resolve '{}': {err}", func.name())); + + let value = ResultBoundary::new(func.result()) + .lower_call(ctx, call.invoke()) + .unwrap_or_else(|err| panic!("Failed to call '{}': {err}", func.name())); + + if let Some(value) = value { + cx.push_value(ctx, value); } }); with_ctx(drain_resource_drops); @@ -144,30 +66,11 @@ impl Interpreter for QjsInterpreter { ) -> u32 { with_ctx(|ctx| { drain_resource_drops(ctx); - let args = cx.stack_into_args(ctx); + let values = cx.take_values(); ctx.task().init(cx); - let globals = ctx.globals(); - - let cqjs: rquickjs::Object = globals.get("__cqjs").expect("__cqjs namespace not found"); - - let async_exports: rquickjs::Object = cqjs - .get("asyncExports") - .expect("__cqjs.asyncExports not found"); - - let wrapper_obj = if let Some(interface) = func.interface() { - async_exports.get(iface_lookup(ctx, interface)).unwrap() - } else { - async_exports - }; - - let func_name = fn_lookup(ctx, func.name()); - let js_func: rquickjs::Function = wrapper_obj - .get(func_name) - .unwrap_or_else(|e| panic!("Failed to get async export '{}': {:?}", func_name, e)); - - let _result = js_func - .call_arg::(args) + prepare_export(ctx, func, values) + .and_then(|call| call.invoke_async(ctx, func)) .catch(ctx) .unwrap_or_else(|e| panic!("Failed to call async '{}': {e}", func.name())); }); diff --git a/crates/runtime/src/lib.rs b/crates/runtime/src/lib.rs index b8ebe09..090a9ec 100644 --- a/crates/runtime/src/lib.rs +++ b/crates/runtime/src/lib.rs @@ -2,6 +2,8 @@ mod abi; mod bindings; mod buffer; mod call; +mod endpoint; +mod exports; mod futures; mod interpreter; mod module; @@ -18,7 +20,7 @@ use std::cell::{Cell, OnceCell, RefCell}; use std::collections::hash_map::DefaultHasher; use rquickjs::runtime::UserDataGuard; -use rquickjs::{Context, JsLifetime, Persistent, Runtime, Value, function}; +use rquickjs::{Context, JsLifetime, Persistent, Runtime, Value}; use smallvec::SmallVec; use task::TaskState; use wit_dylib_ffi::Wit; @@ -39,7 +41,7 @@ pub(crate) type DetIndexMap = indexmap::IndexMap; mod init { wit_bindgen::generate!({ world: "init", - path: "wit/init.wit", + path: "../core/wit/init.wit", generate_all, disable_run_ctors_once_workaround: true, }); @@ -243,8 +245,19 @@ impl QjsCallContext { self.pop_persistent().restore(ctx).expect("stack underflow") } - pub(crate) fn shift_value<'js>(&mut self, ctx: &rquickjs::Ctx<'js>) -> Value<'js> { - self.stack.remove(0).restore(ctx).expect("stack underflow") + /// Drain canonical arguments in call order, retaining storage for the result. + pub(crate) fn drain_values( + &mut self, + ) -> impl ExactSizeIterator>> { + self.stack.drain(..) + } + + /// Take canonical arguments in call order without allocating another vector. + /// + /// The iterator does not borrow this context, so async dispatch can retain its + /// resource guards in the active task before looking up JavaScript exports. + pub(crate) fn take_values(&mut self) -> std::vec::IntoIter>> { + std::mem::take(&mut self.stack).into_iter() } pub(crate) fn pop_persistent(&mut self) -> Persistent> { @@ -263,16 +276,6 @@ impl QjsCallContext { .map(|persistent| persistent.restore(ctx)) .transpose() } - - pub(crate) fn stack_into_args<'js>(&mut self, ctx: &rquickjs::Ctx<'js>) -> function::Args<'js> { - let mut args = function::Args::new(ctx.clone(), self.stack.len()); - for p in self.stack.drain(..) { - p.restore(ctx) - .and_then(|val| args.push_arg(val)) - .expect("Failed to restore arg"); - } - args - } } impl Drop for QjsCallContext { diff --git a/crates/runtime/src/streams.rs b/crates/runtime/src/streams.rs index 8582758..30b2c36 100644 --- a/crates/runtime/src/streams.rs +++ b/crates/runtime/src/streams.rs @@ -4,20 +4,22 @@ //! `StreamWritable`) whose state lives on the Rust side. Methods on the //! shared prototype avoid per-instance closure allocations. #![allow(unsafe_code)] + +mod helpers; + use crate::CtxExt; -use crate::abi::{CopyEnd, CopyResult, CopyState, is_blocked_raw, unpack_copy_result}; +use crate::abi::{CopyResult, is_blocked_raw, unpack_copy_result}; use crate::buffer::BufferGuard; +use crate::endpoint::CopyEnd; use crate::task::Pending; -use crate::typed_array::{copy_typed_array_as, typed_array_len_as}; +use crate::typed_array::copy_typed_array_as; use crate::{QjsCallContext, resolve_promise, symbol_dispose, with_ctx}; use rquickjs::JsLifetime; use rquickjs::class::{Class, JsClass, Trace}; -use rquickjs::function::{self, Opt, Rest, This}; +use rquickjs::function::{self, Opt, This}; use rquickjs::{Ctx, Function, Object, Persistent, Symbol, Value}; -use std::cell::Cell; - const BYTE_ITERATOR_CHUNK_SIZE: usize = 64 * 1024; /// Rust side state for the readable end of a component-model stream. @@ -92,11 +94,7 @@ impl<'js> JsClass<'js> for StreamWritable { let proto = Object::new(ctx.clone())?; proto.set("write", Function::new(ctx.clone(), stream_write)?)?; proto.set("writeOne", Function::new(ctx.clone(), stream_write_one)?)?; - proto.set("writeAll", Function::new(ctx.clone(), stream_write_all)?)?; - proto.set( - "writeIterableItem", - Function::new(ctx.clone(), stream_write_iterable_item)?, - )?; + helpers::register(ctx, &proto)?; proto.set( "cancelWrite", Function::new(ctx.clone(), stream_cancel_write)?, @@ -177,28 +175,6 @@ pub(crate) fn lower_iterable<'js>( Ok(handle) } -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), - wit_dylib_ffi::Type::S16 => typed_array_len_as!(obj, i16), - wit_dylib_ffi::Type::U32 => typed_array_len_as!(obj, u32), - wit_dylib_ffi::Type::S32 => typed_array_len_as!(obj, i32), - wit_dylib_ffi::Type::U64 => typed_array_len_as!(obj, u64), - 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), - _ => Ok(None), - } -} - /// Fast path for `writable.write(typedArray)` fn try_typed_array_to_buffer<'js>( data: &Value<'js>, @@ -241,10 +217,7 @@ fn stream_next<'js>( ) -> rquickjs::Result> { let (type_index, finished) = { let readable = this.0.borrow(); - ( - readable.end.type_index, - readable.end.handle.is_none() || readable.end.state == CopyState::Done, - ) + (readable.end.type_index(), readable.end.is_closed()) }; if finished { @@ -303,29 +276,8 @@ fn stream_read_impl<'js>( }; ctx.task().register(handle, pending); } else { - let (actual_count, copy_result) = - unpack_copy_result(code).expect("non-BLOCKED stream read must decode"); - - let dropped_handle = { - let mut readable = this.0.borrow_mut(); - readable.end.mark_completed(copy_result); - (copy_result == CopyResult::Dropped) - .then(|| readable.end.handle.take()) - .flatten() - }; + let result_val = finish_read(&ctx, &this.0, call, buffer, code, iterator, false)?; - let result_val = lift_stream_read_result( - &ctx, - ty, - call, - buffer, - actual_count as usize, - copy_result, - iterator, - )?; - if let Some(handle) = dropped_handle { - unsafe { ty.drop_readable()(handle) }; - } resolve .call::<_, Value>((result_val,)) .expect("resolve stream read"); @@ -334,6 +286,90 @@ fn stream_read_impl<'js>( Ok(promise.into_value()) } +/// Finish a stream read from an immediate result, callback, or cancellation. +/// +/// Update endpoint ownership before lifting, release closed handles outside the +/// class borrow, and adapt the lifted value to the requested iterator shape. +/// `close` completes an iterator's pending `return()`: drop the readable handle +/// even if the copy was only cancelled, and resolve the read as `done`. +fn finish_read<'js>( + ctx: &Ctx<'js>, + cls: &Class<'js, StreamReadable>, + call: QjsCallContext, + buffer: BufferGuard, + code: u32, + iter: bool, + close: bool, +) -> rquickjs::Result> { + let (progress, result) = + unpack_copy_result(code).expect("stream read completion must not block"); + + let (type_index, dropped_handle) = { + let mut readable = cls.borrow_mut(); + readable.end.mark_completed(result); + + let handle = if close || result == CopyResult::Dropped { + readable.end.begin_drop()? + } else { + None + }; + + (readable.end.type_index(), handle) + }; + + let ty = ctx.wit().stream(type_index as usize); + let value = lift_stream_read_result(ctx, ty, call, buffer, progress as usize, result, iter)?; + + if let Some(handle) = dropped_handle { + unsafe { ty.drop_readable()(handle) }; + } + + if close && iter { + iterator_result(ctx, Value::new_undefined(ctx.clone()), true) + } else { + Ok(value) + } +} + +/// Finish a stream write, committing only the resource groups actually consumed. +/// +/// Keep the conversion context alive through settlement in the caller. Host +/// destructors run only after the endpoint's class borrow has been released. +fn finish_write<'js>( + ctx: &Ctx<'js>, + cls: &Class<'js, StreamWritable>, + call: &mut QjsCallContext, + buffer: BufferGuard, + code: u32, +) -> rquickjs::Result> { + let (progress, result) = + unpack_copy_result(code).expect("stream write completion must not block"); + + drop(buffer); + call.complete_transfers(progress as usize); + + let (type_index, dropped_handle) = { + let mut writable = cls.borrow_mut(); + writable.end.mark_completed(result); + + let handle = if result == CopyResult::Dropped { + writable.end.begin_drop()? + } else { + None + }; + + (writable.end.type_index(), handle) + }; + + if let Some(handle) = dropped_handle { + let ty = ctx.wit().stream(type_index as usize); + unsafe { ty.drop_writable()(handle) }; + } + + Ok(Value::new_number(ctx.clone(), progress as f64)) +} + +/// Lift copied elements before releasing their ABI buffer and resource guards. fn lift_stream_read_result<'js>( ctx: &Ctx<'js>, ty: wit_dylib_ffi::Stream, @@ -341,7 +377,7 @@ fn lift_stream_read_result<'js>( buffer: BufferGuard, progress: usize, copy_result: CopyResult, - iterator: bool, + iter: bool, ) -> rquickjs::Result> { let value = if matches!(ty.ty(), Some(wit_dylib_ffi::Type::U8)) { let vec = unsafe { buffer.into_vec(progress) }; @@ -352,9 +388,10 @@ fn lift_stream_read_result<'js>( unsafe { ty.lift(&mut call, buffer.ptr().add(ty.abi_payload_size() * offset)) }; arr.set(offset, call.pop_value(ctx))?; } + drop(buffer); - if iterator { + if iter { if progress == 0 { Value::new_undefined(ctx.clone()) } else { @@ -365,7 +402,7 @@ fn lift_stream_read_result<'js>( } }; - if iterator { + if iter { iterator_result( ctx, value, @@ -384,6 +421,7 @@ fn iterator_result<'js>( let result = Object::new(ctx.clone())?; result.set("value", value)?; result.set("done", done)?; + Ok(result.into_value()) } @@ -395,6 +433,7 @@ fn resolved_iterator_result<'js>( let (promise, resolve, _reject) = ctx.promise()?; let result = iterator_result(&ctx, value, done)?; resolve.call::<_, Value>((result,))?; + Ok(promise.into_value()) } @@ -408,36 +447,37 @@ fn stream_iterator_return<'js>( this: This>, ctx: Ctx<'js>, ) -> rquickjs::Result> { - let (handle, type_index, state) = { + let (has_handle, pending, cancellable) = { let readable = this.0.borrow(); ( - readable.end.handle, - readable.end.type_index, - readable.end.state, + readable.end.has_handle(), + readable.end.is_pending(), + readable.end.can_cancel(), ) }; - let Some(handle) = handle else { + if !has_handle { return resolved_iterator_result(ctx.clone(), Value::new_undefined(ctx), true); - }; - let ty = ctx.wit().stream(type_index as usize); + } - if matches!(state, CopyState::Idle | CopyState::Done) { - this.0.borrow_mut().end.handle.take(); - unsafe { ty.drop_readable()(handle) }; - this.0.borrow_mut().end.state = CopyState::Done; + if !pending { + stream_drop_readable(this, ctx.clone())?; return resolved_iterator_result(ctx.clone(), Value::new_undefined(ctx), true); } - if state != CopyState::AsyncCopying { + if !cancellable { return Err(rquickjs::Error::new_from_js( "stream", "iterator return while cancellation is in progress", )); } + let (handle, type_index) = this.0.borrow().end.begin_cancel()?; + let ty = ctx.wit().stream(type_index as usize); let (promise, resolve, _reject) = ctx.promise()?; + ctx.task().unjoin(handle); + let code = unsafe { ty.cancel_read()(handle) }; ctx.task() .set_stream_iterator_return(handle, Persistent::save(&ctx, resolve)); @@ -483,8 +523,9 @@ fn stream_drop_readable<'js>( ) -> rquickjs::Result<()> { let (handle, type_index) = { let mut readable = this.0.borrow_mut(); - (readable.end.begin_drop()?, readable.end.type_index) + (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) }; @@ -515,68 +556,6 @@ fn stream_write_one<'js>( stream_write_impl(this, ctx, data, StreamWriteMode::One) } -fn stream_write_iterable_item<'js>( - this: This>, - ctx: Ctx<'js>, - data: Value<'js>, -) -> rquickjs::Result> { - let (type_index, closed) = { - let writable = this.0.borrow(); - ( - writable.end.type_index, - writable.end.handle.is_none() || writable.end.state == CopyState::Done, - ) - }; - if closed { - let result = Value::new_number(ctx.clone(), 0.0); - return map_write_completion(ctx, result, 1); - } - - let ty = ctx.wit().stream(type_index as usize); - let batch_len = typed_array_batch_len(&data, &ty)?; - let expected = batch_len.unwrap_or(1); - - let result = if batch_len.is_some() { - stream_write_all(this, ctx.clone(), data)? - } else { - stream_write_one(this, ctx.clone(), data)? - }; - - map_write_completion(ctx, result, expected) -} - -fn map_write_completion<'js>( - ctx: Ctx<'js>, - result: Value<'js>, - expected: usize, -) -> rquickjs::Result> { - let Some(promise) = result.as_object() else { - let written: usize = result.get()?; - let (promise, resolve, _reject) = ctx.promise()?; - let complete = Value::new_bool(ctx.clone(), written == expected); - resolve.call::<_, Value>((complete,))?; - return Ok(promise.into_value()); - }; - - let then: Function = promise.get("then")?; - let complete = crate::coerce_fn( - move |ctx: Ctx<'_>, args: Rest>| -> rquickjs::Result> { - let written: usize = args - .0 - .into_iter() - .next() - .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "write count"))? - .get()?; - Ok(Value::new_bool(ctx, written == expected)) - }, - ); - let callback = Function::new(ctx.clone(), complete)?; - let mut args = function::Args::new(ctx, 1); - args.this(result)?; - args.push_arg(callback)?; - then.call_arg(args) -} - fn stream_write_impl<'js>( this: This>, ctx: Ctx<'js>, @@ -635,21 +614,8 @@ fn stream_write_impl<'js>( }; ctx.task().register(handle, pending); } 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); - (copy_result == CopyResult::Dropped) - .then(|| writable.end.handle.take()) - .flatten() - }; - if let Some(handle) = dropped_handle { - unsafe { ty.drop_writable()(handle) }; - } + let result = finish_write(&ctx, &this.0, &mut call, buffer, code)?; - let result = Value::new_number(ctx.clone(), progress as f64); resolve .call::<_, Value>((result,)) .expect("resolve stream write"); @@ -658,126 +624,6 @@ fn stream_write_impl<'js>( Ok(promise.into_value()) } -fn stream_write_all<'js>( - this: This>, - ctx: Ctx<'js>, - buffer: Value<'js>, -) -> rquickjs::Result> { - let stream_val = this.0.into_inner().into_value(); - write_all_step(ctx, stream_val, buffer, 0) -} - -fn write_all_buffer_len(value: &Value<'_>) -> rquickjs::Result { - if let Some(array) = value.as_array() { - return Ok(array.len()); - } - - let object = value - .as_object() - .ok_or_else(|| rquickjs::Error::new_from_js(value.type_of().as_str(), "array"))?; - object.get("length") -} - -fn write_all_step<'js>( - ctx: Ctx<'js>, - stream: Value<'js>, - buffer: Value<'js>, - total: usize, -) -> rquickjs::Result> { - // Check termination: buffer empty or stream done. - let class = Class::::from_value(&stream)?; - let state = class.borrow().end.state; - if state == CopyState::Done { - return Ok(Value::new_number(ctx, total as f64)); - } - - let buffer_len = write_all_buffer_len(&buffer)?; - if buffer_len == 0 { - return Ok(Value::new_number(ctx, total as f64)); - } - - // Call stream.write(buffer) with proper `this` binding. - let stream_obj = stream - .as_object() - .ok_or_else(|| rquickjs::Error::new_from_js("value", "stream object"))?; - - let write_fn: Function = stream_obj.get("write")?; - let mut call_args = function::Args::new(ctx.clone(), 1); - call_args.this(stream.clone())?; - call_args.push_arg(buffer.clone())?; - - let write_result: Value = write_fn.call_arg(call_args)?; - - let promise_obj = write_result - .as_object() - .ok_or_else(|| rquickjs::Error::new_from_js("value", "promise"))?; - let then_fn: Function = promise_obj.get("then")?; - - let stream_c = Cell::new(Some(Persistent::save(&ctx, stream))); - let buffer_c = Cell::new(Some(Persistent::save(&ctx, buffer))); - - let next = crate::coerce_fn( - move |ctx: Ctx<'_>, args: Rest>| -> rquickjs::Result> { - let count_val = args - .0 - .into_iter() - .next() - .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "write count"))?; - - let count: usize = count_val.get()?; - let buf = buffer_c - .take() - .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "write buffer"))? - .restore(&ctx)?; - - let s = stream_c - .take() - .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "stream"))? - .restore(&ctx)?; - - let buffer_len = write_all_buffer_len(&buf)?; - if count > buffer_len { - return Err(rquickjs::Exception::throw_range( - &ctx, - &format!("stream write reported {count} items for a {buffer_len}-item buffer"), - )); - } - - let class = Class::::from_value(&s)?; - let state = class.borrow().end.state; - if count == 0 { - if state == CopyState::Done { - return Ok(Value::new_number(ctx, total as f64)); - } - return Err(rquickjs::Exception::throw_range( - &ctx, - "stream write made no progress", - )); - } - if state == CopyState::Done { - return Ok(Value::new_number(ctx, (total + count) as f64)); - } - - let obj = buf - .as_object() - .ok_or_else(|| rquickjs::Error::new_from_js(buf.type_of().as_str(), "array"))?; - - let slice_fn: Function = obj.get("slice")?; - let mut slice_args = function::Args::new(ctx.clone(), 1); - slice_args.this(buf.clone())?; - slice_args.push_arg(count)?; - let sliced = slice_fn.call_arg(slice_args)?; - - write_all_step(ctx, s, sliced, total + count) - }, - ); - let cb = Function::new(ctx.clone(), next)?; - let mut then_args = function::Args::new(ctx.clone(), 1); - then_args.this(write_result)?; - then_args.push_arg(cb)?; - then_fn.call_arg(then_args) -} - pub(crate) fn stream_cancel_write<'js>( this: This>, ctx: Ctx<'js>, @@ -809,7 +655,7 @@ fn stream_drop_writable<'js>( ) -> rquickjs::Result<()> { let (handle, type_index) = { let mut writable = this.0.borrow_mut(); - (writable.end.begin_drop()?, writable.end.type_index) + (writable.end.begin_drop()?, writable.end.type_index()) }; if let Some(handle) = handle { let ty = ctx.wit().stream(type_index as usize); @@ -826,38 +672,17 @@ pub(crate) fn handle_write_event(handle: u32, result: u32) { mut call, resolve, wrapper, - .. + buffer, } = pending else { unreachable!("expected StreamWrite pending"); }; - 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(); - let cls = Class::::from_value(&w).unwrap(); - - let (type_index, dropped_handle) = { - let mut writable = cls.borrow_mut(); - writable.end.mark_completed(copy_result); - ( - writable.end.type_index, - (copy_result == CopyResult::Dropped) - .then(|| writable.end.handle.take()) - .flatten(), - ) - }; - if let Some(handle) = dropped_handle { - let ty = ctx.wit().stream(type_index as usize); - unsafe { ty.drop_writable()(handle) }; - } - - let val = Value::new_number(ctx.clone(), progress as f64); - let res = Persistent::save(ctx, val); - Some(res) + let class = Class::::from_value(&w).unwrap(); + let value = finish_write(ctx, &class, &mut call, buffer, result).unwrap(); + Some(Persistent::save(ctx, value)) }); resolve_promise(resolve, result); } @@ -878,40 +703,12 @@ pub(crate) fn handle_read_event(handle: u32, result: u32) { unreachable!("expected StreamRead pending"); }; - let (progress, copy_result) = - unpack_copy_result(result).expect("StreamRead callback should not be BLOCKED"); - let (result, return_result) = with_ctx(|ctx| { let w = wrapper.restore(ctx).unwrap(); let class = Class::::from_value(&w).unwrap(); let close = iterator_return.is_some(); - let (type_index, dropped_handle) = { - let mut cls = class.borrow_mut(); - cls.end.mark_completed(copy_result); - if close { - cls.end.state = CopyState::Done; - } - ( - cls.end.type_index, - (close || copy_result == CopyResult::Dropped) - .then(|| cls.end.handle.take()) - .flatten(), - ) - }; - - let ty = ctx.wit().stream(type_index as usize); - let progress = progress as usize; - - let mut result_val = - lift_stream_read_result(ctx, ty, call, buffer, progress, copy_result, iterator) - .unwrap(); - if close && iterator { - result_val = iterator_result(ctx, Value::new_undefined(ctx.clone()), true).unwrap(); - } - if let Some(handle) = dropped_handle { - unsafe { ty.drop_readable()(handle) }; - } + let result_val = finish_read(ctx, &class, call, buffer, result, iterator, close).unwrap(); let return_result = iterator_return.map(|resolve| { let result = iterator_result(ctx, Value::new_undefined(ctx.clone()), true).unwrap(); diff --git a/crates/runtime/src/streams/helpers.rs b/crates/runtime/src/streams/helpers.rs new file mode 100644 index 0000000..c4a4de7 --- /dev/null +++ b/crates/runtime/src/streams/helpers.rs @@ -0,0 +1,244 @@ +//! Stream convenience algorithms, separate from native ABI reads and writes. + +use std::cell::Cell; + +use rquickjs::class::Class; +use rquickjs::function::{Args, Rest, This}; +use rquickjs::{Ctx, Function, Object, Persistent, Value}; + +use super::{StreamWritable, stream_write_one}; +use crate::CtxExt; +use crate::typed_array::typed_array_len_as; + +/// Install the public convenience methods on the writable prototype. +pub(super) fn register<'js>(ctx: &Ctx<'js>, prototype: &Object<'js>) -> rquickjs::Result<()> { + prototype.set("writeAll", Function::new(ctx.clone(), stream_write_all)?)?; + prototype.set( + "writeIterableItem", + Function::new(ctx.clone(), stream_write_iterable_item)?, + ) +} + +/// Write matching typed arrays as batches and all other iterable items as scalars. +/// +/// Inspect the endpoint without retaining its borrow across writes, dispatch the +/// appropriate writer, then map its progress to a promise of full completion. +fn stream_write_iterable_item<'js>( + this: This>, + ctx: Ctx<'js>, + data: Value<'js>, +) -> rquickjs::Result> { + let (type_index, closed) = { + let writable = this.0.borrow(); + (writable.end.type_index(), writable.end.is_closed()) + }; + + if closed { + let result = Value::new_number(ctx.clone(), 0.0); + return map_write_completion(ctx, result, 1); + } + + let ty = ctx.wit().stream(type_index as usize); + let batch_len = typed_array_batch_len(&data, &ty)?; + let expected = batch_len.unwrap_or(1); + + let result = if batch_len.is_some() { + stream_write_all(this, ctx.clone(), data)? + } else { + stream_write_one(this, ctx.clone(), data)? + }; + + map_write_completion(ctx, result, expected) +} + +/// Map an immediate or promised write count to a promise of full completion. +fn map_write_completion<'js>( + ctx: Ctx<'js>, + result: Value<'js>, + expected: usize, +) -> rquickjs::Result> { + let Some(promise) = result.as_object() else { + let written: usize = result.get()?; + let (promise, resolve, _reject) = ctx.promise()?; + let complete = Value::new_bool(ctx.clone(), written == expected); + + resolve.call::<_, Value>((complete,))?; + return Ok(promise.into_value()); + }; + + let then: Function = promise.get("then")?; + + let complete = crate::coerce_fn( + move |ctx: Ctx<'_>, args: Rest>| -> rquickjs::Result> { + let written: usize = args + .0 + .into_iter() + .next() + .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "write count"))? + .get()?; + Ok(Value::new_bool(ctx, written == expected)) + }, + ); + + let callback = Function::new(ctx.clone(), complete)?; + let mut args = Args::new(ctx, 1); + + args.this(result)?; + args.push_arg(callback)?; + then.call_arg(args) +} + +/// Write an entire buffer, preserving immediate completion for empty or closed writes. +fn stream_write_all<'js>( + this: This>, + ctx: Ctx<'js>, + buffer: Value<'js>, +) -> rquickjs::Result> { + let stream = this.0.into_inner().into_value(); + write_all_step(ctx, stream, buffer, 0) +} + +/// Write a suffix, retain its JS values, and continue after the promise settles. +/// +/// Validate progress before slicing, stop on closure, and consume captured values +/// only once even when a custom writer supplies its own thenable. +fn write_all_step<'js>( + ctx: Ctx<'js>, + stream: Value<'js>, + buffer: Value<'js>, + total: usize, +) -> rquickjs::Result> { + let class = Class::::from_value(&stream)?; + let closed = class.borrow().end.is_closed(); + + if closed { + return Ok(Value::new_number(ctx, total as f64)); + } + + let buffer_len = write_all_buffer_len(&buffer)?; + + if buffer_len == 0 { + return Ok(Value::new_number(ctx, total as f64)); + } + + let stream_obj = stream + .as_object() + .ok_or_else(|| rquickjs::Error::new_from_js("value", "stream object"))?; + + let write_fn: Function = stream_obj.get("write")?; + let mut call_args = Args::new(ctx.clone(), 1); + call_args.this(stream.clone())?; + call_args.push_arg(buffer.clone())?; + + let write_result: Value = write_fn.call_arg(call_args)?; + let promise_obj = write_result + .as_object() + .ok_or_else(|| rquickjs::Error::new_from_js("value", "promise"))?; + let then_fn: Function = promise_obj.get("then")?; + + let stream_c = Cell::new(Some(Persistent::save(&ctx, stream))); + let buffer_c = Cell::new(Some(Persistent::save(&ctx, buffer))); + + let next = crate::coerce_fn( + move |ctx: Ctx<'_>, args: Rest>| -> rquickjs::Result> { + let count_val = args + .0 + .into_iter() + .next() + .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "write count"))?; + let count: usize = count_val.get()?; + let buf = buffer_c + .take() + .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "write buffer"))? + .restore(&ctx)?; + let stream = stream_c + .take() + .ok_or_else(|| rquickjs::Error::new_from_js("undefined", "stream"))? + .restore(&ctx)?; + + let buffer_len = write_all_buffer_len(&buf)?; + + if count > buffer_len { + return Err(rquickjs::Exception::throw_range( + &ctx, + &format!("stream write reported {count} items for a {buffer_len}-item buffer"), + )); + } + + let class = Class::::from_value(&stream)?; + let closed = class.borrow().end.is_closed(); + + if count == 0 { + if closed { + return Ok(Value::new_number(ctx, total as f64)); + } + + return Err(rquickjs::Exception::throw_range( + &ctx, + "stream write made no progress", + )); + } + + if closed { + return Ok(Value::new_number(ctx, (total + count) as f64)); + } + + let obj = buf + .as_object() + .ok_or_else(|| rquickjs::Error::new_from_js(buf.type_of().as_str(), "array"))?; + + let slice_fn: Function = obj.get("slice")?; + let mut slice_args = Args::new(ctx.clone(), 1); + + slice_args.this(buf.clone())?; + slice_args.push_arg(count)?; + let sliced = slice_fn.call_arg(slice_args)?; + + write_all_step(ctx, stream, sliced, total + count) + }, + ); + + let callback = Function::new(ctx.clone(), next)?; + let mut then_args = Args::new(ctx.clone(), 1); + + then_args.this(write_result)?; + then_args.push_arg(callback)?; + then_fn.call_arg(then_args) +} + +/// Recognize typed arrays whose element type matches the stream payload. +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), + wit_dylib_ffi::Type::S16 => typed_array_len_as!(obj, i16), + wit_dylib_ffi::Type::U32 => typed_array_len_as!(obj, u32), + wit_dylib_ffi::Type::S32 => typed_array_len_as!(obj, i32), + wit_dylib_ffi::Type::U64 => typed_array_len_as!(obj, u64), + 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), + _ => Ok(None), + } +} + +/// Preserve native array lengths and checked numeric conversion for writeAll. +fn write_all_buffer_len(value: &Value<'_>) -> rquickjs::Result { + if let Some(array) = value.as_array() { + return Ok(array.len()); + } + + let object = value + .as_object() + .ok_or_else(|| rquickjs::Error::new_from_js(value.type_of().as_str(), "array"))?; + + object.get("length") +} diff --git a/crates/runtime/src/trivia.rs b/crates/runtime/src/trivia.rs index d14cd93..4219d8f 100644 --- a/crates/runtime/src/trivia.rs +++ b/crates/runtime/src/trivia.rs @@ -1,10 +1,26 @@ use crate::CtxExt; +use crate::endpoint::EndpointError; use crate::with_ctx; use heck::ToLowerCamelCase; use rquickjs::{Atom, Function, Object, Persistent, Symbol}; use rquickjs::{Result, Value, function::Rest}; +impl From for rquickjs::Error { + /// Preserve JS-facing errors without coupling endpoint state to QuickJS. + fn from(error: EndpointError) -> Self { + let (from, message) = match error { + EndpointError::NotIdle(kind) => (kind, "operation requires an idle endpoint"), + EndpointError::WrongType(kind) => (kind, "matching WIT type"), + EndpointError::InFlight(kind) => (kind, "cancel and await before dropping"), + EndpointError::NotCancellable(kind) => (kind, "cancel without active operation"), + EndpointError::Dropped => ("object", "already dropped"), + }; + + Self::new_from_js(from, message) + } +} + /// Coerce closure lifetimes so the returned `Value<'js>` gets the same /// lifetime as the `Ctx<'js>` argument. /// diff --git a/crates/runtime/wit/init.wit b/crates/runtime/wit/init.wit deleted file mode 100644 index ca38264..0000000 --- a/crates/runtime/wit/init.wit +++ /dev/null @@ -1,11 +0,0 @@ -package local:init; - -interface module-loader { - resolve: func(referrer: string, specifier: string) -> result; - load: func(path: string) -> result; -} - -world init { - import module-loader; - export init: func(shim: string, script: string, entry-path: option, disable-gc: bool) -> result<_, string>; -} diff --git a/docs/runtime-intrinsics.md b/docs/runtime-intrinsics.md index af07a05..c497da3 100644 --- a/docs/runtime-intrinsics.md +++ b/docs/runtime-intrinsics.md @@ -64,29 +64,20 @@ Trigger a quickjs garbage collection cycle. - Returns : `undefined` -### `__cqjs.asyncExports` +## Export Dispatch -An object containing wrapper functions for async WIT exports. Each wrapper -calls the user's export function and chains `.then()` to signal `task_return` -back to the component model host. - -Structure mirrors the WIT export layout: - -```js -__cqjs.asyncExports = { - myFunc: Function, // root-scope async export - myInterface: { // interface-scoped exports - anotherFunc: Function, - }, -}; -``` +`crates/runtime/src/exports.rs` resolves and invokes both sync and async exports. +Resource methods take their receiver from the first canonical argument and +preserve the remaining argument order. Async calls activate their task before +export lookup, then attach fulfillment/rejection callbacks to signal +`task.return`; no JavaScript wrapper table is created. --- ## `globalThis.wit` : Public Stream/Future API The user-facing API for creating streams and futures from JavaScript. -Installed by the generated JS shim (see `src/codegen.rs`). +Installed by the generated JS shim (see `crates/core/src/codegen.rs`). ### `wit.Stream(type)` diff --git a/tests/async_types.rs b/tests/async_types.rs index 6f64582..e59db2b 100644 --- a/tests/async_types.rs +++ b/tests/async_types.rs @@ -10,8 +10,8 @@ use std::task::{Context, Poll}; use common::{AsyncComponentInstance, TestCase, WasiCtxState}; use wasmtime::component::{ - Component, Destination, FutureConsumer, FutureReader, Source, StreamConsumer, StreamProducer, - StreamReader, StreamResult, Val, VecBuffer, + Component, Destination, FutureConsumer, FutureReader, Lift, Source, StreamConsumer, + StreamProducer, StreamReader, StreamResult, Val, VecBuffer, }; use wasmtime::{AsContextMut, StoreContextMut}; @@ -113,6 +113,144 @@ async fn test_async_with_await() { assert_eq!(result, Val::U32(100)); } +/// Export getters can start async operations before the function is invoked. +#[tokio::test] +async fn test_async_export_lookup_has_active_task() { + let mut instance = TestCase::new() + .wit( + r#" + package test:async-lookup; + + interface api { + compute: async func(left: u32, right: u32) -> u32; + } + + world async-lookup { + export api; + export unused-future: async func() -> future; + } + "#, + ) + .script( + r#" + export const api = { + get compute() { + if ("asyncExports" in __cqjs) { + throw new Error("legacy async wrapper table is still installed"); + } + + const { readable, writable } = wit.Future(); + const read = readable.read(); + const write = writable.write(1); + + return async (left, right) => { + const offset = await read; + + if (!await write) throw new Error("future write was not consumed"); + + readable.drop(); + writable.drop(); + return left * 10 + right + offset; + }; + }, + }; + + export async function unusedFuture() { + throw new Error("not called"); + } + "#, + ) + .build_async() + .await + .unwrap(); + + let (component, store) = instance.parts(); + let api = component + .get_export_index(&mut *store, None, "test:async-lookup/api") + .expect("api interface not found"); + let compute = component + .get_export_index(&mut *store, Some(&api), "compute") + .expect("compute export not found"); + let compute = component.get_func(&mut *store, compute).unwrap(); + + for (left, right, expected) in [(4, 2, 43), (8, 3, 84)] { + let mut results = [Val::Bool(false)]; + + compute + .call_async( + &mut *store, + &[Val::U32(left), Val::U32(right)], + &mut results, + ) + .await + .unwrap(); + + assert_eq!(results[0], Val::U32(expected)); + } +} + +/// Direct dispatch preserves thenable receivers and synchronous settlement. +#[tokio::test] +async fn test_async_export_thenables() { + let mut instance = TestCase::new() + .wit( + r#" + package test:async-thenables; + + world async-thenables { + export settle: async func(value: u32, fail: bool) -> result; + export settle-void: async func(); + } + "#, + ) + .script( + r#" + export function settle(value, fail) { + return { + value, + fail, + then(resolve, reject) { + if (this.fail) { + reject("expected failure"); + } else { + resolve(this.value); + } + }, + }; + } + + export function settleVoid() { + return { then(resolve) { resolve(); } }; + } + "#, + ) + .build_async() + .await + .unwrap(); + + for fail in [false, true, false] { + let result = instance + .call1_async("settle", &[Val::U32(42), Val::Bool(fail)]) + .await + .unwrap(); + let expected = if fail { + Err(Some(Box::new(Val::String("expected failure".into())))) + } else { + Ok(Some(Box::new(Val::U32(42)))) + }; + + assert_eq!(result, Val::Result(expected)); + } + + assert!( + instance + .call_async("settle-void", &[], 0) + .await + .unwrap() + .is_empty() + ); +} + #[tokio::test] async fn test_async_exported_resource_members() { let mut instance = TestCase::new() @@ -1431,6 +1569,133 @@ async fn test_stream_write_all_rejects_invalid_or_stalled_writes() { ); } +/// Exercise native retry and completion behavior with controlled JavaScript writes. +#[tokio::test] +async fn test_stream_helpers_preserve_completion_shapes() { + let mut instance = TestCase::new() + .wit( + r#" + package test:stream-helper-completions; + world stream-helper-completions { + export verify-helpers: async func(); + export unused-stream: async func() -> stream; + } + "#, + ) + .script(include_str!("js/stream-helpers.js")) + .build_async() + .await + .unwrap(); + + instance.call_async("verify-helpers", &[], 0).await.unwrap(); +} + +/// Preserve iterable completion promises and distinguish batches from tuple values. +/// +/// Tuple items use a host consumer because Wasmtime does not support +/// intra-component stream copies with non-numeric payloads. +#[tokio::test] +async fn test_stream_iterable_item_completion_shapes() { + let mut instance = TestCase::new() + .wit( + r#" + package test:stream-iterable-item; + world stream-iterable-item { + export bytes: async func() -> list; + export tuple-item: async func() -> stream>; + export tuple-completed: async func() -> bool; + export unused-bytes: async func() -> stream; + } + "#, + ) + .script( + r#" + let tupleCompletion; + + export async function bytes() { + const { readable, writable } = wit.Stream(wit.Stream.U8); + const empty = writable.writeIterableItem(new Uint8Array()); + if (!(empty instanceof Promise) || await empty !== true) { + throw new Error("empty batch must resolve true"); + } + + const read = readable.read(2); + const written = writable.writeIterableItem(new Uint8Array([6, 7])); + if (!(written instanceof Promise) || await written !== true) { + throw new Error("complete batch must resolve true"); + } + + const data = await read; + writable.drop(); + readable.drop(); + + const closed = writable.writeIterableItem([1, 2]); + if (!(closed instanceof Promise) || await closed !== false) { + throw new Error("closed stream must resolve false"); + } + + return Array.from(data); + } + + export async function tupleItem() { + const { readable, writable } = wit.Stream(wit.Stream.TUPLE_U8_U8); + tupleCompletion = writable.writeIterableItem([1, 2]).then(complete => { + writable.drop(); + return complete; + }); + return readable; + } + + export async function tupleCompleted() { + return await tupleCompletion; + } + + export async function unusedBytes() { + throw new Error("type declaration only"); + } + + "#, + ) + .build_async() + .await + .unwrap(); + + assert_eq!( + instance.call1_async("bytes", &[]).await.unwrap(), + Val::List(vec![Val::U8(6), Val::U8(7)]) + ); + let output = Arc::new(Mutex::new(Vec::new())); + let (inst, store) = instance.parts(); + let func = inst + .get_typed_func::<(), (StreamReader<(u8, u8)>,)>(&mut *store, "tuple-item") + .unwrap(); + let (reader,) = func.call_async(&mut *store, ()).await.unwrap(); + reader + .pipe( + &mut *store, + StreamCollector { + values: Arc::clone(&output), + expected: 1, + }, + ) + .unwrap(); + store + .as_context_mut() + .run_concurrent(async |_| { + while output.lock().unwrap().is_empty() { + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + + assert_eq!(&*output.lock().unwrap(), &[(1, 2)]); + assert_eq!( + instance.call1_async("tuple-completed", &[]).await.unwrap(), + Val::Bool(true) + ); +} + #[tokio::test] async fn test_stream_write_uint32_array() { // Verify the typed-array fast path handles wider primitive element types @@ -1582,7 +1847,7 @@ async fn test_blocked_stream_write_keeps_lowered_payload_alive() { reader .pipe( &mut *store, - StringStreamConsumer { + StreamCollector { values: Arc::clone(&values), expected: 1, }, @@ -1705,7 +1970,7 @@ async fn test_stream_async_iterable_round_trip_infers_type() { reader .pipe( &mut *store, - StringStreamConsumer { + StreamCollector { values: Arc::clone(&output), expected: 2, }, @@ -1762,7 +2027,7 @@ async fn test_stream_async_iterable_batches_byte_arrays() { reader .pipe( &mut *store, - ByteStreamConsumer { + StreamCollector { values: Arc::clone(&output), expected: 5, }, @@ -1845,6 +2110,78 @@ async fn test_future_create_and_return_u32() { assert_eq!(results.len(), 1); } +/// Immediate, callback, and cancelled reads must share promise/ownership semantics. +#[tokio::test] +async fn test_future_read_completion_and_retry_paths() { + let mut instance = TestCase::new() + .wit( + r#" + package test:future-completion; + world future-completion { + export exercise: async func() -> list; + export unused-future: async func() -> future; + } + "#, + ) + .script( + r#" + export async function exercise() { + const values = []; + + for (const readFirst of [true, false]) { + const { readable, writable } = wit.Future(); + let read, write; + + if (readFirst) { + read = readable.read(); + write = writable.write(10); + } else { + write = writable.write(20); + read = readable.read(); + } + + if (read !== readable.read()) throw new Error("read promise not cached"); + values.push(await read); + if (!await write) throw new Error("write was not consumed"); + if (read !== readable.read()) throw new Error("completed promise not cached"); + readable.drop(); + writable.drop(); + } + + const { readable, writable } = wit.Future(); + const cancelled = readable.read(); + const cancellation = cancelled.catch(reason => String(reason)); + readable.cancelRead(); + + if (await cancellation !== "future read cancelled") { + throw new Error("cancellation did not reject the original read"); + } + + const retry = readable.read(); + if (retry === cancelled) throw new Error("cancelled promise was reused"); + if (!await writable.write(30)) throw new Error("retry write was not consumed"); + values.push(await retry); + readable.drop(); + writable.drop(); + + return values; + } + + export async function unusedFuture() { + throw new Error("not called"); + } + "#, + ) + .build_async() + .await + .unwrap(); + + assert_eq!( + instance.call1_async("exercise", &[]).await.unwrap(), + Val::List(vec![Val::U32(10), Val::U32(20), Val::U32(30)]) + ); +} + #[tokio::test] async fn test_future_create_and_return_string() { let mut instance = TestCase::new() @@ -1981,7 +2318,10 @@ async fn test_async_error_in_promise() { .call_async("might-fail", &[Val::Bool(true)], 1) .await; - assert!(result.is_ok() || result.is_err()); + assert!( + result.is_err(), + "rejected async export unexpectedly succeeded" + ); } #[tokio::test] @@ -2070,14 +2410,16 @@ impl StreamProducer for EmptyProducer } } -struct StringStreamConsumer { - values: Arc>>, +/// Collect typed stream items in host memory until the expected count is reached. +struct StreamCollector { + values: Arc>>, expected: usize, } -impl StreamConsumer for StringStreamConsumer { - type Item = String; +impl StreamConsumer for StreamCollector { + type Item = T; + /// Lift available values and close the consumer once its target count is met. fn poll_consume( self: Pin<&mut Self>, _cx: &mut Context<'_>, @@ -2088,33 +2430,7 @@ impl StreamConsumer for StringStreamConsumer { let mut values = Vec::with_capacity(source.remaining(&mut store)); source.read(&mut store, &mut values)?; self.values.lock().unwrap().extend(values); - let result = if self.values.lock().unwrap().len() >= self.expected { - StreamResult::Dropped - } else { - StreamResult::Completed - }; - Poll::Ready(Ok(result)) - } -} - -struct ByteStreamConsumer { - values: Arc>>, - expected: usize, -} - -impl StreamConsumer for ByteStreamConsumer { - type Item = u8; - fn poll_consume( - self: Pin<&mut Self>, - _cx: &mut Context<'_>, - mut store: StoreContextMut<'_, WasiCtxState>, - mut source: Source<'_, Self::Item>, - _finish: bool, - ) -> Poll> { - let mut values = Vec::with_capacity(source.remaining(&mut store)); - source.read(&mut store, &mut values)?; - self.values.lock().unwrap().extend(values); let result = if self.values.lock().unwrap().len() >= self.expected { StreamResult::Dropped } else { diff --git a/tests/js/stream-helpers.js b/tests/js/stream-helpers.js new file mode 100644 index 0000000..a416ea8 --- /dev/null +++ b/tests/js/stream-helpers.js @@ -0,0 +1,128 @@ +export async function verifyHelpers() { + await withStream(verifyPartialWrites); + await withStream(verifyEmptyAndClosedWrites); + await withStream(verifyClosureAfterProgress); + await withStream(verifyClosureWithoutProgress); + await withStream(verifyIncompleteIterable); + await withStream(verifyThenables); + await withStream(verifyWriteErrors); +} + +async function withStream(verify) { + const { readable, writable } = wit.Stream(wit.Stream.U8); + + try { + await verify(writable); + } finally { + writable.drop(); + readable.drop(); + } +} + +async function verifyPartialWrites(writable) { + const chunks = []; + writable.write = async buffer => { + chunks.push(Array.from(buffer)); + return Math.min(2, buffer.length); + }; + + assert(await writable.writeAll(new Uint8Array([1, 2, 3, 4, 5])) === 5, "partial total"); + assert(JSON.stringify(chunks) === "[[1,2,3,4,5],[3,4,5],[5]]", "partial suffixes"); + + assert(await writable.writeAll([6, 7, 8]) === 3, "plain-array batch"); + assert(await writable.writeIterableItem(new Uint8Array([6, 7, 8])), "typed-array batch"); +} + +async function verifyEmptyAndClosedWrites(writable) { + writable.write = () => { throw new Error("empty or closed stream must not write"); }; + + assert(writable.writeAll([]) === 0, "empty writes retain their immediate return shape"); + + const empty = writable.writeIterableItem(new Uint8Array()); + assert(empty instanceof Promise && await empty, "empty batch resolves true"); + + writable.drop(); + assert(writable.writeAll(42) === 0, "closed writes do not inspect the buffer"); + + const closed = writable.writeIterableItem([1, 2]); + assert(closed instanceof Promise && await closed === false, "closed iterable result"); +} + +async function verifyClosureAfterProgress(writable) { + writable.write = async function () { + this.drop(); + return 1; + }; + + assert(await writable.writeAll([1, 2]) === 1, "partial closure total"); +} + +async function verifyClosureWithoutProgress(writable) { + writable.write = async function () { + this.drop(); + return 0; + }; + + assert(await writable.writeAll([1]) === 0, "closure without progress"); +} + +async function verifyIncompleteIterable(writable) { + writable.write = async function () { + this.drop(); + return 1; + }; + + assert(await writable.writeIterableItem(new Uint8Array([1, 2])) === false, "partial iterable result"); +} + +async function verifyThenables(writable) { + writable.write = buffer => ({ then: next => next(buffer.length) }); + + const mapped = writable.writeIterableItem(new Uint8Array([1, 2])); + assert(mapped instanceof Promise && await mapped, "immediate thenable completion resolves true"); + + writable.write = () => ({ + then(next) { + next(1); + return next(1); + } + }); + + await rejects(() => writable.writeAll([1]), "write buffer"); +} + +async function verifyWriteErrors(writable) { + const write = () => writable.writeAll([1]); + + writable.write = async () => 0; + await rejects(() => writable.writeAll(42), "array"); + await rejects(write, "made no progress"); + + writable.write = async () => 2; + await rejects(write, "2 items for a 1-item buffer"); + + writable.write = () => 1; + await rejects(write, "promise"); + + writable.write = async () => { throw new Error("write failed"); }; + await rejects(write, "write failed"); +} + +function assert(condition, message) { + if (!condition) throw new Error(message); +} + +async function rejects(call, message) { + try { + await call(); + } catch (error) { + assert(error.message.includes(message), error.message); + return; + } + + throw new Error(`expected rejection: ${message}`); +} + +export async function unusedStream() { + throw new Error("type declaration only"); +} diff --git a/xtask/Cargo.toml b/xtask/Cargo.toml new file mode 100644 index 0000000..414a40c --- /dev/null +++ b/xtask/Cargo.toml @@ -0,0 +1,14 @@ +[package] +name = "xtask" +version = "0.0.0" +edition.workspace = true +rust-version.workspace = true +publish = false + +[dependencies] +anyhow.workspace = true +clap.workspace = true +flate2.workspace = true +glob.workspace = true +tar.workspace = true +ureq.workspace = true diff --git a/xtask/src/main.rs b/xtask/src/main.rs new file mode 100644 index 0000000..218e708 --- /dev/null +++ b/xtask/src/main.rs @@ -0,0 +1,52 @@ +//! Prepare runtime Wasm directly, without building the host componentizer. + +#[path = "../../crates/core/build/runtime.rs"] +pub mod runtime; + +use std::path::{Path, PathBuf}; + +use anyhow::Result; +use clap::Parser; +use runtime::{RUNTIME_AUDITABLE_ENV, RUNTIME_BUILDS, RuntimeOptions}; + +/// Portable runtime-artifact preparation options. +#[derive(Parser)] +#[command(about = "Prepare componentize-qjs runtime Wasm artifacts")] +struct Args { + /// Directory receiving the stable runtime*.wasm filenames. + #[arg(long, default_value = "target/runtime")] + output: PathBuf, + /// Directory for tool downloads and disposable nested Cargo targets. + #[arg(long, default_value = "target/runtime-build")] + build_dir: PathBuf, + /// Optimize the runtime with release settings and wasm-opt. + #[arg(long)] + release: bool, + /// Prepare only the two non-async runtime variants. + #[arg(long)] + sync_only: bool, + /// Retain dependency metadata using cargo-auditable. + #[arg(long)] + auditable: bool, +} + +/// Build the requested runtime variants and report their explicit output paths. +fn main() -> Result<()> { + let args = Args::parse(); + let runtime_dir = Path::new(env!("CARGO_MANIFEST_DIR")).join("../crates/runtime"); + let options = RuntimeOptions { + release: args.release, + async_support: !args.sync_only, + auditable: args.auditable || std::env::var_os(RUNTIME_AUDITABLE_ENV).is_some(), + }; + + runtime::prepare(&runtime_dir, &args.output, &args.build_dir, options)?; + + for build in RUNTIME_BUILDS { + if !build.async_support || options.async_support { + println!("{}", args.output.join(build.filename).display()); + } + } + + Ok(()) +}