Skip to main content

frust_reactive/
tracked.rs

1//! [`TrackedScope`]: a custom `reactive_graph` subscriber that runs a frame
2//! rebuild with dependency tracking and bridges "a tracked signal changed" into
3//! a single, coalesced [`FrameWaker`] fire — **without** the `effects` cargo
4//! feature and **without** the executor in the frame path.
5//!
6//! # How it works
7//!
8//! `reactive_graph`'s reactive graph is source/subscriber: reading a signal
9//! under an ambient [`Observer`] records the signal as a *source* of that
10//! observer and the observer as a *subscriber* of the signal (see
11//! `reactive_graph::traits::Track`). [`TrackedScope`] *is* such an observer: it
12//! implements [`Subscriber`]/[`ReactiveNode`], and [`TrackedScope::track`] runs
13//! the caller's closure with itself installed as the observer via
14//! [`WithObserver`]. Every signal read inside `track` therefore subscribes this
15//! scope to that signal.
16//!
17//! When any tracked source later changes it notifies this scope
18//! ([`ReactiveNode::mark_dirty`] for a directly-read signal, or
19//! [`ReactiveNode::mark_check`] for a value reached through a memo — a signal
20//! write marks the memo `mark_dirty`, and the memo relays `mark_check` to *its*
21//! subscribers). We treat **both** as "something changed": set the dirty flag
22//! and, on the clean→dirty edge only, fire the frame waker. That edge-triggering
23//! is the coalescing: N signal writes before the next `track` produce exactly
24//! one wake. `mark_check` waking is deliberately conservative — a check that
25//! resolves to no change wakes spuriously (wasteful but correct); a missed real
26//! change would be a bug.
27//!
28//! This is the *custom subscriber* mechanism, chosen over the `RenderEffect`
29//! fallback: it is synchronous, needs no `any_spawner` executor in the frame
30//! path, and the `reactive_graph::graph` traits are all publicly reachable, so
31//! there is no orphan/visibility wall forcing the fallback.
32
33use std::sync::Mutex;
34use std::sync::atomic::{AtomicBool, Ordering};
35use std::sync::{Arc, Weak};
36
37use reactive_graph::graph::{
38    AnySource, AnySubscriber, ReactiveNode, Source, Subscriber, WithObserver,
39};
40
41use crate::ReactiveRuntime;
42
43/// A reactive observer for a single frame-rebuild pass.
44///
45/// Run the rebuild under [`track`](Self::track); afterwards any change to a
46/// signal read during that pass sets the dirty flag ([`is_dirty`](Self::is_dirty))
47/// and fires the process-wide frame waker once (coalesced). Re-running `track`
48/// re-records dependencies from scratch, so a signal read in one pass but not
49/// the next stops dirtying the scope.
50///
51/// Cheap to clone-around by holding the inner `Arc`; construct one per
52/// long-lived rebuild target (e.g. one per app root / retained component).
53pub struct TrackedScope {
54    inner: Arc<ScopeInner>,
55}
56
57/// The shared, `Send + Sync` heart of a [`TrackedScope`]. Held behind an `Arc`
58/// so a `Weak<dyn Subscriber + Send + Sync>` can be handed to every source's
59/// subscriber list (dirty notifications may arrive from any thread).
60struct ScopeInner {
61    /// The sources (signals/memos) read during the last [`TrackedScope::track`].
62    /// `reactive_graph`'s own `SourceSet` is `pub(crate)`, so we keep our own
63    /// de-duplicated list; on re-track we walk it to unsubscribe from each
64    /// source before recording the new pass's reads.
65    sources: Mutex<Vec<AnySource>>,
66    /// Set when any tracked source notifies; cleared at the start of `track`.
67    /// The clean→dirty edge is what fires the waker (coalescing).
68    dirty: AtomicBool,
69}
70
71impl ScopeInner {
72    fn new() -> Self {
73        Self {
74            sources: Mutex::new(Vec::new()),
75            dirty: AtomicBool::new(false),
76        }
77    }
78
79    /// Marks the scope dirty and, on the clean→dirty edge only, fires the frame
80    /// waker. Callable from any thread (signal writes may originate anywhere).
81    ///
82    /// Also trips the process-wide signals-dirty flag
83    /// (`ReactiveRuntime::mark_signals_dirty`) on *every* call, not just the
84    /// clean→dirty edge below — that flag is a plain bool a shell drains once
85    /// per frame (`ReactiveRuntime::take_signals_dirty`), so any tracked
86    /// scope's invalidation should trip it, coalesced by nature since it has
87    /// no "already set" distinction to preserve.
88    fn notify_dirty(&self) {
89        let rt = ReactiveRuntime::get();
90        if let Some(rt) = rt {
91            rt.mark_signals_dirty();
92        }
93        // `swap` gives us the previous value atomically: only the thread that
94        // observed `false` (the clean→dirty transition) fires the waker, so N
95        // concurrent or sequential writes between tracks coalesce to one wake.
96        let was_dirty = self.dirty.swap(true, Ordering::SeqCst);
97        if !was_dirty && let Some(rt) = rt {
98            rt.wake();
99        }
100    }
101}
102
103impl ReactiveNode for ScopeInner {
104    /// A directly-read signal that changed notifies us here.
105    fn mark_dirty(&self) {
106        self.notify_dirty();
107    }
108
109    /// A value reached through a memo relays `mark_check` to us when the memo's
110    /// upstream changed. We can't cheaply prove the memo's output actually
111    /// changed without recomputing it, so we wake — spurious wakes are correct,
112    /// missed changes are not.
113    fn mark_check(&self) {
114        self.notify_dirty();
115    }
116
117    /// A [`TrackedScope`] is a leaf observer: nothing subscribes to *it*, so it
118    /// has no downstream to propagate a check to.
119    fn mark_subscribers_check(&self) {}
120
121    /// Reports whether the scope has been marked dirty since the last `track`.
122    fn update_if_necessary(&self) -> bool {
123        self.dirty.load(Ordering::SeqCst)
124    }
125}
126
127impl Subscriber for ScopeInner {
128    /// Records a signal/memo read during `track` as a source of this scope.
129    /// De-duplicated so reading the same signal twice in one pass records it
130    /// once.
131    fn add_source(&self, source: AnySource) {
132        let mut sources = self.sources.lock().expect("TrackedScope sources poisoned");
133        if !sources.contains(&source) {
134            sources.push(source);
135        }
136    }
137
138    /// Drops every recorded source, unsubscribing this scope from each so a
139    /// source read in a previous pass but not the current one can no longer
140    /// dirty us.
141    fn clear_sources(&self, subscriber: &AnySubscriber) {
142        let drained: Vec<AnySource> = {
143            let mut sources = self.sources.lock().expect("TrackedScope sources poisoned");
144            std::mem::take(&mut *sources)
145        };
146        for source in drained {
147            source.remove_subscriber(subscriber);
148        }
149    }
150}
151
152impl TrackedScope {
153    /// Creates an empty scope with no tracked sources and a clean dirty flag.
154    pub fn new() -> Self {
155        Self {
156            inner: Arc::new(ScopeInner::new()),
157        }
158    }
159
160    /// Builds the type-erased subscriber handle for this scope.
161    ///
162    /// The `reactive_graph::ToAnySubscriber` trait can't be implemented for
163    /// `Arc<ScopeInner>` here (orphan rule: both `Arc` and the trait are
164    /// foreign), so we construct the `AnySubscriber` directly — the same
165    /// `(ptr-as-usize, Weak<dyn Subscriber + Send + Sync>)` shape
166    /// `reactive_graph` uses internally. The `usize` is the stable allocation
167    /// pointer, which the graph's source/subscriber sets key identity on.
168    fn any_subscriber(&self) -> AnySubscriber {
169        AnySubscriber(
170            Arc::as_ptr(&self.inner) as usize,
171            Arc::downgrade(&self.inner) as Weak<dyn Subscriber + Send + Sync>,
172        )
173    }
174
175    /// Runs `f` with this scope installed as the reactive observer.
176    ///
177    /// Dependency tracking is re-recorded from scratch each call: previously
178    /// recorded sources are unsubscribed first, the dirty flag is cleared, then
179    /// every signal read inside `f` re-subscribes this scope. Returns whatever
180    /// `f` returns.
181    pub fn track<R>(&self, f: impl FnOnce() -> R) -> R {
182        let any = self.any_subscriber();
183        // Re-track from scratch: unsubscribe from the previous pass's sources so
184        // a signal no longer read this pass stops notifying us.
185        any.clear_sources(&any);
186        // Start the pass clean; any write *after* this that we still subscribe
187        // to will re-dirty us.
188        self.inner.dirty.store(false, Ordering::SeqCst);
189        any.with_observer(f)
190    }
191
192    /// Whether any tracked source has notified since the last [`track`](Self::track).
193    pub fn is_dirty(&self) -> bool {
194        self.inner.dirty.load(Ordering::SeqCst)
195    }
196}
197
198impl Default for TrackedScope {
199    fn default() -> Self {
200        Self::new()
201    }
202}
203
204#[cfg(test)]
205mod tests {
206    use super::*;
207    use reactive_graph::computed::Memo;
208    use reactive_graph::signal::RwSignal;
209    use reactive_graph::traits::{Get, Set};
210    use std::sync::atomic::{AtomicUsize, Ordering};
211    use std::time::Duration;
212
213    use crate::FrameWaker;
214
215    /// A recording waker: an `Arc<AtomicUsize>` bumped once per `wake()`.
216    fn recording_waker() -> (FrameWaker, Arc<AtomicUsize>) {
217        let counter = Arc::new(AtomicUsize::new(0));
218        let seen = counter.clone();
219        let waker: FrameWaker = Arc::new(move || {
220            counter.fetch_add(1, Ordering::SeqCst);
221        });
222        (waker, seen)
223    }
224
225    /// The whole tracked-scope scenario runs in one `#[test]` under the shared
226    /// waker lock: the frame waker is process-global and swappable, so any test
227    /// that installs a recording waker and counts wakes must not run concurrently
228    /// with another that swaps it (`runtime`'s end-to-end test does). Each block
229    /// maps to an acceptance criterion.
230    #[test]
231    fn tracked_scope_dirty_and_wake_bridge() {
232        let _guard = crate::WAKER_TEST_LOCK
233            .lock()
234            .unwrap_or_else(|e| e.into_inner());
235
236        let (waker, wakes) = recording_waker();
237        let rt = ReactiveRuntime::init(waker);
238
239        // Criterion 1: track a signal, then multiple writes → dirty + exactly one
240        // coalesced wake before the next track.
241        let scope = TrackedScope::new();
242        let sig = rt.with_owner(|| RwSignal::new(0));
243        scope.track(|| sig.get());
244        assert!(!scope.is_dirty(), "fresh track starts clean");
245        let before = wakes.load(Ordering::SeqCst);
246        sig.set(1);
247        sig.set(2);
248        sig.set(3);
249        assert!(scope.is_dirty(), "a tracked write must dirty the scope");
250        assert_eq!(
251            wakes.load(Ordering::SeqCst) - before,
252            1,
253            "N writes between tracks must coalesce to exactly one wake"
254        );
255
256        // Criterion 2: a signal NOT read inside track does not dirty the scope.
257        let other = rt.with_owner(|| RwSignal::new(0));
258        scope.track(|| sig.get()); // re-track: still only `sig`
259        assert!(!scope.is_dirty());
260        let before = wakes.load(Ordering::SeqCst);
261        other.set(99);
262        assert!(
263            !scope.is_dirty(),
264            "an untracked signal must not dirty the scope"
265        );
266        assert_eq!(
267            wakes.load(Ordering::SeqCst),
268            before,
269            "an untracked write must not wake"
270        );
271
272        // Criterion 3: a signal read in run N but not run N+1 no longer dirties
273        // (re-track clears the old sources).
274        scope.track(|| other.get()); // now tracking `other`, dropped `sig`
275        assert!(!scope.is_dirty());
276        let before = wakes.load(Ordering::SeqCst);
277        sig.set(4); // `sig` was dropped last re-track
278        assert!(
279            !scope.is_dirty(),
280            "a source dropped on re-track must no longer dirty the scope"
281        );
282        assert_eq!(
283            wakes.load(Ordering::SeqCst),
284            before,
285            "dropped source must not wake"
286        );
287        // ...but the newly-tracked `other` still dirties.
288        other.set(100);
289        assert!(
290            scope.is_dirty(),
291            "the freshly-tracked source must still dirty"
292        );
293
294        // Criterion 4: a write from a spawned OS thread dirties the scope and
295        // fires the waker.
296        let scope2 = TrackedScope::new();
297        let cross = rt.with_owner(|| RwSignal::new(0));
298        scope2.track(|| cross.get());
299        assert!(!scope2.is_dirty());
300        let before = wakes.load(Ordering::SeqCst);
301        std::thread::spawn(move || {
302            cross.set(7);
303        })
304        .join()
305        .expect("cross-thread writer panicked");
306        assert!(
307            scope2.is_dirty(),
308            "a write from another thread must dirty the scope"
309        );
310        assert_eq!(
311            wakes.load(Ordering::SeqCst) - before,
312            1,
313            "a cross-thread write must fire the waker once"
314        );
315
316        // Criterion 5: a memo chain (memo over a signal, memo read inside track,
317        // signal written) wakes. The signal write marks the memo dirty, and the
318        // memo relays `mark_check` to this scope.
319        let scope3 = TrackedScope::new();
320        let base = rt.with_owner(|| RwSignal::new(2));
321        let doubled = rt.with_owner(|| Memo::new(move |_| base.get() * 2));
322        let seen = scope3.track(|| doubled.get());
323        assert_eq!(seen, 4, "memo computes from the signal");
324        assert!(!scope3.is_dirty());
325        let before = wakes.load(Ordering::SeqCst);
326        base.set(5);
327        assert!(
328            scope3.is_dirty(),
329            "writing a memo's upstream signal must dirty the tracking scope"
330        );
331        assert!(
332            wakes.load(Ordering::SeqCst) > before,
333            "a memo-chain change must fire the waker"
334        );
335
336        // Sanity: give any late cross-thread notification a beat (there is none
337        // outstanding, but this keeps the test honest about ordering).
338        std::thread::sleep(Duration::from_millis(1));
339    }
340}