use std::task::Poll;
use loom::{future::block_on, thread};
use crate::{Producer, Ref, Shared};
fn equals(n: u32) -> impl FnMut(&Ref<'_, u32>) -> Poll<()> + Unpin {
move |v: &Ref<'_, u32>| if **v == n { Poll::Ready(()) } else { Poll::Pending }
}
#[test]
fn write_wakes_a_parked_consumer() {
loom::model(|| {
let producer = Producer::new(0u32);
let consumer = producer.consume();
let writer = thread::spawn(move || {
*producer.write().ok().expect("open") = 1;
});
assert_eq!(block_on(consumer.wait(equals(1))), Ok(()), "the write was lost");
writer.join().unwrap();
});
}
#[test]
fn concurrent_writes_never_lose_a_wakeup() {
loom::model(|| {
let producer = Producer::new(0u32);
let consumer = producer.consume();
let second = producer.clone();
let a = thread::spawn(move || *producer.write().ok().expect("open") += 1);
let b = thread::spawn(move || *second.write().ok().expect("open") += 1);
assert_eq!(block_on(consumer.wait(equals(2))), Ok(()), "a write was lost");
a.join().unwrap();
b.join().unwrap();
});
}
#[test]
fn racing_last_producer_drops_still_close() {
loom::model(|| {
let producer = Producer::new(0u32);
let second = producer.clone();
let consumer = producer.consume();
let a = thread::spawn(move || drop(producer));
let b = thread::spawn(move || drop(second));
block_on(consumer.closed());
a.join().unwrap();
b.join().unwrap();
});
}
#[test]
fn weak_upgrade_never_resurrects_a_closed_channel() {
loom::model(|| {
let producer = Producer::new(0u32);
let weak = producer.weak();
let closer = thread::spawn(move || drop(producer));
let upgraded = weak.produce();
closer.join().unwrap();
if let Some(upgraded) = upgraded {
assert!(
upgraded.write().is_ok(),
"produce() handed back a producer on a closed channel"
);
}
});
}
#[test]
fn weak_downgrade_upgrade_races_the_last_drop() {
loom::model(|| {
let producer = Producer::new(0u32);
let weak = producer.downgrade();
let closer = thread::spawn(move || drop(producer));
let upgraded = weak.upgrade();
closer.join().unwrap();
if let Some(upgraded) = upgraded {
assert!(
upgraded.write().is_ok(),
"upgrade() handed back a producer on a closed channel"
);
}
});
}
#[test]
fn first_consumer_wakes_used() {
loom::model(|| {
let producer = Producer::new(0u32);
let second = producer.clone();
let maker = thread::spawn(move || second.consume());
assert_eq!(block_on(producer.used()), Ok(()), "the new consumer was missed");
drop(maker.join().unwrap());
});
}
#[test]
fn last_consumer_wakes_unused() {
loom::model(|| {
let producer = Producer::new(0u32);
let consumer = producer.consume();
let dropper = thread::spawn(move || drop(consumer));
assert_eq!(block_on(producer.unused()), Ok(()), "the last drop was missed");
dropper.join().unwrap();
});
}
#[test]
fn consumer_churn_resolves_unused() {
loom::model(|| {
let producer = Producer::new(0u32);
let second = producer.clone();
let churn = thread::spawn(move || drop(second.consume()));
assert_eq!(block_on(producer.unused()), Ok(()), "unused() stalled on churn");
churn.join().unwrap();
});
}
#[test]
fn shared_mutation_wakes_a_parked_handle() {
loom::model(|| {
let shared = Shared::new(0u32);
let other = shared.clone();
let writer = thread::spawn(move || *other.lock() = 1);
drop(block_on(shared.wait(equals(1))));
writer.join().unwrap();
});
}