use alloc::boxed::Box;
use alloc::collections::VecDeque;
use alloc::sync::Arc;
use alloc::task::Wake;
use alloc::vec::Vec;
use core::future::Future;
use core::pin::Pin;
use core::sync::atomic::{AtomicBool, Ordering};
use core::task::{Context as TaskContext, Poll, Waker};
use crate::local::instance_local;
use crate::ui_runtime::Event;
pub const EVENT_QUEUE: usize = 1024;
const ROUNDS: usize = 64;
pub(crate) type BoxFuture = Pin<Box<dyn Future<Output = ()>>>;
pub(crate) struct Flag(AtomicBool);
impl Wake for Flag {
fn wake(self: Arc<Self>) {
self.0.store(true, Ordering::Relaxed);
}
fn wake_by_ref(self: &Arc<Self>) {
self.0.store(true, Ordering::Relaxed);
}
}
pub(crate) struct Task {
future: BoxFuture,
flag: Arc<Flag>,
}
impl Task {
pub(crate) fn new(future: BoxFuture) -> Self {
Self {
future,
flag: Arc::new(Flag(AtomicBool::new(true))),
}
}
}
#[derive(Default)]
pub(crate) struct Reactor {
now_ms: u64,
events: VecDeque<Event>,
event_waiters: Vec<Waker>,
pub(crate) timers: Vec<(u64, Waker)>,
spawned: Vec<BoxFuture>,
#[cfg(all(feature = "wasi", target_arch = "wasm32"))]
pub(crate) io: Vec<(u32, Waker)>,
}
instance_local! {
pub(crate) fn reactor() -> Reactor = Reactor::default();
}
pub(crate) fn now_ms() -> u64 {
reactor(|reactor| reactor.now_ms)
}
pub(crate) fn spawn(task: impl Future<Output = ()> + 'static) {
reactor(|reactor| reactor.spawned.push(Box::pin(task)));
}
pub(crate) fn push_event(event: Event) {
let waiters = reactor(|reactor| {
if reactor.events.len() == EVENT_QUEUE {
reactor.events.pop_front();
}
reactor.events.push_back(event);
core::mem::take(&mut reactor.event_waiters)
});
for waker in waiters {
waker.wake();
}
}
pub(crate) fn run(tasks: &mut Vec<Task>, now_ms: u64) -> Option<u64> {
let due = reactor(|reactor| {
reactor.now_ms = reactor.now_ms.max(now_ms);
let now = reactor.now_ms;
let mut due = Vec::new();
reactor.timers.retain(|(at, waker)| {
if *at <= now {
due.push(waker.clone());
false
} else {
true
}
});
due
});
for waker in due {
waker.wake();
}
#[cfg(all(feature = "wasi", target_arch = "wasm32"))]
super::io::wake_ready();
for _ in 0..ROUNDS {
tasks.extend(
reactor(|reactor| core::mem::take(&mut reactor.spawned))
.into_iter()
.map(Task::new),
);
let mut progressed = false;
tasks.retain_mut(|task| {
if !task.flag.0.swap(false, Ordering::Relaxed) {
return true;
}
progressed = true;
let waker = Waker::from(task.flag.clone());
let mut cx = TaskContext::from_waker(&waker);
task.future.as_mut().poll(&mut cx).is_pending()
});
if !progressed && reactor(|reactor| reactor.spawned.is_empty()) {
break;
}
}
let woken = tasks.iter().any(|task| task.flag.0.load(Ordering::Relaxed))
|| reactor(|reactor| !reactor.spawned.is_empty());
let now = reactor(|reactor| reactor.now_ms);
if woken {
return Some(now);
}
#[cfg(all(feature = "wasi", target_arch = "wasm32"))]
if super::io::pending() {
return Some(super::io::recheck_at(now));
}
reactor(|reactor| reactor.timers.iter().map(|(at, _)| *at).min())
}
pub(crate) fn reset() {
reactor(|reactor| *reactor = Reactor::default());
}
#[derive(Debug)]
#[must_use = "futures do nothing unless awaited"]
pub struct Sleep {
at_ms: u64,
registered: Option<Waker>,
}
impl Sleep {
pub(crate) fn until(at_ms: u64) -> Self {
Self {
at_ms,
registered: None,
}
}
}
impl Future for Sleep {
type Output = ();
fn poll(mut self: Pin<&mut Self>, cx: &mut TaskContext<'_>) -> Poll<()> {
let at = self.at_ms;
let ready = reactor(|reactor| {
if reactor.now_ms >= at {
true
} else {
reactor
.timers
.retain(|(when, waker)| !(*when == at && waker.will_wake(cx.waker())));
reactor.timers.push((at, cx.waker().clone()));
false
}
});
if ready {
self.registered = None;
Poll::Ready(())
} else {
self.registered = Some(cx.waker().clone());
Poll::Pending
}
}
}
impl Drop for Sleep {
fn drop(&mut self) {
let Some(registered) = self.registered.take() else {
return;
};
let at = self.at_ms;
reactor(|reactor| {
reactor
.timers
.retain(|(when, waker)| !(*when == at && waker.will_wake(®istered)));
});
}
}
#[derive(Debug)]
#[must_use = "futures do nothing unless awaited"]
pub struct NextEventUntil {
event: NextEvent,
deadline: Option<Sleep>,
}
impl NextEventUntil {
pub(crate) fn new(at_ms: Option<u64>) -> Self {
Self {
event: NextEvent::new(),
deadline: at_ms.map(Sleep::until),
}
}
}
impl Future for NextEventUntil {
type Output = Option<Event>;
fn poll(mut self: Pin<&mut Self>, cx: &mut TaskContext<'_>) -> Poll<Option<Event>> {
if let Poll::Ready(event) = Pin::new(&mut self.event).poll(cx) {
return Poll::Ready(Some(event));
}
let expired = self
.deadline
.as_mut()
.is_some_and(|deadline| Pin::new(deadline).poll(cx).is_ready());
if expired {
Poll::Ready(None)
} else {
Poll::Pending
}
}
}
#[derive(Debug, Default)]
#[must_use = "futures do nothing unless awaited"]
pub struct NextEvent {
_private: (),
}
impl NextEvent {
pub(crate) fn new() -> Self {
Self::default()
}
}
impl Future for NextEvent {
type Output = Event;
fn poll(self: Pin<&mut Self>, cx: &mut TaskContext<'_>) -> Poll<Event> {
reactor(|reactor| match reactor.events.pop_front() {
Some(event) => Poll::Ready(event),
None => {
if !reactor
.event_waiters
.iter()
.any(|waker| waker.will_wake(cx.waker()))
{
reactor.event_waiters.push(cx.waker().clone());
}
Poll::Pending
}
})
}
}