standard_plugin/daemon/
executor.rs1use 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
21pub const EVENT_QUEUE: usize = 1024;
23const 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 #[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
92pub(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
147pub(crate) fn reset() {
149 reactor(|reactor| *reactor = Reactor::default());
150}
151
152#[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 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(®istered)));
208 });
209 }
210}
211
212#[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#[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}