#![cfg(feature = "scheduler")]
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use parking_lot::Mutex;
use hyphae::{
Cell, Gettable, Materialize, MergeExt, MergeMapExt, Mutable, Signal, SwitchMapExt, Watchable,
batch, scheduler::no_coalesce,
};
const WIDTH: usize = 8;
const ITERATIONS: i64 = 3_000;
static SCHEDULER_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
fn iteration_value(iteration: i64, index: usize, offset: i64) -> i64 {
iteration
.saturating_mul(1000)
.saturating_add(i64::try_from(index).unwrap_or(i64::MAX).saturating_mul(2))
.saturating_add(offset)
}
fn scheduler_test_serial() -> std::sync::MutexGuard<'static, ()> {
hyphae::scheduler::set_wave_threshold_for_test(4);
SCHEDULER_TEST_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
#[test]
fn concurrent_merges_never_drop_or_duplicate_and_complete_once() {
let _serial = scheduler_test_serial();
let mut lefts = Vec::new();
let mut rights = Vec::new();
let mut received = Vec::new();
let mut completes = Vec::new();
let mut guards = Vec::new();
for _ in 0..WIDTH {
let (left, right, merged) = no_coalesce(|| {
let left = Cell::new(0i64);
let right = Cell::new(0i64);
let merged = left.clone().merge(right.clone()).materialize();
(left, right, merged)
});
let got = Arc::new(Mutex::new(Vec::<i64>::new()));
let done = Arc::new(AtomicUsize::new(0));
let (received_buffer, completion_count) = (got.clone(), done.clone());
let guard = merged.subscribe(move |signal| match signal {
Signal::Value(value) => received_buffer.lock().push(**value),
Signal::Complete => {
completion_count.fetch_add(1, Ordering::SeqCst);
}
Signal::Error(_) => {}
});
lefts.push(left);
rights.push(right);
received.push(got);
completes.push(done);
guards.push(guard);
}
for it in 1..=ITERATIONS {
for buf in &received {
buf.lock().clear();
}
batch(|| {
for (index, (left, right)) in lefts.iter().zip(&rights).enumerate() {
let base = iteration_value(it, index, 0);
left.set(base);
right.set(base.saturating_add(1));
}
});
for (index, buffer) in received.iter().enumerate() {
let base = iteration_value(it, index, 0);
let mut got = buffer.lock().clone();
got.sort_unstable();
assert_eq!(
got,
vec![base, base.saturating_add(1)],
"merge dropped/duplicated/mis-forwarded an emission at iteration {it}, merge {index}"
);
}
}
batch(|| {
for (left, right) in lefts.iter().zip(&rights) {
left.complete();
right.complete();
}
});
for (index, complete) in completes.iter().enumerate() {
assert_eq!(
complete.load(Ordering::SeqCst),
1,
"merge {index} fired Complete {} times (expected exactly 1)",
complete.load(Ordering::SeqCst)
);
}
drop(guards);
}
#[test]
fn concurrent_merge_maps_complete_exactly_once() {
let _serial = scheduler_test_serial();
for it in 0..ITERATIONS {
let mut outers = Vec::new();
let mut inner_srcs = Vec::new();
let mut completes = Vec::new();
let mut guards = Vec::new();
no_coalesce(|| {
for _ in 0..WIDTH {
let outer = Cell::new(0i64);
let inner_src = Cell::new(0i64);
let inner_for_merge = inner_src.clone();
let merged = outer
.clone()
.merge_map(move |_: &i64| inner_for_merge.clone().lock())
.materialize();
let done = Arc::new(AtomicUsize::new(0));
let completion_count = done.clone();
let guard = merged.subscribe(move |signal| {
if matches!(signal, Signal::Complete) {
completion_count.fetch_add(1, Ordering::SeqCst);
}
});
outers.push(outer);
inner_srcs.push(inner_src);
completes.push(done);
guards.push(guard);
}
});
batch(|| {
for (outer, inner) in outers.iter().zip(&inner_srcs) {
outer.complete();
inner.complete();
}
});
for (index, complete) in completes.iter().enumerate() {
assert_eq!(
complete.load(Ordering::SeqCst),
1,
"merge_map fired Complete {} times at iteration {it}, unit {index} (expected exactly 1)",
complete.load(Ordering::SeqCst)
);
}
drop(guards);
}
}
#[test]
fn concurrent_switch_map_latest_inner_wins_same_height() {
let _serial = scheduler_test_serial();
for it in 1..=ITERATIONS {
let mut sels = Vec::new();
let mut old_srcs = Vec::new();
let mut results = Vec::new();
for index in 0..WIDTH {
let sel = Cell::new(0i64); let old_src = Cell::new(0i64);
let new_value = iteration_value(it, index, 500);
let new_src = Cell::new(new_value);
let old = old_src.clone().lock();
let new = new_src.clone().lock();
let result = sel
.clone()
.switch_map(move |&selector| {
if selector == 0 {
old.clone()
} else {
new.clone()
}
})
.materialize();
sels.push(sel);
old_srcs.push(old_src);
results.push(result);
}
batch(|| {
for (index, (old_source, selector)) in old_srcs.iter().zip(&sels).enumerate() {
old_source.set(iteration_value(it, index, 0));
selector.set(1); }
});
for (index, result) in results.iter().enumerate() {
let new_value = iteration_value(it, index, 500);
assert_eq!(
result.get(),
new_value,
"switch_map settled on a STALE old-inner value at iteration {it}, unit {index} \
(expected the switched-in inner's {new_value})"
);
}
}
}