use super::helpers::{blocked, last_system_time, pause_before_park, pause_before_rewait};
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_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_sleep_reparks_after_partial_advance() {
for until in [false, true] {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let waiting = thread::spawn({
let clock = clock.clone();
move || {
if until {
clock.sleep_until(deadline);
} else {
clock.sleep(Duration::from_secs(5));
}
clock.now()
}
});
tester.wait_blocked(1);
let (checked, resume) = pause_before_park(&clock);
tester.advance(Duration::from_secs(2));
assert_eq!(checked.recv().unwrap(), None, "{until}");
resume.send(()).unwrap();
tester.wait_blocked(1);
tester.advance_to(deadline);
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);
let seen = waiter.signal.generation();
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!(waiter.signal.generation(), seen, "{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(Duration::from_secs(1));
None::<()>
});
assert_eq!(result, None);
assert_eq!(clock.now(), deadline);
}
#[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());
}