use super::*;
use crate::tests::helpers::{blocked, pause_before_park};
use crossbeam_channel::{bounded, select};
use std::thread;
use std::time::UNIX_EPOCH;
#[test]
fn test_at_fires_at_exact_deadline() {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let timer = clock.at(deadline);
tester.wait_timers(1);
assert_eq!(tester.next_deadline(), Some(deadline));
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty));
tester.advance(Duration::from_secs(4));
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty));
assert_eq!(tester.paused.lock().timers.len(), 1);
tester.advance_to(deadline);
assert_eq!(timer.try_recv(), Ok(deadline));
assert!(tester.paused.lock().timers.is_empty());
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_advance_publishes_time_and_fires_all_due_timers_in_order() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
tester.set_system_time(UNIX_EPOCH);
let late = clock.at(start + Duration::from_secs(3));
let first = clock.at(start + Duration::from_secs(1));
let equal = clock.at(start + Duration::from_secs(3));
drop(clock.at(start + Duration::from_secs(2)));
let future = clock.at(start + Duration::from_secs(7));
tester.wait_timers(5);
let target = start + Duration::from_secs(5);
*tester.paused.after_timer_send.lock().unwrap() = Some(Box::new({
let clock = clock.clone();
let (first, late, equal) = (first.clone(), late.clone(), equal.clone());
move || {
assert_eq!(clock.now(), target);
assert_eq!(clock.system_time(), UNIX_EPOCH + Duration::from_secs(5));
assert_eq!(first.len(), 1);
assert_eq!(late.try_recv(), Err(TryRecvError::Empty));
assert_eq!(equal.try_recv(), Err(TryRecvError::Empty));
}
}));
tester.advance_to(target);
assert_eq!(first.try_recv(), Ok(start + Duration::from_secs(1)));
assert_eq!(late.try_recv(), Ok(start + Duration::from_secs(3)));
assert_eq!(equal.try_recv(), Ok(start + Duration::from_secs(3)));
assert_eq!(future.try_recv(), Err(TryRecvError::Empty));
assert_eq!(tester.paused.lock().fired.len(), 3);
assert_eq!(tester.paused.lock().timers.len(), 1);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(7)));
}
#[test]
fn test_reached_timers_deliver_immediately() {
let mut tester = TestClock::new();
let clock = tester.clock();
let past = clock.now();
tester.advance(Duration::from_secs(2));
let now = clock.now();
for (name, timer, deadline) in [
("past", clock.at(past), past),
("current", clock.at(now), now),
("zero", clock.after(Duration::ZERO), now),
] {
assert_eq!(timer.capacity(), Some(1), "{name}");
assert_eq!(timer.try_recv(), Ok(deadline), "{name}");
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty), "{name}");
}
assert!(tester.paused.lock().timers.is_empty());
assert_eq!(tester.paused.lock().fired.len(), 3);
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_after_uses_time_at_call() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
let first = clock.after(Duration::from_secs(5));
tester.advance(Duration::from_secs(2));
let second = clock.after(Duration::from_secs(5));
tester.advance_to(start + Duration::from_secs(5));
assert_eq!(first.try_recv(), Ok(start + Duration::from_secs(5)));
assert_eq!(second.try_recv(), Err(TryRecvError::Empty));
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(7)));
tester.advance(Duration::from_secs(2));
assert_eq!(second.try_recv(), Ok(start + Duration::from_secs(7)));
assert_eq!(first.try_recv(), Err(TryRecvError::Empty));
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_after_overflow_never_fires_or_counts() {
let mut tester = TestClock::new();
let clock = tester.clock();
assert!(clock.now().checked_add(Duration::MAX).is_none());
let timer = clock.after(Duration::MAX);
assert_eq!(timer.capacity(), Some(0));
assert!(tester.paused.lock().timers.is_empty());
assert_eq!(tester.next_deadline(), None);
tester.advance(Duration::from_secs(60));
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty));
assert!(tester.paused.lock().fired.is_empty());
drop(tester);
drop(clock);
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty));
}
#[test]
fn test_timer_connection_lasts_as_long_as_shared_clock() {
for owner_last in [false, true] {
let mut tester = TestClock::new();
let clock = tester.clock();
let clone = clock.clone();
let deadline = clock.now() + Duration::from_secs(1);
let timer = clock.at(deadline);
tester.advance_to(deadline);
assert_eq!(timer.capacity(), Some(1), "{owner_last}");
assert_eq!(timer.try_recv(), Ok(deadline), "{owner_last}");
drop(clock);
if owner_last {
drop(clone);
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty), "{owner_last}");
drop(tester);
} else {
drop(tester);
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty), "{owner_last}");
drop(clone);
}
assert_eq!(
timer.try_recv(),
Err(TryRecvError::Disconnected),
"{owner_last}"
);
}
}
#[test]
fn test_wall_jumps_and_zero_advances_leave_timers_armed() {
let mut tester = TestClock::new();
let clock = tester.clock();
let now = clock.now();
let deadline = now + Duration::from_secs(1);
let timer = clock.at(deadline);
for seconds in [100, 50] {
tester.set_system_time(UNIX_EPOCH + Duration::from_secs(seconds));
tester.advance(Duration::ZERO);
tester.advance_to(now);
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty), "{seconds}");
assert_eq!(tester.paused.lock().timers.len(), 1, "{seconds}");
assert_eq!(tester.next_deadline(), Some(deadline), "{seconds}");
}
tester.advance_to(deadline);
assert_eq!(timer.try_recv(), Ok(deadline));
}
#[test]
fn test_invalid_advances_leave_timers_armed() {
let mut tester = TestClock::new();
let clock = tester.clock();
let now = clock.now();
let deadline = now + Duration::from_secs(1);
let timer = clock.at(deadline);
let last = crate::tests::helpers::last_system_time();
tester.set_system_time(last);
for case in ["backward", "monotonic overflow", "wall overflow"] {
assert!(
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
match case {
"backward" => tester.advance_to(now - Duration::from_secs(1)),
"monotonic overflow" => tester.advance(Duration::MAX),
_ => tester.advance_to(deadline),
}
}))
.is_err(),
"{case}"
);
assert_eq!(clock.now(), now, "{case}");
assert_eq!(clock.system_time(), last, "{case}");
assert_eq!(timer.try_recv(), Err(TryRecvError::Empty), "{case}");
assert_eq!(tester.paused.lock().timers.len(), 1, "{case}");
}
tester.set_system_time(UNIX_EPOCH);
tester.advance_to(deadline);
assert_eq!(timer.try_recv(), Ok(deadline));
}
#[test]
fn test_wait_timers_counts_only_armed_timers() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
tester.wait_timers(0);
let public = clock.at(start + Duration::from_secs(9));
let fired = clock.at(start + Duration::from_secs(1));
tester.wait_timers(1);
tester.wait_timers(2);
tester.advance(Duration::from_secs(1));
assert_eq!(fired.try_recv(), Ok(start + Duration::from_secs(1)));
assert_eq!(tester.paused.lock().timers.len(), 1);
let (sender, receiver) = bounded(0);
let driver = thread::spawn(move || {
tester.wait_timers(2);
assert_eq!(tester.paused.lock().timers.len(), 2);
sender.send(7).unwrap();
tester
});
assert_eq!(clock.recv_timeout(&receiver, Duration::from_secs(5)), Ok(7));
let mut tester = driver.join().unwrap();
tester.wait_timers(1);
assert_eq!(tester.paused.lock().timers.len(), 1);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(9)));
tester.advance_to(start + Duration::from_secs(9));
assert_eq!(public.try_recv(), Ok(start + Duration::from_secs(9)));
assert!(tester.paused.lock().timers.is_empty());
}
#[test]
fn test_next_deadline_combines_timers_and_parked_waits() {
let mut tester = TestClock::new();
let clock = tester.clock();
let start = clock.now();
let first = clock.at(start + Duration::from_secs(1));
let last = clock.at(start + Duration::from_secs(3));
let waiting = thread::spawn({
let clock = clock.clone();
move || clock.sleep_until(start + Duration::from_secs(2))
});
tester.wait_blocked(1);
tester.wait_timers(2);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(1)));
let (checked, resume) = pause_before_park(&clock);
tester.advance(Duration::from_secs(1));
assert_eq!(first.try_recv(), Ok(start + Duration::from_secs(1)));
assert_eq!(checked.recv().unwrap(), None);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(3)));
resume.send(()).unwrap();
tester.wait_blocked(1);
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(2)));
tester.advance_to(start + Duration::from_secs(2));
waiting.join().unwrap();
assert_eq!(tester.next_deadline(), Some(start + Duration::from_secs(3)));
tester.advance_to(start + Duration::from_secs(3));
assert_eq!(last.try_recv(), Ok(start + Duration::from_secs(3)));
assert_eq!(tester.next_deadline(), None);
}
#[test]
fn test_select_wakes_on_clock_timer() {
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 || {
select! {
recv(clock.at(deadline)) -> result => (result.unwrap(), clock.now()),
recv(crossbeam_channel::never::<()>()) -> _ => unreachable!(),
}
}
});
tester.wait_timers(1);
assert_eq!(blocked(&clock), 0);
tester.advance_to(deadline);
assert_eq!(waiting.join().unwrap(), (deadline, deadline));
assert!(tester.paused.lock().timers.is_empty());
}
#[test]
fn test_receive_checks_channel_before_reached_deadline() {
for relative in [false, true] {
let tester = TestClock::new();
let clock = tester.clock();
let (sender, receiver) = bounded(1);
sender.send(7).unwrap();
drop(sender);
let receive = || {
if relative {
clock.recv_timeout(&receiver, Duration::ZERO)
} else {
clock.recv_deadline(&receiver, clock.now())
}
};
assert_eq!(receive(), Ok(7), "{relative}");
assert_eq!(receive(), Err(RecvTimeoutError::Disconnected), "{relative}");
assert!(tester.paused.lock().timers.is_empty(), "{relative}");
assert!(tester.paused.lock().fired.is_empty(), "{relative}");
}
}
#[test]
fn test_receive_reached_deadline_times_out_without_timer() {
let mut tester = TestClock::new();
let clock = tester.clock();
let past = clock.now();
let (_sender, receiver) = bounded::<()>(1);
tester.advance(Duration::from_secs(1));
for deadline in [past, clock.now()] {
assert_eq!(
clock.recv_deadline(&receiver, deadline),
Err(RecvTimeoutError::Timeout),
"{deadline:?}"
);
}
assert_eq!(
clock.recv_timeout(&receiver, Duration::ZERO),
Err(RecvTimeoutError::Timeout)
);
assert!(tester.paused.lock().timers.is_empty());
assert!(tester.paused.lock().fired.is_empty());
}
#[test]
fn test_receive_message_or_disconnect_disarms_timer() {
for relative in [false, true] {
for disconnect in [false, true] {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let public = clock.at(deadline);
let (sender, receiver) = bounded(0);
let waiting = thread::spawn(move || {
if relative {
clock.recv_timeout(&receiver, Duration::from_secs(5))
} else {
clock.recv_deadline(&receiver, deadline)
}
});
tester.wait_timers(2);
if disconnect {
drop(sender);
} else {
sender.send(7).unwrap();
}
let expected = if disconnect {
Err(RecvTimeoutError::Disconnected)
} else {
Ok(7)
};
assert_eq!(waiting.join().unwrap(), expected, "{relative} {disconnect}");
assert_eq!(
tester.paused.lock().timers.len(),
1,
"{relative} {disconnect}"
);
assert!(
tester.paused.lock().fired.is_empty(),
"{relative} {disconnect}"
);
tester.advance_to(deadline);
assert_eq!(public.try_recv(), Ok(deadline), "{relative} {disconnect}");
assert!(
tester.paused.lock().timers.is_empty(),
"{relative} {disconnect}"
);
assert_eq!(
tester.paused.lock().fired.len(),
1,
"{relative} {disconnect}"
);
assert_eq!(tester.next_deadline(), None, "{relative} {disconnect}");
}
}
}
#[test]
fn test_receive_times_out_at_exact_deadline() {
for relative in [false, true] {
let mut tester = TestClock::new();
tester.advance(Duration::from_secs(2));
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let (_sender, receiver) = bounded::<()>(1);
let waiting = thread::spawn(move || {
let result = if relative {
clock.recv_timeout(&receiver, Duration::from_secs(5))
} else {
clock.recv_deadline(&receiver, deadline)
};
(result, clock.now())
});
tester.wait_timers(1);
assert_eq!(tester.next_deadline(), Some(deadline), "{relative}");
tester.advance(Duration::from_secs(4));
assert_eq!(tester.paused.lock().timers.len(), 1, "{relative}");
assert!(!waiting.is_finished(), "{relative}");
tester.advance_to(deadline);
assert_eq!(
waiting.join().unwrap(),
(Err(RecvTimeoutError::Timeout), deadline),
"{relative}"
);
assert!(tester.paused.lock().timers.is_empty(), "{relative}");
assert!(tester.paused.lock().fired.is_empty(), "{relative}");
assert_eq!(tester.next_deadline(), None, "{relative}");
}
}
#[test]
fn test_receive_rechecks_channel_after_timer_wins() {
for relative in [false, true] {
for disconnect in [false, true] {
let mut tester = TestClock::new();
let clock = tester.clock();
let deadline = clock.now() + Duration::from_secs(5);
let (sender, receiver) = bounded(1);
*tester.paused.after_timer_receive.lock().unwrap() = Some(Box::new({
let clock = clock.clone();
move || {
assert_eq!(clock.now(), deadline);
assert!(clock.paused.as_ref().unwrap().lock().timers.is_empty());
if !disconnect {
sender.send(7).unwrap();
}
drop(sender);
}
}));
let waiting = thread::spawn(move || {
if relative {
clock.recv_timeout(&receiver, Duration::from_secs(5))
} else {
clock.recv_deadline(&receiver, deadline)
}
});
tester.wait_timers(1);
tester.advance_to(deadline);
let expected = if disconnect {
Err(RecvTimeoutError::Disconnected)
} else {
Ok(7)
};
assert_eq!(waiting.join().unwrap(), expected, "{relative} {disconnect}");
assert!(
tester.paused.lock().timers.is_empty(),
"{relative} {disconnect}"
);
assert!(
tester.paused.lock().fired.is_empty(),
"{relative} {disconnect}"
);
assert_eq!(tester.next_deadline(), None, "{relative} {disconnect}");
}
}
}
#[test]
fn test_receive_timeout_overflow_waits_without_timer() {
for disconnect in [false, true] {
let mut tester = TestClock::new();
let clock = tester.clock();
assert!(clock.now().checked_add(Duration::MAX).is_none());
let (sender, receiver) = bounded(0);
let waiting = thread::spawn(move || clock.recv_timeout(&receiver, Duration::MAX));
tester.advance(Duration::from_secs(60));
if disconnect {
drop(sender);
} else {
sender.send(7).unwrap();
}
let expected = if disconnect {
Err(RecvTimeoutError::Disconnected)
} else {
Ok(7)
};
assert_eq!(waiting.join().unwrap(), expected, "{disconnect}");
assert!(tester.paused.lock().timers.is_empty(), "{disconnect}");
assert!(tester.paused.lock().fired.is_empty(), "{disconnect}");
}
}