#![cfg(feature = "scheduler")]
#![allow(clippy::needless_range_loop)]
use std::sync::{Arc, Mutex};
use hyphae::{
BufferCountExt, Cell, CellMutable, FirstExt, Gettable, LastExt, MaterializeDefinite,
MaterializeEmpty, Mutable, PairwiseExt, Signal, SkipExt, SkipWhileExt, TakeExt, TakeUntilExt,
TakeWhileExt, Watchable, WindowExt, 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 WIDE: usize = 16;
const ITERS: i64 = 2_000;
const BURST: i64 = 32;
const ORDER_ITERS: usize = 2_000;
#[test]
fn take_wide_independent_wave_each_instance_exact() {
let _serial = scheduler_test_serial();
let sources: Vec<Cell<i64, CellMutable>> = (0..WIDE).map(|_| Cell::new(0i64)).collect();
let counts: Vec<usize> = (0..WIDE).map(|k| k + 2).collect();
let cells: Vec<_> = sources
.iter()
.zip(&counts)
.map(|(s, &c)| s.clone().take(c).materialize())
.collect();
let recv: Vec<Arc<Mutex<Vec<i64>>>> = (0..WIDE)
.map(|_| Arc::new(Mutex::new(Vec::new())))
.collect();
let mut guards = Vec::new();
for (cell, r) in cells.iter().zip(&recv) {
let rr = r.clone();
guards.push(cell.subscribe(move |sig| {
if let Signal::Value(v) = sig {
rr.lock().unwrap().push(**v);
}
}));
}
for i in 1..=ITERS {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
for k in 0..WIDE {
let expected: Vec<i64> = (0..counts[k] as i64).collect();
assert_eq!(
*recv[k].lock().unwrap(),
expected,
"take instance {k} (count {}) settled wrong under the wide parallel wave",
counts[k]
);
assert!(cells[k].is_complete(), "take instance {k} never completed");
}
drop(guards);
}
#[test]
fn take_while_wide_independent_wave_each_instance_exact() {
let _serial = scheduler_test_serial();
let sources: Vec<Cell<i64, CellMutable>> = (0..WIDE).map(|_| Cell::new(0i64)).collect();
let thresholds: Vec<i64> = (0..WIDE).map(|k| k as i64 + 3).collect();
let cells: Vec<_> = sources
.iter()
.zip(&thresholds)
.map(|(s, &th)| s.clone().take_while(move |x| *x < th).materialize())
.collect();
let recv: Vec<Arc<Mutex<Vec<i64>>>> = (0..WIDE)
.map(|_| Arc::new(Mutex::new(Vec::new())))
.collect();
let mut guards = Vec::new();
for (cell, r) in cells.iter().zip(&recv) {
let rr = r.clone();
guards.push(cell.subscribe(move |sig| {
if let Signal::Value(v) = sig {
if let Some(x) = **v {
rr.lock().unwrap().push(x);
}
}
}));
}
for i in 1..=ITERS {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
for k in 0..WIDE {
let expected: Vec<i64> = (0..thresholds[k]).collect();
assert_eq!(
*recv[k].lock().unwrap(),
expected,
"take_while instance {k} (th {}) settled wrong under the wide parallel wave",
thresholds[k]
);
assert_eq!(
cells[k].get(),
Some(thresholds[k] - 1),
"take_while instance {k} froze on the wrong last-passing value"
);
assert!(
cells[k].is_complete(),
"take_while instance {k} never completed"
);
}
drop(guards);
}
#[test]
fn take_until_wide_independent_wave_each_instance_exact() {
let _serial = scheduler_test_serial();
let sources: Vec<Cell<i64, CellMutable>> = (0..WIDE).map(|_| Cell::new(0i64)).collect();
let stoppers: Vec<Cell<bool, CellMutable>> = (0..WIDE).map(|_| Cell::new(false)).collect();
let cells: Vec<_> = sources
.iter()
.zip(&stoppers)
.map(|(s, st)| s.take_until(st))
.collect();
for i in 1..=ITERS {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
batch(|| {
for st in &stoppers {
st.set(true);
}
});
for i in ITERS + 1..=ITERS + 50 {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
for k in 0..WIDE {
assert_eq!(
cells[k].get(),
ITERS,
"take_until instance {k} did not freeze at the pre-stop value under the wide wave"
);
assert!(
cells[k].is_complete(),
"take_until instance {k} never completed"
);
}
}
#[test]
fn skip_wide_independent_wave_each_instance_exact() {
let _serial = scheduler_test_serial();
let sources: Vec<Cell<i64, CellMutable>> = (0..WIDE).map(|_| Cell::new(0i64)).collect();
let skips: Vec<usize> = (0..WIDE).map(|k| k + 1).collect();
let cells: Vec<_> = sources
.iter()
.zip(&skips)
.map(|(s, &n)| s.clone().skip(n).materialize())
.collect();
let recv: Vec<Arc<Mutex<Vec<i64>>>> = (0..WIDE)
.map(|_| Arc::new(Mutex::new(Vec::new())))
.collect();
let mut guards = Vec::new();
for (cell, r) in cells.iter().zip(&recv) {
let rr = r.clone();
guards.push(cell.subscribe(move |sig| {
if let Signal::Value(v) = sig
&& let Some(x) = (**v)
{
rr.lock().unwrap().push(x);
}
}));
}
for i in 1..=ITERS {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
for k in 0..WIDE {
let n = skips[k] as i64;
let expected: Vec<i64> = (n..=ITERS).collect();
assert_eq!(
*recv[k].lock().unwrap(),
expected,
"skip instance {k} (n {}) settled wrong under the wide parallel wave",
skips[k]
);
}
drop(guards);
}
#[test]
fn skip_while_wide_independent_wave_each_instance_exact() {
let _serial = scheduler_test_serial();
let sources: Vec<Cell<i64, CellMutable>> = (0..WIDE).map(|_| Cell::new(0i64)).collect();
let thresholds: Vec<i64> = (0..WIDE).map(|k| k as i64 + 1).collect();
let cells: Vec<_> = sources
.iter()
.zip(&thresholds)
.map(|(s, &th)| s.clone().skip_while(move |x| *x < th).materialize())
.collect();
let recv: Vec<Arc<Mutex<Vec<i64>>>> = (0..WIDE)
.map(|_| Arc::new(Mutex::new(Vec::new())))
.collect();
let mut guards = Vec::new();
for (cell, r) in cells.iter().zip(&recv) {
let rr = r.clone();
guards.push(cell.subscribe(move |sig| {
if let Signal::Value(v) = sig
&& let Some(x) = (**v)
{
rr.lock().unwrap().push(x);
}
}));
}
for i in 1..=ITERS {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
for k in 0..WIDE {
let th = thresholds[k];
let expected: Vec<i64> = (th..=ITERS).collect();
assert_eq!(
*recv[k].lock().unwrap(),
expected,
"skip_while instance {k} (th {th}) settled wrong under the wide parallel wave"
);
}
drop(guards);
}
#[test]
fn first_wide_independent_wave_each_instance_exact() {
let _serial = scheduler_test_serial();
let sources: Vec<Cell<i64, CellMutable>> = (0..WIDE).map(|_| Cell::new(0i64)).collect();
let cells: Vec<_> = sources
.iter()
.map(|s| s.clone().first().materialize())
.collect();
let recv: Vec<Arc<Mutex<Vec<i64>>>> = (0..WIDE)
.map(|_| Arc::new(Mutex::new(Vec::new())))
.collect();
let mut guards = Vec::new();
for (cell, r) in cells.iter().zip(&recv) {
let rr = r.clone();
guards.push(cell.subscribe(move |sig| {
if let Signal::Value(v) = sig {
rr.lock().unwrap().push(**v);
}
}));
}
for i in 1..=ITERS {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
for k in 0..WIDE {
assert_eq!(
*recv[k].lock().unwrap(),
vec![0i64],
"first instance {k} emitted more than the initial value under the wide wave"
);
assert!(cells[k].is_complete(), "first instance {k} never completed");
}
drop(guards);
}
#[test]
fn last_wide_independent_wave_each_instance_exact() {
let _serial = scheduler_test_serial();
let sources: Vec<Cell<i64, CellMutable>> = (0..WIDE).map(|_| Cell::new(0i64)).collect();
let cells: Vec<_> = sources
.iter()
.map(|s| s.clone().last().materialize())
.collect();
let recv: Vec<Arc<Mutex<Vec<i64>>>> = (0..WIDE)
.map(|_| Arc::new(Mutex::new(Vec::new())))
.collect();
let mut guards = Vec::new();
for (cell, r) in cells.iter().zip(&recv) {
let rr = r.clone();
guards.push(cell.subscribe(move |sig| {
if let Signal::Value(v) = sig
&& let Some(x) = (**v)
{
rr.lock().unwrap().push(x);
}
}));
}
for i in 1..=ITERS {
batch(|| {
for s in &sources {
s.set(i);
}
});
}
batch(|| {
for s in &sources {
s.complete();
}
});
for k in 0..WIDE {
assert_eq!(
*recv[k].lock().unwrap(),
vec![ITERS],
"last instance {k} emitted the wrong most-recent value under the wide wave"
);
assert_eq!(
cells[k].get(),
Some(ITERS),
"last instance {k} settled wrong"
);
assert!(cells[k].is_complete(), "last instance {k} never completed");
}
drop(guards);
}
#[test]
fn pairwise_no_coalesce_burst_keeps_pairs_in_order() {
let _serial = scheduler_test_serial();
let (src, pairs) = no_coalesce(|| {
let src = Cell::<i64, CellMutable>::new(0);
let pairs = src.clone().pairwise().materialize();
(src, pairs)
});
let received: Arc<Mutex<Vec<(i64, i64)>>> = Arc::new(Mutex::new(Vec::new()));
let r = received.clone();
let guard = pairs.subscribe(move |sig| {
if let Signal::Value(v) = sig
&& let Some(pair) = (**v)
{
r.lock().unwrap().push(pair);
}
});
let mut next = 1i64;
for it in 0..ORDER_ITERS {
received.lock().unwrap().clear();
let lo = next;
batch(|| {
for _ in 0..BURST {
src.set(next);
next += 1;
}
});
let hi = next - 1;
let expected: Vec<(i64, i64)> = (lo..=hi).map(|x| (x - 1, x)).collect();
let got = received.lock().unwrap().clone();
assert_eq!(
got, expected,
"iter {it}: pairwise produced REORDERED pairs under the parallel no_coalesce wave"
);
}
std::mem::forget(guard);
}
#[test]
fn buffer_count_no_coalesce_burst_keeps_chunks_in_order() {
let _serial = scheduler_test_serial();
const COUNT: usize = 4;
let (src, buffered) = no_coalesce(|| {
let src = Cell::<i64, CellMutable>::new(0);
let buffered = src.buffer_count(COUNT);
(src, buffered)
});
let received: Arc<Mutex<Vec<Vec<i64>>>> = Arc::new(Mutex::new(Vec::new()));
let r = received.clone();
let guard = buffered.subscribe(move |sig| {
if let Signal::Value(v) = sig {
r.lock().unwrap().push((**v).clone());
}
});
let mut next = 1i64;
for it in 0..ORDER_ITERS {
received.lock().unwrap().clear();
let lo = next;
batch(|| {
for _ in 0..BURST {
src.set(next);
next += 1;
}
});
let hi = next - 1;
let vals: Vec<i64> = (lo..=hi).collect();
let expected: Vec<Vec<i64>> = vals.chunks(COUNT).map(|c| c.to_vec()).collect();
let got = received.lock().unwrap().clone();
assert_eq!(
got, expected,
"iter {it}: buffer_count produced REORDERED/torn chunks under the parallel no_coalesce wave"
);
}
std::mem::forget(guard);
}
#[test]
fn window_no_coalesce_burst_keeps_windows_in_order() {
let _serial = scheduler_test_serial();
const COUNT: i64 = 4;
let (src, windowed) = no_coalesce(|| {
let src = Cell::<i64, CellMutable>::new(0);
let windowed = src.window(COUNT as usize);
(src, windowed)
});
let received: Arc<Mutex<Vec<Vec<i64>>>> = Arc::new(Mutex::new(Vec::new()));
let r = received.clone();
let guard = windowed.subscribe(move |sig| {
if let Signal::Value(v) = sig {
r.lock().unwrap().push((**v).clone());
}
});
let mut next = 1i64;
for it in 0..ORDER_ITERS {
received.lock().unwrap().clear();
let lo = next;
batch(|| {
for _ in 0..BURST {
src.set(next);
next += 1;
}
});
let hi = next - 1;
let expected: Vec<Vec<i64>> = (lo..=hi)
.map(|x| {
let low = (x + 1 - COUNT).max(0);
(low..=x).collect()
})
.collect();
let got = received.lock().unwrap().clone();
assert_eq!(
got, expected,
"iter {it}: window produced REORDERED/torn windows under the parallel no_coalesce wave"
);
}
std::mem::forget(guard);
}