use courierust::courierust_pool::ThreadPool;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
fn main() -> std::io::Result<()> {
let pool = ThreadPool::with_size(4)?;
println!("pool running with {} workers", pool.len());
let counter = Arc::new(AtomicUsize::new(0));
let jobs = 4096;
for _ in 0..jobs {
let c = counter.clone();
pool.spawn(move || {
c.fetch_add(1, Ordering::SeqCst);
});
}
wait_until(|| counter.load(Ordering::SeqCst) == jobs);
assert_eq!(counter.load(Ordering::SeqCst), jobs);
println!("{jobs} independent jobs completed");
let pool = Arc::new(pool);
let nested = Arc::new(AtomicUsize::new(0));
let parent_count = 8;
let per_parent = 64;
for _ in 0..parent_count {
let n = nested.clone();
let p = pool.clone();
pool.spawn(move || {
for _ in 0..per_parent {
let n = n.clone();
p.spawn(move || {
n.fetch_add(1, Ordering::SeqCst);
});
}
});
}
let expected = parent_count * per_parent;
wait_until(|| nested.load(Ordering::SeqCst) == expected);
assert_eq!(nested.load(Ordering::SeqCst), expected);
println!("{parent_count} parent jobs spawned {expected} nested jobs");
let running = Arc::new(AtomicUsize::new(0));
let done = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let start = Instant::now();
for _ in 0..8 {
let running = running.clone();
let done = done.clone();
let peak = peak.clone();
pool.spawn(move || {
let now = running.fetch_add(1, Ordering::SeqCst) + 1;
peak.fetch_max(now, Ordering::SeqCst);
std::thread::sleep(Duration::from_millis(50));
running.fetch_sub(1, Ordering::SeqCst);
done.fetch_add(1, Ordering::SeqCst);
});
}
wait_until(|| done.load(Ordering::SeqCst) == 8);
let elapsed = start.elapsed();
println!(
"8 x 50ms jobs finished in {:?} (peak {}/4 workers busy)",
elapsed,
peak.load(Ordering::SeqCst)
);
assert!(
peak.load(Ordering::SeqCst) >= 2,
"the pool should overlap slow jobs across workers"
);
println!("pool drops cleanly (workers joined on Drop)");
Ok(())
}
fn wait_until(cond: impl Fn() -> bool) {
let deadline = Instant::now() + Duration::from_secs(15);
while !cond() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(2));
}
assert!(cond(), "timed out waiting for pool jobs");
}