Skip to main content

fusor/
reactive.rs

1use std::{
2    cell::{Cell, RefCell},
3    collections::VecDeque,
4    rc::{Rc, Weak},
5};
6
7mod graph;
8mod memo;
9pub mod versions;
10use graph::{Observer, ObserverKind, Source, track};
11pub use memo::{Memo, memo, memo_with_eq};
12
13thread_local! {
14    static CURRENT: RefCell<Option<Weak<Observer>>> = const { RefCell::new(None) };
15    static QUEUE: RefCell<VecDeque<Weak<EffectInner>>> = const { RefCell::new(VecDeque::new()) };
16    static BATCH_DEPTH: Cell<usize> = const { Cell::new(0) };
17    static FLUSHING: Cell<bool> = const { Cell::new(false) };
18    static NEXT_ID: Cell<u64> = const { Cell::new(0) };
19    static NEXT_FLUSH: Cell<u64> = const { Cell::new(0) };
20    static NEXT_WAVE: Cell<u64> = const { Cell::new(0) };
21    static COMPUTING: Cell<usize> = const { Cell::new(0) };
22    #[cfg(feature = "javascript")]
23    static AFTER_FLUSH: RefCell<VecDeque<Box<dyn FnOnce()>>> = const { RefCell::new(VecDeque::new()) };
24}
25
26fn assert_not_computing() {
27    assert!(
28        COMPUTING.with(Cell::get) == 0,
29        "memo computations and equality functions must not write signals or create effects"
30    );
31}
32
33struct SignalInner<T> {
34    value: RefCell<T>,
35    source: Rc<Source>,
36    render_values: RefCell<Vec<(Rc<T>, versions::Versions)>>,
37}
38
39/// Shared, single-threaded reactive state. Cloning shares the same value.
40pub struct Signal<T>(Rc<SignalInner<T>>);
41
42impl<T> Clone for Signal<T> {
43    fn clone(&self) -> Self {
44        Self(self.0.clone())
45    }
46}
47
48/// Create reactive state; reads inside an effect automatically subscribe to it.
49pub fn signal<T>(value: T) -> Signal<T> {
50    Signal(Rc::new(SignalInner {
51        value: RefCell::new(value),
52        source: Rc::new(Source::new(None)),
53        render_values: RefCell::new(Vec::new()),
54    }))
55}
56
57impl<T> Signal<T> {
58    /// Read without cloning the value, tracking this dependency.
59    /// Do not write this signal while the read closure holds its borrow.
60    pub fn with<R>(&self, read: impl FnOnce(&T) -> R) -> R {
61        let render_value = self.0.render_values.borrow().last().cloned();
62        if let Some((value, inputs)) = render_value {
63            versions::Versions::exclude(|| track(&self.0.source));
64            inputs.include();
65            return read(&value);
66        }
67        track(&self.0.source);
68        read(&self.0.value.borrow())
69    }
70
71    /// Read without subscribing the current effect.
72    pub fn with_untracked<R>(&self, read: impl FnOnce(&T) -> R) -> R {
73        let render_value = self.0.render_values.borrow().last().cloned();
74        if let Some((value, _)) = render_value {
75            return read(&value);
76        }
77        read(&self.0.value.borrow())
78    }
79
80    /// A generated keyed row's candidate input. The override is synchronous,
81    /// never published to observers, and validates against the collection's
82    /// actual source versions. This is not historical storage for signals.
83    #[doc(hidden)]
84    pub fn with_render_value<R>(
85        &self,
86        value: Rc<T>,
87        inputs: versions::Versions,
88        render: impl FnOnce() -> R,
89    ) -> R {
90        struct Pop<'a, T>(&'a Signal<T>);
91        impl<T> Drop for Pop<'_, T> {
92            fn drop(&mut self) {
93                let value = self.0.0.render_values.borrow_mut().pop();
94                drop(value);
95            }
96        }
97        self.0.render_values.borrow_mut().push((value, inputs));
98        let _pop = Pop(self);
99        render()
100    }
101
102    /// Mutate state, then notify subscribers after releasing the mutable borrow.
103    /// Always notifies; use `set` to skip unchanged values.
104    /// The mutation closure must not read or write this same signal.
105    pub fn update<R>(&self, update: impl FnOnce(&mut T) -> R) -> R {
106        assert_not_computing();
107        crate::coherence::mutation("signal write");
108        let result = update(&mut self.0.value.borrow_mut());
109        notify(&self.0.source);
110        result
111    }
112
113    /// Replace the value and notify subscribers, returning the previous value.
114    /// Unlike `set`, this needs no `PartialEq` and always notifies. Dropping
115    /// the result retires the old value after notification.
116    pub fn replace(&self, value: T) -> T {
117        self.update(|current| std::mem::replace(current, value))
118    }
119}
120
121impl<T: Clone> Signal<T> {
122    /// Clone the current value and automatically track the read.
123    pub fn get(&self) -> T {
124        self.with(Clone::clone)
125    }
126
127    /// Clone the current value without dependency tracking.
128    pub fn get_untracked(&self) -> T {
129        self.with_untracked(Clone::clone)
130    }
131}
132
133impl<T: PartialEq> Signal<T> {
134    /// Replace the value, notifying only when it actually changes.
135    /// Retire the old value after notification, outside the internal borrow.
136    /// Inside a batch, notification queues observers until the batch completes.
137    pub fn set(&self, value: T) {
138        assert_not_computing();
139        crate::coherence::mutation("signal write");
140        let retired = {
141            let mut current = self.0.value.borrow_mut();
142            if *current == value {
143                None
144            } else {
145                Some(std::mem::replace(&mut *current, value))
146            }
147        };
148        if retired.is_some() {
149            notify(&self.0.source);
150        }
151        drop(retired);
152    }
153}
154
155/// A lazy computation. Reads track its underlying signals in the consuming effect.
156/// It is deliberately not a cache: each `get` evaluates the function once.
157pub struct Derived<T>(Rc<dyn Fn() -> T>);
158
159impl<T> Clone for Derived<T> {
160    fn clone(&self) -> Self {
161        Self(self.0.clone())
162    }
163}
164
165pub fn derived<T>(compute: impl Fn() -> T + 'static) -> Derived<T> {
166    Derived(Rc::new(compute))
167}
168
169impl<T> Derived<T> {
170    pub fn get(&self) -> T {
171        (self.0)()
172    }
173}
174
175struct EffectInner {
176    observer: Rc<Observer>,
177    active: Cell<bool>,
178    queued: Cell<bool>,
179    last_flush: Cell<u64>,
180    flush_runs: Cell<u32>,
181    callback: RefCell<Box<dyn FnMut()>>,
182    lifecycle: RefCell<Vec<crate::Registration>>,
183}
184
185impl EffectInner {
186    fn unsubscribe(&self) {
187        self.observer.unsubscribe();
188    }
189
190    fn run(self: &Rc<Self>, initial: bool) {
191        if !self.active.get() {
192            return;
193        }
194        // Invalidated memos are checked lazily before deciding whether the
195        // effect's observed values actually changed.
196        if !initial && !untrack(|| self.observer.changed()) {
197            return;
198        }
199        // Recollect on every run, so conditional reads shed stale dependencies.
200        let _run = self.observer.begin();
201        let _tracking = TrackingGuard::replace(Some(Rc::downgrade(&self.observer)));
202        (self.callback.borrow_mut())();
203    }
204}
205
206/// Owns a subscription. Dropping it detaches every dependency.
207#[must_use = "retain the effect handle for as long as the subscription should live"]
208pub struct Effect(Rc<EffectInner>);
209
210impl Effect {
211    #[cfg(feature = "dom")]
212    pub(crate) fn initializer(&self) -> impl FnOnce() + 'static {
213        let weak = Rc::downgrade(&self.0);
214        move || {
215            if let Some(inner) = weak.upgrade() {
216                batch(|| inner.run(true));
217            }
218        }
219    }
220
221    /// Stop reacting. Safe to call more than once.
222    pub fn dispose(&self) {
223        self.0.active.set(false);
224        self.0.unsubscribe();
225    }
226}
227
228impl Drop for Effect {
229    fn drop(&mut self) {
230        self.dispose();
231    }
232}
233
234/// Run immediately, then rerun when any signal read by the callback changes.
235/// Effects are synchronous outside a batch; queued reruns are deduplicated.
236fn allocate_effect(callback: impl FnMut() + 'static) -> Effect {
237    assert_not_computing();
238    Effect(Rc::new_cyclic(|weak| EffectInner {
239        observer: Observer::new(ObserverKind::Effect(weak.clone())),
240        active: Cell::new(true),
241        queued: Cell::new(false),
242        last_flush: Cell::new(0),
243        flush_runs: Cell::new(0),
244        callback: RefCell::new(Box::new(callback)),
245        lifecycle: RefCell::new(Vec::new()),
246    }))
247}
248
249#[cfg(feature = "dom")]
250pub(crate) fn prepared_effect(callback: impl FnMut() + 'static) -> Effect {
251    allocate_effect(callback)
252}
253
254pub fn effect(callback: impl FnMut() + 'static) -> Effect {
255    let subscription = allocate_effect(callback);
256    // Initial effects still run immediately inside a batch. Their writes flush
257    // after the initial callback returns, avoiding recursive RefCell borrows.
258    if let Some(owner) = crate::coherence::preparing_owner().filter(|owner| !owner.is_active()) {
259        let weak = Rc::downgrade(&subscription.0);
260        let activation = owner.on_activate(move || {
261            if let Some(inner) = weak.upgrade() {
262                batch(|| inner.run(true));
263            }
264        });
265        let weak = Rc::downgrade(&subscription.0);
266        let cleanup = owner.on_cleanup(move || {
267            if let Some(inner) = weak.upgrade() {
268                inner.active.set(false);
269                inner.unsubscribe();
270            }
271        });
272        subscription
273            .0
274            .lifecycle
275            .borrow_mut()
276            .extend([activation, cleanup]);
277    } else if !crate::coherence::mutation("effect creation") {
278        batch(|| subscription.0.run(true));
279    }
280    subscription
281}
282
283fn notify(source: &Source) {
284    source.advance();
285    // Invalidate the entire reachable graph before executing any effects.
286    source.notify();
287    flush();
288}
289
290fn clear_queue() {
291    #[cfg(feature = "javascript")]
292    AFTER_FLUSH.with(|queue| queue.borrow_mut().clear());
293    QUEUE.with(|queue| {
294        for pending in queue
295            .borrow_mut()
296            .drain(..)
297            .filter_map(|item| item.upgrade())
298        {
299            pending.queued.set(false);
300        }
301    });
302}
303
304struct FlushGuard;
305
306impl Drop for FlushGuard {
307    fn drop(&mut self) {
308        FLUSHING.with(|flushing| flushing.set(false));
309        if std::thread::panicking() {
310            clear_queue();
311        }
312    }
313}
314
315fn flush() {
316    if BATCH_DEPTH.with(Cell::get) > 0 || FLUSHING.with(|flushing| flushing.replace(true)) {
317        return;
318    }
319    let _guard = FlushGuard;
320    let epoch = NEXT_FLUSH.with(|next| {
321        let epoch = next
322            .get()
323            .checked_add(1)
324            .expect("reactive flush ID exhausted");
325        next.set(epoch);
326        epoch
327    });
328    loop {
329        let next = QUEUE.with(|queue| queue.borrow_mut().pop_front());
330        let Some(next) = next else {
331            #[cfg(feature = "javascript")]
332            {
333                let callback = AFTER_FLUSH.with(|queue| queue.borrow_mut().pop_front());
334                if let Some(callback) = callback {
335                    untrack(callback);
336                    continue;
337                }
338            }
339            break;
340        };
341        if let Some(next) = next.upgrade() {
342            next.queued.set(false);
343            // Width and acyclic propagation depth are not feedback loops.
344            // Limit repeated execution of each effect, without a per-flush map.
345            let runs = if next.last_flush.replace(epoch) == epoch {
346                next.flush_runs.get() + 1
347            } else {
348                1
349            };
350            next.flush_runs.set(runs);
351            assert!(
352                runs <= 10_000,
353                "reactive cycle: one effect exceeded 10,000 runs in one flush"
354            );
355            next.run(false);
356        }
357    }
358}
359
360/// Schedule foreign notifications after ordinary reactive work has settled.
361/// The queue never lends a borrow across callbacks; their writes join the next
362/// wave of the current flush rather than recursively running an observer.
363#[cfg(feature = "javascript")]
364pub(crate) fn after_flush(callback: impl FnOnce() + 'static) {
365    AFTER_FLUSH.with(|queue| queue.borrow_mut().push_back(Box::new(callback)));
366}
367
368struct TrackingGuard(Option<Weak<Observer>>);
369
370impl TrackingGuard {
371    fn replace(next: Option<Weak<Observer>>) -> Self {
372        Self(CURRENT.with(|current| current.replace(next)))
373    }
374}
375
376impl Drop for TrackingGuard {
377    fn drop(&mut self) {
378        CURRENT.with(|current| current.replace(self.0.take()));
379    }
380}
381
382/// Evaluate a closure without collecting dependencies in the current effect.
383pub fn untrack<R>(read: impl FnOnce() -> R) -> R {
384    let _guard = TrackingGuard::replace(None);
385    read()
386}
387
388struct BatchGuard;
389
390impl Drop for BatchGuard {
391    fn drop(&mut self) {
392        BATCH_DEPTH.with(|depth| depth.set(depth.get() - 1));
393        if std::thread::panicking() {
394            clear_queue();
395        } else {
396            flush();
397        }
398    }
399}
400
401/// Group synchronous writes into one notification per affected effect.
402/// Nested batches flush at the outer boundary. This is not a rollback transaction.
403pub fn batch<R>(update: impl FnOnce() -> R) -> R {
404    BATCH_DEPTH.with(|depth| depth.set(depth.get() + 1));
405    let _guard = BatchGuard;
406    update()
407}