diff --git a/Cargo.toml b/Cargo.toml index 6db1876..b2e2627 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -37,7 +37,7 @@ serde_json = { workspace = true, optional = true } wasip2.workspace = true [target.'cfg(all(target_os = "wasi", target_env = "p3"))'.dependencies] -wasip3.workspace = true +wasip3 = { workspace = true, features = ["async-spawn"] } [dev-dependencies] anyhow.workspace = true diff --git a/src/lib.rs b/src/lib.rs index ad5aad4..43030ba 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -67,18 +67,7 @@ pub mod iter; pub mod net; #[cfg(target_os = "wasi")] pub mod rand; -#[cfg(all(target_os = "wasi", target_env = "p2"))] pub mod runtime; -#[cfg(all(target_os = "wasi", target_env = "p3"))] -pub mod runtime { - pub fn block_on(fut: F) -> F::Output - where - F: Future, - T: 'static, - { - wasip3::wit_bindgen::block_on(fut) - } -} #[cfg(all(target_os = "wasi", target_env = "p2"))] pub mod task; #[cfg(all(target_os = "wasi", target_env = "p2"))] diff --git a/src/runtime/mod.rs b/src/runtime/mod.rs index 24b9fc2..a92d7fa 100644 --- a/src/runtime/mod.rs +++ b/src/runtime/mod.rs @@ -1,25 +1,36 @@ //! Async event loop support. //! -//! The way to use this is to call [`block_on()`]. Inside the future, [`Reactor::current`] -//! will give an instance of the [`Reactor`] running the event loop, which can be -//! to [`AsyncPollable::wait_for`] instances of +//! On WASI 0.2 the way to use this is to call [`block_on()`]. Inside the +//! future, [`Reactor::current`] will give an instance of the [`Reactor`] +//! running the event loop, which can be used to [`AsyncPollable::wait_for`] +//! instances of //! [`wasip2::Pollable`](https://docs.rs/wasi/latest/wasi/io/poll/struct.Pollable.html). //! This will automatically wait for the futures to resolve, and call the //! necessary wakers to work. +//! +//! On WASI 0.3 [`block_on`] can be used to drive a future, but an async +//! function can also be directly exported and will be driven by the host. #![deny(missing_debug_implementations, nonstandard_style)] #![warn(missing_docs, unreachable_pub)] +pub use ::async_task::Task; + +#[cfg(target_env = "p2")] mod block_on; +#[cfg(target_env = "p2")] mod reactor; -pub use ::async_task::Task; +#[cfg(target_env = "p2")] pub use block_on::block_on; +#[cfg(target_env = "p2")] pub use reactor::{AsyncPollable, Reactor, WaitFor}; +#[cfg(target_env = "p2")] use std::cell::RefCell; // There are no threads in WASI 0.2, so this is just a safe way to thread a single reactor to all // use sites in the background. +#[cfg(target_env = "p2")] std::thread_local! { pub(crate) static REACTOR: RefCell> = const { RefCell::new(None) }; } @@ -27,6 +38,7 @@ pub(crate) static REACTOR: RefCell> = const { RefCell::new(None) /// Spawn a `Future` as a `Task` on the current `Reactor`. /// /// Panics if called from outside `block_on`. +#[cfg(target_env = "p2")] pub fn spawn(fut: F) -> Task where F: std::future::Future + 'static, @@ -34,3 +46,26 @@ where { Reactor::current().spawn(fut) } + +#[cfg(target_env = "p3")] +pub use ::async_task::Runnable; +#[cfg(target_env = "p3")] +pub use wasip3::wit_bindgen::block_on; + +/// Spawn a `Future` as a `Task` on the WASI 0.3 async runtime. +#[cfg(target_env = "p3")] +pub fn spawn(fut: F) -> Task +where + F: std::future::Future + 'static, + T: 'static, +{ + let (runnable, task) = async_task::spawn_local(fut, |runnable: Runnable| { + // Scheduling the task is accomplished by spawning a future which + // executes the `run` method. + wasip3::spawn_local(async move { + let _ = runnable.run(); + }); + }); + runnable.schedule(); + task +} diff --git a/tests/runtime_spawn.rs b/tests/runtime_spawn.rs new file mode 100644 index 0000000..29e9e96 --- /dev/null +++ b/tests/runtime_spawn.rs @@ -0,0 +1,76 @@ +use std::cell::Cell; +use std::future::pending; +use std::rc::Rc; + +use futures_lite::future::yield_now; + +struct DropFlag(Rc>); + +impl Drop for DropFlag { + fn drop(&mut self) { + self.0.set(true); + } +} + +async fn wait_until(flag: &Cell) { + while !flag.get() { + yield_now().await; + } +} + +#[wstd::test] +async fn spawn_detach_cancel_and_drop() { + assert_eq!(wstd::runtime::spawn(async { 42 }).await, 42); + + // Check that a detached task completes eventually. + let detached_completed = Rc::new(Cell::new(false)); + let completed = detached_completed.clone(); + wstd::runtime::spawn(async move { + yield_now().await; + completed.set(true); + }) + .detach(); + + assert!(!detached_completed.get()); + wait_until(&detached_completed).await; + assert!(detached_completed.get()); + + // Check that calling `cancel` on a task cancels and drops the spawned + // future. + let canceled_started = Rc::new(Cell::new(false)); + let canceled_dropped = Rc::new(Cell::new(false)); + let canceled_completed = Rc::new(Cell::new(false)); + let started = canceled_started.clone(); + let dropped = canceled_dropped.clone(); + let completed = canceled_completed.clone(); + let task = wstd::runtime::spawn(async move { + let _drop_flag = DropFlag(dropped); + started.set(true); + pending::<()>().await; + completed.set(true); + }); + + wait_until(&canceled_started).await; + assert_eq!(task.cancel().await, None); + assert!(canceled_dropped.get()); + assert!(!canceled_completed.get()); + + // Check that dropping a task cancels and drops the spawned future. + let dropped_started = Rc::new(Cell::new(false)); + let dropped_dropped = Rc::new(Cell::new(false)); + let dropped_completed = Rc::new(Cell::new(false)); + let started = dropped_started.clone(); + let dropped = dropped_dropped.clone(); + let completed = dropped_completed.clone(); + let task = wstd::runtime::spawn(async move { + let _drop_flag = DropFlag(dropped); + started.set(true); + pending::<()>().await; + completed.set(true); + }); + + wait_until(&dropped_started).await; + drop(task); + wait_until(&dropped_dropped).await; + assert!(!dropped_completed.get()); +}