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}