Skip to main content

frust_reactive/
runtime.rs

1//! [`ReactiveRuntime`]: the process-wide reactive substrate.
2//!
3//! It owns a background tokio runtime, installs the custom [`any_spawner`]
4//! executor, holds the root reactive [`Owner`], and carries a swappable
5//! [`FrameWaker`] the executor fires to nudge the shell into pumping the
6//! UI-thread local task queue.
7//!
8//! **wasm32 arm.** `wasm32-unknown-unknown` has no OS threads and no
9//! mio-backed reactor, so the native background runtime's `rt-multi-thread`
10//! and `net` tokio features fail to compile there (see `Cargo.toml`'s
11//! target-gated `tokio` rows). [`ReactiveRuntime::init`] therefore builds a
12//! bare current-thread runtime instead of the multi-thread pool, and
13//! installs a wasm [`CustomExecutor`](any_spawner::CustomExecutor) (built
14//! inline here, not the native [`ForgeExecutor`](crate::executor::ForgeExecutor))
15//! whose `spawn` **and** `spawn_local` both route straight to
16//! `wasm_bindgen_futures::spawn_local` — a browser has one JS thread, so
17//! "spawn on a thread pool" and "spawn on the local microtask queue" are the
18//! same operation there, and `Send` futures need no different handling than
19//! `!Send` ones. This is the one call site that makes `wasm-bindgen-futures`
20//! (`Cargo.toml`'s wasm-only rows) a real direct dependency rather than a
21//! transitively-resolved one. Because both paths are driven by the browser's
22//! own microtask queue, [`pump_local`](ReactiveRuntime::pump_local) has
23//! nothing of this crate's own to drain on this arm, and `poll_local` is a
24//! no-op. The current-thread runtime is never driven (nothing calls
25//! `block_on`), so it exists only to give
26//! [`handle`](ReactiveRuntime::handle)/[`spawn_blocking`] a real
27//! `tokio::runtime::Handle` to type-check against.
28//!
29//! **What still doesn't work on wasm32.** `frust::spawn`/`frust::spawn_local`
30//! now run for real (above). [`spawn_blocking`] and `use_task`'s background
31//! half do not: `wasm32-unknown-unknown` has no OS threads, so there is no
32//! blocking pool to hand CPU-bound work to. `spawn_blocking` fails loudly
33//! (panics, naming the API and target) on this target rather than silently
34//! swallowing the closure — see its own docs. Treat that as an open gap, not
35//! a proven path.
36
37use std::sync::atomic::{AtomicBool, Ordering};
38use std::sync::{Arc, Mutex, OnceLock};
39
40use any_spawner::Executor;
41#[cfg(target_family = "wasm")]
42use any_spawner::{CustomExecutor, PinnedFuture, PinnedLocalFuture};
43use reactive_graph::owner::Owner;
44use tokio::runtime::{Builder, Handle, Runtime};
45
46use crate::executor;
47#[cfg(not(target_family = "wasm"))]
48use crate::executor::ForgeExecutor;
49
50/// The wasm [`any_spawner`] executor (see the module docs for why it exists
51/// instead of `any_spawner::Executor::init_wasm_bindgen()` and instead of the
52/// native [`ForgeExecutor`](crate::executor::ForgeExecutor)). A unit struct —
53/// nothing to hold, since `wasm_bindgen_futures::spawn_local` is a free
54/// function reachable from anywhere.
55#[cfg(target_family = "wasm")]
56struct WasmExecutor;
57
58#[cfg(target_family = "wasm")]
59impl CustomExecutor for WasmExecutor {
60    /// `Send` futures also go to `wasm_bindgen_futures::spawn_local`: a
61    /// browser has one JS thread, so there is no separate thread-pool
62    /// destination to route `Send` futures to the way the native
63    /// [`ForgeExecutor`](crate::executor::ForgeExecutor) does. This is the
64    /// call site that makes `wasm-bindgen-futures` a real direct dependency
65    /// (see `Cargo.toml`'s wasm-only rows).
66    fn spawn(&self, fut: PinnedFuture<()>) {
67        wasm_bindgen_futures::spawn_local(fut);
68    }
69
70    /// `!Send` futures route the same way — `wasm_bindgen_futures::spawn_local`
71    /// requires only `Future<Output = ()> + 'static`, which `PinnedFuture`
72    /// (this crate's `spawn` path) already satisfies too.
73    fn spawn_local(&self, fut: PinnedLocalFuture<()>) {
74        wasm_bindgen_futures::spawn_local(fut);
75    }
76
77    /// Nothing to drain: every task spawned above is driven by the browser's
78    /// own microtask queue, not a queue this crate owns (see the module
79    /// docs).
80    fn poll_local(&self) {}
81}
82
83/// A thread-safe, cheaply-cloneable "wake up and pump soon" callback. On
84/// desktop this is a winit `EventLoopProxy` send; on mobile (continuous frame
85/// loop) it is a no-op. It must be callable from any thread, since the executor
86/// fires it and background tasks may drive signal writes.
87pub type FrameWaker = Arc<dyn Fn() + Send + Sync>;
88
89/// Builds the one background runtime [`ReactiveRuntime::init`] installs.
90///
91/// Native: the real multi-thread worker pool with the time + IO drivers
92/// enabled — this arm is byte-identical to the pre-wasm code, `cfg`'d only by
93/// which arm compiles for a given target.
94#[cfg(not(target_family = "wasm"))]
95fn build_background_runtime() -> Runtime {
96    Builder::new_multi_thread()
97        .worker_threads(2)
98        .enable_time()
99        // IO driver: `frust::spawn`'s background tasks (tonic/hyper socket
100        // work) need a live reactor, not just the timer driver above.
101        .enable_io()
102        .thread_name("frust-reactive")
103        .build()
104        .expect("frust-reactive: failed to build the background tokio runtime")
105}
106
107/// Wasm: a bare current-thread runtime with no driver enabled (see the module
108/// docs for why) — nothing ever calls `block_on` on it, so it exists purely
109/// to give [`ReactiveRuntime::handle`] a real `tokio::runtime::Handle` to
110/// type-check against, not a functioning executor. [`spawn_blocking`] never
111/// reaches this handle on wasm32 (it panics first — see its own docs).
112#[cfg(target_family = "wasm")]
113fn build_background_runtime() -> Runtime {
114    Builder::new_current_thread()
115        .thread_name("frust-reactive")
116        .build()
117        .expect("frust-reactive: failed to build the wasm placeholder tokio runtime")
118}
119
120/// The one process-wide runtime. Installed on first [`ReactiveRuntime::init`]
121/// and never torn down (process lifetime — the matrix-rust-sdk precedent).
122static RUNTIME: OnceLock<ReactiveRuntime> = OnceLock::new();
123
124/// The process-wide reactive runtime. Construct via [`ReactiveRuntime::init`]
125/// (once, on the UI thread) and reach later calls via [`ReactiveRuntime::get`].
126pub struct ReactiveRuntime {
127    // Kept alive for the process lifetime; the static is never dropped, so the
128    // runtime never shuts down. Read only through `handle`.
129    _runtime: Runtime,
130    handle: Handle,
131    root: Owner,
132    waker: Mutex<FrameWaker>,
133    /// Process-wide "did any tracked signal change since I last asked" flag —
134    /// the reactive input to the mobile frame gate. Set by
135    /// [`crate::tracked::TrackedScope`]'s dirty path (any tracked scope's
136    /// invalidation trips it, coalesced by nature since it's a bool, not a
137    /// counter); drained by [`take_signals_dirty`](Self::take_signals_dirty).
138    signals_dirty: AtomicBool,
139}
140
141impl ReactiveRuntime {
142    /// Process-wide initialization (idempotent).
143    ///
144    /// **Must be called on the UI thread** — the calling thread claims itself as
145    /// the owner of the local task queue that [`pump_local`](Self::pump_local)
146    /// drains. On the first call it builds the background tokio runtime, the
147    /// custom executor, and the root [`Owner`]. A second call (e.g. a relaunched
148    /// shell in the same process) does **not** rebuild anything: it re-marks the
149    /// current thread as the UI thread, **replaces the waker** so the new shell
150    /// re-owns wake-up, and returns the existing runtime. `any_spawner`'s
151    /// `AlreadySet` on a repeated executor install is expected and benign.
152    pub fn init(waker: FrameWaker) -> &'static ReactiveRuntime {
153        // Claim the UI thread on every call: the fast path below must still
154        // (re-)mark a relaunched shell's thread.
155        executor::init_ui_thread();
156
157        if let Some(existing) = RUNTIME.get() {
158            existing.set_waker(waker);
159            return existing;
160        }
161
162        let runtime = build_background_runtime();
163        let handle = runtime.handle().clone();
164        let root = Owner::new();
165
166        let candidate = ReactiveRuntime {
167            _runtime: runtime,
168            handle: handle.clone(),
169            root,
170            waker: Mutex::new(waker),
171            signals_dirty: AtomicBool::new(false),
172        };
173
174        match RUNTIME.set(candidate) {
175            Ok(()) => {
176                let rt = RUNTIME.get().expect("runtime was just installed");
177                // Install the executor AFTER the runtime is reachable, so the
178                // executor's `spawn_local` can fire the waker via `get()`.
179                // `AlreadySet` (a prior shell already installed it) is benign.
180                #[cfg(not(target_family = "wasm"))]
181                {
182                    let _ = Executor::init_custom_executor(ForgeExecutor::new(handle));
183                }
184                // No automatic wasm default exists (see module docs); install
185                // the wasm `CustomExecutor` explicitly. `AlreadySet` on a
186                // repeated install (e.g. a relaunched shell in the same
187                // process) is handled the same benign way as the native arm.
188                #[cfg(target_family = "wasm")]
189                {
190                    let _ = Executor::init_custom_executor(WasmExecutor);
191                }
192                rt
193            }
194            Err(candidate) => {
195                // Lost a concurrent init race (init is documented single-thread,
196                // so this is defensive): adopt the winner, swap our waker in.
197                let winner = RUNTIME.get().expect("runtime is set on the Err path");
198                let waker = candidate
199                    .waker
200                    .into_inner()
201                    .expect("candidate waker mutex is uncontended");
202                winner.set_waker(waker);
203                winner
204            }
205        }
206    }
207
208    /// The installed runtime, if [`init`](Self::init) has run.
209    pub fn get() -> Option<&'static ReactiveRuntime> {
210        RUNTIME.get()
211    }
212
213    /// Runs `f` under the root reactive [`Owner`]. Shells wrap their per-frame
214    /// rebuild in this so signals created during a rebuild are owned by the root
215    /// (and disposed only at process end).
216    pub fn with_owner<R>(&self, f: impl FnOnce() -> R) -> R {
217        self.root.with(f)
218    }
219
220    /// Drains the UI-thread local task queue until it stalls, inside a tokio
221    /// runtime context so `tokio::time::sleep` in a local task registers with
222    /// the background runtime's timer driver. Shells call this each loop turn.
223    ///
224    /// Must be called on the UI thread (the one [`init`](Self::init) ran on);
225    /// off-thread it pumps a distinct, empty queue.
226    ///
227    /// **Wasm:** a harmless no-op drain of an always-empty queue. On this
228    /// target `spawn_local` routes to `wasm_bindgen_futures::spawn_local`
229    /// rather than this crate's own local pool (see the module docs), so
230    /// there is nothing of this crate's own for a shell to pump — it is still
231    /// safe (and, for cross-platform shell code, simplest) to call this every
232    /// loop turn on wasm too.
233    pub fn pump_local(&self) {
234        debug_assert!(
235            executor::is_ui_thread(),
236            "frust-reactive: pump_local called off the UI thread — this pumps \
237             an unrelated empty queue and is a wiring bug"
238        );
239        let _guard = self.handle.enter();
240        executor::run_until_stalled();
241    }
242
243    /// Replaces the frame waker (a relaunched shell re-owns wake-up).
244    pub fn set_waker(&self, waker: FrameWaker) {
245        *self
246            .waker
247            .lock()
248            .expect("frust-reactive: waker mutex poisoned") = waker;
249    }
250
251    /// Fires the current frame waker. Used by the executor's `spawn_local` and
252    /// by the signals dirty-bridge.
253    pub fn wake(&self) {
254        // Clone the Arc out before invoking so the lock is not held across the
255        // callback (which may re-enter the runtime).
256        let waker = self
257            .waker
258            .lock()
259            .expect("frust-reactive: waker mutex poisoned")
260            .clone();
261        waker();
262    }
263
264    /// The background runtime handle. `pub(crate)` so the heavy-work idiom
265    /// ([`crate::task`]) can `spawn`/`spawn_blocking` onto the background
266    /// runtime, and so the test that asserts a repeated executor install is
267    /// benign can rebuild the executor. Deliberately not part of the public
268    /// API — app code routes through `spawn`/`spawn_local`/`spawn_blocking`
269    /// rather than naming a raw `tokio::runtime::Handle`.
270    pub(crate) fn handle(&self) -> Handle {
271        self.handle.clone()
272    }
273
274    /// Trips the process-wide signals-dirty flag. Called from
275    /// [`crate::tracked::TrackedScope`]'s dirty-notification path — any
276    /// tracked scope's invalidation trips this, not just the coalesced
277    /// frame-waker edge, so a shell can ask "did *anything* change" cheaply
278    /// once per frame.
279    pub(crate) fn mark_signals_dirty(&self) {
280        self.signals_dirty.store(true, Ordering::SeqCst);
281    }
282
283    /// Drains the signals-dirty flag: returns whether any tracked signal
284    /// changed since the last `take_signals_dirty` call, and clears it.
285    ///
286    /// **Ordering contract for shells:** pump local tasks
287    /// ([`pump_local`](Self::pump_local)) FIRST, then call this — a
288    /// `spawn_local` continuation that writes a signal during the pump must
289    /// be observed by the *same* frame's dirty check. Calling this before the
290    /// pump can miss a write a just-drained local task makes.
291    pub fn take_signals_dirty(&self) -> bool {
292        self.signals_dirty.swap(false, Ordering::SeqCst)
293    }
294
295    /// Non-draining peek at the signals-dirty flag (does not clear it).
296    /// Prefer [`take_signals_dirty`](Self::take_signals_dirty) for the actual
297    /// once-per-frame gate check; this is for tests/diagnostics that want to
298    /// observe the flag without consuming it.
299    pub fn signals_dirty(&self) -> bool {
300        self.signals_dirty.load(Ordering::SeqCst)
301    }
302}
303
304/// Runs a one-off, blocking CPU workload on the background runtime's blocking
305/// thread pool, returning a [`JoinHandle`](tokio::task::JoinHandle) to `.await`
306/// its result.
307///
308/// This is the CPU-bound entry point of the heavy-work routing convention:
309///
310/// | Call | Use for |
311/// |---|---|
312/// | `frust::spawn` | `Send` async IO-bound work |
313/// | `frust::spawn_local` | `!Send` work that must stay on the UI thread |
314/// | `frust::spawn_blocking` | one-off **CPU-bound** blocking work (JSON parse, decode, hashing) |
315/// | `rayon` | data-parallel compute — an **app-level** choice, deliberately not bundled |
316///
317/// The pool already exists (the runtime is built `rt-multi-thread`), so this
318/// is a thin facade over [`tokio::runtime::Handle::spawn_blocking`]. Compose it
319/// inside a [`use_task`](crate::use_task) fetcher —
320/// `use_task(|| async { spawn_blocking(parse).await })` — to get load/error
321/// states and cancellation for free.
322///
323/// **Cancellation limitation:** dropping/aborting the returned handle stops the
324/// result from being delivered, but a blocking closure *already running*
325/// cannot be interrupted (there is no safe way to unwind arbitrary blocking
326/// code) — the same limitation every runtime has.
327///
328/// **No working wasm equivalent yet — fails loudly.** `wasm32-unknown-unknown`
329/// has no OS threads, so there is no blocking pool for `Handle::spawn_blocking`
330/// to hand work to (see the module docs). Rather than returning an inert
331/// `JoinHandle` that silently never resolves, this function panics up front
332/// on that target (see Panics below), before ever reaching tokio. Treat this
333/// as an open gap on wasm, not a proven path — the eventual fix is a
334/// browser-thread/Web-Worker-backed `JoinHandle`, not yet built.
335///
336/// # Panics
337///
338/// - Panics if [`ReactiveRuntime::init`] has not run yet — the same
339///   wiring-bug-not-runtime-condition contract as `spawn_local` off the UI
340///   thread.
341/// - **Always panics on `wasm32`** (see above), independent of `init` state.
342pub fn spawn_blocking<F, R>(f: F) -> tokio::task::JoinHandle<R>
343where
344    F: FnOnce() -> R + Send + 'static,
345    R: Send + 'static,
346{
347    #[cfg(target_arch = "wasm32")]
348    {
349        let _ = f;
350        panic!(
351            "frust-reactive: spawn_blocking called on wasm32 — \
352             wasm32-unknown-unknown has no OS threads, so there is no \
353             blocking thread pool to hand this closure to. This is an \
354             unimplemented gap on this target, not a wiring bug: route the \
355             work through a different mechanism until a \
356             browser-thread/Web-Worker-backed JoinHandle lands for \
357             `frust::spawn_blocking`."
358        );
359    }
360
361    #[cfg(not(target_arch = "wasm32"))]
362    {
363        let rt = ReactiveRuntime::get().expect(
364            "frust-reactive: spawn_blocking called before ReactiveRuntime::init — \
365             this is a wiring bug: initialize the reactive runtime (the shell does \
366             this on startup) before spawning work",
367        );
368        rt.handle.spawn_blocking(f)
369    }
370}
371
372#[cfg(test)]
373mod tests {
374    use super::*;
375    use crate::tracked::TrackedScope;
376    use any_spawner::Executor;
377    use reactive_graph::signal::RwSignal;
378    use reactive_graph::traits::{Get, Set};
379
380    fn noop_waker() -> FrameWaker {
381        Arc::new(|| {})
382    }
383
384    /// Criterion 1: a write inside a tracked rebuild trips the process-wide
385    /// signals-dirty flag; `take_signals_dirty` drains it (true once, then
386    /// false with no intervening write).
387    #[test]
388    fn signals_dirty_set_by_tracked_write_and_drained_by_take() {
389        let _guard = crate::WAKER_TEST_LOCK
390            .lock()
391            .unwrap_or_else(|e| e.into_inner());
392
393        let rt = ReactiveRuntime::init(noop_waker());
394        // Drain any flag left dirty by a previous test sharing this
395        // process-wide runtime.
396        rt.take_signals_dirty();
397
398        let scope = TrackedScope::new();
399        let sig = rt.with_owner(|| RwSignal::new(0));
400        scope.track(|| sig.get());
401        assert!(!rt.signals_dirty(), "no write yet — flag must be clean");
402
403        sig.set(1);
404        assert!(
405            rt.signals_dirty(),
406            "a write to a tracked signal must trip the process-wide flag"
407        );
408        assert!(
409            rt.take_signals_dirty(),
410            "take must observe the dirty flag and drain it"
411        );
412        assert!(
413            !rt.take_signals_dirty(),
414            "a second take with no intervening write must return false"
415        );
416    }
417
418    /// Criterion: writes from multiple distinct tracked scopes still coalesce
419    /// onto the one process-wide flag (a bool, not a counter) — one drain
420    /// clears every pending scope's contribution at once.
421    #[test]
422    fn signals_dirty_coalesces_across_scopes() {
423        let _guard = crate::WAKER_TEST_LOCK
424            .lock()
425            .unwrap_or_else(|e| e.into_inner());
426
427        let rt = ReactiveRuntime::init(noop_waker());
428        rt.take_signals_dirty();
429
430        let scope_a = TrackedScope::new();
431        let scope_b = TrackedScope::new();
432        let sig_a = rt.with_owner(|| RwSignal::new(0));
433        let sig_b = rt.with_owner(|| RwSignal::new(0));
434        scope_a.track(|| sig_a.get());
435        scope_b.track(|| sig_b.get());
436
437        sig_a.set(1);
438        sig_b.set(1);
439        sig_a.set(2);
440
441        assert!(
442            rt.take_signals_dirty(),
443            "multiple writes across multiple scopes must still trip the flag"
444        );
445        assert!(
446            !rt.take_signals_dirty(),
447            "drain must clear every pending contribution at once"
448        );
449    }
450
451    /// The ordering contract: a spawned local task that writes a tracked
452    /// signal *during* `pump_local` must be observed by a `take_signals_dirty`
453    /// call that runs after the pump — the documented pump-first contract
454    /// shells rely on.
455    #[test]
456    fn signals_dirty_pump_then_take_ordering() {
457        let _guard = crate::WAKER_TEST_LOCK
458            .lock()
459            .unwrap_or_else(|e| e.into_inner());
460
461        let rt = ReactiveRuntime::init(noop_waker());
462        rt.take_signals_dirty();
463
464        let scope = TrackedScope::new();
465        let sig = rt.with_owner(|| RwSignal::new(0));
466        scope.track(|| sig.get());
467
468        Executor::spawn_local(async move {
469            sig.set(1);
470        });
471        assert!(
472            !rt.signals_dirty(),
473            "the local task has not run yet — spawning it must not itself \
474             dirty the flag"
475        );
476
477        rt.pump_local();
478        assert!(
479            rt.take_signals_dirty(),
480            "a signal write from a pumped local task must be observed by \
481             take_signals_dirty called after the pump — the pump-first \
482             ordering contract"
483        );
484    }
485
486    /// The runtime's IO driver must actually be live on the `frust::spawn`
487    /// path (`Executor::spawn` -> `ForgeExecutor::spawn` -> `Handle::spawn`,
488    /// the same route real async IO takes), not just the timer driver the
489    /// other tests in this module exercise: a spawned task must be able to
490    /// bind a socket and complete an accept/connect round-trip.
491    #[test]
492    fn spawn_reaches_io_driver() {
493        let _guard = crate::WAKER_TEST_LOCK
494            .lock()
495            .unwrap_or_else(|e| e.into_inner());
496
497        ReactiveRuntime::init(noop_waker());
498
499        let (tx, rx) = std::sync::mpsc::channel();
500        Executor::spawn(async move {
501            let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
502                .await
503                .expect("bind must succeed with the IO driver enabled");
504            let addr = listener.local_addr().expect("listener has a local addr");
505
506            let accept = tokio::spawn(async move { listener.accept().await });
507            tokio::net::TcpStream::connect(addr)
508                .await
509                .expect("connect must succeed with the IO driver enabled");
510            accept
511                .await
512                .expect("accept task must not panic")
513                .expect("accept must succeed");
514
515            tx.send(()).expect("test receiver must still be alive");
516        });
517
518        rx.recv_timeout(std::time::Duration::from_secs(5)).expect(
519            "a spawned task must complete an accept/connect round-trip \
520             through the IO driver within 5s",
521        );
522    }
523}