diff --git a/src/lib.rs b/src/lib.rs index 43030ba..45c6a85 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -68,9 +68,9 @@ pub mod net; #[cfg(target_os = "wasi")] pub mod rand; pub mod runtime; -#[cfg(all(target_os = "wasi", target_env = "p2"))] +#[cfg(target_os = "wasi")] pub mod task; -#[cfg(all(target_os = "wasi", target_env = "p2"))] +#[cfg(target_os = "wasi")] pub mod time; #[cfg(all(target_os = "wasi", target_env = "p2"))] diff --git a/src/time/duration.rs b/src/time/duration.rs index 7f67ceb..15a6103 100644 --- a/src/time/duration.rs +++ b/src/time/duration.rs @@ -1,7 +1,10 @@ use super::{Instant, Wait}; use std::future::IntoFuture; use std::ops::{Add, AddAssign, Sub, SubAssign}; +#[cfg(target_env = "p2")] use wasip2::clocks::monotonic_clock; +#[cfg(target_env = "p3")] +use wasip3::clocks::monotonic_clock; /// A Duration type to represent a span of time, typically used for system /// timeouts. diff --git a/src/time/instant.rs b/src/time/instant.rs index 6e9cf97..03b7357 100644 --- a/src/time/instant.rs +++ b/src/time/instant.rs @@ -1,7 +1,10 @@ use super::{Duration, Wait}; use std::future::IntoFuture; use std::ops::{Add, AddAssign, Sub, SubAssign}; -use wasip2::clocks::monotonic_clock; +#[cfg(target_env = "p2")] +use wasip2::clocks::monotonic_clock::{self, Instant as WasiInstant}; +#[cfg(target_env = "p3")] +use wasip3::clocks::monotonic_clock::{self, Mark as WasiInstant}; /// A measurement of a monotonically nondecreasing clock. Opaque and useful only /// with Duration. @@ -10,7 +13,7 @@ use wasip2::clocks::monotonic_clock; /// without coherence issues, just like if we were implementing this in the /// stdlib. #[derive(Debug, PartialEq, PartialOrd, Ord, Eq, Hash, Clone, Copy)] -pub struct Instant(pub(crate) monotonic_clock::Instant); +pub struct Instant(pub(crate) WasiInstant); impl Instant { /// Returns an instant corresponding to "now". @@ -24,7 +27,7 @@ impl Instant { /// ``` #[must_use] pub fn now() -> Self { - Instant(wasip2::clocks::monotonic_clock::now()) + Instant(monotonic_clock::now()) } /// Returns the amount of time elapsed from another instant to this one, or zero duration if diff --git a/src/time/mod.rs b/src/time/mod.rs index db0e1b3..b715c11 100644 --- a/src/time/mod.rs +++ b/src/time/mod.rs @@ -1,5 +1,6 @@ //! Async time interfaces. +#[cfg(target_env = "p2")] pub(crate) mod utils; mod duration; @@ -7,26 +8,24 @@ mod instant; pub use duration::Duration; pub use instant::Instant; -use pin_project_lite::pin_project; use std::future::Future; use std::pin::Pin; use std::task::{Context, Poll}; -use wasip2::clocks::{ - monotonic_clock::{subscribe_duration, subscribe_instant}, - wall_clock, -}; +#[cfg(target_env = "p2")] +use wasip2::clocks::wall_clock::{self, Datetime}; +#[cfg(target_env = "p3")] +use wasip3::clocks::system_clock::{self as wall_clock, Instant as Datetime}; -use crate::{ - iter::AsyncIterator, - runtime::{AsyncPollable, Reactor}, -}; +use crate::iter::AsyncIterator; +#[cfg(target_env = "p2")] +use crate::runtime::{AsyncPollable, Reactor}; /// A measurement of the system clock, useful for talking to external entities /// like the file system or other processes. May be converted losslessly to a /// more useful `std::time::SystemTime` to provide more methods. #[derive(Debug, Clone, Copy)] #[allow(dead_code)] -pub struct SystemTime(wall_clock::Datetime); +pub struct SystemTime(Datetime); impl SystemTime { pub fn now() -> Self { @@ -35,11 +34,23 @@ impl SystemTime { } impl From for std::time::SystemTime { + #[cfg(target_env = "p2")] fn from(st: SystemTime) -> Self { std::time::SystemTime::UNIX_EPOCH + std::time::Duration::from_secs(st.0.seconds) + std::time::Duration::from_nanos(st.0.nanoseconds.into()) } + + #[cfg(target_env = "p3")] + fn from(st: SystemTime) -> Self { + let mut result = std::time::SystemTime::UNIX_EPOCH; + if st.0.seconds < 0 { + result -= std::time::Duration::from_secs(st.0.seconds.unsigned_abs()); + } else { + result += std::time::Duration::from_secs(st.0.seconds.unsigned_abs()); + } + result + std::time::Duration::from_nanos(st.0.nanoseconds.into()) + } } /// An async iterator representing notifications at fixed interval. @@ -62,58 +73,191 @@ impl AsyncIterator for Interval { } } -#[derive(Debug)] -pub struct Timer(Option); +#[cfg(target_env = "p2")] +mod timer { + use super::*; + use pin_project_lite::pin_project; + use wasip2::clocks::monotonic_clock::{subscribe_duration, subscribe_instant}; + + #[derive(Debug)] + pub struct Timer(Option); -impl Timer { - pub fn never() -> Timer { - Timer(None) + impl Timer { + pub fn never() -> Timer { + Timer(None) + } + pub fn at(deadline: Instant) -> Timer { + let pollable = Reactor::current().schedule(subscribe_instant(deadline.0)); + Timer(Some(pollable)) + } + pub fn after(duration: Duration) -> Timer { + let pollable = Reactor::current().schedule(subscribe_duration(duration.0)); + Timer(Some(pollable)) + } + pub fn set_after(&mut self, duration: Duration) { + *self = Self::after(duration); + } + pub fn wait(&self) -> Wait { + let wait_for = self.0.as_ref().map(AsyncPollable::wait_for); + Wait { wait_for } + } } - pub fn at(deadline: Instant) -> Timer { - let pollable = Reactor::current().schedule(subscribe_instant(deadline.0)); - Timer(Some(pollable)) + + pin_project! { + /// Future created by [`Timer::wait`] + #[must_use = "futures do nothing unless polled or .awaited"] + pub struct Wait { + #[pin] + wait_for: Option + } } - pub fn after(duration: Duration) -> Timer { - let pollable = Reactor::current().schedule(subscribe_duration(duration.0)); - Timer(Some(pollable)) + + impl Future for Wait { + type Output = Instant; + + fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + let this = self.project(); + match this.wait_for.as_pin_mut() { + None => Poll::Pending, + Some(f) => match f.poll(cx) { + Poll::Pending => Poll::Pending, + Poll::Ready(()) => Poll::Ready(Instant::now()), + }, + } + } } - pub fn set_after(&mut self, duration: Duration) { - *self = Self::after(duration); +} + +#[cfg(target_env = "p3")] +mod timer { + use super::*; + use wasip3::clocks::monotonic_clock::{wait_for, wait_until}; + + #[derive(Debug)] + pub struct Timer(TimerInner); + + enum TimerInner { + Never, + At(Instant), + After(Duration), } - pub fn wait(&self) -> Wait { - let wait_for = self.0.as_ref().map(AsyncPollable::wait_for); - Wait { wait_for } + + impl std::fmt::Debug for TimerInner { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str("Timer") + } } -} -pin_project! { - /// Future created by [`Timer::wait`] - #[must_use = "futures do nothing unless polled or .awaited"] - pub struct Wait { - #[pin] - wait_for: Option + impl Timer { + pub fn never() -> Timer { + Timer(TimerInner::Never) + } + pub fn at(deadline: Instant) -> Timer { + Timer(TimerInner::At(deadline)) + } + pub fn after(duration: Duration) -> Timer { + Timer(TimerInner::After(duration)) + } + pub fn set_after(&mut self, duration: Duration) { + *self = Self::after(duration); + } + pub fn wait(&self) -> Wait { + match &self.0 { + TimerInner::Never => Wait(Box::pin(std::future::pending())), + TimerInner::At(instant) => Wait(Box::pin(wait_until(instant.0))), + TimerInner::After(duration) => Wait(Box::pin(wait_for(duration.0))), + } + } } -} -impl Future for Wait { - type Output = Instant; + /// Future created by [`Timer::wait`]. + #[must_use = "futures do nothing unless polled or .awaited"] + pub struct Wait(Pin>>); + + impl Future for Wait { + type Output = Instant; - fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { - let this = self.project(); - match this.wait_for.as_pin_mut() { - None => Poll::Pending, - Some(f) => match f.poll(cx) { + fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + match self.0.as_mut().poll(cx) { Poll::Pending => Poll::Pending, Poll::Ready(()) => Poll::Ready(Instant::now()), - }, + } } } } +pub use timer::*; + #[cfg(test)] mod test { use super::*; + #[cfg(target_env = "p2")] + fn system_time(seconds: u64, nanoseconds: u32) -> SystemTime { + SystemTime(Datetime { + seconds, + nanoseconds, + }) + } + + #[cfg(target_env = "p3")] + fn system_time(seconds: i64, nanoseconds: u32) -> SystemTime { + SystemTime(Datetime { + seconds, + nanoseconds, + }) + } + + #[test] + fn system_time_conversion_at_epoch() { + let actual: std::time::SystemTime = system_time(0, 0).into(); + + assert_eq!(actual, std::time::SystemTime::UNIX_EPOCH); + } + + #[test] + fn system_time_conversion_after_epoch() { + let actual: std::time::SystemTime = system_time(1, 999_999_999).into(); + let elapsed = actual + .duration_since(std::time::SystemTime::UNIX_EPOCH) + .unwrap(); + + assert_eq!(elapsed, std::time::Duration::new(1, 999_999_999)); + } + + #[cfg(target_env = "p3")] + #[test] + fn system_time_conversion_just_before_epoch() { + let actual: std::time::SystemTime = system_time(-1, 999_999_999).into(); + let before_epoch = std::time::SystemTime::UNIX_EPOCH + .duration_since(actual) + .unwrap(); + + assert_eq!(before_epoch, std::time::Duration::from_nanos(1)); + } + + #[cfg(target_env = "p3")] + #[test] + fn system_time_conversion_at_seconds_limits() { + for seconds in [i64::MIN, i64::MAX] { + let actual: std::time::SystemTime = system_time(seconds, 0).into(); + let elapsed = if seconds < 0 { + std::time::SystemTime::UNIX_EPOCH + .duration_since(actual) + .unwrap() + } else { + actual + .duration_since(std::time::SystemTime::UNIX_EPOCH) + .unwrap() + }; + + assert_eq!( + elapsed, + std::time::Duration::from_secs(seconds.unsigned_abs()) + ); + } + } + async fn debug_duration(what: &str, f: impl Future) { let start = Instant::now(); let now = f.await; diff --git a/tests/sleep.rs b/tests/sleep.rs index c888302..446736f 100644 --- a/tests/sleep.rs +++ b/tests/sleep.rs @@ -1,6 +1,6 @@ -#![cfg(all(target_os = "wasi", target_env = "p2"))] - +use std::cell::RefCell; use std::error::Error; +use std::rc::Rc; use wstd::task::sleep; use wstd::time::Duration; @@ -9,3 +9,23 @@ async fn just_sleep() -> Result<(), Box> { sleep(Duration::from_secs(1)).await; Ok(()) } + +#[wstd::test] +async fn concurrent_sleeps_wake_in_deadline_order() { + let wake_order = Rc::new(RefCell::new(Vec::new())); + let mut tasks = Vec::new(); + + for (delay, task) in [(60, "slow"), (20, "fast"), (40, "medium")] { + let wake_order = wake_order.clone(); + tasks.push(wstd::runtime::spawn(async move { + sleep(Duration::from_millis(delay)).await; + wake_order.borrow_mut().push(task); + })); + } + + for task in tasks { + task.await; + } + + assert_eq!(*wake_order.borrow(), ["fast", "medium", "slow"]); +}