emscripten-futures 0.7.6

A local async executor and awaitable Emscripten operations for Rust WebAssembly applications
Documentation
use emscripten_rs_sys::em_asm::{SignatureBuilder, emscripten_asm_const_int, js_asm};
use futures::{
    FutureExt,
    future::{CatchUnwind, LocalBoxFuture},
    task::{ArcWake, waker_ref},
};
use send_wrapper::SendWrapper;
use std::{
    ffi::c_void,
    panic::AssertUnwindSafe,
    sync::Arc,
    task::{Context, Poll},
};

use crate::task::spawn::{keepalive::EmscriptenKeepalive, microtask_waker::MicroTaskWaker};

#[cfg(test)]
mod tests;

mod keepalive {
    use emscripten_rs_sys::{emscripten_runtime_keepalive_pop, emscripten_runtime_keepalive_push};

    pub struct EmscriptenKeepalive {}

    impl Default for EmscriptenKeepalive {
        fn default() -> Self {
            unsafe { emscripten_runtime_keepalive_push() };
            Self {}
        }
    }

    impl Drop for EmscriptenKeepalive {
        fn drop(&mut self) {
            unsafe { emscripten_runtime_keepalive_pop() };
        }
    }
}

struct SpawnedTask {
    future: CatchUnwind<AssertUnwindSafe<LocalBoxFuture<'static, ()>>>,
    mwaker: Arc<MicroTaskWakerWrapper>,
    _keepalive: EmscriptenKeepalive,
}

// Records wakes during polling; schedules sleeping tasks with queueMicrotask.
struct MicroTaskWakerWrapper {
    wrapper: SendWrapper<microtask_waker::MicroTaskWaker>,
}

mod microtask_waker {
    use crate::task::spawn::{SpawnedTask, schedule};
    use bitflags::bitflags;
    use std::cell::Cell;

    bitflags! {
        #[derive(Clone, Copy, PartialEq, Eq)]
        struct TaskStateFlags : u8 {
            const NOTIFIED = 1 << 0;
            const SLEEPING = 1 << 1;
        }
    }

    pub struct MicroTaskWaker {
        task: Cell<*mut SpawnedTask>,
        task_state: Cell<TaskStateFlags>,
    }

    impl Default for MicroTaskWaker {
        fn default() -> Self {
            Self {
                task: Default::default(),
                task_state: Cell::new(TaskStateFlags::all()),
            }
        }
    }

    impl MicroTaskWaker {
        #[inline]
        pub fn init(&self, task: *mut SpawnedTask) {
            self.task.replace(task);
        }

        #[inline]
        pub fn start_poll(&self) {
            self.task_state.set(TaskStateFlags::empty());
        }

        #[inline]
        pub fn try_wake(&self) {
            self.update_state::<{ TaskStateFlags::NOTIFIED.bits() }, { TaskStateFlags::SLEEPING.bits() }>();
        }

        #[inline]
        pub fn try_sleep(&self) {
            self.update_state::<{ TaskStateFlags::SLEEPING.bits() }, { TaskStateFlags::NOTIFIED.bits() }>();
        }

        #[inline]
        fn update_state<const SET: u8, const SCHEDULE_IF: u8>(&self) {
            let current_state = self.task_state.get().bits();
            self.task_state
                .set(TaskStateFlags::from_bits_retain(current_state | SET));
            // Decide whether to schedule from the state before setting the flag.
            if current_state == SCHEDULE_IF {
                schedule(self.task.get());
            }
        }
    }
}

impl ArcWake for MicroTaskWakerWrapper {
    fn wake_by_ref(arc_self: &Arc<Self>) {
        arc_self.wrapper.try_wake();
    }
}

#[unsafe(no_mangle)]
pub extern "C" fn emscripten_futures_poll_task(task: *mut c_void) {
    let task_ptr = task.cast::<SpawnedTask>();
    // SAFETY: at most one callback is queued, and the Box stays allocated
    // until its final poll. Wakes during polling cannot queue a callback.
    let mut task = unsafe { Box::from_raw(task_ptr) };
    let state = &*task.mwaker.wrapper;

    state.start_poll();
    let result = {
        let waker = waker_ref(&task.mwaker);
        let mut cx = Context::from_waker(&waker);
        task.future.poll_unpin(&mut cx)
    };

    match result {
        Poll::Ready(_) => {
            // Leave SLEEPING clear so retained wakers cannot schedule the
            // freed task, including wakes from the future's destructor.
            drop(task);
        }
        Poll::Pending => {
            state.try_sleep();
            let _ = Box::into_raw(task);
        }
    }
}

fn schedule(raw_task: *mut SpawnedTask) {
    js_asm! {
        |raw_task| {
            queueMicrotask(() => {
                _emscripten_futures_poll_task(raw_task);
            });
        }
    }
}

/// Spawns a detached future on the calling thread (requires the `spawn` feature).
///
/// The first poll is deferred to `queueMicrotask`. Each callback polls once;
/// repeated wakes are coalesced. A wake during polling schedules another
/// microtask only if the future returns `Pending`. Wakes after completion on
/// the originating thread are ignored. No running [`LocalPool`](crate::executor::LocalPool)
/// is required. The future may hold non-`Send` values and keeps the Emscripten
/// runtime alive until it completes. There is no cancellation handle.
///
/// Cross-thread wakeups are not supported and panic before accessing task state.
/// Wakers may be cloned and dropped on any thread, but must only be woken on
/// the originating thread.
/// Await asynchronous operations inside the future rather than calling
/// `block_on` or otherwise suspending a poll with JSPI.
///
/// # Panics
///
/// With `panic = "unwind"`, panics while polling are caught and the failed future
/// is dropped, releasing its runtime keepalive. The panic hook still runs, and
/// other tasks can continue. Retained wakers cannot reschedule the failed task.
/// Panics while dropping the future or panic payload are not caught, and
/// `panic = "abort"` builds cannot recover from panics.
///
/// To receive a result or propagate a polling panic to an awaiter, use
/// [`FutureExt::remote_handle`] and spawn its `Remote` future. Keep the handle on
/// the originating thread: dropping it requests cancellation by waking the task.
/// Calling the handle's `forget` method instead lets the task continue detached.
pub fn spawn_local(f: impl Future<Output = ()> + 'static) {
    spawn_local_boxed(f.boxed_local());
}

fn spawn_local_boxed(f: LocalBoxFuture<'static, ()>) {
    let waker_state = SendWrapper::new(MicroTaskWaker::default());
    let mwaker = Arc::new(MicroTaskWakerWrapper {
        wrapper: waker_state,
    });
    let mut spawned_task: Box<SpawnedTask> = Box::new(SpawnedTask {
        future: AssertUnwindSafe(f).catch_unwind(),
        mwaker,
        _keepalive: Default::default(),
    });
    let task_ptr = Box::as_mut_ptr(&mut spawned_task);
    spawned_task.mwaker.wrapper.init(task_ptr);
    schedule(Box::into_raw(spawned_task));
}