use crate::{Computed, Effect, Signal, Transaction, flush_effects};
use std::collections::HashSet;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
struct TestCircuit {
voltage_signal: Signal,
}
impl TestCircuit {
fn new() -> Self {
Self {
voltage_signal: Signal::new(),
}
}
}
#[test]
fn computed_caches_until_invalidated() {
let _circuit = Arc::new(TestCircuit::new());
let computation_count = Arc::new(AtomicUsize::new(0));
let computation_count_clone = computation_count.clone();
let power = Computed::new(move || {
computation_count_clone.fetch_add(1, Ordering::Relaxed);
100.0
});
let _p1 = power.get();
assert_eq!(computation_count.load(Ordering::Relaxed), 1);
let _p2 = power.get();
assert_eq!(computation_count.load(Ordering::Relaxed), 1);
power.invalidate();
flush_effects();
let _p3 = power.get();
assert_eq!(computation_count.load(Ordering::Relaxed), 2);
}
#[test]
fn effect_debouncing_batches_invalidations() {
let run_count = Arc::new(AtomicUsize::new(0));
let run_count_clone = run_count.clone();
let effect = Effect::new(move || {
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
for _ in 0..20 {
effect.invalidate();
}
Effect::process_all();
assert_eq!(run_count.load(Ordering::Relaxed), 2);
}
#[test]
fn signal_emit_origin_tracking() {
let circuit = TestCircuit::new();
circuit.voltage_signal.emit_from_api();
circuit.voltage_signal.emit_from_ui();
}
#[test]
fn transaction_suppresses_ui_updates() {
let circuit = TestCircuit::new();
let run_count = Arc::new(AtomicUsize::new(0));
let run_count_clone = run_count.clone();
let signal_id = circuit.voltage_signal.id();
let _effect = Effect::new(move || {
signal_id.track_dependency();
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
Transaction::no_ui_updates(|| {
circuit.voltage_signal.emit();
});
flush_effects();
assert_eq!(run_count.load(Ordering::Relaxed), 2);
}
#[test]
fn debouncing_batches_rapid_emissions() {
let effect_runs = Arc::new(AtomicUsize::new(0));
let effect_runs_clone = effect_runs.clone();
let circuit = Arc::new(TestCircuit::new());
let voltage_signal_id = circuit.voltage_signal.id();
let effect = Effect::new(move || {
voltage_signal_id.track_dependency();
effect_runs_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(effect_runs.load(Ordering::Relaxed), 1);
for _ in 0..20 {
effect.invalidate();
}
Effect::process_all();
assert_eq!(effect_runs.load(Ordering::Relaxed), 2);
}
#[test]
fn multiple_effects_debounce_independently() {
let run_count_1 = Arc::new(AtomicUsize::new(0));
let run_count_2 = Arc::new(AtomicUsize::new(0));
let rc1_clone = run_count_1.clone();
let effect1 = Effect::new(move || {
rc1_clone.fetch_add(1, Ordering::Relaxed);
});
let rc2_clone = run_count_2.clone();
let effect2 = Effect::new(move || {
rc2_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count_1.load(Ordering::Relaxed), 1);
assert_eq!(run_count_2.load(Ordering::Relaxed), 1);
for _ in 0..10 {
effect1.invalidate();
}
for _ in 0..10 {
effect2.invalidate();
}
Effect::process_all();
assert_eq!(run_count_1.load(Ordering::Relaxed), 2);
assert_eq!(run_count_2.load(Ordering::Relaxed), 2);
}
#[test]
fn transaction_nesting_preserves_state() {
let circuit = TestCircuit::new();
Transaction::run(|| {
Transaction::no_ui_updates(|| {
circuit.voltage_signal.emit();
});
circuit.voltage_signal.emit();
});
}
#[test]
fn computed_recomputes_on_dependency_change() {
let circuit = Arc::new(TestCircuit::new());
let voltage_signal_id = circuit.voltage_signal.id();
let computation_count = Arc::new(AtomicUsize::new(0));
let cc_clone = computation_count.clone();
let power = Computed::new(move || {
cc_clone.fetch_add(1, Ordering::Relaxed);
voltage_signal_id.track_dependency();
100.0
});
let _p1 = power.get();
assert_eq!(computation_count.load(Ordering::Relaxed), 1);
let _p2 = power.get();
assert_eq!(computation_count.load(Ordering::Relaxed), 1);
power.invalidate();
flush_effects();
let _p3 = power.get();
assert_eq!(computation_count.load(Ordering::Relaxed), 2);
}
#[test]
fn signal_emit_in_transaction_defers_effect() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
Transaction::run(|| {
signal.emit();
assert_eq!(counter.load(Ordering::Relaxed), 1);
});
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
#[test]
fn transaction_batches_multiple_emissions() {
use std::sync::atomic::AtomicI32;
let signal1 = Signal::new();
let signal2 = Signal::new();
let signal1_id = signal1.id();
let signal2_id = signal2.id();
let run_count = Arc::new(AtomicI32::new(0));
let run_count_clone = run_count.clone();
let _effect = Effect::new(move || {
signal1_id.track_dependency();
signal2_id.track_dependency();
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
Transaction::run(|| {
signal1.emit();
signal2.emit();
signal1.emit();
signal2.emit();
});
assert_eq!(run_count.load(Ordering::Relaxed), 2);
}
#[test]
fn nested_transactions_defer_until_outermost_exits() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let run_count = Arc::new(AtomicI32::new(0));
let run_count_clone = run_count.clone();
let _effect = Effect::new(move || {
signal_id.track_dependency();
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
Transaction::run(|| {
signal.emit();
Transaction::run(|| {
signal.emit();
});
signal.emit();
});
assert_eq!(run_count.load(Ordering::Relaxed), 2);
}
#[test]
fn transaction_active_flag_works() {
use crate::transaction::is_transaction_active;
assert!(!is_transaction_active());
Transaction::run(|| {
assert!(is_transaction_active());
Transaction::run(|| {
assert!(is_transaction_active());
});
assert!(is_transaction_active());
});
assert!(!is_transaction_active());
}
#[test]
fn effect_created_in_effect_becomes_child() {
use std::sync::atomic::AtomicI32;
let parent_runs = Arc::new(AtomicI32::new(0));
let child_runs = Arc::new(AtomicI32::new(0));
let parent_runs_clone = parent_runs.clone();
let child_runs_clone = child_runs.clone();
let parent_effect = Effect::new(move || {
parent_runs_clone.fetch_add(1, Ordering::Relaxed);
let child_runs_inner = child_runs_clone.clone();
let _child = Effect::new(move || {
child_runs_inner.fetch_add(1, Ordering::Relaxed);
});
});
assert_eq!(parent_runs.load(Ordering::Relaxed), 1);
assert_eq!(child_runs.load(Ordering::Relaxed), 1);
let child_count = parent_effect.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(child_count, 1);
}
#[test]
fn children_destroyed_when_parent_reruns() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let parent_runs = Arc::new(AtomicI32::new(0));
let child_creates = Arc::new(AtomicI32::new(0));
let parent_runs_clone = parent_runs.clone();
let child_creates_clone = child_creates.clone();
let parent_effect = Effect::new(move || {
signal_id.track_dependency();
parent_runs_clone.fetch_add(1, Ordering::Relaxed);
let cc = child_creates_clone.clone();
let _child = Effect::new(move || {
cc.fetch_add(1, Ordering::Relaxed);
});
});
assert_eq!(parent_runs.load(Ordering::Relaxed), 1);
assert_eq!(child_creates.load(Ordering::Relaxed), 1);
signal.emit();
flush_effects();
assert_eq!(parent_runs.load(Ordering::Relaxed), 2);
assert_eq!(child_creates.load(Ordering::Relaxed), 2);
let child_count = parent_effect.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(child_count, 1);
}
#[test]
fn many_effects_process_correctly() {
use std::sync::atomic::AtomicI32;
let counters: Vec<_> = (0..100).map(|_| Arc::new(AtomicI32::new(0))).collect();
let effects: Vec<_> = counters
.iter()
.map(|counter| {
let counter_clone = counter.clone();
Effect::new(move || {
counter_clone.fetch_add(1, Ordering::Relaxed);
})
})
.collect();
for counter in &counters {
assert_eq!(counter.load(Ordering::Relaxed), 1);
}
for effect in &effects {
effect.invalidate();
}
Effect::process_all();
for counter in &counters {
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
}
#[test]
fn signal_emit_schedules_processing_when_not_in_transaction() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let run_count = Arc::new(AtomicI32::new(0));
let run_count_clone = run_count.clone();
let _effect = Effect::new(move || {
signal_id.track_dependency();
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
signal.emit();
assert_eq!(run_count.load(Ordering::Relaxed), 1);
flush_effects();
assert_eq!(run_count.load(Ordering::Relaxed), 2);
}
#[test]
fn deeply_nested_hierarchy_tracked_correctly() {
use std::sync::atomic::AtomicI32;
let level0_runs = Arc::new(AtomicI32::new(0));
let level1_runs = Arc::new(AtomicI32::new(0));
let level2_runs = Arc::new(AtomicI32::new(0));
let l0 = level0_runs.clone();
let l1 = level1_runs.clone();
let l2 = level2_runs.clone();
let root_effect = Effect::new(move || {
l0.fetch_add(1, Ordering::Relaxed);
let l1_inner = l1.clone();
let l2_inner = l2.clone();
let _child = Effect::new(move || {
l1_inner.fetch_add(1, Ordering::Relaxed);
let l2_deep = l2_inner.clone();
let _grandchild = Effect::new(move || {
l2_deep.fetch_add(1, Ordering::Relaxed);
});
});
});
assert_eq!(level0_runs.load(Ordering::Relaxed), 1);
assert_eq!(level1_runs.load(Ordering::Relaxed), 1);
assert_eq!(level2_runs.load(Ordering::Relaxed), 1);
let child_count = root_effect.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(child_count, 1);
let first_child_id = root_effect.id().with_children(|c| c[0]).unwrap();
let grandchild_count = first_child_id.with_children(Vec::len).unwrap_or(0);
assert_eq!(grandchild_count, 1);
}
#[test]
fn children_destroyed_immediately_on_parent_invalidation() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let parent_runs = Arc::new(AtomicI32::new(0));
let child_created = Arc::new(AtomicI32::new(0));
let parent_runs_clone = parent_runs.clone();
let child_created_clone = child_created.clone();
let parent_effect = Effect::new(move || {
signal_id.track_dependency();
parent_runs_clone.fetch_add(1, Ordering::Relaxed);
let cc = child_created_clone.clone();
let _child = Effect::new(move || {
cc.fetch_add(1, Ordering::Relaxed);
});
});
assert_eq!(parent_runs.load(Ordering::Relaxed), 1);
assert_eq!(child_created.load(Ordering::Relaxed), 1);
let child_count = parent_effect.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(child_count, 1);
Transaction::run(|| {
parent_effect.invalidate();
let child_count_after = parent_effect.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(
child_count_after, 0,
"Children should be destroyed immediately when parent is marked dirty"
);
assert_eq!(parent_runs.load(Ordering::Relaxed), 1);
});
assert_eq!(parent_runs.load(Ordering::Relaxed), 2);
assert_eq!(child_created.load(Ordering::Relaxed), 2);
let new_child_count = parent_effect.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(new_child_count, 1);
}
#[test]
fn debouncing_batches_rapid_signal_changes() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let run_count = Arc::new(AtomicI32::new(0));
let rc = run_count.clone();
let _effect = Effect::new(move || {
signal_id.track_dependency();
rc.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
signal.emit();
signal.emit();
signal.emit();
signal.emit();
signal.emit();
assert_eq!(run_count.load(Ordering::Relaxed), 1);
flush_effects();
assert_eq!(run_count.load(Ordering::Relaxed), 2);
}
#[test]
fn complex_chain_propagates_correctly() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let computed1_count = Arc::new(AtomicI32::new(0));
let cc1 = computed1_count.clone();
let computed1 = Computed::new(move || {
cc1.fetch_add(1, Ordering::Relaxed);
signal_id.track_dependency();
10
});
let computed2_count = Arc::new(AtomicI32::new(0));
let cc2 = computed2_count.clone();
let c1_clone = computed1.clone();
let _computed2 = Computed::new(move || {
cc2.fetch_add(1, Ordering::Relaxed);
c1_clone.get() * 2
});
assert_eq!(computed1.get(), 10);
assert_eq!(computed1_count.load(Ordering::Relaxed), 1);
assert_eq!(computed1.get(), 10);
assert_eq!(computed1_count.load(Ordering::Relaxed), 1);
computed1.invalidate();
flush_effects();
assert_eq!(computed1.get(), 10);
assert_eq!(computed1_count.load(Ordering::Relaxed), 2);
}
#[test]
fn diamond_dependency_updates_correctly() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let c1 = counter.clone();
let _effect1 = Effect::new(move || {
signal_id.track_dependency();
c1.fetch_add(1, Ordering::Relaxed);
});
let c2 = counter.clone();
let _effect2 = Effect::new(move || {
signal_id.track_dependency();
c2.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 2);
signal.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 4);
}
#[test]
fn effect_resubscribes_on_each_run() {
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicI32;
let signal_a = Signal::new();
let signal_b = Signal::new();
let signal_a_id = signal_a.id();
let signal_b_id = signal_b.id();
let use_b = Arc::new(AtomicBool::new(false));
let counter = Arc::new(AtomicI32::new(0));
let use_b_clone = use_b.clone();
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
if use_b_clone.load(Ordering::Relaxed) {
signal_b_id.track_dependency();
} else {
signal_a_id.track_dependency();
}
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal_a.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
use_b.store(true, Ordering::Relaxed);
signal_a.emit(); flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 3);
signal_a.emit();
flush_effects();
assert_eq!(
counter.load(Ordering::Relaxed),
3,
"Effect should not react to A anymore"
);
signal_b.emit();
flush_effects();
assert_eq!(
counter.load(Ordering::Relaxed),
4,
"Effect should react to B now"
);
}
#[test]
fn stress_test_many_signals() {
use std::sync::atomic::AtomicI32;
let signals: Vec<Signal> = (0..100).map(|_| Signal::new()).collect();
let signal_ids: Vec<_> = signals.iter().map(super::signal::Signal::id).collect();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let ids_clone = signal_ids.clone();
let _effect = Effect::new(move || {
for id in &ids_clone {
id.track_dependency();
}
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signals[50].emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
#[test]
fn stress_test_many_effects() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let _effects: Vec<Effect> = (0..100)
.map(|_| {
let counter_clone = counter.clone();
Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
})
})
.collect();
assert_eq!(counter.load(Ordering::Relaxed), 100);
signal.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 200);
}
#[test]
fn rapid_emissions_debounced() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
for _ in 0..100 {
signal.emit();
}
assert_eq!(counter.load(Ordering::Relaxed), 1);
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
#[test]
fn deeply_nested_hierarchy_stress() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let levels_created = Arc::new(AtomicI32::new(0));
fn create_nested_effects(
levels_created: Arc<AtomicI32>,
signal_id: crate::arena::SignalId,
depth: usize,
max_depth: usize,
) {
levels_created.fetch_add(1, Ordering::Relaxed);
if depth < max_depth {
let lc = levels_created.clone();
let _child = Effect::new(move || {
signal_id.track_dependency();
create_nested_effects(lc.clone(), signal_id, depth + 1, max_depth);
});
}
}
let lc = levels_created.clone();
let root_effect = Effect::new(move || {
signal_id.track_dependency();
create_nested_effects(lc.clone(), signal_id, 1, 10);
});
assert_eq!(levels_created.load(Ordering::Relaxed), 10);
root_effect.invalidate();
let child_count = root_effect.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(
child_count, 0,
"Children should be destroyed when parent is marked pending"
);
flush_effects();
assert_eq!(levels_created.load(Ordering::Relaxed), 20);
}
#[test]
fn signal_drop_cleanup_is_safe() {
use std::sync::atomic::AtomicI32;
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let signal = Signal::new();
let signal_id = signal.id();
let effect = Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
drop(signal);
effect.invalidate();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
#[test]
fn signal_drop_with_pending_effect_cleans_up() {
use std::sync::atomic::AtomicI32;
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let signal = Signal::new();
let signal_id = signal.id();
let _effect = Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal.emit();
drop(signal);
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
#[test]
fn effect_drop_cleanup_is_safe() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
{
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
}
signal.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 1);
}
#[test]
fn effect_drop_while_pending_does_not_run() {
use std::sync::atomic::AtomicI32;
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let signal = Signal::new();
let signal_id = signal.id();
let effect = Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
effect.invalidate();
drop(effect);
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 1);
}
#[test]
fn stale_signal_id_returns_none() {
let signal = Signal::new();
let signal_id = signal.id();
assert!(signal_id.with_subscribers(|_| ()).is_some());
drop(signal);
let _ = signal_id.with_subscribers(|_| ());
}
#[test]
fn stale_effect_id_does_not_panic() {
let effect = Effect::new(|| {});
let effect_id = effect.id();
assert!(effect_id.has_callback());
drop(effect);
let _ = effect_id.has_callback();
}
#[test]
fn ui_origin_filtering_with_track_ui_only() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let _effect = Effect::new_ui_updates_only(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal.emit_from_api();
flush_effects();
assert_eq!(
counter.load(Ordering::Relaxed),
2,
"Effect should run on API updates"
);
signal.emit_from_ui();
flush_effects();
assert_eq!(
counter.load(Ordering::Relaxed),
2,
"Effect with ui_updates_only should NOT run on UI updates"
);
}
#[test]
fn ui_origin_filtering_behavior() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let data_only_counter = Arc::new(AtomicI32::new(0));
let dc = data_only_counter.clone();
let all_counter = Arc::new(AtomicI32::new(0));
let ac = all_counter.clone();
let _data_only_effect = Effect::new_ui_updates_only(move || {
signal_id.track_dependency();
dc.fetch_add(1, Ordering::Relaxed);
});
let _all_effect = Effect::new(move || {
signal_id.track_dependency();
ac.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(data_only_counter.load(Ordering::Relaxed), 1);
assert_eq!(all_counter.load(Ordering::Relaxed), 1);
signal.emit_from_api();
flush_effects();
assert_eq!(
all_counter.load(Ordering::Relaxed),
2,
"all_effect should run on API update"
);
assert_eq!(
data_only_counter.load(Ordering::Relaxed),
2,
"data_only_effect should run on API update"
);
signal.emit_from_ui();
flush_effects();
assert_eq!(
all_counter.load(Ordering::Relaxed),
3,
"all_effect should run on UI update"
);
assert_eq!(
data_only_counter.load(Ordering::Relaxed),
2,
"data_only_effect should NOT run on UI update"
);
}
#[test]
fn transaction_cleans_up_on_panic() {
use std::panic::catch_unwind;
let initial_active = crate::transaction::is_transaction_active();
assert!(!initial_active, "Should not be in transaction initially");
let result = catch_unwind(|| {
Transaction::run(|| {
assert!(crate::transaction::is_transaction_active());
panic!("Test panic");
});
});
assert!(result.is_err(), "Transaction should have panicked");
assert!(
!crate::transaction::is_transaction_active(),
"Transaction should be cleaned up after panic"
);
}
#[test]
fn transaction_returns_value_correctly() {
let int_result = Transaction::run(|| 42i32);
assert_eq!(int_result, 42);
let string_result = Transaction::run(|| String::from("hello"));
assert_eq!(string_result, "hello");
let option_result: Option<i32> = Transaction::run(|| Some(42));
assert_eq!(option_result, Some(42));
let tuple_result = Transaction::run(|| (1, "two", 3.0));
assert_eq!(tuple_result, (1, "two", 3.0));
}
#[test]
fn computed_chain_propagates() {
use std::sync::atomic::AtomicI32;
let c1_count = Arc::new(AtomicI32::new(0));
let c2_count = Arc::new(AtomicI32::new(0));
let c3_count = Arc::new(AtomicI32::new(0));
let cc1 = c1_count.clone();
let c1 = Computed::new(move || {
cc1.fetch_add(1, Ordering::Relaxed);
10
});
let cc2 = c2_count.clone();
let c1_for_c2 = c1.clone();
let c2 = Computed::new(move || {
cc2.fetch_add(1, Ordering::Relaxed);
c1_for_c2.get() * 2
});
let cc3 = c3_count.clone();
let c2_for_c3 = c2.clone();
let c3 = Computed::new(move || {
cc3.fetch_add(1, Ordering::Relaxed);
c2_for_c3.get() + 5
});
assert_eq!(c3.get(), 25); assert_eq!(c1_count.load(Ordering::Relaxed), 1);
assert_eq!(c2_count.load(Ordering::Relaxed), 1);
assert_eq!(c3_count.load(Ordering::Relaxed), 1);
assert_eq!(c3.get(), 25);
assert_eq!(c3_count.load(Ordering::Relaxed), 1);
c1.invalidate();
flush_effects();
assert_eq!(c3.get(), 25);
assert_eq!(c1_count.load(Ordering::Relaxed), 2); assert_eq!(c3_count.load(Ordering::Relaxed), 1); }
#[test]
fn effect_creates_sibling_effects() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let child_count = Arc::new(AtomicI32::new(0));
let cc = child_count.clone();
let parent = Effect::new(move || {
signal_id.track_dependency();
for _ in 0..3 {
let cc_inner = cc.clone();
let _child = Effect::new(move || {
cc_inner.fetch_add(1, Ordering::Relaxed);
});
}
});
assert_eq!(child_count.load(Ordering::Relaxed), 3);
let child_count_check = parent.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(child_count_check, 3);
signal.emit();
let child_count_after_emit = parent.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(
child_count_after_emit, 0,
"Children destroyed when parent marked pending"
);
flush_effects();
assert_eq!(child_count.load(Ordering::Relaxed), 6);
let new_child_count = parent.id().with_children(Vec::len).unwrap_or(0);
assert_eq!(new_child_count, 3);
}
#[test]
fn effect_writing_tracked_as_output() {
use std::sync::atomic::AtomicI32;
let input_signal = Signal::new();
let output_signal = Signal::new();
let input_signal_id = input_signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let effect = Effect::new(move || {
input_signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
output_signal.emit();
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
let output_count = effect
.id()
.with_outputs(|outputs| outputs.count())
.unwrap_or(0);
assert_eq!(output_count, 1);
}
#[test]
fn effect_emit_tracked_as_output() {
use crate::arena::signal_arena::with_signal_writers;
let input_signal = Signal::new();
let output_signal = Signal::new();
let input_signal_id = input_signal.id();
let output_signal_id = output_signal.id();
let effect = Effect::new(move || {
input_signal_id.track_dependency();
});
let source_count = effect
.id()
.with_sources(|sources| sources.count())
.unwrap_or(0);
assert_eq!(source_count, 1);
let has_input_source = effect
.id()
.with_sources(|sources| {
for s in sources {
if s == input_signal_id {
return true;
}
}
false
})
.unwrap_or(false);
assert!(has_input_source);
let output_count = effect
.id()
.with_outputs(|outputs| outputs.count())
.unwrap_or(0);
assert_eq!(output_count, 0);
let has_writers = with_signal_writers(output_signal_id, |writers| !writers.is_empty());
assert!(!has_writers.unwrap_or(false));
}
#[test]
fn untracked_prevents_dependency() {
use crate::untracked;
let signal = Signal::new();
let signal_id = signal.id();
let effect = Effect::new(move || {
untracked(|| {
signal_id.track_dependency();
});
});
let source_count = effect
.id()
.with_sources(|sources| sources.count())
.unwrap_or(0);
assert_eq!(
source_count, 0,
"untracked() should prevent dependency tracking"
);
let subscriber_count = signal_id.with_subscribers(HashSet::len).unwrap_or(0);
assert_eq!(
subscriber_count, 0,
"untracked() should prevent subscription"
);
}
#[test]
fn untracked_returns_value() {
use crate::untracked;
let result = untracked(|| 42);
assert_eq!(result, 42);
let result = untracked(|| "hello".to_string());
assert_eq!(result, "hello");
let result: Option<i32> = untracked(|| Some(123));
assert_eq!(result, Some(123));
}
#[test]
fn effect_rerun_clears_outputs() {
use std::sync::atomic::AtomicBool;
let input_signal = Signal::new();
let _output_signal_a = Signal::new();
let _output_signal_b = Signal::new();
let input_signal_id = input_signal.id();
let use_b = Arc::new(AtomicBool::new(false));
let use_b_clone = use_b.clone();
let effect = Effect::new(move || {
input_signal_id.track_dependency();
let _ = use_b_clone.load(Ordering::Relaxed);
});
assert!(effect.id().with_sources(|_| ()).is_some());
use_b.store(true, Ordering::Relaxed);
effect.invalidate();
flush_effects();
assert!(effect.id().with_sources(|_| ()).is_some());
}
#[test]
fn pull_ensures_fresh_values() {
let signal = Signal::new();
let signal_id = signal.id();
signal_id.pull(false);
assert!(signal_id.with_subscribers(|_| ()).is_some());
}
#[test]
fn diamond_with_pull_gets_fresh_values() {
use std::sync::atomic::AtomicI32;
let input_signal = Signal::new();
let input_signal_id = input_signal.id();
let effect1_runs = Arc::new(AtomicI32::new(0));
let effect2_runs = Arc::new(AtomicI32::new(0));
let e1 = effect1_runs.clone();
let _effect1 = Effect::new(move || {
input_signal_id.track_dependency();
e1.fetch_add(1, Ordering::Relaxed);
});
let e2 = effect2_runs.clone();
let _effect2 = Effect::new(move || {
input_signal_id.track_dependency();
e2.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(effect1_runs.load(Ordering::Relaxed), 1);
assert_eq!(effect2_runs.load(Ordering::Relaxed), 1);
input_signal.emit();
flush_effects();
assert_eq!(effect1_runs.load(Ordering::Relaxed), 2);
assert_eq!(effect2_runs.load(Ordering::Relaxed), 2);
}
#[test]
fn fixed_point_iteration_processes_cascades() {
use std::sync::atomic::AtomicI32;
let counter = Arc::new(AtomicI32::new(0));
let signal = Signal::new();
let signal_id = signal.id();
let counter_clone = counter.clone();
let effect = Effect::new(move || {
signal_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
effect.invalidate();
effect.invalidate();
effect.invalidate();
Effect::process_all();
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
#[test]
fn read_write_cycle_detection() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let run_count = Arc::new(AtomicI32::new(0));
let run_count_clone = run_count.clone();
let _effect = Effect::new(move || {
signal_id.track_dependency();
signal.emit();
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
flush_effects();
assert_eq!(run_count.load(Ordering::Relaxed), 1);
}
#[test]
fn read_write_different_signals_allowed() {
use std::sync::atomic::AtomicI32;
let signal_read = Signal::new();
let signal_write = Signal::new();
let signal_read_id = signal_read.id();
let run_count = Arc::new(AtomicI32::new(0));
let run_count_clone = run_count.clone();
let _effect = Effect::new(move || {
signal_read_id.track_dependency();
signal_write.emit();
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
signal_read.emit();
flush_effects();
assert_eq!(run_count.load(Ordering::Relaxed), 2);
signal_read.emit();
flush_effects();
assert_eq!(run_count.load(Ordering::Relaxed), 3);
}
#[test]
fn skippable_effect_flag_set_correctly() {
let effect = Effect::new_skippable(|| {});
assert!(effect.id().is_skippable());
let regular_effect = Effect::new(|| {});
assert!(!regular_effect.id().is_skippable());
}
#[test]
fn flush_with_budget_defers_skippable() {
use std::time::Duration;
let skippable_count = Arc::new(AtomicUsize::new(0));
let sc = skippable_count.clone();
let skippable = Effect::new_skippable(move || {
sc.fetch_add(1, Ordering::Relaxed);
});
let non_skippable_count = Arc::new(AtomicUsize::new(0));
let nsc = non_skippable_count.clone();
let non_skippable = Effect::new(move || {
nsc.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(skippable_count.load(Ordering::Relaxed), 1);
assert_eq!(non_skippable_count.load(Ordering::Relaxed), 1);
skippable.invalidate();
non_skippable.invalidate();
let budget = Duration::from_millis(0);
let mut must_run_buf = Vec::new();
let mut skippable_buf = Vec::new();
crate::effect::flush_effects_with_budget(budget, &mut must_run_buf, &mut skippable_buf);
assert_eq!(non_skippable_count.load(Ordering::Relaxed), 2);
assert_eq!(skippable_count.load(Ordering::Relaxed), 1);
assert_eq!(skippable_buf.len(), 1);
assert!(skippable_buf[0].is_skippable());
}
#[test]
fn pull_skips_skippable_when_flag_set() {
use std::sync::atomic::AtomicI32;
let signal = Signal::new();
let signal_id = signal.id();
let skippable_count = Arc::new(AtomicI32::new(0));
let sc = skippable_count.clone();
let _skippable_computed = Computed::new_skippable(move || {
sc.fetch_add(1, Ordering::Relaxed);
42
});
assert_eq!(skippable_count.load(Ordering::Relaxed), 1);
signal_id.pull(true);
assert_eq!(skippable_count.load(Ordering::Relaxed), 1);
signal_id.pull(false);
}
#[test]
fn skippable_effects_integrate_with_budget_system() {
use std::time::Duration;
let producer_count = Arc::new(AtomicUsize::new(0));
let mut producers = Vec::new();
for _ in 0..5 {
let pc = producer_count.clone();
producers.push(Effect::new_skippable(move || {
pc.fetch_add(1, Ordering::Relaxed);
}));
}
let consumer_count = Arc::new(AtomicUsize::new(0));
let cc = consumer_count.clone();
let _consumer = Effect::new(move || {
cc.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(producer_count.load(Ordering::Relaxed), 5);
assert_eq!(consumer_count.load(Ordering::Relaxed), 1);
for producer in &producers {
producer.invalidate();
}
let budget = Duration::from_millis(1);
let mut must_run_buf = Vec::new();
let mut skippable_buf = Vec::new();
crate::effect::flush_effects_with_budget(budget, &mut must_run_buf, &mut skippable_buf);
for deferred_id in &skippable_buf {
assert!(deferred_id.is_skippable());
}
}
#[test]
fn check_state_effects_added_to_pending() {
use crate::arena::current_effect;
cov_mark::check!(check_effect_added_to_pending);
let trigger = Signal::new();
let output = Signal::new();
let trigger_id = trigger.id();
let output_id = output.id();
let _effect_a = Effect::new(move || {
trigger_id.track_dependency();
if let Some(effect_id) = current_effect() {
effect_id.add_output(output_id);
}
output_id.notify_subscribers();
});
let b_count = Arc::new(AtomicUsize::new(0));
let b_count_clone = b_count.clone();
let _effect_b = Effect::new(move || {
output_id.track_dependency();
b_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(b_count.load(Ordering::Relaxed), 1);
trigger.emit();
flush_effects();
assert_eq!(b_count.load(Ordering::Relaxed), 2);
for i in 3..=5 {
trigger.emit();
flush_effects();
assert_eq!(
b_count.load(Ordering::Relaxed),
i,
"Effect B should run on emission {i}. \
If this fails, Check-state effects may not be added to the pending set."
);
}
}
#[test]
fn check_state_propagation_through_chain() {
use crate::arena::current_effect;
cov_mark::check!(check_effect_added_to_pending);
let source = Signal::new();
let middle = Signal::new();
let end = Signal::new();
let source_id = source.id();
let middle_id = middle.id();
let end_id = end.id();
let _effect_a = Effect::new(move || {
source_id.track_dependency();
if let Some(effect_id) = current_effect() {
effect_id.add_output(middle_id);
}
middle_id.notify_subscribers();
});
let _effect_b = Effect::new(move || {
middle_id.track_dependency();
if let Some(effect_id) = current_effect() {
effect_id.add_output(end_id);
}
end_id.notify_subscribers();
});
let c_count = Arc::new(AtomicUsize::new(0));
let c_count_clone = c_count.clone();
let _effect_c = Effect::new(move || {
end_id.track_dependency();
c_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(c_count.load(Ordering::Relaxed), 1);
for i in 2..=5 {
source.emit();
flush_effects();
assert_eq!(
c_count.load(Ordering::Relaxed),
i,
"Effect C should run on emission {i}"
);
}
}