Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"))]
Expand Down
3 changes: 3 additions & 0 deletions src/time/duration.rs
Original file line number Diff line number Diff line change
@@ -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.
Expand Down
9 changes: 6 additions & 3 deletions src/time/instant.rs
Original file line number Diff line number Diff line change
@@ -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.
Expand All @@ -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".
Expand All @@ -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
Expand Down
228 changes: 186 additions & 42 deletions src/time/mod.rs
Original file line number Diff line number Diff line change
@@ -1,32 +1,31 @@
//! Async time interfaces.

#[cfg(target_env = "p2")]
pub(crate) mod utils;

mod duration;
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 {
Expand All @@ -35,11 +34,23 @@ impl SystemTime {
}

impl From<SystemTime> 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.
Expand All @@ -62,58 +73,191 @@ impl AsyncIterator for Interval {
}
}

#[derive(Debug)]
pub struct Timer(Option<AsyncPollable>);
#[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<AsyncPollable>);

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<crate::runtime::WaitFor>
}
}
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<Self::Output> {
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<crate::runtime::WaitFor>
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<Box<dyn Future<Output = ()>>>);

impl Future for Wait {
type Output = Instant;

fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
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<Self::Output> {
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<Output = Instant>) {
let start = Instant::now();
let now = f.await;
Expand Down
24 changes: 22 additions & 2 deletions tests/sleep.rs
Original file line number Diff line number Diff line change
@@ -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;

Expand All @@ -9,3 +9,23 @@ async fn just_sleep() -> Result<(), Box<dyn Error>> {
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"]);
}
Loading