use super::helpers::{
blocked, last_system_time, pause_before_park, pause_before_rewait, signals, wakes,
};
use crate::{Clock, TestClock, Waiter};
use std::fmt;
use std::panic::{self, AssertUnwindSafe};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
fn parked(waiter: &Waiter, count: usize) {
let mut state = waiter.signal.lock();
while state.parks < count {
state = waiter.signal.parked.wait(state).unwrap();
}
}
#[test]
fn test_clock_identity_and_shared_time() {
const REAL: Clock = Clock::real();
let before = Instant::now();
let mut other = TestClock::new();
let mut tester = TestClock::new();
let real_now = REAL.now();
let after = Instant::now();
let clock = tester.clock();
let clone = clock.clone();
let other_start = other.clock().now();
let start = clock.now();
assert!((before..=after).contains(&start));
assert!((before..=after).contains(&real_now));
tester.advance(Duration::from_secs(1));
assert_eq!(clone.now(), start + Duration::from_secs(1));
assert_eq!(other.clock().now(), other_start);
tester.advance_to(start + Duration::from_secs(3));
assert_eq!(clock.now(), start + Duration::from_secs(3));
other.advance_to(clock.now());
other.set_system_time(clock.system_time());
assert_eq!(other.clock().now(), clock.now());
assert_eq!(other.clock().system_time(), clock.system_time());
assert_eq!(clock, clone);
assert_eq!(clock, tester.clock());
assert_ne!(clock, other.clock());
assert_ne!(clock, REAL);
assert_ne!(REAL, clock);
assert_eq!(REAL, Clock::real());
}
#[test]
fn test_debug_writer_can_read_clock() {
struct ClockWriter {
clock: Clock,
output: String,
}
impl fmt::Write for ClockWriter {
fn write_str(&mut self, text: &str) -> fmt::Result {
let _ = self.clock.now();
self.output.push_str(text);
Ok(())
}
}
let tester = TestClock::new();
#[cfg(feature = "crossbeam")]
let _timer = tester.clock().after(Duration::from_secs(5));
let mut writer = ClockWriter {
clock: tester.clock(),
output: String::new(),
};
fmt::write(&mut writer, format_args!("{tester:?}")).unwrap();
for field in ["advanced:", "system_time:", "blocked:"] {
assert!(writer.output.contains(field), "{field}");
}
#[cfg(feature = "crossbeam")]
assert!(writer.output.contains("timers: 1"));
}
#[test]
fn test_clock_debug_uses_one_snapshot() {
struct AdvancingWriter {
tester: TestClock,
output: String,
}
impl fmt::Write for AdvancingWriter {
fn write_str(&mut self, text: &str) -> fmt::Result {
self.tester.advance(Duration::from_secs(1));
self.output.push_str(text);
Ok(())
}
}
let mut tester = TestClock::new();
tester.set_system_time(UNIX_EPOCH);
let clock = tester.clock();
let mut writer = AdvancingWriter {
tester,
output: String::new(),
};
fmt::write(&mut writer, format_args!("{clock:?}")).unwrap();
assert_eq!(
writer.output,
format!("Clock {{ paused: true, advanced: 0ns, system_time: {UNIX_EPOCH:?} }}")
);
assert!(clock.system_time() > UNIX_EPOCH);
}
#[test]
fn test_elapsed_saturates() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
tester.advance(Duration::from_secs(3));
assert_eq!(clock.elapsed(start), Duration::from_secs(3));
assert_eq!(clock.elapsed(clock.now()), Duration::ZERO);
assert_eq!(
clock.elapsed(clock.now() + Duration::from_secs(1)),
Duration::ZERO
);
}
#[test]
fn test_system_time_follows_advances_and_jumps_alone() {
let before = SystemTime::now();
let mut tester = TestClock::default();
let real_time = Clock::real().system_time();
let after = SystemTime::now();
let clock = tester.clock();
let start = clock.now();
assert!((before..=after).contains(&clock.system_time()));
assert!((before..=after).contains(&real_time));
tester.set_system_time(UNIX_EPOCH + Duration::from_secs(100));
tester.advance(Duration::from_secs(7));
assert_eq!(clock.now(), start + Duration::from_secs(7));
assert_eq!(clock.system_time(), UNIX_EPOCH + Duration::from_secs(107));
tester.advance_to(start + Duration::from_secs(10));
assert_eq!(clock.system_time(), UNIX_EPOCH + Duration::from_secs(110));
for seconds in [200, 50] {
let time = UNIX_EPOCH + Duration::from_secs(seconds);
tester.set_system_time(time);
assert_eq!(clock.system_time(), time, "{seconds}");
assert_eq!(clock.now(), start + Duration::from_secs(10), "{seconds}");
}
tester.advance(Duration::from_secs(1));
assert_eq!(clock.system_time(), UNIX_EPOCH + Duration::from_secs(51));
assert_eq!(clock.now(), start + Duration::from_secs(11));
}
#[test]
fn test_failed_monotonic_advance_keeps_both_times() {
let mut tester = TestClock::new();
let clock = tester.clock();
tester.advance(Duration::from_secs(1));
for backwards in [true, false] {
let now = clock.now();
let wall = clock.system_time();
assert!(
panic::catch_unwind(AssertUnwindSafe(|| {
if backwards {
tester.advance_to(now - Duration::from_secs(1));
} else {
tester.advance(Duration::MAX);
}
}))
.is_err(),
"{backwards}"
);
assert_eq!(clock.now(), now, "{backwards}");
assert_eq!(clock.system_time(), wall, "{backwards}");
tester.advance(Duration::from_secs(1));
assert_eq!(clock.now(), now + Duration::from_secs(1), "{backwards}");
assert_eq!(
clock.system_time(),
wall + Duration::from_secs(1),
"{backwards}"
);
}
}
#[test]
fn test_failed_wall_advance_keeps_both_times() {
let mut tester = TestClock::new();
let clock = tester.clock();
let last = last_system_time();
assert!(last.checked_add(Duration::from_secs(1)).is_none());
for absolute in [false, true] {
let now = clock.now();
tester.set_system_time(last);
assert!(
panic::catch_unwind(AssertUnwindSafe(|| {
if absolute {
tester.advance_to(now + Duration::from_secs(1));
} else {
tester.advance(Duration::from_secs(1));
}
}))
.is_err(),
"{absolute}"
);
assert_eq!(clock.now(), now, "{absolute}");
assert_eq!(clock.system_time(), last, "{absolute}");
tester.set_system_time(UNIX_EPOCH);
tester.advance(Duration::from_secs(1));
assert_eq!(clock.now(), now + Duration::from_secs(1), "{absolute}");
assert_eq!(
clock.system_time(),
UNIX_EPOCH + Duration::from_secs(1),
"{absolute}"
);
}
}
#[test]
fn test_sleep_observes_advance_before_parking() {
for until in [false, true] {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(60);
let (checked, resume) = pause_before_park(&clock);
let waiting = thread::spawn(move || {
if until {
clock.sleep_until(deadline);
} else {
clock.sleep(Duration::from_secs(60));
}
clock.now()
});
assert_eq!(checked.recv().unwrap(), None, "{until}");
tester.advance_to(deadline);
resume.send(()).unwrap();
assert_eq!(waiting.join().unwrap(), deadline, "{until}");
}
}
#[test]
fn test_advance_wakes_all_parked_sleepers() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(60);
let threads: Vec<_> = (0..6)
.map(|index| {
let clock = clock.clone();
thread::spawn(move || {
if index % 2 == 0 {
clock.sleep(Duration::from_secs(60));
} else {
clock.sleep_until(deadline);
}
clock.now()
})
})
.collect();
tester.wait_blocked(6);
assert_eq!(blocked(&clock), 6);
tester.advance(Duration::from_secs(60));
for waiting in threads {
assert_eq!(waiting.join().unwrap(), deadline);
}
let state = clock.paused.as_ref().unwrap().state.lock().unwrap();
assert_eq!(state.blocked, 0);
assert!(state.signals.is_empty());
}
#[test]
fn test_sleep_returns_for_reached_deadlines() {
let mut tester = TestClock::new();
let start = tester.clock().now();
tester.advance(Duration::from_secs(1));
for clock in [Clock::real(), tester.clock()] {
clock.sleep_until(start);
clock.sleep_until(clock.now());
clock.sleep(Duration::ZERO);
}
}
#[test]
fn test_wall_jumps_and_zero_advances_do_not_wake() {
let mut tester = TestClock::new();
let clock = tester.clock();
let now = clock.now();
let waiter = Arc::new(clock.waiter());
let waiting = thread::spawn({
let waiter = waiter.clone();
move || waiter.wait_until(Some(now + Duration::from_secs(3)), || None::<()>)
});
tester.wait_blocked(1);
for seconds in [100, 50] {
let wall = UNIX_EPOCH + Duration::from_secs(seconds);
tester.set_system_time(wall);
tester.advance(Duration::ZERO);
tester.advance_to(now);
assert_eq!(clock.now(), now, "{seconds}");
assert_eq!(clock.system_time(), wall, "{seconds}");
assert_eq!(wakes(&waiter.signal), 0, "{seconds}");
}
tester.advance(Duration::from_secs(3));
assert_eq!(waiting.join().unwrap(), None);
}
#[test]
fn test_spurious_wakeup_keeps_park_counted() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(1);
let waiter = Arc::new(clock.waiter());
let waiting = thread::spawn({
let waiter = waiter.clone();
move || waiter.wait_until(Some(deadline), || None::<()>)
});
tester.wait_blocked(1);
let (checked, resume) = pause_before_rewait(&clock);
{
let _state = waiter.signal.lock();
waiter.signal.changed.notify_all();
}
assert_eq!(checked.recv().unwrap(), None);
assert_eq!(blocked(&clock), 1);
assert_eq!(tester.next_deadline(), Some(deadline));
resume.send(()).unwrap();
tester.advance_to(deadline);
assert_eq!(waiting.join().unwrap(), None);
}
#[test]
fn test_poisoned_clock_remains_usable() {
let mut tester = TestClock::new();
let clock = tester.clock();
let now = clock.now();
let wall = clock.system_time();
let paused = clock.paused.as_ref().unwrap();
assert!(
panic::catch_unwind(AssertUnwindSafe(|| {
let _state = paused.state.lock().unwrap();
panic!("poison the shared clock lock");
}))
.is_err()
);
assert!(paused.state.is_poisoned());
assert_eq!(clock.now(), now);
assert_eq!(clock.system_time(), wall);
tester.set_system_time(UNIX_EPOCH);
let waiting = thread::spawn({
let clock = clock.clone();
move || clock.sleep(Duration::from_secs(1))
});
tester.wait_blocked(1);
tester.advance_to(now + Duration::from_secs(1));
waiting.join().unwrap();
assert_eq!(clock.now(), now + Duration::from_secs(1));
assert_eq!(clock.system_time(), UNIX_EPOCH + Duration::from_secs(1));
}
#[test]
fn test_wait_checks_ready_first() {
for clock in [Clock::real(), TestClock::new().clock()] {
let waiter = clock.waiter();
let now = clock.now();
assert_eq!(
waiter.wait_until(Some(now), || Some(1)),
Some(1),
"{clock:?}"
);
assert_eq!(
waiter.wait_until(Some(now), || None::<u8>),
None,
"{clock:?}"
);
}
}
#[test]
fn test_wait_observes_notification_during_check() {
for clock in [Clock::real(), TestClock::new().clock()] {
let waiter = clock.waiter();
let mut notified = false;
let result = waiter.wait_until(None, || {
if notified {
Some(42)
} else {
waiter.signal.notify_all();
notified = true;
None
}
});
assert_eq!(result, Some(42), "{clock:?}");
}
}
#[test]
fn test_wait_observes_advance_during_check() {
let mut tester = TestClock::new();
let clock = tester.clock();
let waiter = clock.waiter();
let deadline = clock.now() + Duration::from_secs(2);
let result = waiter.wait_until(Some(deadline), || {
tester.advance_to(deadline);
None::<()>
});
assert_eq!(result, None);
assert_eq!(clock.now(), deadline);
assert_eq!(waiter.signal.lock().parks, 0);
}
#[test]
fn test_ready_value_survives_advance() {
let mut tester = TestClock::new();
let clock = tester.clock();
let waiter = Arc::new(clock.waiter());
let deadline = clock.now() + Duration::from_secs(1);
let ready = Arc::new(AtomicBool::new(false));
let waiting = thread::spawn({
let (waiter, ready) = (waiter.clone(), ready.clone());
move || {
waiter.wait_until(Some(deadline), || {
ready.load(Ordering::SeqCst).then_some(())
})
}
});
tester.wait_blocked(1);
ready.store(true, Ordering::SeqCst);
tester.advance(Duration::from_secs(1));
assert_eq!(waiting.join().unwrap(), Some(()));
}
#[test]
fn test_notify_wakes_all_parked_waiters() {
for clock in [Clock::real(), TestClock::new().clock()] {
let waiter = Arc::new(clock.waiter());
let open = Arc::new(AtomicBool::new(false));
let threads: Vec<_> = (0..2)
.map(|_| {
let (waiter, open) = (waiter.clone(), open.clone());
thread::spawn(move || {
waiter.wait_until(None, || open.load(Ordering::SeqCst).then_some(()))
})
})
.collect();
parked(&waiter, 2);
open.store(true, Ordering::SeqCst);
waiter.signal.notify_all();
for waiting in threads {
assert_eq!(waiting.join().unwrap(), Some(()), "{clock:?}");
}
}
}
#[test]
fn test_waiter_registry_tracks_live_waiters() {
let tester = TestClock::new();
let clock = tester.clock();
let paused = clock.paused.as_ref().unwrap();
let mut waiters: Vec<_> = (0..128).map(|_| clock.waiter()).collect();
let signal = waiters.last().unwrap().signal.clone();
waiters.truncate(3);
assert_eq!(paused.state.lock().unwrap().signals.len(), 3);
for _ in 0..256 {
drop(clock.waiter());
}
assert_eq!(paused.state.lock().unwrap().signals.len(), 3);
drop(waiters);
drop(signal);
assert!(paused.state.lock().unwrap().signals.is_empty());
}
#[test]
fn test_dropping_a_waiter_keeps_other_registrations() {
let tester = TestClock::new();
let clock = tester.clock();
let mut waiters = [clock.waiter(), clock.waiter()];
waiters.sort_by_key(|waiter| waiter.signal.key());
let [survivor, dropped] = waiters;
drop(dropped);
let state = clock.paused.as_ref().unwrap().state.lock().unwrap();
assert_eq!(
state.signals.keys().copied().collect::<Vec<_>>(),
[survivor.signal.key()]
);
}
#[test]
fn test_wall_time_accumulates_fractional_advances() {
let mut tester = TestClock::new();
let clock = tester.clock();
tester.set_system_time(UNIX_EPOCH);
tester.advance(Duration::from_nanos(150));
tester.advance(Duration::from_nanos(150));
assert_eq!(clock.system_time(), UNIX_EPOCH + Duration::from_nanos(300));
tester.set_system_time(UNIX_EPOCH + Duration::from_secs(1));
tester.advance(Duration::from_nanos(150));
tester.advance(Duration::from_nanos(150));
assert_eq!(
clock.system_time(),
UNIX_EPOCH + Duration::from_nanos(1_000_000_300)
);
}
#[test]
fn test_next_deadline_drives_sleeps_in_order() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
let mut sleeps: Vec<_> = [3, 1, 2]
.into_iter()
.map(|seconds| {
let clock = clock.clone();
let deadline = start + Duration::from_secs(seconds);
(deadline, thread::spawn(move || clock.sleep_until(deadline)))
})
.collect();
tester.wait_blocked(3);
sleeps.sort_by_key(|(deadline, _)| *deadline);
let signals = signals(&clock);
for (index, (deadline, sleep)) in sleeps.into_iter().enumerate() {
let next = tester.next_deadline().unwrap();
assert_eq!(next, deadline, "{index}");
tester.advance_to(next);
sleep.join().unwrap();
let woken: usize = signals.iter().map(|signal| wakes(signal)).sum();
assert_eq!(woken, index + 1, "{index}");
}
assert_eq!(tester.next_deadline(), None);
assert_eq!(blocked(&clock), 0);
}
#[test]
fn test_real_sleep_reaches_duration() {
let clock = Clock::real();
let start = Instant::now();
clock.sleep(Duration::from_millis(1));
assert!(clock.elapsed(start) >= Duration::from_millis(1));
}
#[test]
fn test_real_sleep_until_reaches_deadline() {
let clock = Clock::real();
let deadline = Instant::now() + Duration::from_millis(1);
clock.sleep_until(deadline);
assert!(Instant::now() >= deadline);
}