Skip to main content

frust_reactive/
lib.rs

1//! Reactive-programming substrate for Frust.
2//!
3//! This crate owns the process-wide [`ReactiveRuntime`]: a background tokio
4//! runtime, a custom [`any_spawner`] executor that routes `spawn` to that
5//! runtime and `spawn_local` to a UI-thread local task queue, a swappable
6//! [`FrameWaker`], and the root reactive [`Owner`]. It also re-exports the
7//! `reactive_graph` types the later programming-model tasks build on, and
8//! owns the process-wide deep-link source shells write platform deep links
9//! into (see [`push_deep_link`]/[`deep_links`]), the process-wide Android
10//! back-press source next to it (see [`push_back_press`]/[`back_presses`] and
11//! the `handles_back` flag), and the process-wide desktop menu-activation
12//! source beside those (see [`push_menu_event`]/[`menu_events`]). A process-wide
13//! signals-dirty flag (`ReactiveRuntime::take_signals_dirty`), tripped by
14//! [`TrackedScope`]'s dirty path, lets a shell ask synchronously and cheaply
15//! once per frame "did any tracked signal change since I last asked" — the
16//! reactive input to the mobile frame gate.
17//!
18//! It is a leaf substrate: no `winit`, no `vello`/`wgpu`, no `frust-core`
19//! dependency. Shells own the wake-up wiring and call [`ReactiveRuntime::init`]
20//! (once, on the UI thread) and [`ReactiveRuntime::pump_local`] each frame.
21
22mod back;
23mod deep_link;
24mod executor;
25mod menu;
26mod runtime;
27mod task;
28mod tracked;
29
30pub use back::{
31    BackPresses, CanPopRegistration, back_presses, clear_can_pop_provider, handles_back,
32    push_back_press, set_can_pop_provider, set_handles_back,
33};
34pub use deep_link::{DeepLink, DeepLinks, deep_links, push_deep_link};
35pub use menu::{MenuEvent, MenuEvents, menu_events, push_menu_event};
36pub use runtime::{FrameWaker, ReactiveRuntime, spawn_blocking};
37pub use task::{AsyncValue, TaskError, UseTask, use_task};
38pub use tracked::TrackedScope;
39
40pub use reactive_graph::owner::{Owner, on_cleanup, provide_context, use_context};
41pub use reactive_graph::signal::RwSignal;
42
43/// Serializes every test that installs a recording [`FrameWaker`] and asserts on
44/// wake counts. The frame waker is process-global and swappable, so two such
45/// tests running on the parallel test-runner's separate threads would clobber
46/// each other's waker mid-assertion. Both the runtime end-to-end test and the
47/// tracked-scope test acquire this before touching the waker.
48#[cfg(test)]
49pub(crate) static WAKER_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
50
51#[cfg(test)]
52mod tests {
53    use super::*;
54    use reactive_graph::traits::{Get, Set};
55    use std::sync::Arc;
56    use std::sync::atomic::{AtomicUsize, Ordering};
57    use std::time::{Duration, Instant};
58
59    use any_spawner::{Executor, ExecutorError};
60
61    #[test]
62    fn owner_and_signal_round_trip() {
63        let owner = Owner::new();
64        owner.set();
65
66        let signal = RwSignal::new(1);
67        assert_eq!(signal.get(), 1);
68
69        signal.set(2);
70        assert_eq!(signal.get(), 2);
71    }
72
73    /// A recording waker: an `Arc<AtomicUsize>` bumped once per `wake()`.
74    fn recording_waker() -> (FrameWaker, Arc<AtomicUsize>) {
75        let counter = Arc::new(AtomicUsize::new(0));
76        let seen = counter.clone();
77        let waker: FrameWaker = Arc::new(move || {
78            counter.fetch_add(1, Ordering::SeqCst);
79        });
80        (waker, seen)
81    }
82
83    /// The runtime, executor, and UI-thread markers are all process-global, and
84    /// the UI-thread markers are thread-local — so the whole scenario runs in one
85    /// `#[test]` on one thread to stay independent of the parallel test runner.
86    /// Each block maps to an acceptance criterion.
87    #[test]
88    fn reactive_runtime_end_to_end() {
89        // Serialize with the tracked-scope test: both swap the global waker.
90        let _guard = crate::WAKER_TEST_LOCK
91            .lock()
92            .unwrap_or_else(|e| e.into_inner());
93
94        let (waker1, waker1_count) = recording_waker();
95        let rt = ReactiveRuntime::init(waker1);
96
97        // `with_owner` runs its closure under the root owner, so a signal
98        // created there is valid and readable afterward.
99        let scoped = rt.with_owner(|| {
100            let s = RwSignal::new(41);
101            s.set(42);
102            s
103        });
104        assert_eq!(scoped.get(), 42);
105
106        // Criterion 2: `spawn` runs on the background runtime.
107        {
108            let (tx, rx) = std::sync::mpsc::channel();
109            let ui_thread = std::thread::current().id();
110            Executor::spawn(async move {
111                let _ = tx.send(std::thread::current().id());
112            });
113            let ran_on = rx
114                .recv_timeout(Duration::from_secs(5))
115                .expect("spawned background task did not run");
116            assert_ne!(
117                ran_on, ui_thread,
118                "Executor::spawn should run on a background worker thread"
119            );
120        }
121
122        // Criterion 1 (immediate): a `spawn_local` future runs to completion via
123        // `pump_local`. Criterion 3: `spawn_local` fires the waker.
124        {
125            let ran = Arc::new(AtomicUsize::new(0));
126            let flag = ran.clone();
127            let before = waker1_count.load(Ordering::SeqCst);
128            Executor::spawn_local(async move {
129                flag.fetch_add(1, Ordering::SeqCst);
130            });
131            assert_eq!(
132                waker1_count.load(Ordering::SeqCst),
133                before + 1,
134                "spawn_local should fire the frame waker exactly once"
135            );
136            assert_eq!(
137                ran.load(Ordering::SeqCst),
138                0,
139                "task should not run before pump"
140            );
141            rt.pump_local();
142            assert_eq!(
143                ran.load(Ordering::SeqCst),
144                1,
145                "spawn_local future should complete after pump_local"
146            );
147        }
148
149        // Criterion 1 (timer): a local future awaiting tokio time completes after
150        // pumping past the deadline.
151        {
152            let done = Arc::new(AtomicUsize::new(0));
153            let flag = done.clone();
154            Executor::spawn_local(async move {
155                tokio::time::sleep(Duration::from_millis(10)).await;
156                flag.fetch_add(1, Ordering::SeqCst);
157            });
158            let start = Instant::now();
159            while done.load(Ordering::SeqCst) == 0 && start.elapsed() < Duration::from_secs(2) {
160                rt.pump_local();
161                std::thread::sleep(Duration::from_millis(1));
162            }
163            assert_eq!(
164                done.load(Ordering::SeqCst),
165                1,
166                "local task awaiting tokio::time::sleep should complete after pumping"
167            );
168        }
169
170        // Criterion 4: a second `init` is benign, returns the same runtime, and
171        // swaps the waker (the new one now fires on spawn_local).
172        {
173            let (waker2, waker2_count) = recording_waker();
174            let rt2 = ReactiveRuntime::init(waker2);
175            assert!(
176                std::ptr::eq(rt, rt2),
177                "second init must return the existing runtime"
178            );
179            Executor::spawn_local(async {});
180            assert_eq!(
181                waker2_count.load(Ordering::SeqCst),
182                1,
183                "second init must swap in the new waker"
184            );
185            rt.pump_local();
186
187            // Criterion 4 (executor double-init): a repeated executor install
188            // returns AlreadySet rather than panicking.
189            let repeated =
190                Executor::init_custom_executor(crate::executor::ForgeExecutor::new(rt.handle()));
191            assert!(
192                matches!(repeated, Err(ExecutorError::AlreadySet)),
193                "repeated executor install should be a benign AlreadySet"
194            );
195        }
196
197        // Criterion 5: `spawn_local` from a non-UI OS thread panics with the
198        // wiring-bug message (the executor is already installed above, so this is
199        // our panic, not any_spawner's uninitialized-executor panic).
200        {
201            let joined = std::thread::spawn(|| {
202                Executor::spawn_local(async {});
203            })
204            .join();
205            let payload = joined.expect_err("spawn_local off the UI thread must panic");
206            let msg = payload
207                .downcast_ref::<&str>()
208                .map(|s| (*s).to_string())
209                .or_else(|| payload.downcast_ref::<String>().cloned())
210                .unwrap_or_default();
211            assert!(
212                msg.contains("UI thread") && msg.contains("wiring bug"),
213                "panic should name the UI-thread wiring bug, got: {msg}"
214            );
215        }
216    }
217}