use std::cell::{Cell, RefCell};
use std::future::Future;
use std::pin::Pin;
use std::rc::Rc;
use std::sync::Arc;
use std::sync::OnceLock;
use std::sync::atomic::{AtomicBool, Ordering};
use std::task::{Context, Poll, Wake, Waker};
use teksilo_core::{
AppEventPoster, AsyncCompletionHandle, AsyncCompletionPayload, EventContext, TeksiloWindowId,
};
struct AsyncWake;
struct ExecWaker {
woken: AtomicBool,
poster: OnceLock<Arc<dyn AppEventPoster>>,
}
impl Wake for ExecWaker {
fn wake(self: Arc<Self>) {
self.wake_by_ref();
}
fn wake_by_ref(self: &Arc<Self>) {
self.woken.store(true, Ordering::SeqCst);
if let Some(poster) = self.poster.get() {
poster.post_external(Box::new(AsyncWake));
}
}
}
type BoxFuture = Pin<Box<dyn Future<Output = ()>>>;
struct Task {
future: BoxFuture,
cancelled: Rc<Cell<bool>>,
}
#[must_use = "dropping the TaskHandle cancels the task — call `.detach()` to let it keep running"]
pub struct TaskHandle {
cancelled: Option<Rc<Cell<bool>>>,
}
impl TaskHandle {
pub fn detach(mut self) {
self.cancelled = None;
}
}
impl Drop for TaskHandle {
fn drop(&mut self) {
if let Some(flag) = &self.cancelled {
flag.set(true);
}
}
}
struct ExecInner {
tasks: RefCell<Vec<Task>>,
spawn_queue: RefCell<Vec<Task>>,
wake: Arc<ExecWaker>,
waker: Waker,
poll_source: Rc<Cell<bool>>,
completions: AsyncCompletionHandle,
ticking: Cell<bool>,
}
struct TickGuard<'a>(&'a Cell<bool>);
impl Drop for TickGuard<'_> {
fn drop(&mut self) {
self.0.set(false);
}
}
impl ExecInner {
fn flush_spawns(&self) {
let mut queued = self.spawn_queue.borrow_mut();
if !queued.is_empty() {
self.tasks.borrow_mut().append(&mut queued);
}
}
fn spawn(&self, future: BoxFuture) -> TaskHandle {
let cancelled = Rc::new(Cell::new(false));
self.spawn_queue.borrow_mut().push(Task {
future,
cancelled: cancelled.clone(),
});
self.wake.woken.store(true, Ordering::SeqCst);
TaskHandle {
cancelled: Some(cancelled),
}
}
fn tick(&self) -> bool {
if self.ticking.get() {
return false;
}
self.ticking.set(true);
let _reset = TickGuard(&self.ticking);
self.flush_spawns();
if !self.wake.woken.swap(false, Ordering::SeqCst) {
self.poll_source.set(false);
return false;
}
let mut cx = Context::from_waker(&self.waker);
let taken = std::mem::take(&mut *self.tasks.borrow_mut());
let mut survivors = Vec::with_capacity(taken.len());
for mut task in taken {
if task.cancelled.get() {
continue; }
match task.future.as_mut().poll(&mut cx) {
Poll::Ready(()) => {} Poll::Pending => survivors.push(task),
}
}
*self.tasks.borrow_mut() = survivors;
self.flush_spawns();
self.poll_source.set(self.wake.woken.load(Ordering::SeqCst));
true
}
}
#[derive(Clone)]
pub struct AsyncRuntimeHandle {
inner: Rc<ExecInner>,
}
impl Default for AsyncRuntimeHandle {
fn default() -> Self {
Self::new()
}
}
impl AsyncRuntimeHandle {
pub fn new() -> Self {
let wake = Arc::new(ExecWaker {
woken: AtomicBool::new(false),
poster: OnceLock::new(),
});
let waker = Waker::from(wake.clone());
Self {
inner: Rc::new(ExecInner {
tasks: RefCell::new(Vec::new()),
spawn_queue: RefCell::new(Vec::new()),
wake,
waker,
poll_source: Rc::new(Cell::new(false)),
completions: AsyncCompletionHandle::new(),
ticking: Cell::new(false),
}),
}
}
pub fn poll_source(&self) -> Rc<Cell<bool>> {
self.inner.poll_source.clone()
}
pub fn completions(&self) -> AsyncCompletionHandle {
self.inner.completions.clone()
}
pub fn set_poster(&self, poster: Arc<dyn AppEventPoster>) {
let _ = self.inner.wake.poster.set(poster);
}
pub fn tick(&self) -> bool {
self.inner.tick()
}
pub fn spawn_local(&self, future: impl Future<Output = ()> + 'static) -> TaskHandle {
self.inner.spawn(Box::pin(future))
}
pub fn spawn_local_with<R: 'static>(
&self,
window_id: TeksiloWindowId,
future: impl Future<Output = R> + 'static,
on_complete: impl FnOnce(R, &mut EventContext) + 'static,
) -> TaskHandle {
let completions = self.inner.completions.clone();
let poster = self.inner.wake.poster.get().cloned();
let wrapper = async move {
let result = future.await;
let callback: Box<dyn FnOnce(&mut EventContext)> =
Box::new(move |ctx| on_complete(result, ctx));
let id = completions.register(window_id, callback);
if let Some(poster) = &poster {
poster.post_external(Box::new(AsyncCompletionPayload { id, window_id }));
}
};
self.inner.spawn(Box::pin(wrapper))
}
}