standard-plugin-sdk 0.1.1

Write Standard Code plugins in Rust: wasm components against standard:plugin@2.0.0
Documentation
//! A single-threaded executor driven by the component's `drive` export.
//!
//! Tasks are woken by timers ([`Sleep`]), events ([`NextEvent`]) and, with
//! the `wasi` feature, WASI pollables (`daemon::io`). `drive(now)` fires due
//! timers, polls woken tasks until none is woken (bounded, so one busy task
//! cannot hold the call), and reports when it next needs to run.

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;

/// Events kept while nothing awaits them.
pub const EVENT_QUEUE: usize = 1024;
/// Rounds of polling woken tasks per `drive` call before yielding to the
/// host with an immediate re-poll.
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>,
    /// Tasks waiting on WASI pollables (`daemon::io`).
    #[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();
    }
}

/// Runs the tasks once for `now_ms`; returns when the executor next wants
/// `drive` (none: only on an event), and drops finished tasks.
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())
}

/// Forgets every waiter and queued event (at deactivate).
pub(crate) fn reset() {
    reactor(|reactor| *reactor = Reactor::default());
}

/// Completes at an instant on the daemon's clock. Dropped before it
/// completes (it lost a race with an event), it takes its timer with it, so
/// the host is not asked to drive the plugin at an instant nobody waits for.
#[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 {
                // Polled again while pending: keep one timer per waiter.
                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(&registered)));
        });
    }
}

/// The next event, or `None` once the daemon's clock reaches a deadline
/// ([`Context::next_event_until`](super::Context::next_event_until)).
#[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
        }
    }
}

/// The next event the host delivers.
#[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
            }
        })
    }
}