standard-plugin-sdk 0.1.0

Write Standard Code plugins in Rust: wasm components against standard:plugin@2.0.0
Documentation
//! Awaiting WASI pollables on the daemon executor (`wasi` feature).
//!
//! A task awaits [`wait`] on any `wasi:io/poll` pollable (a stream, an HTTP
//! response, a timer). While tasks wait on I/O and nothing else is ready,
//! `drive` blocks in `wasi:io/poll.poll` on their pollables for at most
//! [`IO_SLICE_MS`] (or until the next timer), then asks the host to call
//! again at once, so events still interleave with long I/O.

use alloc::vec::Vec;
use core::future::Future;
use core::mem::ManuallyDrop;
use core::pin::Pin;
use core::task::{Context, Poll};

use wasip2::io::poll::Pollable;

use super::executor::reactor;

/// The longest `drive` blocks on pending I/O.
pub const IO_SLICE_MS: u64 = 10;

/// Completes when `pollable` is ready.
pub fn wait(pollable: &Pollable) -> Wait<'_> {
    Wait { pollable }
}

#[must_use = "futures do nothing unless awaited"]
pub struct Wait<'a> {
    pollable: &'a Pollable,
}

impl Future for Wait<'_> {
    type Output = ();

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
        let handle = self.pollable.handle();
        if self.pollable.ready() {
            reactor(|reactor| reactor.io.retain(|(h, _)| *h != handle));
            return Poll::Ready(());
        }
        reactor(|reactor| {
            reactor.io.retain(|(h, _)| *h != handle);
            reactor.io.push((handle, cx.waker().clone()));
        });
        Poll::Pending
    }
}

impl Drop for Wait<'_> {
    fn drop(&mut self) {
        let handle = self.pollable.handle();
        reactor(|reactor| reactor.io.retain(|(h, _)| *h != handle));
    }
}

/// Completes when any of `pollables` is ready, with its index.
pub fn wait_any<'a>(pollables: &'a [&'a Pollable]) -> WaitAny<'a> {
    WaitAny { pollables }
}

#[must_use = "futures do nothing unless awaited"]
pub struct WaitAny<'a> {
    pollables: &'a [&'a Pollable],
}

impl WaitAny<'_> {
    fn forget(&self) {
        let handles: Vec<u32> = self
            .pollables
            .iter()
            .map(|pollable| pollable.handle())
            .collect();
        reactor(|reactor| reactor.io.retain(|(h, _)| !handles.contains(h)));
    }
}

impl Future for WaitAny<'_> {
    type Output = usize;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<usize> {
        self.forget();
        if let Some(index) = self.pollables.iter().position(|pollable| pollable.ready()) {
            return Poll::Ready(index);
        }
        reactor(|reactor| {
            for pollable in self.pollables {
                reactor.io.push((pollable.handle(), cx.waker().clone()));
            }
        });
        Poll::Pending
    }
}

impl Drop for WaitAny<'_> {
    fn drop(&mut self) {
        self.forget();
    }
}

/// A pollable borrowed by handle: the awaiting future owns it.
fn borrowed(handle: u32) -> ManuallyDrop<Pollable> {
    // SAFETY: the handle belongs to a pollable a pending `Wait` borrows;
    // `Wait` deregisters it before that pollable can be dropped, and the
    // `ManuallyDrop` never closes it.
    ManuallyDrop::new(unsafe { Pollable::from_handle(handle) })
}

pub(crate) fn wake_ready() {
    let ready: Vec<_> = reactor(|reactor| {
        let mut ready = Vec::new();
        reactor.io.retain(|(handle, waker)| {
            if borrowed(*handle).ready() {
                ready.push(waker.clone());
                false
            } else {
                true
            }
        });
        ready
    });
    for waker in ready {
        waker.wake();
    }
}

pub(crate) fn pending() -> bool {
    reactor(|reactor| !reactor.io.is_empty())
}

/// Blocks on the pending pollables for at most a slice, then asks for an
/// immediate re-poll.
pub(crate) fn recheck_at(now_ms: u64) -> u64 {
    let (handles, next_timer) = reactor(|reactor| {
        (
            reactor
                .io
                .iter()
                .map(|(handle, _)| *handle)
                .collect::<Vec<_>>(),
            reactor.timers.iter().map(|(at, _)| *at).min(),
        )
    });
    let slice = next_timer
        .map_or(IO_SLICE_MS, |at| at.saturating_sub(now_ms))
        .min(IO_SLICE_MS);
    let timeout = wasip2::clocks::monotonic_clock::subscribe_duration(slice * 1_000_000);
    let pollables: Vec<ManuallyDrop<Pollable>> = handles.into_iter().map(borrowed).collect();
    let mut list: Vec<&Pollable> = pollables.iter().map(|pollable| &**pollable).collect();
    list.push(&timeout);
    let _ = wasip2::io::poll::poll(&list);
    now_ms
}