diff --git a/.cargo/config.toml b/.cargo/config.toml index e9a242c..aa3925b 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -3,3 +3,9 @@ # * AWS auth environment variables, for running the wstd-aws integration tests. # * . directory is available at . runner = "wasmtime run -Shttp --env AWS_ACCESS_KEY_ID --env AWS_SECRET_ACCESS_KEY --env AWS_SESSION_TOKEN --dir .::." + +[target.wasm32-wasip3] +# wasmtime is given: +# * AWS auth environment variables, for running the wstd-aws integration tests. +# * . directory is available at . +runner = "wasmtime run -Shttp -Sp3 --env AWS_ACCESS_KEY_ID --env AWS_SECRET_ACCESS_KEY --env AWS_SESSION_TOKEN --dir .::." diff --git a/.github/actions/install-rust/action.yml b/.github/actions/install-rust/action.yml index 0d11907..3c16f0c 100644 --- a/.github/actions/install-rust/action.yml +++ b/.github/actions/install-rust/action.yml @@ -29,7 +29,15 @@ runs: run: | rustup set profile minimal rustup update "${{ steps.select.outputs.version }}" --no-self-update - rustup default "${{ steps.select.outputs.version }}" + # `rustup default` gets overwritten by rust-toolchain.toml, `rustup + # override` has higher precedence though. + rustup override set "${{ steps.select.outputs.version }}" + + rustup target add wasm32-wasip2 + + if [ "${{ inputs.toolchain }}" = "nightly" ]; then + rustup target add wasm32-wasip3 + fi # Save disk space by avoiding incremental compilation. Also turn down # debuginfo from 2 to 0 to help save disk space. diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index 1362609..598dd5f 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -48,12 +48,19 @@ jobs: role-to-assume: arn:aws:iam::313377415443:role/wstd-aws-ci-role role-session-name: github-ci + - name: rust version + run: cargo --version + - name: check run: cargo check --workspace --all --bins --examples - - name: wstd tests + - name: wstd p2 tests run: cargo test -p wstd -p wstd-axum --target wasm32-wasip2 -- --nocapture + - name: wstd p3 tests + if: matrix.rust == 'nightly' + run: cargo test -p wstd -p wstd-axum --target wasm32-wasip3 -- --nocapture + - name: test-programs tests run: cargo test -p test-programs -- --nocapture if: steps.creds.outcome == 'success' @@ -90,4 +97,3 @@ jobs: - run: ./publish bump # Make sure the tree is publish-able as-is - run: ./publish verify - diff --git a/Cargo.toml b/Cargo.toml index fd399e4..2a7e28b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -27,13 +27,18 @@ http.workspace = true itoa.workspace = true pin-project-lite.workspace = true slab.workspace = true -wasip2.workspace = true wstd-macro.workspace = true # optional serde = { workspace = true, optional = true } serde_json = { workspace = true, optional = true } +[target.'cfg(all(target_os = "wasi", target_env = "p2"))'.dependencies] +wasip2.workspace = true + +[target.'cfg(all(target_os = "wasi", target_env = "p3"))'.dependencies] +wasip3.workspace = true + [dev-dependencies] anyhow.workspace = true clap.workspace = true @@ -59,7 +64,7 @@ license = "Apache-2.0 WITH LLVM-exception" repository = "https://github.com/bytecodealliance/wstd" keywords = ["WebAssembly", "async", "stdlib", "Components"] categories = ["wasm", "asynchronous"] -rust-version = "1.91.1" +rust-version = "1.92.0" authors = [ "Yoshua Wuyts ", "Pat Hickey ", @@ -95,13 +100,13 @@ test-programs = { path = "test-programs" } tower-service = "0.3.3" ureq = { version = "3.1", default-features = false, features = ["json"] } wasip2 = "1.0" -wstd = { path = ".", version = "=0.6.8" } +wasip3 = "0.8" +wstd = { path = ".", version = "=0.6.8", default-features = false } wstd-axum = { path = "./axum", version = "=0.6.8" } wstd-axum-macro = { path = "./axum/macro", version = "=0.6.8" } wstd-macro = { path = "./macro", version = "=0.6.8" } [package.metadata.docs.rs] -all-features = true targets = [ "wasm32-wasip2" ] diff --git a/axum/Cargo.toml b/axum/Cargo.toml index 154a021..74141e9 100644 --- a/axum/Cargo.toml +++ b/axum/Cargo.toml @@ -13,7 +13,7 @@ rust-version.workspace = true [dependencies] axum.workspace = true tower-service.workspace = true -wstd.workspace = true +wstd = { workspace = true, features = ["json"] } wstd-axum-macro.workspace = true [dev-dependencies] diff --git a/axum/examples/hello_world.rs b/axum/examples/hello_world.rs index 90a9cbd..9931344 100644 --- a/axum/examples/hello_world.rs +++ b/axum/examples/hello_world.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + //! Run with //! //! ```sh diff --git a/axum/examples/hello_world_nomacro.rs b/axum/examples/hello_world_nomacro.rs index f3068d1..bc0a39b 100644 --- a/axum/examples/hello_world_nomacro.rs +++ b/axum/examples/hello_world_nomacro.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + //! Run with //! //! ```sh diff --git a/axum/examples/weather.rs b/axum/examples/weather.rs index 5163fcc..0df8f97 100644 --- a/axum/examples/weather.rs +++ b/axum/examples/weather.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + //! This demo app shows a Axum based wasi-http server making an arbitrary //! number of http requests as part of serving a single response. //! diff --git a/axum/src/lib.rs b/axum/src/lib.rs index 3272b91..c5a2af1 100644 --- a/axum/src/lib.rs +++ b/axum/src/lib.rs @@ -1,3 +1,4 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] //! Support for the [`axum`] web server framework in wasi-http components, via //! [`wstd`]. //! diff --git a/examples/complex_http_client.rs b/examples/complex_http_client.rs index 703e931..7c3e470 100644 --- a/examples/complex_http_client.rs +++ b/examples/complex_http_client.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use anyhow::{Result, anyhow}; use clap::{ArgAction, Parser}; use std::str::FromStr; diff --git a/examples/http_client.rs b/examples/http_client.rs index f4465a8..178fbeb 100644 --- a/examples/http_client.rs +++ b/examples/http_client.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use anyhow::{Result, anyhow}; use clap::{ArgAction, Parser}; use wstd::http::{Body, BodyExt, Client, Method, Request, Uri}; diff --git a/examples/http_server.rs b/examples/http_server.rs index 8449f83..fa67518 100644 --- a/examples/http_server.rs +++ b/examples/http_server.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use anyhow::{Context, Result}; use futures_lite::stream::{once_future, unfold}; use http_body_util::{BodyExt, StreamBody}; diff --git a/examples/http_server_proxy.rs b/examples/http_server_proxy.rs index ccad608..8d0cf01 100644 --- a/examples/http_server_proxy.rs +++ b/examples/http_server_proxy.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + //! Run the example with: //! ```sh //! cargo build --example http_server_proxy --target=wasm32-wasip2 diff --git a/examples/macro_main.rs b/examples/macro_main.rs new file mode 100644 index 0000000..d301547 --- /dev/null +++ b/examples/macro_main.rs @@ -0,0 +1,9 @@ +#![cfg_attr(not(target_os = "wasi"), no_main)] +#![cfg(target_os = "wasi")] + +//! Verifies that the `main` macro is compiling. + +#[wstd::main] +async fn main() { + println!("Hello world"); +} diff --git a/examples/tcp_echo_server.rs b/examples/tcp_echo_server.rs index f1dd895..224222a 100644 --- a/examples/tcp_echo_server.rs +++ b/examples/tcp_echo_server.rs @@ -1,10 +1,13 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use wstd::io; use wstd::iter::AsyncIterator; use wstd::net::TcpListener; #[wstd::main] async fn main() -> io::Result<()> { - let listener = TcpListener::bind("127.0.0.1:8080").await?; + let mut listener = TcpListener::bind("127.0.0.1:8080").await?; println!("Listening on {}", listener.local_addr()?); println!("type `nc localhost 8080` to create a TCP client"); @@ -14,7 +17,9 @@ async fn main() -> io::Result<()> { println!("Accepted from: {}", stream.peer_addr()?); wstd::runtime::spawn(async move { // If echo copy fails, we can ignore it. - let _ = io::copy(&stream, &stream).await; + let mut stream = stream; + let (mut read_half, mut write_half) = stream.split(); + let _ = io::copy(&mut read_half, &mut write_half).await; }) .detach(); } diff --git a/examples/tcp_stream_client.rs b/examples/tcp_stream_client.rs index 2e93faf..a269b8c 100644 --- a/examples/tcp_stream_client.rs +++ b/examples/tcp_stream_client.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use wstd::io::{self, AsyncRead, AsyncWrite}; use wstd::net::TcpStream; diff --git a/examples/udp_echo_server.rs b/examples/udp_echo_server.rs index f840107..c441d87 100644 --- a/examples/udp_echo_server.rs +++ b/examples/udp_echo_server.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use wstd::io; use wstd::net::UdpSocket; diff --git a/examples/udp_stream_client.rs b/examples/udp_stream_client.rs index 013caee..f26f5d5 100644 --- a/examples/udp_stream_client.rs +++ b/examples/udp_stream_client.rs @@ -1,3 +1,6 @@ +#![cfg_attr(not(all(target_os = "wasi", target_env = "p2")), no_main)] +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use wstd::io; use wstd::net::{UdpSocket, UdpStream}; diff --git a/src/http/body.rs b/src/http/body.rs index 95e8a2e..b4de743 100644 --- a/src/http/body.rs +++ b/src/http/body.rs @@ -3,7 +3,7 @@ use crate::http::{ error::Context as _, fields::{header_map_from_wasi, header_map_to_wasi}, }; -use crate::io::{AsyncInputStream, AsyncOutputStream}; +use crate::io::{AsyncInputStream, AsyncOutputStream, AsyncWrite}; use crate::runtime::{AsyncPollable, Reactor, WaitFor}; pub use ::http_body::{Body as HttpBody, Frame, SizeHint}; @@ -77,7 +77,7 @@ impl Body { match self.0 { BodyInner::Incoming(incoming) => incoming.send(outgoing_body).await, BodyInner::Boxed(box_body) => { - let out_stream = AsyncOutputStream::new( + let mut out_stream = AsyncOutputStream::new( outgoing_body .write() .expect("outgoing body already written"), @@ -108,7 +108,7 @@ impl Body { } } BodyInner::Complete { data, trailers } => { - let out_stream = AsyncOutputStream::new( + let mut out_stream = AsyncOutputStream::new( outgoing_body .write() .expect("outgoing body already written"), @@ -348,14 +348,14 @@ impl Incoming { } async fn send(self, outgoing_body: WasiOutgoingBody) -> Result<(), Error> { let in_body = self.body; - let in_stream = + let mut in_stream = AsyncInputStream::new(in_body.stream().expect("incoming body already read")); - let out_stream = AsyncOutputStream::new( + let mut out_stream = AsyncOutputStream::new( outgoing_body .write() .expect("outgoing body already written"), ); - in_stream.copy_to(&out_stream).await.map_err(|e| { + in_stream.copy_to(&mut out_stream).await.map_err(|e| { Error::from(e).context("copying incoming body stream to outgoing body stream") })?; drop(in_stream); diff --git a/src/io/read.rs b/src/io/read.rs index a6a95da..474188e 100644 --- a/src/io/read.rs +++ b/src/io/read.rs @@ -28,7 +28,7 @@ pub trait AsyncRead { // If the `AsyncRead` implementation is an unbuffered wrapper around an // `AsyncInputStream`, some I/O operations can be more efficient. #[inline] - fn as_async_input_stream(&self) -> Option<&io::AsyncInputStream> { + fn as_async_input_stream(&mut self) -> Option<&mut io::AsyncInputStream> { None } } @@ -45,7 +45,7 @@ impl AsyncRead for &mut R { } #[inline] - fn as_async_input_stream(&self) -> Option<&io::AsyncInputStream> { + fn as_async_input_stream(&mut self) -> Option<&mut io::AsyncInputStream> { (**self).as_async_input_stream() } } diff --git a/src/io/stdio.rs b/src/io/stdio.rs index b2ac153..e403986 100644 --- a/src/io/stdio.rs +++ b/src/io/stdio.rs @@ -43,8 +43,8 @@ impl AsyncRead for Stdin { } #[inline] - fn as_async_input_stream(&self) -> Option<&AsyncInputStream> { - Some(&self.stream) + fn as_async_input_stream(&mut self) -> Option<&mut AsyncInputStream> { + Some(&mut self.stream) } } @@ -93,7 +93,7 @@ impl AsyncWrite for Stdout { } #[inline] - fn as_async_output_stream(&self) -> Option<&AsyncOutputStream> { + fn as_async_output_stream(&mut self) -> Option<&mut AsyncOutputStream> { self.stream.as_async_output_stream() } } @@ -143,7 +143,7 @@ impl AsyncWrite for Stderr { } #[inline] - fn as_async_output_stream(&self) -> Option<&AsyncOutputStream> { + fn as_async_output_stream(&mut self) -> Option<&mut AsyncOutputStream> { self.stream.as_async_output_stream() } } diff --git a/src/io/streams.rs b/src/io/streams.rs index 3676d21..34fd3b0 100644 --- a/src/io/streams.rs +++ b/src/io/streams.rs @@ -26,7 +26,7 @@ impl AsyncInputStream { stream, } } - fn poll_ready(&self, cx: &mut Context<'_>) -> Poll<()> { + fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<()> { // Lazily initialize the AsyncPollable let subscription = self .subscription @@ -43,40 +43,13 @@ impl AsyncInputStream { } } /// Await for read readiness. - async fn ready(&self) { + async fn ready(&mut self) { poll_fn(|cx| self.poll_ready(cx)).await } - /// Asynchronously read from the input stream. - /// This method is the same as [`AsyncRead::read`], but doesn't require a `&mut self`. - pub async fn read(&self, buf: &mut [u8]) -> std::io::Result { - let read = loop { - self.ready().await; - // Ideally, the ABI would be able to read directly into buf. - // However, with the default generated bindings, it returns a - // newly allocated vec, which we need to copy into buf. - match self.stream.read(buf.len() as u64) { - // A read of 0 bytes from WASI's `read` doesn't mean - // end-of-stream as it does in Rust. However, `self.ready()` - // cannot guarantee that at least one byte is ready for - // reading, so in this case we try again. - Ok(r) if r.is_empty() => continue, - Ok(r) => break r, - // 0 bytes from Rust's `read` means end-of-stream. - Err(StreamError::Closed) => return Ok(0), - Err(StreamError::LastOperationFailed(err)) => { - return Err(std::io::Error::other(err.to_debug_string())); - } - } - }; - let len = read.len(); - buf[0..len].copy_from_slice(&read); - Ok(len) - } - /// Move the entire contents of an input stream directly into an output /// stream, until the input stream has closed. This operation is optimized /// to avoid copying stream contents into and out of memory. - pub async fn copy_to(&self, writer: &AsyncOutputStream) -> std::io::Result { + pub async fn copy_to(&mut self, writer: &mut AsyncOutputStream) -> std::io::Result { let mut written = 0; loop { self.ready().await; @@ -124,11 +97,32 @@ impl AsyncInputStream { impl AsyncRead for AsyncInputStream { async fn read(&mut self, buf: &mut [u8]) -> std::io::Result { - Self::read(self, buf).await + let read = loop { + self.ready().await; + // Ideally, the ABI would be able to read directly into buf. + // However, with the default generated bindings, it returns a + // newly allocated vec, which we need to copy into buf. + match self.stream.read(buf.len() as u64) { + // A read of 0 bytes from WASI's `read` doesn't mean + // end-of-stream as it does in Rust. However, `self.ready()` + // cannot guarantee that at least one byte is ready for + // reading, so in this case we try again. + Ok(r) if r.is_empty() => continue, + Ok(r) => break r, + // 0 bytes from Rust's `read` means end-of-stream. + Err(StreamError::Closed) => return Ok(0), + Err(StreamError::LastOperationFailed(err)) => { + return Err(std::io::Error::other(err.to_debug_string())); + } + } + }; + let len = read.len(); + buf[0..len].copy_from_slice(&read); + Ok(len) } #[inline] - fn as_async_input_stream(&self) -> Option<&AsyncInputStream> { + fn as_async_input_stream(&mut self) -> Option<&mut AsyncInputStream> { Some(self) } } @@ -150,9 +144,10 @@ impl AsyncInputChunkStream { impl futures_lite::stream::Stream for AsyncInputChunkStream { type Item = Result, std::io::Error>; fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - match self.stream.poll_ready(cx) { + let this = self.get_mut(); + match this.stream.poll_ready(cx) { Poll::Pending => Poll::Pending, - Poll::Ready(()) => match self.stream.stream.read(self.chunk_size as u64) { + Poll::Ready(()) => match this.stream.stream.read(this.chunk_size as u64) { Ok(r) if r.is_empty() => Poll::Pending, Ok(r) => Poll::Ready(Some(Ok(r))), Err(StreamError::LastOperationFailed(err)) => { @@ -233,7 +228,7 @@ impl AsyncOutputStream { } } /// Await write readiness. - async fn ready(&self) { + async fn ready(&mut self) { // Lazily initialize the AsyncPollable let subscription = self .subscription @@ -241,15 +236,11 @@ impl AsyncOutputStream { // Wait on readiness subscription.wait_for().await; } - /// Asynchronously write to the output stream. This method is the same as - /// [`AsyncWrite::write`], but doesn't require a `&mut self`. - /// - /// Awaits for write readiness, and then performs at most one write to the - /// output stream. Returns how much of the argument `buf` was written, or - /// a `std::io::Error` indicating either an error returned by the stream write - /// using the debug string provided by the WASI error, or else that the, - /// indicated by `std::io::ErrorKind::ConnectionReset`. - pub async fn write(&self, buf: &[u8]) -> std::io::Result { +} + +impl AsyncWrite for AsyncOutputStream { + // Required methods + async fn write(&mut self, buf: &[u8]) -> std::io::Result { // Loops at most twice. loop { match self.stream.check_write() { @@ -279,32 +270,7 @@ impl AsyncOutputStream { } } } - - /// Asynchronously write to the output stream. This method is the same as - /// [`AsyncWrite::write_all`], but doesn't require a `&mut self`. - pub async fn write_all(&self, buf: &[u8]) -> std::io::Result<()> { - let mut to_write = &buf[0..]; - loop { - let bytes_written = self.write(to_write).await?; - to_write = &to_write[bytes_written..]; - if to_write.is_empty() { - return Ok(()); - } - } - } - - /// Asyncronously flush the output stream. Initiates a flush, and then - /// awaits until the flush is complete and the output stream is ready for - /// writing again. - /// - /// This method is the same as [`AsyncWrite::flush`], but doesn't require - /// a `&mut self`. - /// - /// Fails with a `std::io::Error` indicating either an error returned by - /// the stream flush, using the debug string provided by the WASI error, - /// or else that the stream is closed, indicated by - /// `std::io::ErrorKind::ConnectionReset`. - pub async fn flush(&self) -> std::io::Result<()> { + async fn flush(&mut self) -> std::io::Result<()> { match self.stream.flush() { Ok(()) => { self.ready().await; @@ -318,19 +284,9 @@ impl AsyncOutputStream { } } } -} - -impl AsyncWrite for AsyncOutputStream { - // Required methods - async fn write(&mut self, buf: &[u8]) -> std::io::Result { - Self::write(self, buf).await - } - async fn flush(&mut self) -> std::io::Result<()> { - Self::flush(self).await - } #[inline] - fn as_async_output_stream(&self) -> Option<&AsyncOutputStream> { + fn as_async_output_stream(&mut self) -> Option<&mut AsyncOutputStream> { Some(self) } } diff --git a/src/io/write.rs b/src/io/write.rs index 79cf0d9..12fd6ae 100644 --- a/src/io/write.rs +++ b/src/io/write.rs @@ -20,7 +20,7 @@ pub trait AsyncWrite { // If the `AsyncWrite` implementation is an unbuffered wrapper around an // `AsyncOutputStream`, some I/O operations can be more efficient. #[inline] - fn as_async_output_stream(&self) -> Option<&io::AsyncOutputStream> { + fn as_async_output_stream(&mut self) -> Option<&mut io::AsyncOutputStream> { None } } @@ -42,7 +42,7 @@ impl AsyncWrite for &mut W { } #[inline] - fn as_async_output_stream(&self) -> Option<&io::AsyncOutputStream> { + fn as_async_output_stream(&mut self) -> Option<&mut io::AsyncOutputStream> { (**self).as_async_output_stream() } } diff --git a/src/lib.rs b/src/lib.rs index ebc673d..dca9cd9 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -55,30 +55,53 @@ //! These are unique capabilities provided by WASI 0.2, and because this library //! is specific to that are exposed from here. +#[cfg(all(target_os = "wasi", target_env = "p2"))] pub mod future; +#[cfg(all(target_os = "wasi", target_env = "p2"))] #[macro_use] pub mod http; +#[cfg(all(target_os = "wasi", target_env = "p2"))] pub mod io; pub mod iter; +#[cfg(all(target_os = "wasi", target_env = "p2"))] pub mod net; +#[cfg(all(target_os = "wasi", target_env = "p2"))] 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"))] pub mod time; -pub use wstd_macro::attr_macro_http_server as http_server; -pub use wstd_macro::attr_macro_main as main; -pub use wstd_macro::attr_macro_test as test; +pub use wstd_macro::{ + attr_macro_http_server as http_server, attr_macro_main as main, attr_macro_test as test, +}; -// Re-export the wasip2 crate for use only by `wstd-macro` macros. The proc -// macros need to generate code that uses these definitions, but we don't want -// to treat it as part of our public API with regards to semver, so we keep it -// under `__internal` as well as doc(hidden) to indicate it is private. +// Re-export the active WASI backend crate for use only by `wstd-macro` macros. +// The proc macros need to generate code that uses these definitions, but we +// don't want to treat it as part of our public API with regards to semver, so +// we keep it under `__internal` as well as doc(hidden) to indicate it is +// private. #[doc(hidden)] pub mod __internal { + #[cfg(all(target_os = "wasi", target_env = "p2"))] pub use wasip2; + #[cfg(all(target_os = "wasi", target_env = "p3"))] + pub use wasip3; } +#[cfg(all(target_os = "wasi", target_env = "p2"))] pub mod prelude { pub use crate::future::FutureExt as _; pub use crate::io::AsyncRead as _; diff --git a/src/net/tcp_listener.rs b/src/net/tcp_listener.rs index 3ee9007..69a70f3 100644 --- a/src/net/tcp_listener.rs +++ b/src/net/tcp_listener.rs @@ -55,7 +55,7 @@ impl TcpListener { } /// Returns an iterator over the connections being received on this listener. - pub fn incoming(&self) -> Incoming<'_> { + pub fn incoming(&mut self) -> Incoming<'_> { Incoming { listener: self } } } @@ -63,7 +63,7 @@ impl TcpListener { /// An iterator that infinitely accepts connections on a TcpListener. #[derive(Debug)] pub struct Incoming<'a> { - listener: &'a TcpListener, + listener: &'a mut TcpListener, } impl<'a> AsyncIterator for Incoming<'a> { diff --git a/src/net/tcp_stream.rs b/src/net/tcp_stream.rs index af3674a..977fb29 100644 --- a/src/net/tcp_stream.rs +++ b/src/net/tcp_stream.rs @@ -86,8 +86,17 @@ impl TcpStream { Ok(format!("{addr:?}")) } - pub fn split(&self) -> (ReadHalf<'_>, WriteHalf<'_>) { - (ReadHalf(self), WriteHalf(self)) + pub fn split(&mut self) -> (ReadHalf<'_>, WriteHalf<'_>) { + ( + ReadHalf { + stream: &mut self.input, + socket: &self.socket, + }, + WriteHalf { + stream: &mut self.output, + socket: &self.socket, + }, + ) } } @@ -104,18 +113,8 @@ impl io::AsyncRead for TcpStream { self.input.read(buf).await } - fn as_async_input_stream(&self) -> Option<&AsyncInputStream> { - Some(&self.input) - } -} - -impl io::AsyncRead for &TcpStream { - async fn read(&mut self, buf: &mut [u8]) -> io::Result { - self.input.read(buf).await - } - - fn as_async_input_stream(&self) -> Option<&AsyncInputStream> { - (**self).as_async_input_stream() + fn as_async_input_stream(&mut self) -> Option<&mut AsyncInputStream> { + Some(&mut self.input) } } @@ -128,64 +127,56 @@ impl io::AsyncWrite for TcpStream { self.output.flush().await } - fn as_async_output_stream(&self) -> Option<&AsyncOutputStream> { - Some(&self.output) + fn as_async_output_stream(&mut self) -> Option<&mut AsyncOutputStream> { + Some(&mut self.output) } } -impl io::AsyncWrite for &TcpStream { - async fn write(&mut self, buf: &[u8]) -> io::Result { - self.output.write(buf).await - } - - async fn flush(&mut self) -> io::Result<()> { - self.output.flush().await - } +pub struct ReadHalf<'a> { + stream: &'a mut AsyncInputStream, + socket: &'a TcpSocket, +} - fn as_async_output_stream(&self) -> Option<&AsyncOutputStream> { - (**self).as_async_output_stream() +impl<'a> Drop for ReadHalf<'a> { + fn drop(&mut self) { + let _ = self + .socket + .shutdown(wasip2::sockets::tcp::ShutdownType::Receive); } } -pub struct ReadHalf<'a>(&'a TcpStream); impl<'a> io::AsyncRead for ReadHalf<'a> { async fn read(&mut self, buf: &mut [u8]) -> io::Result { - self.0.read(buf).await + self.stream.read(buf).await } - fn as_async_input_stream(&self) -> Option<&AsyncInputStream> { - self.0.as_async_input_stream() + fn as_async_input_stream(&mut self) -> Option<&mut AsyncInputStream> { + self.stream.as_async_input_stream() } } -impl<'a> Drop for ReadHalf<'a> { - fn drop(&mut self) { - let _ = self - .0 - .socket - .shutdown(wasip2::sockets::tcp::ShutdownType::Receive); - } +pub struct WriteHalf<'a> { + stream: &'a mut AsyncOutputStream, + socket: &'a TcpSocket, } -pub struct WriteHalf<'a>(&'a TcpStream); impl<'a> io::AsyncWrite for WriteHalf<'a> { async fn write(&mut self, buf: &[u8]) -> io::Result { - self.0.write(buf).await + self.stream.write(buf).await } async fn flush(&mut self) -> io::Result<()> { - self.0.flush().await + self.stream.flush().await } - fn as_async_output_stream(&self) -> Option<&AsyncOutputStream> { - self.0.as_async_output_stream() + fn as_async_output_stream(&mut self) -> Option<&mut AsyncOutputStream> { + self.stream.as_async_output_stream() } } impl<'a> Drop for WriteHalf<'a> { fn drop(&mut self) { let _ = self - .0 .socket .shutdown(wasip2::sockets::tcp::ShutdownType::Send); } diff --git a/tests/http_first_byte_timeout.rs b/tests/http_first_byte_timeout.rs index 5fd47be..11d381a 100644 --- a/tests/http_first_byte_timeout.rs +++ b/tests/http_first_byte_timeout.rs @@ -1,6 +1,8 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use wstd::http::{Body, Client, Request, error::ErrorCode}; -#[wstd::main] +#[wstd::test] async fn main() -> Result<(), Box> { // Set first byte timeout to 1/2 second. let mut client = Client::new(); diff --git a/tests/http_get.rs b/tests/http_get.rs index e7a3a5a..1735d22 100644 --- a/tests/http_get.rs +++ b/tests/http_get.rs @@ -1,3 +1,5 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use std::error::Error; use wstd::http::{Body, Client, HeaderValue, Request}; diff --git a/tests/http_get_json.rs b/tests/http_get_json.rs index a4f42b6..641c2b7 100644 --- a/tests/http_get_json.rs +++ b/tests/http_get_json.rs @@ -1,3 +1,5 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use serde::Deserialize; use std::error::Error; use wstd::http::{Body, Client, Request}; diff --git a/tests/http_handle_error_code.rs b/tests/http_handle_error_code.rs index 6affb90..fa72e0c 100644 --- a/tests/http_handle_error_code.rs +++ b/tests/http_handle_error_code.rs @@ -1,3 +1,5 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use wstd::http::{Body, Client, Request, error::ErrorCode}; /// Test that `outgoing_handler::handle` errors are properly propagated. diff --git a/tests/http_post.rs b/tests/http_post.rs index 5de184d..1082a40 100644 --- a/tests/http_post.rs +++ b/tests/http_post.rs @@ -1,3 +1,5 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use std::error::Error; use wstd::http::{Client, HeaderValue, Request}; diff --git a/tests/http_post_json.rs b/tests/http_post_json.rs index f9fcf07..f67d050 100644 --- a/tests/http_post_json.rs +++ b/tests/http_post_json.rs @@ -1,3 +1,5 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use serde::{Deserialize, Serialize}; use std::error::Error; use wstd::http::{Body, Client, HeaderValue, Request}; diff --git a/tests/http_timeout.rs b/tests/http_timeout.rs index 96d40de..6a156d9 100644 --- a/tests/http_timeout.rs +++ b/tests/http_timeout.rs @@ -1,3 +1,5 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use wstd::future::FutureExt; use wstd::http::{Body, Client, Request}; use wstd::time::Duration; diff --git a/tests/macro_test.rs b/tests/macro_test.rs new file mode 100644 index 0000000..5dd0a22 --- /dev/null +++ b/tests/macro_test.rs @@ -0,0 +1,12 @@ +// Verify that the test macro is acutally running. + +#[wstd::test] +async fn computation() -> Result<(), String> { + assert_eq!(1 + 1, 2); + Ok(()) +} + +#[wstd::test] +async fn second_computation() { + assert_eq!(1 + 1, 2); +} diff --git a/tests/sleep.rs b/tests/sleep.rs index eb55ea0..c888302 100644 --- a/tests/sleep.rs +++ b/tests/sleep.rs @@ -1,3 +1,5 @@ +#![cfg(all(target_os = "wasi", target_env = "p2"))] + use std::error::Error; use wstd::task::sleep; use wstd::time::Duration;