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}