Skip to main content

standard_plugin/daemon/
executor.rs

1//! A single-threaded executor driven by the component's `drive` export.
2//!
3//! Tasks are woken by timers ([`Sleep`]), events ([`NextEvent`]) and, with
4//! the `wasi` feature, WASI pollables (`daemon::io`). `drive(now)` fires due
5//! timers, polls woken tasks until none is woken (bounded, so one busy task
6//! cannot hold the call), and reports when it next needs to run.
7
8use alloc::boxed::Box;
9use alloc::collections::VecDeque;
10use alloc::sync::Arc;
11use alloc::task::Wake;
12use alloc::vec::Vec;
13use core::future::Future;
14use core::pin::Pin;
15use core::sync::atomic::{AtomicBool, Ordering};
16use core::task::{Context as TaskContext, Poll, Waker};
17
18use crate::local::instance_local;
19use crate::ui_runtime::Event;
20
21/// Events kept while nothing awaits them.
22pub const EVENT_QUEUE: usize = 1024;
23/// Rounds of polling woken tasks per `drive` call before yielding to the
24/// host with an immediate re-poll.
25const ROUNDS: usize = 64;
26
27pub(crate) type BoxFuture = Pin<Box<dyn Future<Output = ()>>>;
28
29pub(crate) struct Flag(AtomicBool);
30
31impl Wake for Flag {
32    fn wake(self: Arc<Self>) {
33        self.0.store(true, Ordering::Relaxed);
34    }
35
36    fn wake_by_ref(self: &Arc<Self>) {
37        self.0.store(true, Ordering::Relaxed);
38    }
39}
40
41pub(crate) struct Task {
42    future: BoxFuture,
43    flag: Arc<Flag>,
44}
45
46impl Task {
47    pub(crate) fn new(future: BoxFuture) -> Self {
48        Self {
49            future,
50            flag: Arc::new(Flag(AtomicBool::new(true))),
51        }
52    }
53}
54
55#[derive(Default)]
56pub(crate) struct Reactor {
57    now_ms: u64,
58    events: VecDeque<Event>,
59    event_waiters: Vec<Waker>,
60    pub(crate) timers: Vec<(u64, Waker)>,
61    spawned: Vec<BoxFuture>,
62    /// Tasks waiting on WASI pollables (`daemon::io`).
63    #[cfg(all(feature = "wasi", target_arch = "wasm32"))]
64    pub(crate) io: Vec<(u32, Waker)>,
65}
66
67instance_local! {
68    pub(crate) fn reactor() -> Reactor = Reactor::default();
69}
70
71pub(crate) fn now_ms() -> u64 {
72    reactor(|reactor| reactor.now_ms)
73}
74
75pub(crate) fn spawn(task: impl Future<Output = ()> + 'static) {
76    reactor(|reactor| reactor.spawned.push(Box::pin(task)));
77}
78
79pub(crate) fn push_event(event: Event) {
80    let waiters = reactor(|reactor| {
81        if reactor.events.len() == EVENT_QUEUE {
82            reactor.events.pop_front();
83        }
84        reactor.events.push_back(event);
85        core::mem::take(&mut reactor.event_waiters)
86    });
87    for waker in waiters {
88        waker.wake();
89    }
90}
91
92/// Runs the tasks once for `now_ms`; returns when the executor next wants
93/// `drive` (none: only on an event), and drops finished tasks.
94pub(crate) fn run(tasks: &mut Vec<Task>, now_ms: u64) -> Option<u64> {
95    let due = reactor(|reactor| {
96        reactor.now_ms = reactor.now_ms.max(now_ms);
97        let now = reactor.now_ms;
98        let mut due = Vec::new();
99        reactor.timers.retain(|(at, waker)| {
100            if *at <= now {
101                due.push(waker.clone());
102                false
103            } else {
104                true
105            }
106        });
107        due
108    });
109    for waker in due {
110        waker.wake();
111    }
112    #[cfg(all(feature = "wasi", target_arch = "wasm32"))]
113    super::io::wake_ready();
114    for _ in 0..ROUNDS {
115        tasks.extend(
116            reactor(|reactor| core::mem::take(&mut reactor.spawned))
117                .into_iter()
118                .map(Task::new),
119        );
120        let mut progressed = false;
121        tasks.retain_mut(|task| {
122            if !task.flag.0.swap(false, Ordering::Relaxed) {
123                return true;
124            }
125            progressed = true;
126            let waker = Waker::from(task.flag.clone());
127            let mut cx = TaskContext::from_waker(&waker);
128            task.future.as_mut().poll(&mut cx).is_pending()
129        });
130        if !progressed && reactor(|reactor| reactor.spawned.is_empty()) {
131            break;
132        }
133    }
134    let woken = tasks.iter().any(|task| task.flag.0.load(Ordering::Relaxed))
135        || reactor(|reactor| !reactor.spawned.is_empty());
136    let now = reactor(|reactor| reactor.now_ms);
137    if woken {
138        return Some(now);
139    }
140    #[cfg(all(feature = "wasi", target_arch = "wasm32"))]
141    if super::io::pending() {
142        return Some(super::io::recheck_at(now));
143    }
144    reactor(|reactor| reactor.timers.iter().map(|(at, _)| *at).min())
145}
146
147/// Forgets every waiter and queued event (at deactivate).
148pub(crate) fn reset() {
149    reactor(|reactor| *reactor = Reactor::default());
150}
151
152/// Completes at an instant on the daemon's clock. Dropped before it
153/// completes (it lost a race with an event), it takes its timer with it, so
154/// the host is not asked to drive the plugin at an instant nobody waits for.
155#[derive(Debug)]
156#[must_use = "futures do nothing unless awaited"]
157pub struct Sleep {
158    at_ms: u64,
159    registered: Option<Waker>,
160}
161
162impl Sleep {
163    pub(crate) fn until(at_ms: u64) -> Self {
164        Self {
165            at_ms,
166            registered: None,
167        }
168    }
169}
170
171impl Future for Sleep {
172    type Output = ();
173
174    fn poll(mut self: Pin<&mut Self>, cx: &mut TaskContext<'_>) -> Poll<()> {
175        let at = self.at_ms;
176        let ready = reactor(|reactor| {
177            if reactor.now_ms >= at {
178                true
179            } else {
180                // Polled again while pending: keep one timer per waiter.
181                reactor
182                    .timers
183                    .retain(|(when, waker)| !(*when == at && waker.will_wake(cx.waker())));
184                reactor.timers.push((at, cx.waker().clone()));
185                false
186            }
187        });
188        if ready {
189            self.registered = None;
190            Poll::Ready(())
191        } else {
192            self.registered = Some(cx.waker().clone());
193            Poll::Pending
194        }
195    }
196}
197
198impl Drop for Sleep {
199    fn drop(&mut self) {
200        let Some(registered) = self.registered.take() else {
201            return;
202        };
203        let at = self.at_ms;
204        reactor(|reactor| {
205            reactor
206                .timers
207                .retain(|(when, waker)| !(*when == at && waker.will_wake(&registered)));
208        });
209    }
210}
211
212/// The next event, or `None` once the daemon's clock reaches a deadline
213/// ([`Context::next_event_until`](super::Context::next_event_until)).
214#[derive(Debug)]
215#[must_use = "futures do nothing unless awaited"]
216pub struct NextEventUntil {
217    event: NextEvent,
218    deadline: Option<Sleep>,
219}
220
221impl NextEventUntil {
222    pub(crate) fn new(at_ms: Option<u64>) -> Self {
223        Self {
224            event: NextEvent::new(),
225            deadline: at_ms.map(Sleep::until),
226        }
227    }
228}
229
230impl Future for NextEventUntil {
231    type Output = Option<Event>;
232
233    fn poll(mut self: Pin<&mut Self>, cx: &mut TaskContext<'_>) -> Poll<Option<Event>> {
234        if let Poll::Ready(event) = Pin::new(&mut self.event).poll(cx) {
235            return Poll::Ready(Some(event));
236        }
237        let expired = self
238            .deadline
239            .as_mut()
240            .is_some_and(|deadline| Pin::new(deadline).poll(cx).is_ready());
241        if expired {
242            Poll::Ready(None)
243        } else {
244            Poll::Pending
245        }
246    }
247}
248
249/// The next event the host delivers.
250#[derive(Debug, Default)]
251#[must_use = "futures do nothing unless awaited"]
252pub struct NextEvent {
253    _private: (),
254}
255
256impl NextEvent {
257    pub(crate) fn new() -> Self {
258        Self::default()
259    }
260}
261
262impl Future for NextEvent {
263    type Output = Event;
264
265    fn poll(self: Pin<&mut Self>, cx: &mut TaskContext<'_>) -> Poll<Event> {
266        reactor(|reactor| match reactor.events.pop_front() {
267            Some(event) => Poll::Ready(event),
268            None => {
269                if !reactor
270                    .event_waiters
271                    .iter()
272                    .any(|waker| waker.will_wake(cx.waker()))
273                {
274                    reactor.event_waiters.push(cx.waker().clone());
275                }
276                Poll::Pending
277            }
278        })
279    }
280}