use std::cell::{Cell, RefCell};
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll, Wake, Waker};
use any_spawner::{CustomExecutor, PinnedFuture, PinnedLocalFuture};
use futures::executor::{LocalPool, LocalSpawner};
use futures::task::LocalSpawnExt;
use tokio::runtime::Handle;
#[allow(clippy::missing_const_for_thread_local)]
mod tls {
use super::*;
thread_local! {
pub(super) static LOCAL_POOL: RefCell<LocalPool> = RefCell::new(LocalPool::new());
pub(super) static LOCAL_SPAWNER: LocalSpawner =
LOCAL_POOL.with(|pool| pool.borrow().spawner());
pub(super) static IS_UI_THREAD: Cell<bool> = const { Cell::new(false) };
}
}
use tls::{IS_UI_THREAD, LOCAL_POOL, LOCAL_SPAWNER};
pub(crate) fn init_ui_thread() {
IS_UI_THREAD.with(|flag| flag.set(true));
}
pub(crate) fn is_ui_thread() -> bool {
IS_UI_THREAD.with(Cell::get)
}
pub(crate) fn run_until_stalled() {
LOCAL_POOL.with(|pool| {
if let Ok(mut pool) = pool.try_borrow_mut() {
pool.run_until_stalled();
}
});
}
struct CompositeWaker {
inner: Waker,
}
impl CompositeWaker {
fn wake_shell() {
if let Some(rt) = crate::ReactiveRuntime::get() {
rt.wake();
}
}
}
impl Wake for CompositeWaker {
fn wake(self: Arc<Self>) {
Self::wake_shell();
self.inner.wake_by_ref();
}
fn wake_by_ref(self: &Arc<Self>) {
Self::wake_shell();
self.inner.wake_by_ref();
}
}
struct WakeBridge {
inner: PinnedLocalFuture<()>,
}
impl Future for WakeBridge {
type Output = ();
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
let this = self.get_mut();
let composite = Waker::from(Arc::new(CompositeWaker {
inner: cx.waker().clone(),
}));
let mut inner_cx = Context::from_waker(&composite);
this.inner.as_mut().poll(&mut inner_cx)
}
}
pub(crate) struct ForgeExecutor {
handle: Handle,
}
impl ForgeExecutor {
pub(crate) fn new(handle: Handle) -> Self {
Self { handle }
}
}
impl CustomExecutor for ForgeExecutor {
fn spawn(&self, fut: PinnedFuture<()>) {
self.handle.spawn(fut);
}
fn spawn_local(&self, fut: PinnedLocalFuture<()>) {
if !is_ui_thread() {
panic!(
"frust-reactive: Executor::spawn_local was called off the UI \
thread. `!Send` local futures can only be spawned on the UI \
thread (the one `ReactiveRuntime::init` ran on). This is a \
wiring bug: route the work through `Executor::spawn` instead, or \
hand it back to the UI thread before spawning it locally."
);
}
let fut = WakeBridge { inner: fut };
LOCAL_SPAWNER.with(|spawner| {
spawner
.spawn_local(fut)
.expect("frust-reactive: UI-thread local task queue rejected a future");
});
if let Some(rt) = crate::ReactiveRuntime::get() {
rt.wake();
}
}
fn poll_local(&self) {
let _guard = self.handle.enter();
run_until_stalled();
}
}
#[cfg(test)]
mod tests {
use crate::{FrameWaker, ReactiveRuntime};
use any_spawner::Executor;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::{Duration, Instant};
fn recording_waker() -> (FrameWaker, Arc<AtomicUsize>) {
let counter = Arc::new(AtomicUsize::new(0));
let seen = counter.clone();
let waker: FrameWaker = Arc::new(move || {
counter.fetch_add(1, Ordering::SeqCst);
});
(waker, seen)
}
#[test]
fn spawn_local_rewake_from_timer_thread_fires_frame_waker() {
let _guard = crate::WAKER_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let (waker, count) = recording_waker();
let rt = ReactiveRuntime::init(waker);
let done = Arc::new(AtomicUsize::new(0));
let flag = done.clone();
let before = count.load(Ordering::SeqCst);
Executor::spawn_local(async move {
tokio::time::sleep(Duration::from_millis(50)).await;
flag.fetch_add(1, Ordering::SeqCst);
});
assert_eq!(
count.load(Ordering::SeqCst),
before + 1,
"spawn_local should fire the frame waker once at spawn"
);
rt.pump_local();
assert_eq!(
done.load(Ordering::SeqCst),
0,
"future must still be pending — the 50ms timer has not fired"
);
let after_pump = count.load(Ordering::SeqCst);
let start = Instant::now();
while count.load(Ordering::SeqCst) == after_pump && start.elapsed() < Duration::from_secs(5)
{
std::thread::sleep(Duration::from_millis(1));
}
assert!(
count.load(Ordering::SeqCst) > after_pump,
"a timer-driven re-wake must fire the FrameWaker with no pump in \
between (this is the desktop wake gap the composite waker closes)"
);
let start = Instant::now();
while done.load(Ordering::SeqCst) == 0 && start.elapsed() < Duration::from_secs(5) {
rt.pump_local();
std::thread::sleep(Duration::from_millis(1));
}
assert_eq!(
done.load(Ordering::SeqCst),
1,
"the local future should complete after the re-wake + pump"
);
}
}