Skip to main content

telar_reactive_core/
lib.rs

1//! Factory functions (`signal`, `effect`, `memo`) create
2//! nodes in the reactive graph. Struct constructors (`Runtime::new`, etc.) own
3//! their state. Free functions (`batch`, `set_flush_notify`) operate on the
4//! thread-local runtime.
5
6mod effect;
7mod memo;
8pub use reactive_local::reentry;
9mod runtime;
10mod signal;
11mod source;
12#[macro_use]
13
14mod task;
15
16pub use effect::{Effect, effect};
17pub use memo::{Memo, memo};
18pub use reactive_local::{SurfaceSlot, surface_local};
19pub use runtime::{
20    FlushNotifyHandle, SurfaceEnterGuard, SurfaceHandle, batch, begin_batch, current_surface,
21    end_batch, reset_runtime, set_current_surface, set_flush_notify, set_surface_enter_hook,
22};
23pub use signal::{ReadSignal, RwSignal, signal};
24pub use source::{Source, derive, derive_pair};
25pub use task::{
26    Emitter, Task, cancel_tasks_for, drain_tasks, reset_tasks, set_task_waker, spawn_stream,
27    spawn_task,
28};
29
30#[cfg(test)]
31mod tests {
32    use std::cell::RefCell;
33    use std::rc::Rc;
34
35    use super::*;
36
37    #[test]
38    fn signal_get_set() {
39        let count = signal(0i32);
40        assert_eq!(count.get(), 0);
41        count.set(42);
42        assert_eq!(count.get(), 42);
43    }
44
45    // M3 owner-scope: the shared runtime stamps each effect with the surface active at registration, and the
46    // flush re-enters that surface before running the effect — even when the write that scheduled it happened
47    // under a different active surface. Here the enter hook records which surface each run resolved against.
48    #[test]
49    fn effect_runs_under_its_own_surface_context() {
50        use std::cell::RefCell;
51        use std::rc::Rc;
52
53        let entered: Rc<RefCell<Vec<u64>>> = Rc::new(RefCell::new(Vec::new()));
54        let entered_hook = Rc::clone(&entered);
55        set_surface_enter_hook(move |handle| {
56            let prev = set_current_surface(handle);
57            entered_hook.borrow_mut().push(handle.0);
58            SurfaceEnterGuard::new(move || {
59                set_current_surface(prev);
60            })
61        });
62
63        // Build an effect "owned by" surface A: A is active while it registers, so it captures A.
64        let trigger = signal(0i32);
65        let read = trigger.read_only();
66        let _guard_a = SurfaceHandle(1).enter();
67        let _e = effect(move || {
68            read.get();
69        });
70        drop(_guard_a);
71
72        // Back on the ambient surface, no A entry has been recorded beyond the build itself; clear the log so
73        // we observe only what the *flush-triggered* run enters.
74        entered.borrow_mut().clear();
75
76        // Write the signal while surface B is active. The scheduled effect belongs to A, so the flush must
77        // enter A (1), not B (2), before running it.
78        let _guard_b = SurfaceHandle(2).enter();
79        trigger.set(1);
80        drop(_guard_b);
81
82        assert!(
83            entered.borrow().contains(&1),
84            "flush must re-enter the effect's own surface (A=1): {:?}",
85            entered.borrow()
86        );
87    }
88
89    // A panic inside `batch` must leave the shared runtime consistent (batch_depth/flushing reset), so a
90    // later write still schedules and flushes effects. Without the RAII guards this would wedge the runtime.
91    #[test]
92    fn runtime_recovers_after_panic_in_batch() {
93        use std::cell::RefCell;
94        use std::panic::{AssertUnwindSafe, catch_unwind};
95        use std::rc::Rc;
96
97        let count = signal(0i32);
98        let read = count.read_only();
99        let seen: Rc<RefCell<Vec<i32>>> = Rc::new(RefCell::new(Vec::new()));
100        let seen_c = Rc::clone(&seen);
101        let _e = effect(move || {
102            seen_c.borrow_mut().push(read.get());
103        });
104
105        let result = catch_unwind(AssertUnwindSafe(|| {
106            batch(|| {
107                count.set(1);
108                panic!("boom");
109            });
110        }));
111        assert!(result.is_err(), "the batch closure should have panicked");
112
113        count.set(2);
114        assert!(
115            seen.borrow().contains(&2),
116            "runtime wedged after panic-in-batch; effect never re-ran: {:?}",
117            seen.borrow()
118        );
119    }
120
121    // A cascade inside one flush: `writer` runs after `reader` and writes the signal `reader` depends on, so
122    // `reader` has to run again. It is scheduled either way — the regression was the flush-wide epoch stamp,
123    // which turned that second run into a no-op with nothing left to reschedule it, so the reader kept the
124    // stale value until some later, unrelated write opened a fresh flush.
125    #[test]
126    fn an_effect_reruns_when_a_later_effect_in_the_same_flush_writes_its_source() {
127        let trigger = signal(0i32);
128        let source = signal(0i32);
129        let seen: Rc<RefCell<Vec<i32>>> = Rc::new(RefCell::new(Vec::new()));
130
131        // The reader is scheduled by the same write as the writer, and runs first — so by the time the writer
132        // moves `source`, the reader has already run once in this flush.
133        let read_trigger = trigger.read_only();
134        let read_source = source.read_only();
135        let seen_c = Rc::clone(&seen);
136        let _reader = effect(move || {
137            read_trigger.get();
138            seen_c.borrow_mut().push(read_source.get());
139        });
140
141        let read_trigger = trigger.read_only();
142        let write_source = source.clone();
143        let _writer = effect(move || {
144            let v = read_trigger.get();
145            if v > 0 {
146                write_source.set(v);
147            }
148        });
149
150        seen.borrow_mut().clear();
151        trigger.set(7);
152        assert_eq!(
153            seen.borrow().last().copied(),
154            Some(7),
155            "the reader never saw a write made later in the same flush: {:?}",
156            seen.borrow()
157        );
158    }
159
160    #[test]
161    fn signal_update() {
162        let count = signal(10i32);
163        count.update(|v| *v *= 2);
164        assert_eq!(count.get(), 20);
165    }
166
167    #[test]
168    fn signal_with() {
169        let name = signal(String::from("rsx"));
170        let len = name.with(|s| s.len());
171        assert_eq!(len, 3);
172    }
173
174    #[test]
175    fn rw_signal() {
176        let count = signal(0i32);
177        count.set(10);
178        assert_eq!(count.get(), 10);
179        count.update(|v| *v += 5);
180        assert_eq!(count.get(), 15);
181    }
182
183    #[test]
184    fn rw_signal_read_only() {
185        let sig = signal(0i32);
186        let read = sig.read_only();
187        sig.set(7);
188        assert_eq!(read.get(), 7);
189    }
190
191    #[test]
192    fn bool_signal_toggle() {
193        let flag = signal(false);
194        flag.toggle();
195        assert!(flag.get());
196        flag.toggle();
197        assert!(!flag.get());
198    }
199
200    #[test]
201    fn effect_runs_immediately() {
202        let ran = Rc::new(RefCell::new(false));
203        let ran_clone = Rc::clone(&ran);
204        let _e = effect(move || {
205            *ran_clone.borrow_mut() = true;
206        });
207        assert!(*ran.borrow());
208    }
209
210    #[test]
211    fn effect_reruns_on_signal_change() {
212        let count = signal(0i32);
213        let log: Rc<RefCell<Vec<i32>>> = Rc::new(RefCell::new(Vec::new()));
214        let log_clone = Rc::clone(&log);
215
216        let read = count.read_only();
217        let _e = effect(move || {
218            log_clone.borrow_mut().push(read.get());
219        });
220
221        count.set(1);
222        count.set(2);
223
224        assert_eq!(*log.borrow(), vec![0, 1, 2]);
225    }
226
227    #[test]
228    fn memo_derives_value() {
229        let count = signal(2i32);
230        let read = count.read_only();
231        let doubled = memo(move || read.get() * 2);
232
233        assert_eq!(doubled.get(), 4);
234        count.set(5);
235        assert_eq!(doubled.get(), 10);
236    }
237
238    #[test]
239    fn effect_reruns_when_memo_changes() {
240        let count = signal(0i32);
241        let read = count.read_only();
242        let doubled = memo(move || read.get() * 2);
243        let log: Rc<RefCell<Vec<i32>>> = Rc::new(RefCell::new(Vec::new()));
244        let log_clone = Rc::clone(&log);
245        let doubled_read = doubled.clone();
246        let _e = effect(move || {
247            log_clone.borrow_mut().push(doubled_read.get());
248        });
249        assert_eq!(*log.borrow(), vec![0]);
250        count.set(3);
251        assert_eq!(*log.borrow(), vec![0, 6]);
252    }
253
254    // Regression: an effect tracking BOTH a signal and a memo must re-run when only the memo's source changes — run_effect's signal-version shortcut used to skip it because memo deps are invisible to `sources` (the sandbox counter's frozen "Double:" text).
255    #[test]
256    fn effect_with_signal_and_memo_sources_reruns_on_memo_change() {
257        let unrelated = signal(0i32);
258        let count = signal(0i32);
259        let read = count.read_only();
260        let doubled = memo(move || read.get() * 2);
261        let log: Rc<RefCell<Vec<i32>>> = Rc::new(RefCell::new(Vec::new()));
262        let log_clone = Rc::clone(&log);
263        let unrelated_read = unrelated.read_only();
264        let doubled_read = doubled.clone();
265        let _e = effect(move || {
266            unrelated_read.get();
267            log_clone.borrow_mut().push(doubled_read.get());
268        });
269        assert_eq!(*log.borrow(), vec![0]);
270        count.set(3);
271        assert_eq!(*log.borrow(), vec![0, 6]);
272    }
273
274    #[test]
275    fn effect_reruns_when_memo_changes_inside_batch() {
276        let count = signal(0i32);
277        let read = count.read_only();
278        let doubled = memo(move || read.get() * 2);
279        let log: Rc<RefCell<Vec<i32>>> = Rc::new(RefCell::new(Vec::new()));
280        let log_clone = Rc::clone(&log);
281        let doubled_read = doubled.clone();
282        let _e = effect(move || {
283            log_clone.borrow_mut().push(doubled_read.get());
284        });
285        batch(|| count.set(3));
286        assert_eq!(*log.borrow(), vec![0, 6]);
287    }
288
289    #[test]
290    fn memo_chains() {
291        let n = signal(3i32);
292        let read = n.read_only();
293        let doubled = memo(move || read.get() * 2);
294        let doubled_for_quad = doubled.clone();
295        let quadrupled = memo(move || doubled_for_quad.get() * 2);
296
297        assert_eq!(quadrupled.get(), 12);
298        n.set(5);
299        assert_eq!(quadrupled.get(), 20);
300    }
301
302    #[test]
303    fn batch_fires_effect_once() {
304        let a = signal(0i32);
305        let b = signal(0i32);
306        let runs = Rc::new(RefCell::new(0usize));
307        let runs_clone = Rc::clone(&runs);
308
309        let a_read = a.read_only();
310        let b_read = b.read_only();
311        let _e = effect(move || {
312            let _ = a_read.get() + b_read.get();
313            *runs_clone.borrow_mut() += 1;
314        });
315
316        assert_eq!(*runs.borrow(), 1);
317
318        batch(|| {
319            a.set(1);
320            b.set(2);
321        });
322
323        assert_eq!(*runs.borrow(), 2);
324    }
325
326    #[test]
327    fn dropping_one_subscriber_keeps_others_consistent() {
328        // Three effects subscribe to the same signal. Dropping the middle one and then writing repeatedly must keep the survivors firing and must not panic — exercising the in-place subscriber-list cleanup that runs alongside notify_signal's reused scratch buffer (the cleanup must stay correct now that subscribers are no longer cloned per write).
329        let count = signal(0i32);
330        let a = Rc::new(RefCell::new(0i32));
331        let b = Rc::new(RefCell::new(0i32));
332        let c = Rc::new(RefCell::new(0i32));
333
334        let mk = |sink: &Rc<RefCell<i32>>, sig: &RwSignal<i32>| {
335            let read = sig.read_only();
336            let sink = Rc::clone(sink);
337            effect(move || {
338                *sink.borrow_mut() = read.get();
339            })
340        };
341
342        let _ea = mk(&a, &count);
343        let eb = mk(&b, &count);
344        let _ec = mk(&c, &count);
345
346        count.set(1);
347        assert_eq!((*a.borrow(), *b.borrow(), *c.borrow()), (1, 1, 1));
348
349        drop(eb);
350        count.set(2);
351        count.set(3);
352
353        // Survivors tracked every write; the dropped one is frozen at its last value, with no panic and no lost survivor.
354        assert_eq!(*a.borrow(), 3);
355        assert_eq!(*c.borrow(), 3);
356        assert_eq!(*b.borrow(), 1);
357    }
358
359    // A signal whose value owns other signal handles must drop cleanly: dropping the outer signal removes its
360    // storage and then drops the value, which re-enters `drop_signal` for the inner handles. If the value were
361    // dropped while the runtime borrow was still held, that re-entry would abort during teardown.
362    #[test]
363    fn dropping_a_signal_whose_value_holds_signals_does_not_double_borrow() {
364        let inner = signal(1i32);
365        let outer = signal(vec![inner.clone()]);
366        drop(inner);
367        drop(outer);
368        // Reaching here without aborting is the assertion; the runtime is still usable afterwards.
369        let after = signal(5i32);
370        assert_eq!(after.get(), 5);
371    }
372}