Skip to main content

ridl_rt/
task.rs

1//! Waiting on a future from a thread, a waker that does nothing, and a waker
2//! that sets a flag.
3//!
4//! This module is gated by the `std` cargo feature, the one feature of this
5//! crate that links the standard library. It is for three callers:
6//!
7//! - **A blocking client built over an async one.** A client that returns a
8//!   future is the one source of truth for a call's behaviour; a blocking
9//!   variant of the same call is [`block_on`] over that future, so the two
10//!   cannot diverge. The `deadline` bounds the whole wait.
11//! - **A frame loop that polls a future once per frame.** Such a loop needs a
12//!   `Context`, and a `Context` needs a `Waker`, but the loop polls on its own
13//!   schedule and has no use for a wake. [`noop_waker`] is that waker: created
14//!   once when the loop starts and cloned into each `Context`.
15//! - **A frame loop that polls again when the future asks for it.** A future
16//!   that does a bounded amount of work per poll and wakes itself when work is
17//!   left — the generated `Serve` takes at most 32 claims per poll — makes
18//!   progress under [`noop_waker`] only at the loop's next frame, because that
19//!   wake is discarded. [`flag_waker`] records the wake in a flag the loop
20//!   reads: poll, then poll again while the flag was set, up to the loop's own
21//!   limit of polls per frame.
22//!
23//! All three are written without `unsafe`, which this crate forbids, over
24//! `std::task::Wake` on an `Arc`. That is why the module is under `std`:
25//! `Waker::noop()` needs Rust 1.85, above this crate's minimum, and building a
26//! `Waker` from a raw vtable needs `unsafe`. The module takes no dependency:
27//! the wait is `std::thread::park_timeout`, whose unpark token cannot be lost
28//! between a poll and the park that follows it.
29//!
30//! The feature compiles for `wasm32-unknown-unknown`, because `just wasm-check`
31//! requires every feature to, and [`block_on`] is not usable on that target:
32//! `Instant::now()` panics there, and a park does not block the thread. A
33//! frame loop on wasm polls with [`noop_waker`] or [`flag_waker`] and never
34//! calls [`block_on`]. A `no_std` frame loop, where this module does not
35//! exist, writes the same small waker over `alloc::task::Wake` on an `Arc`
36//! when it has an allocator; one without an allocator needs a hand-written
37//! `RawWaker`, which needs `unsafe`, as the paragraph above explains.
38//!
39//! Nothing here names a port, an interface or a generated type; the module
40//! knows only `core::future::Future`.
41
42use core::future::Future;
43use core::pin::pin;
44use core::sync::atomic::{AtomicBool, Ordering};
45use core::task::{Context, Poll, Waker};
46use std::sync::Arc;
47use std::task::Wake;
48use std::thread::{self, Thread};
49use std::time::Instant;
50
51/// Wake by unparking the thread that is blocked in [`block_on`].
52struct Unpark(Thread);
53
54impl Wake for Unpark {
55    fn wake(self: Arc<Self>) {
56        self.0.unpark();
57    }
58
59    fn wake_by_ref(self: &Arc<Self>) {
60        self.0.unpark();
61    }
62}
63
64/// A waker whose wake does nothing.
65struct Noop;
66
67impl Wake for Noop {
68    fn wake(self: Arc<Self>) {}
69
70    fn wake_by_ref(self: &Arc<Self>) {}
71}
72
73/// Run `fut` to completion on the current thread, parking the thread between
74/// polls, and give up at `deadline`.
75///
76/// The future is polled once immediately, then once more after every wake of
77/// the waker it was polled with. Between polls the thread is parked, so the
78/// wait costs no CPU. The waker unparks the thread through `std::task::Wake`
79/// on an `Arc`; an unpark that arrives while the thread is not parked is kept
80/// as a token that ends the next park at once, so a wake between a poll and
81/// the park that follows it is not lost. A wake that is not from the waker —
82/// `Thread::unpark` called by someone else, or a park ending on its own — is
83/// harmless: it causes one extra poll, and the wait continues.
84///
85/// # The deadline
86///
87/// - `None`: wait until the future is ready, however long that takes.
88/// - `Some(deadline)`: return `Some(output)` from the first poll that returns
89///   `Ready`, whenever that poll runs — also after a park that ended late —
90///   and `None` once a poll returns `Pending` when
91///   `Instant::now() >= deadline`. The future is polled at least once even
92///   when the deadline has already passed at entry, so a future that is
93///   already ready returns `Some` whatever the deadline. On `None`, the future
94///   is dropped without another poll.
95///
96/// The deadline is checked after each `Pending` poll, so `None` can be
97/// returned only after a poll, and each park is given at most the time left
98/// until the deadline.
99///
100/// ```
101/// use std::time::{Duration, Instant};
102///
103/// let out = ridl_rt::task::block_on(std::future::ready(7), None);
104/// assert_eq!(out, Some(7));
105///
106/// let pending = std::future::pending::<u8>();
107/// let out = ridl_rt::task::block_on(pending, Some(Instant::now() + Duration::from_millis(10)));
108/// assert_eq!(out, None);
109/// ```
110pub fn block_on<F: Future>(fut: F, deadline: Option<Instant>) -> Option<F::Output> {
111    let mut fut = pin!(fut);
112    let waker = Waker::from(Arc::new(Unpark(thread::current())));
113    let mut cx = Context::from_waker(&waker);
114    loop {
115        if let Poll::Ready(out) = fut.as_mut().poll(&mut cx) {
116            return Some(out);
117        }
118        match deadline {
119            None => thread::park(),
120            Some(deadline) => {
121                let now = Instant::now();
122                if now >= deadline {
123                    return None;
124                }
125                thread::park_timeout(deadline - now);
126            }
127        }
128    }
129}
130
131/// A waker whose wake does nothing, for a loop that polls on its own schedule.
132///
133/// Waking it, by value or by reference, has no effect and never panics; a
134/// future polled with it makes progress only when the caller polls it again.
135/// The waker is built as `Waker::from(Arc<Noop>)`, so each call allocates one
136/// `Arc`: create it once per loop and clone it into each `Context`, rather
137/// than calling this function per poll. Two clones of one waker report
138/// `will_wake` as `true`; two wakers from two calls do not.
139///
140/// ```
141/// use std::future::Future;
142/// use std::task::Context;
143///
144/// let waker = ridl_rt::task::noop_waker();
145/// let mut cx = Context::from_waker(&waker);
146/// let mut fut = std::pin::pin!(std::future::ready(1));
147/// assert!(fut.as_mut().poll(&mut cx).is_ready());
148/// ```
149pub fn noop_waker() -> Waker {
150    Waker::from(Arc::new(Noop))
151}
152
153/// A waker whose wake sets a flag, and the handle that reads and clears it.
154#[derive(Debug)]
155struct Flag(AtomicBool);
156
157impl Wake for Flag {
158    fn wake(self: Arc<Self>) {
159        self.0.store(true, Ordering::Release);
160    }
161
162    fn wake_by_ref(self: &Arc<Self>) {
163        self.0.store(true, Ordering::Release);
164    }
165}
166
167/// The read side of a [`flag_waker`]: whether the waker was woken since the
168/// last [`take`](WakeFlag::take).
169#[derive(Debug)]
170pub struct WakeFlag(Arc<Flag>);
171
172impl WakeFlag {
173    /// Return `true` if the waker was woken, by value or by reference, from
174    /// any thread, since the flag was created or last taken, and clear the
175    /// flag.
176    pub fn take(&self) -> bool {
177        self.0 .0.swap(false, Ordering::Acquire)
178    }
179}
180
181/// A waker whose wake sets a flag, for a frame loop that polls again in the
182/// same frame when the future asks for it.
183///
184/// The flag starts clear. A wake of the waker or of any clone of it, by value
185/// or by reference, from any thread, sets it; [`WakeFlag::take`] reads and
186/// clears it. Nothing is polled by the wake itself. The waker is built as
187/// `Waker::from(Arc<Flag>)`, so each call allocates one `Arc`, which the
188/// waker and the flag share: create the pair once per loop.
189///
190/// The frame-loop pattern: poll, then poll again while the flag was set, up to
191/// the loop's own limit of polls per frame, so that a future that wakes itself
192/// on every poll cannot hold the frame. A wake that arrived between two frames
193/// leaves the flag set, which costs at most one extra poll in the next frame.
194///
195/// ```
196/// use std::future::Future;
197/// use std::task::Context;
198///
199/// const POLLS_PER_FRAME: usize = 8;
200///
201/// let (waker, woken) = ridl_rt::task::flag_waker();
202/// let mut cx = Context::from_waker(&waker);
203/// let mut fut = std::pin::pin!(std::future::ready(1));
204///
205/// // One frame.
206/// for _ in 0..POLLS_PER_FRAME {
207///     if fut.as_mut().poll(&mut cx).is_ready() || !woken.take() {
208///         break;
209///     }
210/// }
211///
212/// waker.wake_by_ref();
213/// assert!(woken.take());
214/// assert!(!woken.take());
215/// ```
216pub fn flag_waker() -> (Waker, WakeFlag) {
217    let flag = Arc::new(Flag(AtomicBool::new(false)));
218    (Waker::from(Arc::clone(&flag)), WakeFlag(flag))
219}