#![cfg(feature = "scheduler")]
#![allow(clippy::needless_range_loop)]
use std::sync::{
Arc, Mutex,
atomic::{AtomicI64, AtomicUsize, Ordering},
};
use hyphae::{
Cell, CellImmutable, CellMutable, DedupedExt, DistinctExt, DistinctUntilChangedByExt,
FilterExt, Gettable, MapErrExt, MapExt, MapOkExt, MaterializeDefinite, MaterializeEmpty,
Mutable, ScanExt, Signal, StateTransitionExt, TapExt, TryMapExt, Watchable, batch,
scheduler::no_coalesce,
};
fn scheduler_test_serial() -> std::sync::MutexGuard<'static, ()> {
hyphae::scheduler::set_wave_threshold_for_test(4);
static LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
LOCK.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
}
const N_B: usize = 16;
const ITERS_B: i64 = 2_000;
const N_C: i64 = 32;
const ITERS_C: usize = 2_000;
#[test]
fn map_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for k in 0..N_B {
let kk = k as i64;
let s = Cell::new(0i64);
let d = s.clone().map(move |x| x * (kk + 1) + 7).materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let x = base + k as i64 + 1;
assert_eq!(
derived[k].get(),
x * (k as i64 + 1) + 7,
"map instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn filter_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000; let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for _ in 0..N_B {
let s = Cell::new(0i64);
let d = s.clone().filter(|x| x % 2 == 0).materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let v = base + k as i64 + 1;
let expected = if v % 2 == 0 { Some(v) } else { Some(0) };
assert_eq!(
derived[k].get(),
expected,
"filter instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn try_map_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for _ in 0..N_B {
let s = Cell::new(0i64);
let d = s
.clone()
.try_map(|x| {
if x % 2 == 0 {
Ok(x * 2)
} else {
Err(format!("odd:{x}"))
}
})
.materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let v = base + k as i64 + 1;
let expected: Result<i64, String> = if v % 2 == 0 {
Ok(v * 2)
} else {
Err(format!("odd:{v}"))
};
assert_eq!(
derived[k].get(),
expected,
"try_map instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn tap_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
let side: Vec<Arc<AtomicI64>> = (0..N_B).map(|_| Arc::new(AtomicI64::new(-1))).collect();
for k in 0..N_B {
let s = Cell::new(0i64);
let sc = side[k].clone();
let d = s
.clone()
.tap(move |v| sc.store(*v, Ordering::SeqCst))
.materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let v = base + k as i64 + 1;
assert_eq!(derived[k].get(), v, "tap value {k} wrong at iter {iter}");
assert_eq!(
side[k].load(Ordering::SeqCst),
v,
"tap side effect {k} wrong at iter {iter}"
);
}
}
}
#[test]
fn scan_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for _ in 0..N_B {
let s = Cell::new(0i64);
let d = s.clone().scan(100i64, |acc, x| acc + x).materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let v = base + k as i64 + 1;
assert_eq!(
derived[k].get(),
100 + v,
"scan instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn distinct_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for _ in 0..N_B {
let s = Cell::new(0i64);
let d = s.clone().distinct().materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let v = base + k as i64 + 1;
assert_eq!(
derived[k].get(),
v,
"distinct instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn distinct_until_changed_by_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for _ in 0..N_B {
let s = Cell::new(0i64);
let d = s
.clone()
.distinct_until_changed_by(|a, b| a == b)
.materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let v = base + k as i64 + 1;
assert_eq!(
derived[k].get(),
v,
"distinct_until_changed_by instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn deduped_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for _ in 0..N_B {
let s = Cell::new(0i64);
let d = s.clone().deduped().materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
s.set(base + k as i64 + 1);
}
});
for k in 0..N_B {
let v = base + k as i64 + 1;
assert_eq!(
derived[k].get(),
v,
"deduped instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn state_transition_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let mut sources = Vec::with_capacity(N_B);
let mut machines = Vec::with_capacity(N_B);
for k in 0..N_B {
let kk = k as i64;
let s = Cell::new(0i64);
let sm: Cell<i64, CellImmutable> = s.state_transition(move |sm| {
sm.on(0i64, 1i64, move |_, _| (kk + 1) * 10);
});
sources.push(s);
machines.push(sm);
}
batch(|| {
for s in &sources {
s.set(1);
}
});
for k in 0..N_B {
assert_eq!(
machines[k].get(),
(k as i64 + 1) * 10,
"state_transition instance {k} emitted wrong at iter {iter}"
);
}
}
}
#[test]
fn map_ok_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for k in 0..N_B {
let kk = k as i64;
let s = Cell::new(Ok::<i64, String>(0));
let d = s.clone().map_ok(move |v| v * (kk + 1)).materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
if k % 2 == 0 {
s.set(Ok(base + k as i64 + 1));
} else {
s.set(Err(format!("e{k}")));
}
}
});
for k in 0..N_B {
let expected: Result<i64, String> = if k % 2 == 0 {
Ok((base + k as i64 + 1) * (k as i64 + 1))
} else {
Err(format!("e{k}"))
};
assert_eq!(
derived[k].get(),
expected,
"map_ok instance {k} settled wrong at iter {iter}"
);
}
}
}
#[test]
fn map_err_wide_wave_settles_each_instance_correctly() {
let _serial = scheduler_test_serial();
for iter in 0..ITERS_B {
let base = iter * 1_000;
let mut sources = Vec::with_capacity(N_B);
let mut derived = Vec::with_capacity(N_B);
for k in 0..N_B {
let kk = k as i64;
let s = Cell::new(Err::<i64, String>("init".into()));
let d = s
.clone()
.map_err(move |e| format!("{e}!{kk}"))
.materialize();
sources.push(s);
derived.push(d);
}
batch(|| {
for (k, s) in sources.iter().enumerate() {
if k % 2 == 0 {
s.set(Ok(base + k as i64 + 1));
} else {
s.set(Err(format!("e{k}")));
}
}
});
for k in 0..N_B {
let expected: Result<i64, String> = if k % 2 == 0 {
Ok(base + k as i64 + 1)
} else {
Err(format!("e{k}!{k}"))
};
assert_eq!(
derived[k].get(),
expected,
"map_err instance {k} settled wrong at iter {iter}"
);
}
}
}
fn assert_ordered_event_delivery(
label: &str,
src: Cell<i64, CellMutable>,
derived: Cell<i64, CellImmutable>,
) {
let received: Arc<Mutex<Vec<i64>>> = Arc::new(Mutex::new(Vec::new()));
let in_flight = Arc::new(AtomicUsize::new(0));
let max_in_flight = Arc::new(AtomicUsize::new(0));
let r = received.clone();
let inf = in_flight.clone();
let maxf = max_in_flight.clone();
let guard = derived.subscribe(move |sig| {
if let Signal::Value(v) = sig {
let cur = inf.fetch_add(1, Ordering::SeqCst) + 1;
maxf.fetch_max(cur, Ordering::SeqCst);
for _ in 0..64 {
std::hint::spin_loop();
}
r.lock().unwrap().push(**v);
inf.fetch_sub(1, Ordering::SeqCst);
}
});
let expected: Vec<i64> = (1..=N_C).collect();
for it in 0..ITERS_C {
received.lock().unwrap().clear();
batch(|| {
for i in 1..=N_C {
src.set(i);
}
});
let got = received.lock().unwrap().clone();
assert_eq!(
got.len(),
N_C as usize,
"{label} iter {it}: a no_coalesce event was dropped/coalesced"
);
assert_eq!(
got, expected,
"{label} iter {it}: no_coalesce events delivered OUT OF ORDER under a parallel wave"
);
}
std::mem::forget(guard);
let peak = max_in_flight.load(Ordering::SeqCst);
assert_eq!(
peak, 1,
"{label}: derived event subscriber ran concurrently with itself (peak in-flight = {peak})"
);
}
#[test]
fn scan_event_order_under_parallel_wave_regression() {
let _serial = scheduler_test_serial();
let (src, derived) = no_coalesce(|| {
let src = Cell::<i64, CellMutable>::new(0);
let derived = src.clone().scan(0i64, |_acc, x| *x).materialize();
(src, derived)
});
assert_ordered_event_delivery("scan", src, derived);
}
#[test]
fn distinct_until_changed_by_event_order_under_parallel_wave_regression() {
let _serial = scheduler_test_serial();
let (src, derived) = no_coalesce(|| {
let src = Cell::<i64, CellMutable>::new(0);
let derived = src
.clone()
.distinct_until_changed_by(|a, b| a == b)
.materialize();
(src, derived)
});
assert_ordered_event_delivery("distinct_until_changed_by", src, derived);
}
#[test]
fn deduped_event_order_under_parallel_wave_regression() {
let _serial = scheduler_test_serial();
let (src, derived) = no_coalesce(|| {
let src = Cell::<i64, CellMutable>::new(0);
let derived = src.clone().deduped().materialize();
(src, derived)
});
assert_ordered_event_delivery("deduped", src, derived);
}
#[test]
fn state_transition_event_order_under_parallel_wave_regression() {
let _serial = scheduler_test_serial();
let (src, derived) = no_coalesce(|| {
let src = Cell::<i64, CellMutable>::new(0);
let derived: Cell<i64, CellImmutable> = src.state_transition(|sm| {
for k in 0..N_C {
sm.on(k, k + 1, move |_, to| *to);
}
sm.on(N_C, 1i64, move |_, to| *to);
});
(src, derived)
});
assert_ordered_event_delivery("state_transition", src, derived);
}