use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use rand::Rng;
use takeaway::{Queue, Worker, util::block_on};
struct Global {
queue: Queue<MyTask>,
generated: AtomicUsize,
}
impl Global {
fn finishing(&self) -> bool {
self.generated.load(Ordering::Relaxed) >= 4096
}
}
struct MyTask {
counter: Arc<()>,
}
impl takeaway::Task for MyTask {
type Priority = ();
fn priority(&self) -> Self::Priority {}
}
async fn worker(global: &Global, id: usize) -> Box<[Arc<()>]> {
let mut worker = Worker::new(&global.queue, id);
let mut slots = Vec::new();
let mut rng = rand::rng();
let counter = Arc::new(());
slots.push(counter.clone());
let task = MyTask { counter };
worker.enqueuer().add(task);
while let Some(task) = worker.next().await {
assert_eq!(Arc::strong_count(&task.counter), 2);
std::mem::drop(task.counter);
if !global.finishing() {
let num = rng.random_range(1..4);
global.generated.fetch_add(num, Ordering::Relaxed);
for _ in 0..num {
let counter = Arc::new(());
slots.push(counter.clone());
let task = MyTask { counter };
worker.enqueue_one(task);
}
}
}
slots.into_boxed_slice()
}
#[test]
fn simple() {
let config = takeaway::Config::default().with_oneshot(true);
let num_workers = config.num_workers().get();
let global = Global {
queue: config.build(),
generated: AtomicUsize::new(num_workers),
};
std::thread::scope(|s| {
let global = &global;
let handles = (0..num_workers)
.map(|id| s.spawn(move || block_on(worker(global, id))))
.collect::<Box<[_]>>();
let slots = handles
.into_iter()
.flat_map(|handle| handle.join().unwrap())
.collect::<Box<[_]>>();
assert!(global.finishing());
for slot in slots {
assert_eq!(Arc::strong_count(&slot), 1);
}
println!(
"Executed {} tasks",
global.generated.load(Ordering::Relaxed)
);
})
}