mod back;
mod deep_link;
mod executor;
mod menu;
mod runtime;
mod task;
mod tracked;
pub use back::{
BackPresses, CanPopRegistration, back_presses, clear_can_pop_provider, handles_back,
push_back_press, set_can_pop_provider, set_handles_back,
};
pub use deep_link::{DeepLink, DeepLinks, deep_links, push_deep_link};
pub use menu::{MenuEvent, MenuEvents, menu_events, push_menu_event};
pub use runtime::{FrameWaker, ReactiveRuntime, spawn_blocking};
pub use task::{AsyncValue, TaskError, UseTask, use_task};
pub use tracked::TrackedScope;
pub use reactive_graph::owner::{Owner, on_cleanup, provide_context, use_context};
pub use reactive_graph::signal::RwSignal;
#[cfg(test)]
pub(crate) static WAKER_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
#[cfg(test)]
mod tests {
use super::*;
use reactive_graph::traits::{Get, Set};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::{Duration, Instant};
use any_spawner::{Executor, ExecutorError};
#[test]
fn owner_and_signal_round_trip() {
let owner = Owner::new();
owner.set();
let signal = RwSignal::new(1);
assert_eq!(signal.get(), 1);
signal.set(2);
assert_eq!(signal.get(), 2);
}
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 reactive_runtime_end_to_end() {
let _guard = crate::WAKER_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let (waker1, waker1_count) = recording_waker();
let rt = ReactiveRuntime::init(waker1);
let scoped = rt.with_owner(|| {
let s = RwSignal::new(41);
s.set(42);
s
});
assert_eq!(scoped.get(), 42);
{
let (tx, rx) = std::sync::mpsc::channel();
let ui_thread = std::thread::current().id();
Executor::spawn(async move {
let _ = tx.send(std::thread::current().id());
});
let ran_on = rx
.recv_timeout(Duration::from_secs(5))
.expect("spawned background task did not run");
assert_ne!(
ran_on, ui_thread,
"Executor::spawn should run on a background worker thread"
);
}
{
let ran = Arc::new(AtomicUsize::new(0));
let flag = ran.clone();
let before = waker1_count.load(Ordering::SeqCst);
Executor::spawn_local(async move {
flag.fetch_add(1, Ordering::SeqCst);
});
assert_eq!(
waker1_count.load(Ordering::SeqCst),
before + 1,
"spawn_local should fire the frame waker exactly once"
);
assert_eq!(
ran.load(Ordering::SeqCst),
0,
"task should not run before pump"
);
rt.pump_local();
assert_eq!(
ran.load(Ordering::SeqCst),
1,
"spawn_local future should complete after pump_local"
);
}
{
let done = Arc::new(AtomicUsize::new(0));
let flag = done.clone();
Executor::spawn_local(async move {
tokio::time::sleep(Duration::from_millis(10)).await;
flag.fetch_add(1, Ordering::SeqCst);
});
let start = Instant::now();
while done.load(Ordering::SeqCst) == 0 && start.elapsed() < Duration::from_secs(2) {
rt.pump_local();
std::thread::sleep(Duration::from_millis(1));
}
assert_eq!(
done.load(Ordering::SeqCst),
1,
"local task awaiting tokio::time::sleep should complete after pumping"
);
}
{
let (waker2, waker2_count) = recording_waker();
let rt2 = ReactiveRuntime::init(waker2);
assert!(
std::ptr::eq(rt, rt2),
"second init must return the existing runtime"
);
Executor::spawn_local(async {});
assert_eq!(
waker2_count.load(Ordering::SeqCst),
1,
"second init must swap in the new waker"
);
rt.pump_local();
let repeated =
Executor::init_custom_executor(crate::executor::ForgeExecutor::new(rt.handle()));
assert!(
matches!(repeated, Err(ExecutorError::AlreadySet)),
"repeated executor install should be a benign AlreadySet"
);
}
{
let joined = std::thread::spawn(|| {
Executor::spawn_local(async {});
})
.join();
let payload = joined.expect_err("spawn_local off the UI thread must panic");
let msg = payload
.downcast_ref::<&str>()
.map(|s| (*s).to_string())
.or_else(|| payload.downcast_ref::<String>().cloned())
.unwrap_or_default();
assert!(
msg.contains("UI thread") && msg.contains("wiring bug"),
"panic should name the UI-thread wiring bug, got: {msg}"
);
}
}
}