use super::*;
use crate::TestClock;
use crate::tests::helpers::{blocked, pause_before_park, pause_before_rewait, wakes};
use std::sync::{self, Arc, mpsc};
use std::thread::{self, JoinHandle};
use std::time::Duration;
fn poison<G>(lock: impl FnOnce() -> G + Send) {
thread::scope(|scope| {
let panicked = scope
.spawn(|| {
let _guard = lock();
panic!("poison the mutex");
})
.join();
assert!(panicked.is_err());
});
}
fn start_wait<T, R>(
pair: &Arc<(Mutex<T>, Condvar)>,
wait: impl for<'a> FnOnce(&'a Condvar, MutexGuard<'a, T>) -> R + Send + 'static,
) -> JoinHandle<R>
where
T: Send + 'static,
R: Send + 'static,
{
let (locked, locks) = mpsc::sync_channel(0);
let waiting = thread::spawn({
let pair = pair.clone();
move || {
let guard = pair.0.lock().unwrap();
locked.send(()).unwrap();
wait(&pair.1, guard)
}
});
locks.recv().unwrap();
drop(pair.0.lock().unwrap());
waiting
}
fn wait_unblocked(clock: &Clock) {
let paused = clock.paused.as_ref().unwrap();
let mut state = paused.state.lock().unwrap();
while state.blocked != 0 {
state = paused.changed.wait(state).unwrap();
}
}
fn hold_park(clock: &Clock, condvar: &Condvar) -> mpsc::SyncSender<()> {
let (checked, resume) = pause_before_rewait(clock);
{
let _state = condvar.waiter.signal.lock();
condvar.waiter.signal.changed.notify_all();
}
assert_eq!(checked.recv().unwrap(), None);
resume
}
#[test]
fn test_mutex_try_lock_reports_contention() {
let mutex = Mutex::new(3);
let mut guard = mutex.lock().unwrap();
*guard = 7;
assert!(matches!(mutex.try_lock(), Err(TryLockError::WouldBlock)));
drop(guard);
assert_eq!(*mutex.try_lock().unwrap(), 7);
}
#[test]
fn test_mutex_poisoning_preserves_guard_and_recovers() {
let mutex = Mutex::new(3);
poison(|| mutex.lock().unwrap());
assert!(mutex.is_poisoned());
let mut guard = mutex.lock().unwrap_err().into_inner();
assert_eq!(*guard, 3);
*guard = 7;
drop(guard);
match mutex.try_lock() {
Err(TryLockError::Poisoned(err)) => assert_eq!(*err.into_inner(), 7),
result => panic!("expected a poison error, got {result:?}"),
}
mutex.clear_poison();
assert!(!mutex.is_poisoned());
assert_eq!(*mutex.lock().unwrap(), 7);
}
#[test]
fn test_mutex_accessors_preserve_poisoned_values() {
let mut healthy = Mutex::<Vec<u8>>::default();
healthy.get_mut().unwrap().push(7);
assert_eq!(healthy.into_inner().unwrap(), vec![7]);
let mut poisoned = Mutex::from(vec![3]);
poison(|| poisoned.lock().unwrap());
poisoned.get_mut().unwrap_err().into_inner().push(7);
assert!(poisoned.is_poisoned());
assert_eq!(poisoned.into_inner().unwrap_err().into_inner(), vec![3, 7]);
}
#[test]
fn test_mutex_and_guards_format_like_std_without_blocking() {
let mutex = Mutex::new(7);
let standard = sync::Mutex::new(7);
assert_eq!(format!("{mutex:?}"), format!("{standard:?}"));
let guard = mutex.lock().unwrap();
let standard_guard = standard.lock().unwrap();
assert_eq!(format!("{mutex:?}"), format!("{standard:?}"));
assert_eq!(format!("{guard:?}"), format!("{standard_guard:?}"));
assert_eq!(format!("{guard}"), format!("{standard_guard}"));
drop((guard, standard_guard));
poison(|| mutex.lock().unwrap());
poison(|| standard.lock().unwrap());
assert_eq!(format!("{mutex:?}"), format!("{standard:?}"));
}
#[test]
fn test_notifications_end_waits() {
for all in [false, true] {
for clock in [Clock::real(), TestClock::new().clock()] {
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let (done, returned) = mpsc::channel();
let threads: Vec<_> = (0..2)
.map(|_| {
let done = done.clone();
start_wait(&pair, move |condvar, guard| {
drop(condvar.wait(guard).unwrap());
done.send(()).unwrap();
})
})
.collect();
if all {
pair.1.notify_all();
} else {
for _ in 0..2 {
pair.1.notify_one();
returned.recv().unwrap();
}
}
for waiting in threads {
waiting.join().unwrap();
}
}
}
}
#[test]
fn test_wait_while_rechecks_predicate_after_notifications() {
let tester = TestClock::new();
let pair = Arc::new((Mutex::new(false), Condvar::new(&tester.clock())));
let (checked, checks) = mpsc::channel();
let waiting = start_wait(&pair, move |condvar, guard| {
let guard = condvar
.wait_while(guard, |open| {
checked.send(*open).unwrap();
!*open
})
.unwrap();
assert!(*guard);
});
assert!(!checks.recv().unwrap());
tester.wait_blocked(1);
pair.1.notify_one();
assert!(!checks.recv().unwrap());
tester.wait_blocked(1);
*pair.0.lock().unwrap() = true;
pair.1.notify_one();
assert!(checks.recv().unwrap());
waiting.join().unwrap();
}
#[test]
fn test_advances_never_wake_an_untimed_wait() {
let mut tester = TestClock::new();
let clock = tester.clock();
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let waiting = start_wait(&pair, |condvar, guard| drop(condvar.wait(guard).unwrap()));
tester.wait_blocked(1);
tester.advance(Duration::from_secs(30));
tester.advance_to(clock.now() + Duration::from_secs(30));
assert_eq!(wakes(&pair.1.waiter.signal), 0);
assert_eq!(blocked(&clock), 1);
assert_eq!(tester.next_deadline(), None);
pair.1.notify_one();
waiting.join().unwrap();
}
#[test]
fn test_advances_wake_only_the_waits_they_reach() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
let sleeps: Vec<_> = [false, true]
.into_iter()
.map(|until| {
let clock = clock.clone();
thread::spawn(move || {
if until {
clock.sleep_until(start + Duration::from_secs(3));
} else {
clock.sleep(Duration::from_secs(3));
}
})
})
.collect();
tester.wait_blocked(2);
let deadline = start + Duration::from_secs(5);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let waiting = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
let untimed = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let untimed_thread = start_wait(&untimed, |condvar, guard| {
drop(condvar.wait(guard).unwrap())
});
tester.wait_blocked(4);
tester.advance(Duration::from_secs(2));
for sleep in &sleeps {
assert!(!sleep.is_finished());
}
assert_eq!(wakes(&pair.1.waiter.signal), 0);
assert_eq!(wakes(&untimed.1.waiter.signal), 0);
assert_eq!(blocked(&clock), 4);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(3)));
tester.advance_to(start + Duration::from_secs(3));
for sleep in sleeps {
sleep.join().unwrap();
}
assert_eq!(wakes(&pair.1.waiter.signal), 0);
assert_eq!(wakes(&untimed.1.waiter.signal), 0);
assert_eq!(blocked(&clock), 2);
assert_eq!(tester.next_deadline(), Some(deadline));
tester.advance_to(deadline);
assert_eq!(wakes(&pair.1.waiter.signal), 1);
assert_eq!(wakes(&untimed.1.waiter.signal), 0);
assert!(waiting.join().unwrap().timed_out());
assert_eq!(clock.now(), deadline);
assert_eq!(tester.next_deadline(), None);
assert_eq!(blocked(&clock), 1);
untimed.1.notify_one();
untimed_thread.join().unwrap();
}
#[test]
fn test_wait_deadline_returns_early_on_notification() {
let tester = TestClock::new();
for clock in [Clock::real(), tester.clock()] {
let deadline = clock.now() + Duration::from_secs(86400);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let waiting = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
pair.1.notify_one();
assert!(!waiting.join().unwrap().timed_out(), "{clock:?}");
assert_eq!(tester.next_deadline(), None, "{clock:?}");
}
}
#[test]
fn test_wait_deadline_returns_for_reached_deadlines() {
let mut tester = TestClock::new();
let past = tester.clock().now();
tester.advance(Duration::from_secs(1));
for clock in [Clock::real(), tester.clock()] {
let mutex = Mutex::new(7);
let condvar = Condvar::new(&clock);
for deadline in [past, clock.now()] {
let (guard, result) = condvar
.wait_deadline(mutex.lock().unwrap(), deadline)
.unwrap();
assert!(result.timed_out(), "{clock:?} {deadline:?}");
assert_eq!(*guard, 7, "{clock:?} {deadline:?}");
}
assert_eq!(condvar.waiter.signal.lock().parks, 0, "{clock:?}");
}
}
#[test]
fn test_wait_deadline_observes_advance_before_parking() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let (checked, resume) = pause_before_park(&clock);
let waiting = thread::spawn({
let clock = clock.clone();
move || {
let condvar = Condvar::new(&clock);
let mutex = Mutex::new(());
condvar
.wait_deadline(mutex.lock().unwrap(), deadline)
.unwrap()
.1
}
});
assert_eq!(checked.recv().unwrap(), None);
assert_eq!(tester.next_deadline(), None);
tester.advance_to(deadline);
resume.send(()).unwrap();
assert!(waiting.join().unwrap().timed_out());
assert_eq!(clock.now(), deadline);
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_deadline_worker_recomputes_earliest_deadline() {
let mut tester = TestClock::new();
let clock = tester.clock();
let later = clock.now() + Duration::from_secs(10);
let earlier = clock.now() + Duration::from_secs(3);
let pair = Arc::new((Mutex::new(vec![later]), Condvar::new(&clock)));
let (picked, picks) = mpsc::channel();
let waiting = start_wait(&pair, move |condvar, mut guard| {
loop {
let deadline = *guard.iter().min().unwrap();
picked.send(deadline).unwrap();
let (next, result) = condvar.wait_deadline(guard, deadline).unwrap();
guard = next;
if result.timed_out() {
return deadline;
}
}
});
assert_eq!(picks.recv().unwrap(), later);
tester.wait_blocked(1);
assert_eq!(tester.next_deadline(), Some(later));
pair.0.lock().unwrap().push(earlier);
pair.1.notify_one();
assert_eq!(picks.recv().unwrap(), earlier);
tester.wait_blocked(1);
assert_eq!(tester.next_deadline(), Some(earlier));
tester.advance_to(earlier);
assert_eq!(waiting.join().unwrap(), earlier);
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_notification_between_predicate_and_park_is_observed() {
let tester = TestClock::new();
let clock = tester.clock();
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let (checked, resume) = pause_before_park(&clock);
let waiting = thread::spawn({
let pair = pair.clone();
move || {
let mut checks = 0;
drop(
pair.1
.wait_while(pair.0.lock().unwrap(), |_| {
checks += 1;
checks == 1
})
.unwrap(),
);
checks
}
});
assert_eq!(checked.recv().unwrap(), None);
assert!(matches!(pair.0.try_lock(), Err(TryLockError::WouldBlock)));
pair.1.notify_one();
resume.send(()).unwrap();
assert_eq!(waiting.join().unwrap(), 2);
assert_eq!(pair.1.waiter.signal.lock().parks, 0);
}
#[test]
fn test_wait_relocks_poisoned_mutex() {
for timed in [false, true] {
let tester = TestClock::new();
let deadline = tester.clock().now() + Duration::from_secs(5);
let pair = Arc::new((Mutex::new(7), Condvar::new(&tester.clock())));
let waiting = start_wait(&pair, move |condvar, guard| {
if timed {
let (guard, result) = condvar
.wait_deadline(guard, deadline)
.unwrap_err()
.into_inner();
assert!(!result.timed_out());
*guard
} else {
*condvar.wait(guard).unwrap_err().into_inner()
}
});
tester.wait_blocked(1);
poison(|| pair.0.lock().unwrap());
pair.1.notify_one();
assert_eq!(waiting.join().unwrap(), 7, "{timed}");
assert!(pair.0.is_poisoned(), "{timed}");
assert_eq!(tester.next_deadline(), None, "{timed}");
}
}
#[test]
fn test_notified_wait_does_not_time_out_on_late_relock() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let waiting = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(1);
let guard = pair.0.lock().unwrap();
pair.1.notify_one();
wait_unblocked(&clock);
assert_eq!(tester.next_deadline(), None);
tester.advance_to(deadline);
drop(guard);
assert!(!waiting.join().unwrap().timed_out());
}
#[test]
fn test_consumed_notification_leaves_sibling_to_time_out() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let later = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(1);
let resume = hold_park(&clock, &pair.1);
let first = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(2);
pair.1.notify_one();
assert!(!first.join().unwrap().timed_out());
assert_eq!(blocked(&clock), 1);
assert_eq!(tester.next_deadline(), Some(deadline));
tester.advance_to(deadline);
resume.send(()).unwrap();
assert!(later.join().unwrap().timed_out());
let (_guard, result) = pair
.1
.wait_deadline(pair.0.lock().unwrap(), deadline)
.unwrap();
assert!(result.timed_out());
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_consumed_notification_keeps_unreached_sibling_listed() {
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 report = report.clone();
start_wait(&pair, move |condvar, guard| {
let (guard, result) = condvar.wait_deadline(guard, later).unwrap();
report.send(Some(result.timed_out())).unwrap();
if !result.timed_out() {
let (_guard, result) = condvar.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!(blocked(&clock), 1);
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_eq!(pair.1.waiter.signal.lock().waiting, 0);
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_broadcast_followed_by_single_notification_strands_nothing() {
for pending in [false, true] {
let tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let first = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(1);
let resume = hold_park(&clock, &pair.1);
if pending {
pair.1.notify_one();
}
pair.1.notify_all();
pair.1.notify_one();
{
let state = pair.1.waiter.signal.lock();
assert_eq!(state.waiting, 0, "{pending}");
assert!(state.pending.is_empty(), "{pending}");
}
let next = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(2);
resume.send(()).unwrap();
assert!(!first.join().unwrap().timed_out(), "{pending}");
assert_eq!(pair.1.waiter.signal.lock().waiting, 1, "{pending}");
pair.1.notify_one();
assert!(!next.join().unwrap().timed_out(), "{pending}");
let state = pair.1.waiter.signal.lock();
assert_eq!(state.waiting, 0, "{pending}");
assert!(state.pending.is_empty(), "{pending}");
assert_eq!(tester.next_deadline(), None, "{pending}");
}
}
#[test]
fn test_single_notification_releases_exactly_one_reached_wait() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let mut waits = Vec::new();
let mut resumes = Vec::new();
for count in 1..=2 {
waits.push(start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
}));
tester.wait_blocked(count);
resumes.push(hold_park(&clock, &pair.1));
}
pair.1.notify_one();
tester.advance_to(deadline);
for resume in resumes {
resume.send(()).unwrap();
}
let timeouts = waits
.into_iter()
.map(|wait| usize::from(wait.join().unwrap().timed_out()))
.sum::<usize>();
assert_eq!(timeouts, 1);
let state = pair.1.waiter.signal.lock();
assert_eq!(state.waiting, 0);
assert!(state.pending.is_empty());
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_later_wait_cannot_claim_earlier_notification() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let first = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(1);
let resume = hold_park(&clock, &pair.1);
pair.1.notify_one();
pair.1.notify_one();
assert_eq!(pair.1.waiter.signal.lock().pending.len(), 1);
let next = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.advance_to(deadline);
assert!(next.join().unwrap().timed_out());
resume.send(()).unwrap();
assert!(!first.join().unwrap().timed_out());
let state = pair.1.waiter.signal.lock();
assert_eq!(state.waiting, 0);
assert!(state.pending.is_empty());
}
#[test]
fn test_wait_claims_oldest_eligible_notification() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let first = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(1);
let first_resume = hold_park(&clock, &pair.1);
pair.1.notify_one();
let second = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(2);
let second_resume = hold_park(&clock, &pair.1);
pair.1.notify_one();
first_resume.send(()).unwrap();
assert!(!first.join().unwrap().timed_out());
tester.advance_to(deadline);
second_resume.send(()).unwrap();
assert!(!second.join().unwrap().timed_out());
let state = pair.1.waiter.signal.lock();
assert_eq!(state.waiting, 0);
assert!(state.pending.is_empty());
}
#[test]
fn test_next_deadline_clamps_reached_park_to_current_time() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let waiting = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
tester.wait_blocked(1);
let resume = hold_park(&clock, &pair.1);
tester.advance(Duration::from_secs(2));
let next = tester.next_deadline().unwrap();
assert_eq!(next, clock.now());
tester.advance_to(next);
assert_eq!(blocked(&clock), 1);
resume.send(()).unwrap();
assert!(waiting.join().unwrap().timed_out());
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_later_wait_on_a_reached_condvar_stays_listed() {
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 last = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, later).unwrap().1
});
tester.wait_blocked(1);
let resume = hold_park(&clock, &pair.1);
let first = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, earlier).unwrap().1
});
tester.wait_blocked(2);
tester.advance_to(earlier);
assert!(first.join().unwrap().timed_out());
assert_eq!(tester.next_deadline(), Some(later));
assert_eq!(blocked(&clock), 1);
let paused = clock.paused.as_ref().unwrap();
let (checked, checks) = mpsc::channel();
for hook in [&paused.before_park, &paused.before_rewait] {
let (clock, checked) = (clock.clone(), checked.clone());
*hook.lock().unwrap() = Some(Box::new(move |_| checked.send(blocked(&clock)).unwrap()));
}
resume.send(()).unwrap();
assert_eq!(checks.recv().unwrap(), 1);
tester.advance_to(later);
assert!(last.join().unwrap().timed_out());
assert_eq!(tester.next_deadline(), None);
assert_eq!(blocked(&clock), 0);
}
#[test]
fn test_advance_ends_every_reached_wait_on_a_condvar() {
for same_deadline in [false, true] {
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<_> = [if same_deadline { later } else { earlier }, later]
.into_iter()
.map(|deadline| {
start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
})
})
.collect();
tester.wait_blocked(2);
tester.advance_to(later);
for wait in waits {
assert!(wait.join().unwrap().timed_out(), "{same_deadline}");
}
assert_eq!(blocked(&clock), 0, "{same_deadline}");
assert_eq!(tester.next_deadline(), None, "{same_deadline}");
}
}
#[test]
fn test_advance_deduplicates_nonadjacent_signal_deadlines() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
let first = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let second = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let waits: Vec<_> = [(&first, 1), (&second, 2), (&first, 3)]
.into_iter()
.map(|(pair, seconds)| {
start_wait(pair, move |condvar, guard| {
condvar
.wait_deadline(guard, start + Duration::from_secs(seconds))
.unwrap()
.1
})
})
.collect();
tester.wait_blocked(3);
tester.advance(Duration::from_secs(3));
assert_eq!(wakes(&first.1.waiter.signal), 1);
assert_eq!(wakes(&second.1.waiter.signal), 1);
for wait in waits {
assert!(wait.join().unwrap().timed_out());
}
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_next_deadline_tracks_only_parked_timed_waits() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
assert_eq!(tester.next_deadline(), None);
let untimed = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let untimed_thread = start_wait(&untimed, |condvar, guard| {
drop(condvar.wait(guard).unwrap())
});
tester.wait_blocked(1);
assert_eq!(tester.next_deadline(), None);
let sleep_thread = thread::spawn({
let clock = clock.clone();
move || clock.sleep(Duration::from_secs(3))
});
let until_thread = thread::spawn({
let clock = clock.clone();
move || clock.sleep_until(start + Duration::from_secs(7))
});
let waits: Vec<_> = [1, 3]
.into_iter()
.map(|seconds| {
let pair = Arc::new((Mutex::new(()), Condvar::new(&clock)));
let deadline = start + Duration::from_secs(seconds);
let waiting = start_wait(&pair, move |condvar, guard| {
condvar.wait_deadline(guard, deadline).unwrap().1
});
(pair, waiting)
})
.collect();
tester.wait_blocked(5);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(1)));
let mut waits = waits.into_iter();
let (pair, waiting) = waits.next().unwrap();
pair.1.notify_one();
assert!(!waiting.join().unwrap().timed_out());
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(3)));
let (pair, waiting) = waits.next().unwrap();
pair.1.notify_one();
assert!(!waiting.join().unwrap().timed_out());
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(3)));
untimed.1.notify_one();
untimed_thread.join().unwrap();
tester.advance(Duration::from_secs(3));
sleep_thread.join().unwrap();
assert!(!until_thread.is_finished());
assert_eq!(blocked(&clock), 1);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(7)));
tester.advance_to(start + Duration::from_secs(7));
until_thread.join().unwrap();
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_real_wait_deadline_expires() {
let clock = Clock::real();
let condvar = Condvar::new(&clock);
let mutex = Mutex::new(());
let deadline = clock.now() + Duration::from_millis(1);
let (guard, result) = condvar
.wait_deadline(mutex.lock().unwrap(), deadline)
.unwrap();
assert!(result.timed_out());
assert!(clock.now() >= deadline);
drop(guard);
}