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
114impl<T: Clone> Signal<T> {
115    /// Clone the current value and automatically track the read.
116    pub fn get(&self) -> T {
117        self.with(Clone::clone)
118    }
119
120    /// Clone the current value without dependency tracking.
121    pub fn get_untracked(&self) -> T {
122        self.with_untracked(Clone::clone)
123    }
124}
125
126impl<T: PartialEq> Signal<T> {
127    /// Replace the value, notifying only when it actually changes.
128    /// Retire the old value after notification, outside the internal borrow.
129    /// Inside a batch, notification queues observers until the batch completes.
130    pub fn set(&self, value: T) {
131        assert_not_computing();
132        crate::coherence::mutation("signal write");
133        let retired = {
134            let mut current = self.0.value.borrow_mut();
135            if *current == value {
136                None
137            } else {
138                Some(std::mem::replace(&mut *current, value))
139            }
140        };
141        if retired.is_some() {
142            notify(&self.0.source);
143        }
144        drop(retired);
145    }
146}
147
148/// A lazy computation. Reads track its underlying signals in the consuming effect.
149/// It is deliberately not a cache: each `get` evaluates the function once.
150pub struct Derived<T>(Rc<dyn Fn() -> T>);
151
152impl<T> Clone for Derived<T> {
153    fn clone(&self) -> Self {
154        Self(self.0.clone())
155    }
156}
157
158pub fn derived<T>(compute: impl Fn() -> T + 'static) -> Derived<T> {
159    Derived(Rc::new(compute))
160}
161
162impl<T> Derived<T> {
163    pub fn get(&self) -> T {
164        (self.0)()
165    }
166}
167
168struct EffectInner {
169    observer: Rc<Observer>,
170    active: Cell<bool>,
171    queued: Cell<bool>,
172    last_flush: Cell<u64>,
173    flush_runs: Cell<u32>,
174    callback: RefCell<Box<dyn FnMut()>>,
175    lifecycle: RefCell<Vec<crate::Registration>>,
176}
177
178impl EffectInner {
179    fn unsubscribe(&self) {
180        self.observer.unsubscribe();
181    }
182
183    fn run(self: &Rc<Self>, initial: bool) {
184        if !self.active.get() {
185            return;
186        }
187        // Invalidated memos are checked lazily before deciding whether the
188        // effect's observed values actually changed.
189        if !initial && !untrack(|| self.observer.changed()) {
190            return;
191        }
192        // Recollect on every run, so conditional reads shed stale dependencies.
193        self.unsubscribe();
194        let _tracking = TrackingGuard::replace(Some(Rc::downgrade(&self.observer)));
195        (self.callback.borrow_mut())();
196    }
197}
198
199/// Owns a subscription. Dropping it detaches every dependency.
200#[must_use = "retain the effect handle for as long as the subscription should live"]
201pub struct Effect(Rc<EffectInner>);
202
203impl Effect {
204    #[cfg(feature = "dom")]
205    pub(crate) fn initializer(&self) -> impl FnOnce() + 'static {
206        let weak = Rc::downgrade(&self.0);
207        move || {
208            if let Some(inner) = weak.upgrade() {
209                batch(|| inner.run(true));
210            }
211        }
212    }
213
214    /// Stop reacting. Safe to call more than once.
215    pub fn dispose(&self) {
216        self.0.active.set(false);
217        self.0.unsubscribe();
218    }
219}
220
221impl Drop for Effect {
222    fn drop(&mut self) {
223        self.dispose();
224    }
225}
226
227/// Run immediately, then rerun when any signal read by the callback changes.
228/// Effects are synchronous outside a batch; queued reruns are deduplicated.
229fn allocate_effect(callback: impl FnMut() + 'static) -> Effect {
230    assert_not_computing();
231    Effect(Rc::new_cyclic(|weak| EffectInner {
232        observer: Observer::new(ObserverKind::Effect(weak.clone())),
233        active: Cell::new(true),
234        queued: Cell::new(false),
235        last_flush: Cell::new(0),
236        flush_runs: Cell::new(0),
237        callback: RefCell::new(Box::new(callback)),
238        lifecycle: RefCell::new(Vec::new()),
239    }))
240}
241
242#[cfg(feature = "dom")]
243pub(crate) fn prepared_effect(callback: impl FnMut() + 'static) -> Effect {
244    allocate_effect(callback)
245}
246
247pub fn effect(callback: impl FnMut() + 'static) -> Effect {
248    let subscription = allocate_effect(callback);
249    // Initial effects still run immediately inside a batch. Their writes flush
250    // after the initial callback returns, avoiding recursive RefCell borrows.
251    if let Some(owner) = crate::coherence::preparing_owner().filter(|owner| !owner.is_active()) {
252        let weak = Rc::downgrade(&subscription.0);
253        let activation = owner.on_activate(move || {
254            if let Some(inner) = weak.upgrade() {
255                batch(|| inner.run(true));
256            }
257        });
258        let weak = Rc::downgrade(&subscription.0);
259        let cleanup = owner.on_cleanup(move || {
260            if let Some(inner) = weak.upgrade() {
261                inner.active.set(false);
262                inner.unsubscribe();
263            }
264        });
265        subscription
266            .0
267            .lifecycle
268            .borrow_mut()
269            .extend([activation, cleanup]);
270    } else if !crate::coherence::mutation("effect creation") {
271        batch(|| subscription.0.run(true));
272    }
273    subscription
274}
275
276fn notify(source: &Source) {
277    source.advance();
278    // Invalidate the entire reachable graph before executing any effects.
279    source.notify();
280    flush();
281}
282
283fn clear_queue() {
284    #[cfg(feature = "javascript")]
285    AFTER_FLUSH.with(|queue| queue.borrow_mut().clear());
286    QUEUE.with(|queue| {
287        for pending in queue
288            .borrow_mut()
289            .drain(..)
290            .filter_map(|item| item.upgrade())
291        {
292            pending.queued.set(false);
293        }
294    });
295}
296
297struct FlushGuard;
298
299impl Drop for FlushGuard {
300    fn drop(&mut self) {
301        FLUSHING.with(|flushing| flushing.set(false));
302        if std::thread::panicking() {
303            clear_queue();
304        }
305    }
306}
307
308fn flush() {
309    if BATCH_DEPTH.with(Cell::get) > 0 || FLUSHING.with(|flushing| flushing.replace(true)) {
310        return;
311    }
312    let _guard = FlushGuard;
313    let epoch = NEXT_FLUSH.with(|next| {
314        let epoch = next
315            .get()
316            .checked_add(1)
317            .expect("reactive flush ID exhausted");
318        next.set(epoch);
319        epoch
320    });
321    loop {
322        let next = QUEUE.with(|queue| queue.borrow_mut().pop_front());
323        let Some(next) = next else {
324            #[cfg(feature = "javascript")]
325            {
326                let callback = AFTER_FLUSH.with(|queue| queue.borrow_mut().pop_front());
327                if let Some(callback) = callback {
328                    untrack(callback);
329                    continue;
330                }
331            }
332            break;
333        };
334        if let Some(next) = next.upgrade() {
335            next.queued.set(false);
336            // Width and acyclic propagation depth are not feedback loops.
337            // Limit repeated execution of each effect, without a per-flush map.
338            let runs = if next.last_flush.replace(epoch) == epoch {
339                next.flush_runs.get() + 1
340            } else {
341                1
342            };
343            next.flush_runs.set(runs);
344            assert!(
345                runs <= 10_000,
346                "reactive cycle: one effect exceeded 10,000 runs in one flush"
347            );
348            next.run(false);
349        }
350    }
351}
352
353/// Schedule foreign notifications after ordinary reactive work has settled.
354/// The queue never lends a borrow across callbacks; their writes join the next
355/// wave of the current flush rather than recursively running an observer.
356#[cfg(feature = "javascript")]
357pub(crate) fn after_flush(callback: impl FnOnce() + 'static) {
358    AFTER_FLUSH.with(|queue| queue.borrow_mut().push_back(Box::new(callback)));
359}
360
361struct TrackingGuard(Option<Weak<Observer>>);
362
363impl TrackingGuard {
364    fn replace(next: Option<Weak<Observer>>) -> Self {
365        Self(CURRENT.with(|current| current.replace(next)))
366    }
367}
368
369impl Drop for TrackingGuard {
370    fn drop(&mut self) {
371        CURRENT.with(|current| current.replace(self.0.take()));
372    }
373}
374
375/// Evaluate a closure without collecting dependencies in the current effect.
376pub fn untrack<R>(read: impl FnOnce() -> R) -> R {
377    let _guard = TrackingGuard::replace(None);
378    read()
379}
380
381struct BatchGuard;
382
383impl Drop for BatchGuard {
384    fn drop(&mut self) {
385        BATCH_DEPTH.with(|depth| depth.set(depth.get() - 1));
386        if std::thread::panicking() {
387            clear_queue();
388        } else {
389            flush();
390        }
391    }
392}
393
394/// Group synchronous writes into one notification per affected effect.
395/// Nested batches flush at the outer boundary. This is not a rollback transaction.
396pub fn batch<R>(update: impl FnOnce() -> R) -> R {
397    BATCH_DEPTH.with(|depth| depth.set(depth.get() + 1));
398    let _guard = BatchGuard;
399    update()
400}