use crate::sync::{Condvar, Mutex};
use crate::{Clock, TestClock, Waiter};
use loom::sync::mpsc;
use loom::thread;
use std::sync::Arc;
use std::time::{Duration, UNIX_EPOCH};
fn hold_park(clock: &Clock, waiter: &Waiter) -> mpsc::Sender<()> {
let (checked, checks) = mpsc::channel();
let (resume, resumed) = mpsc::channel();
*clock.paused.as_ref().unwrap().before_rewait.lock().unwrap() = Some(Box::new(move |_| {
checked.send(()).unwrap();
resumed.recv().unwrap();
}));
waiter.signal.wake();
checks.recv().unwrap();
resume
}
fn assert_idle(tester: &TestClock, clock: &Clock) {
assert_eq!(tester.next_deadline(), None);
let state = clock.paused.as_ref().unwrap().state.lock().unwrap();
assert_eq!(state.blocked, 0);
}
fn notification_model(holding: bool) {
loom::model(move || {
let tester = TestClock::new();
let clock = tester.clock();
let pair = Arc::new((Mutex::new(false), Condvar::new(&clock)));
let waiting = thread::spawn({
let pair = pair.clone();
move || {
let ready = pair.1.wait_while(pair.0.lock().unwrap(), |ready| !*ready);
assert!(*ready.unwrap());
}
});
let mut ready = pair.0.lock().unwrap();
*ready = true;
if holding {
pair.1.notify_one();
drop(ready);
} else {
drop(ready);
pair.1.notify_one();
}
waiting.join().unwrap();
drop(pair);
assert_idle(&tester, &clock);
});
}
#[test]
fn test_wait_observes_notification_under_mutex() {
notification_model(true);
}
#[test]
fn test_wait_observes_notification_after_unlock() {
notification_model(false);
}
fn two_waiters_model(broadcast: bool) {
loom::model(move || {
let tester = TestClock::new();
let clock = tester.clock();
let pair = Arc::new((Mutex::new(false), Condvar::new(&clock)));
let waiters: Vec<_> = (0..2)
.map(|_| {
let pair = pair.clone();
thread::spawn(move || {
let ready = pair.1.wait_while(pair.0.lock().unwrap(), |ready| !*ready);
assert!(*ready.unwrap());
})
})
.collect();
tester.wait_blocked(2);
*pair.0.lock().unwrap() = true;
if broadcast {
pair.1.notify_all();
} else {
pair.1.notify_one();
pair.1.notify_one();
}
for waiter in waiters {
waiter.join().unwrap();
}
drop(pair);
assert_idle(&tester, &clock);
});
}
#[test]
fn test_two_notify_one_calls_release_two_waiters() {
two_waiters_model(false);
}
#[test]
fn test_notify_all_releases_two_waiters() {
two_waiters_model(true);
}
#[test]
fn test_wait_deadline_observes_exact_advance() {
loom::model(|| {
let mut tester = TestClock::new();
tester.set_system_time(UNIX_EPOCH);
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let waiting = thread::spawn({
let clock = clock.clone();
move || {
let condvar = Condvar::new(&clock);
let mutex = Mutex::new(());
let (_guard, result) = condvar
.wait_deadline(mutex.lock().unwrap(), deadline)
.unwrap();
assert!(result.timed_out());
assert_eq!(clock.now(), deadline);
assert_eq!(clock.system_time(), UNIX_EPOCH + Duration::from_secs(1));
}
});
tester.advance_to(deadline);
waiting.join().unwrap();
assert_idle(&tester, &clock);
});
}
#[test]
fn test_wait_deadline_observes_notification_after_short_advance() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
let deadline = start + Duration::from_secs(2);
let pair = Arc::new((Mutex::new(false), Condvar::new(&clock)));
let waiting = thread::spawn({
let pair = pair.clone();
move || {
let (ready, result) = pair
.1
.wait_deadline(pair.0.lock().unwrap(), deadline)
.unwrap();
assert!(!result.timed_out());
assert!(*ready);
}
});
tester.wait_blocked(1);
tester.advance(Duration::from_secs(1));
*pair.0.lock().unwrap() = true;
pair.1.notify_one();
waiting.join().unwrap();
assert_eq!(clock.now(), start + Duration::from_secs(1));
drop(pair);
assert_idle(&tester, &clock);
});
}
#[test]
fn test_sleep_until_observes_exact_advance() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let sleeper = thread::spawn({
let clock = clock.clone();
move || {
clock.sleep_until(deadline);
assert_eq!(clock.now(), deadline);
}
});
tester.advance_to(deadline);
sleeper.join().unwrap();
assert_idle(&tester, &clock);
});
}
#[test]
fn test_wait_blocked_exposes_sleep_deadline() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let sleeper = thread::spawn({
let clock = clock.clone();
move || {
clock.sleep_until(deadline);
assert_eq!(clock.now(), deadline);
}
});
tester.wait_blocked(1);
assert_eq!(tester.next_deadline(), Some(deadline));
tester.advance_to(deadline);
sleeper.join().unwrap();
assert_idle(&tester, &clock);
});
}
#[test]
fn test_notify_all_races_advance_with_mixed_waiters() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let pair = Arc::new((Mutex::new((false, false)), Condvar::new(&clock)));
let timed = thread::spawn({
let (pair, clock) = (pair.clone(), clock.clone());
move || {
let mut ready = pair.0.lock().unwrap();
loop {
let (next, result) = pair.1.wait_deadline(ready, deadline).unwrap();
ready = next;
if result.timed_out() {
assert_eq!(clock.now(), deadline);
break;
}
if ready.1 {
break;
}
}
}
});
let untimed = thread::spawn({
let pair = pair.clone();
move || {
let mut ready = pair
.1
.wait_while(pair.0.lock().unwrap(), |ready| !ready.0)
.unwrap();
assert!(ready.0);
ready.1 = true;
drop(ready);
pair.1.notify_all();
}
});
tester.wait_blocked(2);
pair.0.lock().unwrap().0 = true;
pair.1.notify_all();
tester.advance_to(deadline);
timed.join().unwrap();
untimed.join().unwrap();
drop(pair);
assert_idle(&tester, &clock);
});
}
#[test]
fn test_partial_advance_leaves_sleep_listed() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(2);
let sleeper = thread::spawn({
let clock = clock.clone();
move || {
clock.sleep_until(deadline);
assert_eq!(clock.now(), deadline);
}
});
tester.advance(Duration::from_secs(1));
tester.wait_blocked(1);
assert_eq!(tester.next_deadline(), Some(deadline));
tester.advance_to(deadline);
sleeper.join().unwrap();
assert_idle(&tester, &clock);
});
}
#[test]
fn test_reached_wait_keeps_later_wait_on_condvar_listed() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let earlier = clock.now() + Duration::from_secs(1);
let later = clock.now() + Duration::from_secs(2);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let waits: Vec<_> = [earlier, later]
.into_iter()
.map(|deadline| {
let pair = pair.clone();
thread::spawn(move || {
let (_guard, result) = pair
.1
.wait_deadline(pair.0.lock().unwrap(), deadline)
.unwrap();
assert!(result.timed_out());
})
})
.collect();
let mut waits = waits.into_iter();
tester.advance_to(earlier);
waits.next().unwrap().join().unwrap();
tester.wait_blocked(1);
assert_eq!(tester.next_deadline(), Some(later));
tester.advance_to(later);
waits.next().unwrap().join().unwrap();
drop(pair);
assert_idle(&tester, &clock);
});
}
#[test]
fn test_notification_survives_advance_with_mixed_waiters() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let pair = Arc::new((Mutex::new(false), Condvar::new(&clock)));
let timed = thread::spawn({
let (pair, clock) = (pair.clone(), clock.clone());
move || {
let (ready, result) = pair
.1
.wait_deadline(pair.0.lock().unwrap(), deadline)
.unwrap();
if result.timed_out() {
assert_eq!(clock.now(), deadline);
} else {
assert!(*ready);
}
}
});
let untimed = thread::spawn({
let pair = pair.clone();
move || {
let ready = pair.1.wait_while(pair.0.lock().unwrap(), |ready| !*ready);
assert!(*ready.unwrap());
}
});
tester.wait_blocked(2);
let notifying = thread::spawn({
let pair = pair.clone();
move || {
*pair.0.lock().unwrap() = true;
pair.1.notify_all();
}
});
tester.advance_to(deadline);
notifying.join().unwrap();
timed.join().unwrap();
untimed.join().unwrap();
drop(pair);
assert_idle(&tester, &clock);
});
}
#[test]
fn test_consumed_notification_keeps_unreached_sibling_listed() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let earlier = clock.now() + Duration::from_secs(1);
let later = clock.now() + Duration::from_secs(2);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let (report, reports) = mpsc::channel();
let waits: Vec<_> = (0..2)
.map(|_| {
let (pair, report) = (pair.clone(), report.clone());
thread::spawn(move || {
let (guard, result) =
pair.1.wait_deadline(pair.0.lock().unwrap(), later).unwrap();
report.send(Some(result.timed_out())).unwrap();
if !result.timed_out() {
let (_guard, result) = pair.1.wait_deadline(guard, earlier).unwrap();
report.send(Some(result.timed_out())).unwrap();
}
})
})
.collect();
tester.wait_blocked(2);
pair.1.notify_one();
assert_eq!(reports.recv().unwrap(), Some(false));
tester.wait_blocked(2);
*clock.paused.as_ref().unwrap().before_rewait.lock().unwrap() = Some(Box::new(move |_| {
report.send(None).unwrap();
}));
tester.advance_to(earlier);
let events = [reports.recv().unwrap(), reports.recv().unwrap()];
assert!(events.contains(&None));
assert!(events.contains(&Some(true)));
assert_eq!(tester.next_deadline(), Some(later));
tester.advance_to(later);
assert_eq!(reports.recv().unwrap(), Some(true));
for wait in waits {
wait.join().unwrap();
}
assert_idle(&tester, &clock);
});
}
#[test]
fn test_later_wait_cannot_take_pending_notification() {
loom::model(|| {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let waiter = Arc::new(clock.waiter());
let first = thread::spawn({
let waiter = waiter.clone();
move || waiter.wait(Some(deadline))
});
tester.wait_blocked(1);
let resume = hold_park(&clock, &waiter);
waiter.signal.notify_one();
let next = thread::spawn({
let waiter = waiter.clone();
move || waiter.wait(Some(deadline))
});
tester.wait_blocked(2);
tester.advance_to(deadline);
assert!(!next.join().unwrap());
resume.send(()).unwrap();
assert!(first.join().unwrap());
assert_idle(&tester, &clock);
});
}
#[test]
fn test_broadcast_then_single_notification_preserves_later_wait() {
loom::model(|| {
let tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let waiter = Arc::new(clock.waiter());
let first = thread::spawn({
let waiter = waiter.clone();
move || waiter.wait(Some(deadline))
});
tester.wait_blocked(1);
let resume = hold_park(&clock, &waiter);
waiter.signal.notify_all();
waiter.signal.notify_one();
resume.send(()).unwrap();
assert!(first.join().unwrap());
{
let state = waiter.signal.lock();
assert_eq!(state.waiting, 0);
assert!(state.pending.is_empty());
}
let next = thread::spawn({
let waiter = waiter.clone();
move || waiter.wait(Some(deadline))
});
tester.wait_blocked(1);
waiter.signal.notify_one();
assert!(next.join().unwrap());
assert_idle(&tester, &clock);
});
}
fn reached_signals_model(equal: bool, condvars: bool) {
loom::model(move || {
let mut tester = TestClock::new();
let clock = tester.clock();
let earlier = clock.now() + Duration::from_secs(1);
let later = clock.now() + Duration::from_secs(2);
let waits: Vec<_> = [if equal { later } else { earlier }, later]
.into_iter()
.map(|deadline| {
let clock = clock.clone();
thread::spawn(move || {
if condvars {
let condvar = Condvar::new(&clock);
let mutex = Mutex::new(());
let (_guard, result) = condvar
.wait_deadline(mutex.lock().unwrap(), deadline)
.unwrap();
assert!(result.timed_out());
} else {
clock.sleep_until(deadline);
}
assert_eq!(clock.now(), later);
})
})
.collect();
tester.wait_blocked(2);
tester.advance_to(later);
for wait in waits {
wait.join().unwrap();
}
assert_idle(&tester, &clock);
});
}
#[test]
fn test_advance_reaches_two_sleep_signals() {
reached_signals_model(false, false);
}
#[test]
fn test_advance_reaches_equal_sleep_signals() {
reached_signals_model(true, false);
}
#[test]
fn test_advance_reaches_two_condvar_signals() {
reached_signals_model(false, true);
}