use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Duration;
use crate::callback_env::CallbackEnviron;
use crate::work::ThreadpoolWork;
#[test]
fn new_with_no_env_succeeds() {
let work = ThreadpoolWork::new(|| {}, None);
assert!(work.is_ok());
}
#[test]
fn new_with_default_env_succeeds() {
let mut env = CallbackEnviron::new();
let work = ThreadpoolWork::new(|| {}, Some(&mut env));
assert!(work.is_ok());
}
#[test]
fn new_env_with_runs_long_succeeds() {
let mut env = CallbackEnviron::new();
env.set_runs_long();
let work = ThreadpoolWork::new(|| {}, Some(&mut env));
assert!(work.is_ok());
}
#[test]
fn submit_once_callback_runs() {
let ran = Arc::new(AtomicBool::new(false));
let r = Arc::clone(&ran);
let work = ThreadpoolWork::new(move || r.store(true, Ordering::SeqCst), None).unwrap();
work.submit();
work.wait();
assert!(ran.load(Ordering::SeqCst));
}
#[test]
fn submit_increments_counter_once() {
let count = Arc::new(AtomicUsize::new(0));
let c = Arc::clone(&count);
let work = ThreadpoolWork::new(
move || {
c.fetch_add(1, Ordering::SeqCst);
},
None,
)
.unwrap();
work.submit();
work.wait();
assert_eq!(count.load(Ordering::SeqCst), 1);
}
#[test]
fn submit_five_times_counter_is_five() {
let count = Arc::new(AtomicUsize::new(0));
let c = Arc::clone(&count);
let work = ThreadpoolWork::new(
move || {
c.fetch_add(1, Ordering::SeqCst);
},
None,
)
.unwrap();
for _ in 0..5 {
work.submit();
}
work.wait();
assert_eq!(count.load(Ordering::SeqCst), 5);
}
#[test]
fn submit_ten_times_counter_is_ten() {
let count = Arc::new(AtomicUsize::new(0));
let c = Arc::clone(&count);
let work = ThreadpoolWork::new(
move || {
c.fetch_add(1, Ordering::SeqCst);
},
None,
)
.unwrap();
for _ in 0..10 {
work.submit();
}
work.wait();
assert_eq!(count.load(Ordering::SeqCst), 10);
}
#[test]
fn drop_waits_for_in_flight_callback() {
let done = Arc::new(AtomicBool::new(false));
let d = Arc::clone(&done);
{
let work = ThreadpoolWork::new(
move || {
std::thread::sleep(Duration::from_millis(5));
d.store(true, Ordering::SeqCst);
},
None,
)
.unwrap();
work.submit();
}
assert!(done.load(Ordering::SeqCst));
}
#[test]
fn callback_can_own_arc_data() {
let data = Arc::new(vec![1u64, 2, 3, 4, 5]);
let d = Arc::clone(&data);
let sum = Arc::new(AtomicUsize::new(0));
let s = Arc::clone(&sum);
let work = ThreadpoolWork::new(
move || {
s.fetch_add(
d.iter().map(|x| *x as usize).sum::<usize>(),
Ordering::SeqCst,
);
},
None,
)
.unwrap();
work.submit();
work.wait();
assert_eq!(sum.load(Ordering::SeqCst), 15);
}
#[test]
fn context_ref_count_valid_during_callback() {
let data = Arc::new(AtomicUsize::new(0));
let d = Arc::clone(&data);
let inner_count = Arc::new(AtomicUsize::new(0));
let ic = Arc::clone(&inner_count);
let work = ThreadpoolWork::new(
move || {
ic.store(Arc::strong_count(&d), Ordering::SeqCst);
},
None,
)
.unwrap();
work.submit();
work.wait();
assert!(inner_count.load(Ordering::SeqCst) >= 2);
}
#[test]
fn cancel_pending_does_not_panic() {
let count = Arc::new(AtomicUsize::new(0));
let c = Arc::clone(&count);
let work = ThreadpoolWork::new(
move || {
c.fetch_add(1, Ordering::SeqCst);
},
None,
)
.unwrap();
for _ in 0..10 {
work.submit();
}
work.cancel_pending();
assert!(count.load(Ordering::SeqCst) <= 10);
}
#[test]
fn resubmit_after_cancel_runs_once() {
let count = Arc::new(AtomicUsize::new(0));
let c = Arc::clone(&count);
let work = ThreadpoolWork::new(
move || {
c.fetch_add(1, Ordering::SeqCst);
},
None,
)
.unwrap();
for _ in 0..5 {
work.submit();
}
work.cancel_pending();
count.store(0, Ordering::SeqCst);
work.submit();
work.wait();
assert_eq!(count.load(Ordering::SeqCst), 1);
}
#[test]
fn wait_then_resubmit_counts_independently() {
let count = Arc::new(AtomicUsize::new(0));
let c = Arc::clone(&count);
let work = ThreadpoolWork::new(
move || {
c.fetch_add(1, Ordering::SeqCst);
},
None,
)
.unwrap();
work.submit();
work.wait();
assert_eq!(count.load(Ordering::SeqCst), 1);
work.submit();
work.wait();
assert_eq!(count.load(Ordering::SeqCst), 2);
}
#[test]
fn no_submit_drop_is_safe() {
let _work = ThreadpoolWork::new(|| {}, None).unwrap();
}
#[test]
fn work_with_env_callback_runs() {
let ran = Arc::new(AtomicBool::new(false));
let r = Arc::clone(&ran);
let mut env = CallbackEnviron::new();
let work =
ThreadpoolWork::new(move || r.store(true, Ordering::SeqCst), Some(&mut env)).unwrap();
work.submit();
work.wait();
assert!(ran.load(Ordering::SeqCst));
}
#[test]
fn work_with_runs_long_env_callback_runs() {
let count = Arc::new(AtomicUsize::new(0));
let c = Arc::clone(&count);
let mut env = CallbackEnviron::new();
env.set_runs_long();
let work = ThreadpoolWork::new(
move || {
c.fetch_add(1, Ordering::SeqCst);
},
Some(&mut env),
)
.unwrap();
work.submit();
work.wait();
assert_eq!(count.load(Ordering::SeqCst), 1);
}