use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use any_spawner::Executor;
#[cfg(target_family = "wasm")]
use any_spawner::{CustomExecutor, PinnedFuture, PinnedLocalFuture};
use reactive_graph::owner::Owner;
use tokio::runtime::{Builder, Handle, Runtime};
use crate::executor;
#[cfg(not(target_family = "wasm"))]
use crate::executor::ForgeExecutor;
#[cfg(target_family = "wasm")]
struct WasmExecutor;
#[cfg(target_family = "wasm")]
impl CustomExecutor for WasmExecutor {
fn spawn(&self, fut: PinnedFuture<()>) {
wasm_bindgen_futures::spawn_local(fut);
}
fn spawn_local(&self, fut: PinnedLocalFuture<()>) {
wasm_bindgen_futures::spawn_local(fut);
}
fn poll_local(&self) {}
}
pub type FrameWaker = Arc<dyn Fn() + Send + Sync>;
#[cfg(not(target_family = "wasm"))]
fn build_background_runtime() -> Runtime {
Builder::new_multi_thread()
.worker_threads(2)
.enable_time()
.enable_io()
.thread_name("frust-reactive")
.build()
.expect("frust-reactive: failed to build the background tokio runtime")
}
#[cfg(target_family = "wasm")]
fn build_background_runtime() -> Runtime {
Builder::new_current_thread()
.thread_name("frust-reactive")
.build()
.expect("frust-reactive: failed to build the wasm placeholder tokio runtime")
}
static RUNTIME: OnceLock<ReactiveRuntime> = OnceLock::new();
pub struct ReactiveRuntime {
_runtime: Runtime,
handle: Handle,
root: Owner,
waker: Mutex<FrameWaker>,
signals_dirty: AtomicBool,
}
impl ReactiveRuntime {
pub fn init(waker: FrameWaker) -> &'static ReactiveRuntime {
executor::init_ui_thread();
if let Some(existing) = RUNTIME.get() {
existing.set_waker(waker);
return existing;
}
let runtime = build_background_runtime();
let handle = runtime.handle().clone();
let root = Owner::new();
let candidate = ReactiveRuntime {
_runtime: runtime,
handle: handle.clone(),
root,
waker: Mutex::new(waker),
signals_dirty: AtomicBool::new(false),
};
match RUNTIME.set(candidate) {
Ok(()) => {
let rt = RUNTIME.get().expect("runtime was just installed");
#[cfg(not(target_family = "wasm"))]
{
let _ = Executor::init_custom_executor(ForgeExecutor::new(handle));
}
#[cfg(target_family = "wasm")]
{
let _ = Executor::init_custom_executor(WasmExecutor);
}
rt
}
Err(candidate) => {
let winner = RUNTIME.get().expect("runtime is set on the Err path");
let waker = candidate
.waker
.into_inner()
.expect("candidate waker mutex is uncontended");
winner.set_waker(waker);
winner
}
}
}
pub fn get() -> Option<&'static ReactiveRuntime> {
RUNTIME.get()
}
pub fn with_owner<R>(&self, f: impl FnOnce() -> R) -> R {
self.root.with(f)
}
pub fn pump_local(&self) {
debug_assert!(
executor::is_ui_thread(),
"frust-reactive: pump_local called off the UI thread — this pumps \
an unrelated empty queue and is a wiring bug"
);
let _guard = self.handle.enter();
executor::run_until_stalled();
}
pub fn set_waker(&self, waker: FrameWaker) {
*self
.waker
.lock()
.expect("frust-reactive: waker mutex poisoned") = waker;
}
pub fn wake(&self) {
let waker = self
.waker
.lock()
.expect("frust-reactive: waker mutex poisoned")
.clone();
waker();
}
pub(crate) fn handle(&self) -> Handle {
self.handle.clone()
}
pub(crate) fn mark_signals_dirty(&self) {
self.signals_dirty.store(true, Ordering::SeqCst);
}
pub fn take_signals_dirty(&self) -> bool {
self.signals_dirty.swap(false, Ordering::SeqCst)
}
pub fn signals_dirty(&self) -> bool {
self.signals_dirty.load(Ordering::SeqCst)
}
}
pub fn spawn_blocking<F, R>(f: F) -> tokio::task::JoinHandle<R>
where
F: FnOnce() -> R + Send + 'static,
R: Send + 'static,
{
#[cfg(target_arch = "wasm32")]
{
let _ = f;
panic!(
"frust-reactive: spawn_blocking called on wasm32 — \
wasm32-unknown-unknown has no OS threads, so there is no \
blocking thread pool to hand this closure to. This is an \
unimplemented gap on this target, not a wiring bug: route the \
work through a different mechanism until a \
browser-thread/Web-Worker-backed JoinHandle lands for \
`frust::spawn_blocking`."
);
}
#[cfg(not(target_arch = "wasm32"))]
{
let rt = ReactiveRuntime::get().expect(
"frust-reactive: spawn_blocking called before ReactiveRuntime::init — \
this is a wiring bug: initialize the reactive runtime (the shell does \
this on startup) before spawning work",
);
rt.handle.spawn_blocking(f)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::tracked::TrackedScope;
use any_spawner::Executor;
use reactive_graph::signal::RwSignal;
use reactive_graph::traits::{Get, Set};
fn noop_waker() -> FrameWaker {
Arc::new(|| {})
}
#[test]
fn signals_dirty_set_by_tracked_write_and_drained_by_take() {
let _guard = crate::WAKER_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let rt = ReactiveRuntime::init(noop_waker());
rt.take_signals_dirty();
let scope = TrackedScope::new();
let sig = rt.with_owner(|| RwSignal::new(0));
scope.track(|| sig.get());
assert!(!rt.signals_dirty(), "no write yet — flag must be clean");
sig.set(1);
assert!(
rt.signals_dirty(),
"a write to a tracked signal must trip the process-wide flag"
);
assert!(
rt.take_signals_dirty(),
"take must observe the dirty flag and drain it"
);
assert!(
!rt.take_signals_dirty(),
"a second take with no intervening write must return false"
);
}
#[test]
fn signals_dirty_coalesces_across_scopes() {
let _guard = crate::WAKER_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let rt = ReactiveRuntime::init(noop_waker());
rt.take_signals_dirty();
let scope_a = TrackedScope::new();
let scope_b = TrackedScope::new();
let sig_a = rt.with_owner(|| RwSignal::new(0));
let sig_b = rt.with_owner(|| RwSignal::new(0));
scope_a.track(|| sig_a.get());
scope_b.track(|| sig_b.get());
sig_a.set(1);
sig_b.set(1);
sig_a.set(2);
assert!(
rt.take_signals_dirty(),
"multiple writes across multiple scopes must still trip the flag"
);
assert!(
!rt.take_signals_dirty(),
"drain must clear every pending contribution at once"
);
}
#[test]
fn signals_dirty_pump_then_take_ordering() {
let _guard = crate::WAKER_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let rt = ReactiveRuntime::init(noop_waker());
rt.take_signals_dirty();
let scope = TrackedScope::new();
let sig = rt.with_owner(|| RwSignal::new(0));
scope.track(|| sig.get());
Executor::spawn_local(async move {
sig.set(1);
});
assert!(
!rt.signals_dirty(),
"the local task has not run yet — spawning it must not itself \
dirty the flag"
);
rt.pump_local();
assert!(
rt.take_signals_dirty(),
"a signal write from a pumped local task must be observed by \
take_signals_dirty called after the pump — the pump-first \
ordering contract"
);
}
#[test]
fn spawn_reaches_io_driver() {
let _guard = crate::WAKER_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
ReactiveRuntime::init(noop_waker());
let (tx, rx) = std::sync::mpsc::channel();
Executor::spawn(async move {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind must succeed with the IO driver enabled");
let addr = listener.local_addr().expect("listener has a local addr");
let accept = tokio::spawn(async move { listener.accept().await });
tokio::net::TcpStream::connect(addr)
.await
.expect("connect must succeed with the IO driver enabled");
accept
.await
.expect("accept task must not panic")
.expect("accept must succeed");
tx.send(()).expect("test receiver must still be alive");
});
rx.recv_timeout(std::time::Duration::from_secs(5)).expect(
"a spawned task must complete an accept/connect round-trip \
through the IO driver within 5s",
);
}
}