1use std::{
23 any::Any,
24 cell::RefCell,
25 collections::{BTreeMap, BinaryHeap},
26 fmt::Debug,
27 ops::Deref,
28 time::Duration,
29};
30
31use ahash::AHashMap;
32use jiff::Timestamp;
33use nautilus_core::{
34 AtomicTime, DurationNanos, UUID4, UnixNanos,
35 correctness::{check_positive_u64, check_predicate_true, check_valid_string_utf8},
36 datetime::{NANOSECONDS_IN_SECOND, try_datetime_to_unix_nanos},
37 string::formatting::Separable,
38};
39use ustr::Ustr;
40
41use crate::timer::{
42 ScheduledTimeEvent, TestTimer, TimeEvent, TimeEventCallback, TimeEventHandler, Timer,
43 create_valid_interval,
44};
45
46pub trait Clock: Debug + Any {
50 fn utc_now(&self) -> Timestamp {
52 self.timestamp_ns().to_datetime_utc()
53 }
54
55 fn timestamp_ns(&self) -> UnixNanos;
57
58 fn timestamp_us(&self) -> u64;
60
61 fn timestamp_ms(&self) -> u64;
63
64 fn timestamp(&self) -> f64;
66
67 fn timer_names(&self) -> Vec<&str>;
69
70 fn timer_count(&self) -> usize;
72
73 fn timer_exists(&self, name: &Ustr) -> bool;
75
76 fn register_default_handler(&mut self, callback: TimeEventCallback);
78
79 fn cancel_default_handler(&mut self);
86
87 fn cancel_callbacks(&mut self);
94
95 fn set_time_alert(
112 &mut self,
113 name: &str,
114 alert_time: Timestamp,
115 callback: Option<TimeEventCallback>,
116 allow_past: Option<bool>,
117 ) -> anyhow::Result<()> {
118 self.set_time_alert_ns(
119 name,
120 try_datetime_to_unix_nanos(alert_time)?,
121 callback,
122 allow_past,
123 )
124 }
125
126 fn set_time_alert_ns(
150 &mut self,
151 name: &str,
152 alert_time_ns: UnixNanos,
153 callback: Option<TimeEventCallback>,
154 allow_past: Option<bool>,
155 ) -> anyhow::Result<()>;
156
157 #[expect(clippy::too_many_arguments)]
181 fn set_timer(
182 &mut self,
183 name: &str,
184 interval: Duration,
185 start_time: Option<Timestamp>,
186 stop_time: Option<Timestamp>,
187 callback: Option<TimeEventCallback>,
188 allow_past: Option<bool>,
189 fire_immediately: Option<bool>,
190 ) -> anyhow::Result<()> {
191 self.set_timer_ns(
192 name,
193 duration_to_nanos(interval)?,
194 start_time.map(try_datetime_to_unix_nanos).transpose()?,
195 stop_time.map(try_datetime_to_unix_nanos).transpose()?,
196 callback,
197 allow_past,
198 fire_immediately,
199 )
200 }
201
202 #[expect(clippy::too_many_arguments)]
238 fn set_timer_ns(
239 &mut self,
240 name: &str,
241 interval_ns: DurationNanos,
242 start_time_ns: Option<UnixNanos>,
243 stop_time_ns: Option<UnixNanos>,
244 callback: Option<TimeEventCallback>,
245 allow_past: Option<bool>,
246 fire_immediately: Option<bool>,
247 ) -> anyhow::Result<()>;
248
249 fn next_time_ns(&self, name: &str) -> Option<UnixNanos>;
253
254 fn cancel_timer(&mut self, name: &str);
256
257 fn cancel_timers(&mut self);
259
260 fn reset(&mut self);
264}
265
266impl dyn Clock {
267 pub fn as_any(&self) -> &dyn std::any::Any {
269 self
270 }
271
272 pub fn as_any_mut(&mut self) -> &mut dyn std::any::Any {
274 self
275 }
276}
277
278#[derive(Debug)]
283pub struct ClockApi<'a> {
284 backing: ClockApiBacking<'a>,
285}
286
287impl<'a> ClockApi<'a> {
288 pub(crate) fn new(clock: &'a RefCell<dyn Clock>) -> Self {
289 Self {
290 backing: ClockApiBacking::Native(clock),
291 }
292 }
293
294 #[doc(hidden)]
300 #[must_use]
301 #[expect(
302 clippy::too_many_arguments,
303 reason = "clock API backing mirrors the full ClockApi surface"
304 )]
305 pub fn from_handlers<
306 TimestampNs,
307 SetTimeAlertNs,
308 SetTimerNs,
309 TimerNames,
310 TimerCount,
311 TimerExists,
312 NextTimeNs,
313 CancelTimer,
314 CancelTimers,
315 >(
316 timestamp_ns: TimestampNs,
317 set_time_alert_ns: SetTimeAlertNs,
318 set_timer_ns: SetTimerNs,
319 timer_names: TimerNames,
320 timer_count: TimerCount,
321 timer_exists: TimerExists,
322 next_time_ns: NextTimeNs,
323 cancel_timer: CancelTimer,
324 cancel_timers: CancelTimers,
325 ) -> Self
326 where
327 TimestampNs: Fn() -> UnixNanos + 'a,
328 SetTimeAlertNs:
329 Fn(&str, UnixNanos, Option<TimeEventCallback>, Option<bool>) -> anyhow::Result<()> + 'a,
330 SetTimerNs: Fn(
331 &str,
332 DurationNanos,
333 Option<UnixNanos>,
334 Option<UnixNanos>,
335 Option<TimeEventCallback>,
336 Option<bool>,
337 Option<bool>,
338 ) -> anyhow::Result<()>
339 + 'a,
340 TimerNames: Fn() -> Vec<String> + 'a,
341 TimerCount: Fn() -> usize + 'a,
342 TimerExists: Fn(&str) -> bool + 'a,
343 NextTimeNs: Fn(&str) -> Option<UnixNanos> + 'a,
344 CancelTimer: Fn(&str) + 'a,
345 CancelTimers: Fn() + 'a,
346 {
347 Self {
348 backing: ClockApiBacking::Handlers(ClockApiHandlers {
349 timestamp_ns: Box::new(timestamp_ns),
350 set_time_alert_ns: Box::new(set_time_alert_ns),
351 set_timer_ns: Box::new(set_timer_ns),
352 timer_names: Box::new(timer_names),
353 timer_count: Box::new(timer_count),
354 timer_exists: Box::new(timer_exists),
355 next_time_ns: Box::new(next_time_ns),
356 cancel_timer: Box::new(cancel_timer),
357 cancel_timers: Box::new(cancel_timers),
358 }),
359 }
360 }
361
362 #[must_use]
368 pub fn timestamp_ns(&self) -> UnixNanos {
369 match &self.backing {
370 ClockApiBacking::Native(clock) => clock.borrow().timestamp_ns(),
371 ClockApiBacking::Handlers(handlers) => (handlers.timestamp_ns)(),
372 }
373 }
374
375 #[must_use]
381 pub fn timestamp_us(&self) -> u64 {
382 match &self.backing {
383 ClockApiBacking::Native(clock) => clock.borrow().timestamp_us(),
384 ClockApiBacking::Handlers(handlers) => (handlers.timestamp_ns)().as_micros(),
385 }
386 }
387
388 #[must_use]
394 pub fn timestamp_ms(&self) -> u64 {
395 match &self.backing {
396 ClockApiBacking::Native(clock) => clock.borrow().timestamp_ms(),
397 ClockApiBacking::Handlers(handlers) => (handlers.timestamp_ns)().as_millis(),
398 }
399 }
400
401 #[must_use]
407 pub fn timestamp(&self) -> f64 {
408 match &self.backing {
409 ClockApiBacking::Native(clock) => clock.borrow().timestamp(),
410 ClockApiBacking::Handlers(handlers) => {
411 (handlers.timestamp_ns)().as_f64() / (NANOSECONDS_IN_SECOND as f64)
412 }
413 }
414 }
415
416 #[must_use]
422 pub fn utc_now(&self) -> Timestamp {
423 match &self.backing {
424 ClockApiBacking::Native(clock) => clock.borrow().utc_now(),
425 ClockApiBacking::Handlers(handlers) => (handlers.timestamp_ns)().to_datetime_utc(),
426 }
427 }
428
429 pub fn set_time_alert(
443 &self,
444 name: &str,
445 alert_time: Timestamp,
446 callback: Option<TimeEventCallback>,
447 allow_past: Option<bool>,
448 ) -> anyhow::Result<()> {
449 match &self.backing {
450 ClockApiBacking::Native(clock) => clock
451 .borrow_mut()
452 .set_time_alert(name, alert_time, callback, allow_past),
453 ClockApiBacking::Handlers(handlers) => (handlers.set_time_alert_ns)(
454 name,
455 try_datetime_to_unix_nanos(alert_time)?,
456 callback,
457 allow_past,
458 ),
459 }
460 }
461
462 pub fn set_time_alert_ns(
475 &self,
476 name: &str,
477 alert_time_ns: UnixNanos,
478 callback: Option<TimeEventCallback>,
479 allow_past: Option<bool>,
480 ) -> anyhow::Result<()> {
481 match &self.backing {
482 ClockApiBacking::Native(clock) => {
483 clock
484 .borrow_mut()
485 .set_time_alert_ns(name, alert_time_ns, callback, allow_past)
486 }
487 ClockApiBacking::Handlers(handlers) => {
488 (handlers.set_time_alert_ns)(name, alert_time_ns, callback, allow_past)
489 }
490 }
491 }
492
493 #[expect(clippy::too_many_arguments, reason = "timer scheduling mirrors Clock")]
507 pub fn set_timer(
508 &self,
509 name: &str,
510 interval: Duration,
511 start_time: Option<Timestamp>,
512 stop_time: Option<Timestamp>,
513 callback: Option<TimeEventCallback>,
514 allow_past: Option<bool>,
515 fire_immediately: Option<bool>,
516 ) -> anyhow::Result<()> {
517 match &self.backing {
518 ClockApiBacking::Native(clock) => clock.borrow_mut().set_timer(
519 name,
520 interval,
521 start_time,
522 stop_time,
523 callback,
524 allow_past,
525 fire_immediately,
526 ),
527 ClockApiBacking::Handlers(handlers) => (handlers.set_timer_ns)(
528 name,
529 duration_to_nanos(interval)?,
530 start_time.map(try_datetime_to_unix_nanos).transpose()?,
531 stop_time.map(try_datetime_to_unix_nanos).transpose()?,
532 callback,
533 allow_past,
534 fire_immediately,
535 ),
536 }
537 }
538
539 #[expect(clippy::too_many_arguments, reason = "timer scheduling mirrors Clock")]
552 pub fn set_timer_ns(
553 &self,
554 name: &str,
555 interval_ns: DurationNanos,
556 start_time_ns: Option<UnixNanos>,
557 stop_time_ns: Option<UnixNanos>,
558 callback: Option<TimeEventCallback>,
559 allow_past: Option<bool>,
560 fire_immediately: Option<bool>,
561 ) -> anyhow::Result<()> {
562 match &self.backing {
563 ClockApiBacking::Native(clock) => clock.borrow_mut().set_timer_ns(
564 name,
565 interval_ns,
566 start_time_ns,
567 stop_time_ns,
568 callback,
569 allow_past,
570 fire_immediately,
571 ),
572 ClockApiBacking::Handlers(handlers) => (handlers.set_timer_ns)(
573 name,
574 interval_ns,
575 start_time_ns,
576 stop_time_ns,
577 callback,
578 allow_past,
579 fire_immediately,
580 ),
581 }
582 }
583
584 #[must_use]
590 pub fn timer_names(&self) -> Vec<String> {
591 match &self.backing {
592 ClockApiBacking::Native(clock) => clock
593 .borrow()
594 .timer_names()
595 .into_iter()
596 .map(str::to_string)
597 .collect(),
598 ClockApiBacking::Handlers(handlers) => (handlers.timer_names)(),
599 }
600 }
601
602 #[must_use]
608 pub fn timer_count(&self) -> usize {
609 match &self.backing {
610 ClockApiBacking::Native(clock) => clock.borrow().timer_count(),
611 ClockApiBacking::Handlers(handlers) => (handlers.timer_count)(),
612 }
613 }
614
615 #[must_use]
621 pub fn timer_exists(&self, name: &str) -> bool {
622 match &self.backing {
623 ClockApiBacking::Native(clock) => clock.borrow().timer_exists(&Ustr::from(name)),
624 ClockApiBacking::Handlers(handlers) => (handlers.timer_exists)(name),
625 }
626 }
627
628 #[must_use]
636 pub fn next_time_ns(&self, name: &str) -> Option<UnixNanos> {
637 match &self.backing {
638 ClockApiBacking::Native(clock) => clock.borrow().next_time_ns(name),
639 ClockApiBacking::Handlers(handlers) => (handlers.next_time_ns)(name),
640 }
641 }
642
643 pub fn cancel_timer(&self, name: &str) {
649 match &self.backing {
650 ClockApiBacking::Native(clock) => clock.borrow_mut().cancel_timer(name),
651 ClockApiBacking::Handlers(handlers) => (handlers.cancel_timer)(name),
652 }
653 }
654
655 pub fn cancel_timers(&self) {
661 match &self.backing {
662 ClockApiBacking::Native(clock) => clock.borrow_mut().cancel_timers(),
663 ClockApiBacking::Handlers(handlers) => (handlers.cancel_timers)(),
664 }
665 }
666}
667
668enum ClockApiBacking<'a> {
669 Native(&'a RefCell<dyn Clock>),
670 Handlers(ClockApiHandlers<'a>),
671}
672
673struct ClockApiHandlers<'a> {
674 timestamp_ns: Box<dyn Fn() -> UnixNanos + 'a>,
675 set_time_alert_ns: Box<SetTimeAlertNsHandler<'a>>,
676 set_timer_ns: Box<SetTimerNsHandler<'a>>,
677 timer_names: Box<dyn Fn() -> Vec<String> + 'a>,
678 timer_count: Box<dyn Fn() -> usize + 'a>,
679 timer_exists: Box<dyn Fn(&str) -> bool + 'a>,
680 next_time_ns: Box<NextTimeNsHandler<'a>>,
681 cancel_timer: Box<dyn Fn(&str) + 'a>,
682 cancel_timers: Box<dyn Fn() + 'a>,
683}
684
685impl Debug for ClockApiBacking<'_> {
686 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
687 match self {
688 Self::Native(_) => f.write_str("Native"),
689 Self::Handlers(_) => f.write_str("Handlers"),
690 }
691 }
692}
693
694type SetTimeAlertNsHandler<'a> =
695 dyn Fn(&str, UnixNanos, Option<TimeEventCallback>, Option<bool>) -> anyhow::Result<()> + 'a;
696type NextTimeNsHandler<'a> = dyn Fn(&str) -> Option<UnixNanos> + 'a;
697type SetTimerNsHandler<'a> = dyn Fn(
698 &str,
699 DurationNanos,
700 Option<UnixNanos>,
701 Option<UnixNanos>,
702 Option<TimeEventCallback>,
703 Option<bool>,
704 Option<bool>,
705 ) -> anyhow::Result<()>
706 + 'a;
707
708fn duration_to_nanos(duration: Duration) -> anyhow::Result<DurationNanos> {
709 DurationNanos::try_from(duration)
710 .map_err(|_| anyhow::anyhow!("Interval exceeds u64 nanoseconds"))
711}
712
713#[derive(Debug, Default)]
718pub struct CallbackRegistry {
719 default_callback: Option<TimeEventCallback>,
720 callbacks: AHashMap<Ustr, TimeEventCallback>,
721}
722
723impl CallbackRegistry {
724 #[must_use]
726 pub fn new() -> Self {
727 Self::default()
728 }
729
730 pub fn register_default_handler(&mut self, callback: TimeEventCallback) {
732 self.default_callback = Some(callback);
733 }
734
735 pub fn cancel_default_handler(&mut self) {
737 self.default_callback = None;
738 }
739
740 pub fn register_callback(&mut self, name: Ustr, callback: TimeEventCallback) {
742 self.callbacks.insert(name, callback);
743 }
744
745 #[must_use]
747 pub fn has_any_callback(&self, name: &Ustr) -> bool {
748 self.callbacks.contains_key(name) || self.default_callback.is_some()
749 }
750
751 #[must_use]
753 pub fn get_callback(&self, name: &Ustr) -> Option<TimeEventCallback> {
754 self.callbacks
755 .get(name)
756 .cloned()
757 .or_else(|| self.default_callback.clone())
758 }
759
760 #[must_use]
766 pub fn get_handler(&self, event: TimeEvent) -> TimeEventHandler {
767 let callback = self
768 .get_callback(&event.name)
769 .unwrap_or_else(|| panic!("Event '{}' should have associated handler", event.name));
770
771 TimeEventHandler::new(event, callback)
772 }
773
774 pub fn clear(&mut self) {
776 self.callbacks.clear();
777 }
778}
779
780pub fn validate_and_prepare_time_alert(
790 name: &str,
791 mut alert_time_ns: UnixNanos,
792 allow_past: Option<bool>,
793 ts_now: UnixNanos,
794) -> anyhow::Result<(Ustr, UnixNanos)> {
795 check_valid_string_utf8(name, stringify!(name))?;
796
797 let name = Ustr::from(name);
798 let allow_past = allow_past.unwrap_or(true);
799
800 if alert_time_ns < ts_now {
801 if allow_past {
802 log::warn!(
803 "Timer '{name}' alert time {} was in the past, adjusted to current time for immediate firing",
804 alert_time_ns.to_rfc3339(),
805 );
806 alert_time_ns = ts_now;
807 } else {
808 anyhow::bail!(
809 "Timer '{name}' alert time {} was in the past (current time is {ts_now})",
810 alert_time_ns.to_rfc3339(),
811 );
812 }
813 }
814
815 Ok((name, alert_time_ns))
816}
817
818pub fn validate_and_prepare_timer(
834 name: &str,
835 interval_ns: DurationNanos,
836 start_time_ns: Option<UnixNanos>,
837 stop_time_ns: Option<UnixNanos>,
838 allow_past: Option<bool>,
839 fire_immediately: Option<bool>,
840 ts_now: UnixNanos,
841) -> anyhow::Result<(Ustr, UnixNanos, Option<UnixNanos>, bool, bool)> {
842 check_valid_string_utf8(name, stringify!(name))?;
843 check_positive_u64(interval_ns.as_u64(), stringify!(interval_ns))?;
844
845 let name = Ustr::from(name);
846 let allow_past = allow_past.unwrap_or(true);
847 let fire_immediately = fire_immediately.unwrap_or(false);
848
849 let start_time_ns = start_time_ns
850 .filter(|start_time_ns| *start_time_ns != 0)
851 .unwrap_or(ts_now);
852
853 let next_event_time = if fire_immediately {
854 start_time_ns
855 } else {
856 start_time_ns.checked_add(interval_ns).ok_or_else(|| {
857 anyhow::anyhow!("Timer '{name}' first event time exceeds UnixNanos range")
858 })?
859 };
860
861 if !allow_past && next_event_time < ts_now {
862 anyhow::bail!(
863 "Timer '{name}' next event time {} would be in the past (current time is {ts_now})",
864 next_event_time.to_rfc3339(),
865 );
866 }
867
868 if let Some(stop_time) = stop_time_ns {
869 if stop_time <= start_time_ns {
870 anyhow::bail!(
871 "Timer '{name}' stop time {} must be after start time {}",
872 stop_time.to_rfc3339(),
873 start_time_ns.to_rfc3339(),
874 );
875 }
876
877 if !allow_past && stop_time <= ts_now {
878 anyhow::bail!(
879 "Timer '{name}' stop time {} is in the past (current time is {ts_now})",
880 stop_time.to_rfc3339(),
881 );
882 }
883 }
884
885 Ok((
886 name,
887 start_time_ns,
888 stop_time_ns,
889 allow_past,
890 fire_immediately,
891 ))
892}
893
894#[derive(Debug)]
903pub struct TestClock {
904 time: AtomicTime,
905 timers: BTreeMap<Ustr, TestTimer>,
906 timer_queue: BinaryHeap<ScheduledTimeEvent>,
907 callbacks: CallbackRegistry,
908}
909
910impl TestClock {
911 #[must_use]
913 pub fn new() -> Self {
914 Self {
915 time: AtomicTime::new(false, UnixNanos::default()),
916 timers: BTreeMap::new(),
917 timer_queue: BinaryHeap::new(),
918 callbacks: CallbackRegistry::new(),
919 }
920 }
921
922 pub fn advance_time(&mut self, to_time_ns: UnixNanos, set_time: bool) -> Vec<TimeEvent> {
936 const WARN_TIME_EVENTS_THRESHOLD: usize = 1_000_000;
937
938 let from_time_ns = self.time.get_time_ns();
939
940 assert!(
941 to_time_ns >= from_time_ns,
942 "Invariant: time must be non-decreasing, `to_time_ns` {to_time_ns} < `from_time_ns` {from_time_ns}"
943 );
944
945 if set_time {
946 self.time.set_time(to_time_ns);
947 }
948
949 let mut events: Vec<TimeEvent> = Vec::new();
950
951 while self
952 .timer_queue
953 .peek()
954 .is_some_and(|entry| entry.0.ts_event <= to_time_ns)
955 {
956 let entry = self
957 .timer_queue
958 .pop()
959 .expect("timer queue peeked Some but pop returned None");
960
961 let Some((event, next_event)) = self.advance_timer_from_entry(&entry.0) else {
962 continue;
963 };
964
965 events.push(event);
966
967 if let Some(next_event) = next_event {
968 self.timer_queue.push(next_event);
969 }
970 }
971
972 self.compact_timer_queue_if_needed();
973
974 if events.len() >= WARN_TIME_EVENTS_THRESHOLD {
975 log::warn!(
976 "Allocated {} time events during clock advancement from {} to {}, \
977 consider stopping the timer between large time ranges with no data points",
978 events.len().separate_with_commas(),
979 from_time_ns,
980 to_time_ns
981 );
982 }
983
984 events.sort_by(|a, b| {
985 a.ts_event
986 .cmp(&b.ts_event)
987 .then_with(|| a.name.cmp(&b.name))
988 });
989
990 events
991 }
992
993 #[must_use]
1001 pub fn match_handlers(&self, events: Vec<TimeEvent>) -> Vec<TimeEventHandler> {
1002 events
1003 .into_iter()
1004 .map(|event| self.callbacks.get_handler(event))
1005 .collect()
1006 }
1007
1008 fn replace_existing_timer_if_needed(&mut self, name: &Ustr) {
1009 replace_existing_timer(&mut self.timers, name);
1010 self.compact_timer_queue_if_needed();
1011 }
1012
1013 fn insert_timer(&mut self, timer: TestTimer) {
1014 self.timer_queue.push(Self::scheduled_event(&timer));
1015 self.timers.insert(timer.name, timer);
1016 self.compact_timer_queue_if_needed();
1017 }
1018
1019 fn advance_timer_from_entry(
1020 &mut self,
1021 entry: &TimeEvent,
1022 ) -> Option<(TimeEvent, Option<ScheduledTimeEvent>)> {
1023 let timer = self.timers.get_mut(&entry.name)?;
1024 if timer.next_time_ns() != entry.ts_event {
1025 return None;
1026 }
1027
1028 let Some((event, _)) = timer.next() else {
1029 self.timers.remove(&entry.name);
1030 return None;
1031 };
1032
1033 let next_entry = if timer.is_expired() {
1034 self.timers.remove(&entry.name);
1035 None
1036 } else {
1037 Some(Self::scheduled_event(timer))
1038 };
1039
1040 Some((event, next_entry))
1041 }
1042
1043 fn compact_timer_queue_if_needed(&mut self) {
1044 if self.timer_queue.len() > self.timers.len().saturating_mul(2) {
1045 self.compact_timer_queue();
1046 }
1047 }
1048
1049 fn compact_timer_queue(&mut self) {
1050 self.timer_queue = self.timers.values().map(Self::scheduled_event).collect();
1051 }
1052
1053 fn scheduled_event(timer: &TestTimer) -> ScheduledTimeEvent {
1054 ScheduledTimeEvent::new(TimeEvent::new(
1055 timer.name,
1056 UUID4::new(),
1057 timer.next_time_ns(),
1058 timer.next_time_ns(),
1059 ))
1060 }
1061}
1062
1063impl Default for TestClock {
1064 fn default() -> Self {
1066 Self::new()
1067 }
1068}
1069
1070impl Deref for TestClock {
1071 type Target = AtomicTime;
1072
1073 fn deref(&self) -> &Self::Target {
1074 &self.time
1075 }
1076}
1077
1078impl Clock for TestClock {
1079 fn timestamp_ns(&self) -> UnixNanos {
1080 self.time.get_time_ns()
1081 }
1082
1083 fn timestamp_us(&self) -> u64 {
1084 self.time.get_time_us()
1085 }
1086
1087 fn timestamp_ms(&self) -> u64 {
1088 self.time.get_time_ms()
1089 }
1090
1091 fn timestamp(&self) -> f64 {
1092 self.time.get_time()
1093 }
1094
1095 fn timer_names(&self) -> Vec<&str> {
1096 self.timers
1097 .iter()
1098 .filter(|(_, timer)| !timer.is_expired())
1099 .map(|(k, _)| k.as_str())
1100 .collect()
1101 }
1102
1103 fn timer_count(&self) -> usize {
1104 self.timers
1105 .iter()
1106 .filter(|(_, timer)| !timer.is_expired())
1107 .count()
1108 }
1109
1110 fn timer_exists(&self, name: &Ustr) -> bool {
1111 self.timers
1112 .get(name)
1113 .is_some_and(|timer| !timer.is_expired())
1114 }
1115
1116 fn register_default_handler(&mut self, callback: TimeEventCallback) {
1117 self.callbacks.register_default_handler(callback);
1118 }
1119
1120 fn cancel_default_handler(&mut self) {
1121 self.callbacks.cancel_default_handler();
1122 }
1123
1124 fn cancel_callbacks(&mut self) {
1125 self.callbacks.clear();
1126 }
1127
1128 fn set_time_alert_ns(
1129 &mut self,
1130 name: &str,
1131 alert_time_ns: UnixNanos,
1132 callback: Option<TimeEventCallback>,
1133 allow_past: Option<bool>,
1134 ) -> anyhow::Result<()> {
1135 let ts_now = self.get_time_ns();
1136 let (name, alert_time_ns) =
1137 validate_and_prepare_time_alert(name, alert_time_ns, allow_past, ts_now)?;
1138
1139 check_predicate_true(
1140 callback.is_some() | self.callbacks.has_any_callback(&name),
1141 "No callbacks provided",
1142 )?;
1143
1144 self.replace_existing_timer_if_needed(&name);
1145
1146 if let Some(callback) = callback {
1147 self.callbacks.register_callback(name, callback);
1148 }
1149
1150 let interval_ns = create_valid_interval(alert_time_ns - ts_now);
1152 let fire_immediately = alert_time_ns == ts_now;
1153
1154 let timer = TestTimer::new(
1155 name,
1156 interval_ns,
1157 ts_now,
1158 Some(alert_time_ns),
1159 fire_immediately,
1160 );
1161 self.insert_timer(timer);
1162
1163 Ok(())
1164 }
1165
1166 fn set_timer_ns(
1167 &mut self,
1168 name: &str,
1169 interval_ns: DurationNanos,
1170 start_time_ns: Option<UnixNanos>,
1171 stop_time_ns: Option<UnixNanos>,
1172 callback: Option<TimeEventCallback>,
1173 allow_past: Option<bool>,
1174 fire_immediately: Option<bool>,
1175 ) -> anyhow::Result<()> {
1176 let ts_now = self.get_time_ns();
1177 let (name, start_time_ns, stop_time_ns, _allow_past, fire_immediately) =
1178 validate_and_prepare_timer(
1179 name,
1180 interval_ns,
1181 start_time_ns,
1182 stop_time_ns,
1183 allow_past,
1184 fire_immediately,
1185 ts_now,
1186 )?;
1187
1188 check_predicate_true(
1189 callback.is_some() | self.callbacks.has_any_callback(&name),
1190 "No callbacks provided",
1191 )?;
1192
1193 self.replace_existing_timer_if_needed(&name);
1194
1195 if let Some(callback) = callback {
1196 self.callbacks.register_callback(name, callback);
1197 }
1198
1199 let interval_ns = create_valid_interval(interval_ns);
1200
1201 let timer = TestTimer::new(
1202 name,
1203 interval_ns,
1204 start_time_ns,
1205 stop_time_ns,
1206 fire_immediately,
1207 );
1208 self.insert_timer(timer);
1209
1210 Ok(())
1211 }
1212
1213 fn next_time_ns(&self, name: &str) -> Option<UnixNanos> {
1214 self.timers
1215 .get(&Ustr::from(name))
1216 .filter(|timer| !timer.is_expired())
1217 .map(TestTimer::next_time_ns)
1218 }
1219
1220 fn cancel_timer(&mut self, name: &str) {
1221 let timer = self.timers.remove(&Ustr::from(name));
1222 if let Some(mut timer) = timer {
1223 timer.cancel();
1224 }
1225
1226 self.compact_timer_queue_if_needed();
1227 }
1228
1229 fn cancel_timers(&mut self) {
1230 for timer in &mut self.timers.values_mut() {
1231 timer.cancel();
1232 }
1233
1234 self.timers.clear();
1235 self.timer_queue.clear();
1236 }
1237
1238 fn reset(&mut self) {
1239 self.time = AtomicTime::new(false, UnixNanos::default());
1240 self.timers = BTreeMap::new();
1241 self.timer_queue = BinaryHeap::new();
1242 self.callbacks.clear();
1243 }
1244}
1245
1246pub(crate) fn replace_existing_timer<T: Timer>(timers: &mut BTreeMap<Ustr, T>, name: &Ustr) {
1247 let Some(mut timer) = timers.remove(name) else {
1248 return;
1249 };
1250
1251 if timer.is_expired() {
1252 return;
1253 }
1254
1255 timer.cancel();
1256 log::warn!("Timer '{name}' replaced");
1257}
1258
1259#[cfg(test)]
1260mod tests {
1261 use std::{cell::RefCell, collections::BTreeMap, sync::Arc, time::Duration};
1262
1263 use nautilus_core::{DurationNanos, UnixNanos};
1264 use parking_lot::Mutex;
1265 use proptest::{prelude::*, test_runner::TestCaseResult};
1266 use rstest::{fixture, rstest};
1267 use ustr::Ustr;
1268
1269 use super::*;
1270 use crate::timer::{TimeEvent, TimeEventCallback};
1271
1272 #[derive(Debug, Default)]
1273 struct TestCallback {
1274 called: Arc<Mutex<bool>>,
1276 }
1277
1278 impl TestCallback {
1279 fn new(called: Arc<Mutex<bool>>) -> Self {
1280 Self { called }
1281 }
1282 }
1283
1284 impl From<TestCallback> for TimeEventCallback {
1285 fn from(callback: TestCallback) -> Self {
1286 Self::from(move |_event: TimeEvent| {
1287 *callback.called.lock() = true;
1288 })
1289 }
1290 }
1291
1292 #[fixture]
1293 pub fn test_clock() -> TestClock {
1294 let mut clock = TestClock::new();
1295 clock.register_default_handler(TestCallback::default().into());
1296 clock
1297 }
1298
1299 #[rstest]
1300 fn test_time_monotonicity(mut test_clock: TestClock) {
1301 let initial_time = test_clock.timestamp_ns();
1302 test_clock.advance_time(initial_time + DurationNanos::new(1000), true);
1303 assert!(test_clock.timestamp_ns() > initial_time);
1304 }
1305
1306 #[rstest]
1307 fn test_timer_registration(mut test_clock: TestClock) {
1308 test_clock
1309 .set_time_alert_ns(
1310 "test_timer",
1311 test_clock.timestamp_ns() + DurationNanos::new(1000),
1312 None,
1313 None,
1314 )
1315 .unwrap();
1316 assert_eq!(test_clock.timer_count(), 1);
1317 assert_eq!(test_clock.timer_names(), vec!["test_timer"]);
1318 }
1319
1320 #[rstest]
1321 fn test_timer_expiration(mut test_clock: TestClock) {
1322 let alert_time = test_clock.timestamp_ns() + DurationNanos::new(1000);
1323 test_clock
1324 .set_time_alert_ns("test_timer", alert_time, None, None)
1325 .unwrap();
1326 let events = test_clock.advance_time(alert_time, true);
1327 assert_eq!(events.len(), 1);
1328 assert_eq!(events[0].name, "test_timer");
1329 }
1330
1331 #[rstest]
1332 fn test_timer_cancellation(mut test_clock: TestClock) {
1333 test_clock
1334 .set_time_alert_ns(
1335 "test_timer",
1336 test_clock.timestamp_ns() + DurationNanos::new(1000),
1337 None,
1338 None,
1339 )
1340 .unwrap();
1341 assert_eq!(test_clock.timer_count(), 1);
1342 test_clock.cancel_timer("test_timer");
1343 assert_eq!(test_clock.timer_count(), 0);
1344 }
1345
1346 #[rstest]
1347 fn test_time_advancement(mut test_clock: TestClock) {
1348 let start_time = test_clock.timestamp_ns();
1349 test_clock
1350 .set_timer_ns(
1351 "test_timer",
1352 DurationNanos::new(1000),
1353 Some(start_time),
1354 None,
1355 None,
1356 None,
1357 None,
1358 )
1359 .unwrap();
1360 let events = test_clock.advance_time(start_time + DurationNanos::new(2500), true);
1361 assert_eq!(events.len(), 2);
1362 assert_eq!(events[0].ts_event, start_time + DurationNanos::new(1000));
1363 assert_eq!(events[1].ts_event, start_time + DurationNanos::new(2000));
1364 }
1365
1366 #[rstest]
1367 fn test_default_and_custom_callbacks() {
1368 let mut clock = TestClock::new();
1369 let default_called = Arc::new(Mutex::new(false));
1370 let custom_called = Arc::new(Mutex::new(false));
1371
1372 let default_callback = TestCallback::new(Arc::clone(&default_called));
1373 let custom_callback = TestCallback::new(Arc::clone(&custom_called));
1374
1375 clock.register_default_handler(TimeEventCallback::from(default_callback));
1376 clock
1377 .set_time_alert_ns(
1378 "default_timer",
1379 clock.timestamp_ns() + DurationNanos::new(1000),
1380 None,
1381 None,
1382 )
1383 .unwrap();
1384 clock
1385 .set_time_alert_ns(
1386 "custom_timer",
1387 clock.timestamp_ns() + DurationNanos::new(1000),
1388 Some(TimeEventCallback::from(custom_callback)),
1389 None,
1390 )
1391 .unwrap();
1392
1393 let events = clock.advance_time(clock.timestamp_ns() + DurationNanos::new(1000), true);
1394 let handlers = clock.match_handlers(events);
1395
1396 for handler in handlers {
1397 handler.callback.call(handler.event);
1398 }
1399
1400 assert!(*default_called.lock());
1401 assert!(*custom_called.lock());
1402 }
1403
1404 #[rstest]
1405 fn test_timer_with_rust_local_callback() {
1406 use std::{cell::RefCell, rc::Rc};
1407
1408 let mut clock = TestClock::new();
1409 let call_count = Rc::new(RefCell::new(0_u32));
1410 let call_count_clone = Rc::clone(&call_count);
1411
1412 let callback: Rc<dyn Fn(TimeEvent)> = Rc::new(move |_event: TimeEvent| {
1414 *call_count_clone.borrow_mut() += 1;
1415 });
1416
1417 clock
1418 .set_time_alert_ns(
1419 "local_timer",
1420 clock.timestamp_ns() + DurationNanos::new(1000),
1421 Some(TimeEventCallback::from(callback)),
1422 None,
1423 )
1424 .unwrap();
1425
1426 let events = clock.advance_time(clock.timestamp_ns() + DurationNanos::new(1000), true);
1427 let handlers = clock.match_handlers(events);
1428
1429 for handler in handlers {
1430 handler.callback.call(handler.event);
1431 }
1432
1433 assert_eq!(*call_count.borrow(), 1);
1434 }
1435
1436 #[rstest]
1437 fn test_multiple_timers(mut test_clock: TestClock) {
1438 let start_time = test_clock.timestamp_ns();
1439 test_clock
1440 .set_timer_ns(
1441 "timer1",
1442 DurationNanos::new(1000),
1443 Some(start_time),
1444 None,
1445 None,
1446 None,
1447 None,
1448 )
1449 .unwrap();
1450 test_clock
1451 .set_timer_ns(
1452 "timer2",
1453 DurationNanos::new(2000),
1454 Some(start_time),
1455 None,
1456 None,
1457 None,
1458 None,
1459 )
1460 .unwrap();
1461 let events = test_clock.advance_time(start_time + DurationNanos::new(2000), true);
1462 assert_eq!(events.len(), 3);
1463 assert_eq!(events[0].name, "timer1");
1464 assert_eq!(events[1].name, "timer1");
1465 assert_eq!(events[2].name, "timer2");
1466 }
1467
1468 #[rstest]
1469 fn test_allow_past_parameter_true(mut test_clock: TestClock) {
1470 test_clock.set_time(UnixNanos::from(2000));
1471 let current_time = test_clock.timestamp_ns();
1472 let past_time = current_time - DurationNanos::new(1000);
1473
1474 test_clock
1476 .set_time_alert_ns("past_timer", past_time, None, Some(true))
1477 .unwrap();
1478
1479 assert_eq!(test_clock.timer_count(), 1);
1481 assert_eq!(test_clock.timer_names(), vec!["past_timer"]);
1482
1483 let next_time = test_clock.next_time_ns("past_timer").unwrap();
1485 assert!(next_time >= current_time);
1486 }
1487
1488 #[rstest]
1489 fn test_allow_past_parameter_false(mut test_clock: TestClock) {
1490 test_clock.set_time(UnixNanos::from(2000));
1491 let current_time = test_clock.timestamp_ns();
1492 let past_time = current_time - DurationNanos::new(1000);
1493
1494 let result = test_clock.set_time_alert_ns("past_timer", past_time, None, Some(false));
1496
1497 assert!(result.is_err());
1499 assert!(format!("{}", result.unwrap_err()).contains("was in the past"));
1500
1501 assert_eq!(test_clock.timer_count(), 0);
1503 assert!(test_clock.timer_names().is_empty());
1504 }
1505
1506 #[rstest]
1507 fn test_invalid_stop_time_validation(mut test_clock: TestClock) {
1508 test_clock.set_time(UnixNanos::from(2000));
1509 let current_time = test_clock.timestamp_ns();
1510 let start_time = current_time + DurationNanos::new(1000);
1511 let stop_time = current_time + DurationNanos::new(500); let result = test_clock.set_timer_ns(
1515 "invalid_timer",
1516 DurationNanos::new(100),
1517 Some(start_time),
1518 Some(stop_time),
1519 None,
1520 None,
1521 None,
1522 );
1523
1524 assert!(result.is_err());
1526 assert!(format!("{}", result.unwrap_err()).contains("must be after start time"));
1527
1528 assert_eq!(test_clock.timer_count(), 0);
1530 }
1531
1532 #[rstest]
1533 fn test_set_timer_ns_fire_immediately_true(mut test_clock: TestClock) {
1534 let start_time = test_clock.timestamp_ns();
1535 let interval_ns = DurationNanos::new(1000);
1536
1537 test_clock
1538 .set_timer_ns(
1539 "fire_immediately_timer",
1540 interval_ns,
1541 Some(start_time),
1542 None,
1543 None,
1544 None,
1545 Some(true),
1546 )
1547 .unwrap();
1548
1549 let events = test_clock.advance_time(start_time + DurationNanos::new(2500), true);
1551
1552 assert_eq!(events.len(), 3);
1554 assert_eq!(*events[0].ts_event, *start_time); assert_eq!(events[1].ts_event, start_time + DurationNanos::new(1000)); assert_eq!(events[2].ts_event, start_time + DurationNanos::new(2000)); }
1558
1559 #[rstest]
1560 fn test_set_timer_ns_fire_immediately_false(mut test_clock: TestClock) {
1561 let start_time = test_clock.timestamp_ns();
1562 let interval_ns = DurationNanos::new(1000);
1563
1564 test_clock
1565 .set_timer_ns(
1566 "normal_timer",
1567 interval_ns,
1568 Some(start_time),
1569 None,
1570 None,
1571 None,
1572 Some(false),
1573 )
1574 .unwrap();
1575
1576 let events = test_clock.advance_time(start_time + DurationNanos::new(2500), true);
1578
1579 assert_eq!(events.len(), 2);
1581 assert_eq!(events[0].ts_event, start_time + DurationNanos::new(1000)); assert_eq!(events[1].ts_event, start_time + DurationNanos::new(2000)); }
1584
1585 #[rstest]
1586 fn test_set_timer_ns_fire_immediately_default_is_false(mut test_clock: TestClock) {
1587 let start_time = test_clock.timestamp_ns();
1588 let interval_ns = DurationNanos::new(1000);
1589
1590 test_clock
1592 .set_timer_ns(
1593 "default_timer",
1594 interval_ns,
1595 Some(start_time),
1596 None,
1597 None,
1598 None,
1599 None,
1600 )
1601 .unwrap();
1602
1603 let events = test_clock.advance_time(start_time + DurationNanos::new(1500), true);
1604
1605 assert_eq!(events.len(), 1);
1607 assert_eq!(events[0].ts_event, start_time + DurationNanos::new(1000)); }
1609
1610 #[rstest]
1611 fn test_set_timer_ns_fire_immediately_with_zero_start_time(mut test_clock: TestClock) {
1612 test_clock.set_time(5000.into());
1613 let interval_ns = DurationNanos::new(1000);
1614
1615 test_clock
1616 .set_timer_ns(
1617 "zero_start_timer",
1618 interval_ns,
1619 None,
1620 None,
1621 None,
1622 None,
1623 Some(true),
1624 )
1625 .unwrap();
1626
1627 let events = test_clock.advance_time(UnixNanos::from(7000), true);
1628
1629 assert_eq!(events.len(), 3);
1632 assert_eq!(*events[0].ts_event, 5000); assert_eq!(*events[1].ts_event, 6000);
1634 assert_eq!(*events[2].ts_event, 7000);
1635 }
1636
1637 #[rstest]
1638 fn test_multiple_timers_different_fire_immediately_settings(mut test_clock: TestClock) {
1639 let start_time = test_clock.timestamp_ns();
1640 let interval_ns = DurationNanos::new(1000);
1641
1642 test_clock
1644 .set_timer_ns(
1645 "immediate_timer",
1646 interval_ns,
1647 Some(start_time),
1648 None,
1649 None,
1650 None,
1651 Some(true),
1652 )
1653 .unwrap();
1654
1655 test_clock
1657 .set_timer_ns(
1658 "normal_timer",
1659 interval_ns,
1660 Some(start_time),
1661 None,
1662 None,
1663 None,
1664 Some(false),
1665 )
1666 .unwrap();
1667
1668 let events = test_clock.advance_time(start_time + DurationNanos::new(1500), true);
1669
1670 assert_eq!(events.len(), 3);
1672
1673 let mut event_times: Vec<u64> = events.iter().map(|e| e.ts_event.as_u64()).collect();
1675 event_times.sort_unstable();
1676
1677 assert_eq!(event_times[0], start_time.as_u64()); assert_eq!(event_times[1], start_time.as_u64() + 1000); assert_eq!(event_times[2], start_time.as_u64() + 1000); }
1681
1682 #[rstest]
1683 fn test_timer_name_collision_overwrites(mut test_clock: TestClock) {
1684 let start_time = test_clock.timestamp_ns();
1685
1686 test_clock
1688 .set_timer_ns(
1689 "collision_timer",
1690 DurationNanos::new(1000),
1691 Some(start_time),
1692 None,
1693 None,
1694 None,
1695 None,
1696 )
1697 .unwrap();
1698
1699 let result = test_clock.set_timer_ns(
1701 "collision_timer",
1702 DurationNanos::new(2000),
1703 Some(start_time),
1704 None,
1705 None,
1706 None,
1707 None,
1708 );
1709
1710 assert!(result.is_ok());
1711 assert_eq!(test_clock.timer_count(), 1);
1713
1714 let next_time = test_clock.next_time_ns("collision_timer").unwrap();
1716 assert_eq!(next_time, start_time + DurationNanos::new(2000));
1718 }
1719
1720 #[rstest]
1721 fn test_timer_zero_interval_error(mut test_clock: TestClock) {
1722 let start_time = test_clock.timestamp_ns();
1723
1724 let result = test_clock.set_timer_ns(
1726 "zero_interval",
1727 DurationNanos::default(),
1728 Some(start_time),
1729 None,
1730 None,
1731 None,
1732 None,
1733 );
1734
1735 assert!(result.is_err());
1736 assert_eq!(test_clock.timer_count(), 0);
1737 }
1738
1739 #[rstest]
1740 fn test_timer_empty_name_error(mut test_clock: TestClock) {
1741 let start_time = test_clock.timestamp_ns();
1742
1743 let result = test_clock.set_timer_ns(
1745 "",
1746 DurationNanos::new(1000),
1747 Some(start_time),
1748 None,
1749 None,
1750 None,
1751 None,
1752 );
1753
1754 assert!(result.is_err());
1755 assert_eq!(test_clock.timer_count(), 0);
1756 }
1757
1758 #[rstest]
1759 fn test_timer_exists(mut test_clock: TestClock) {
1760 let name = Ustr::from("exists_timer");
1761 assert!(!test_clock.timer_exists(&name));
1762
1763 test_clock
1764 .set_time_alert_ns(
1765 name.as_str(),
1766 test_clock.timestamp_ns() + DurationNanos::new(1_000),
1767 None,
1768 None,
1769 )
1770 .unwrap();
1771
1772 assert!(test_clock.timer_exists(&name));
1773 }
1774
1775 #[rstest]
1776 fn test_timer_exists_consistent_with_names_and_count_after_expiry(mut test_clock: TestClock) {
1777 let name = Ustr::from("expiring_timer");
1778 let start_time = test_clock.timestamp_ns();
1779
1780 test_clock
1781 .set_timer_ns(
1782 name.as_str(),
1783 DurationNanos::new(1_000),
1784 Some(start_time),
1785 Some(start_time + DurationNanos::new(2_500)),
1786 None,
1787 None,
1788 None,
1789 )
1790 .unwrap();
1791
1792 assert!(test_clock.timer_exists(&name));
1793 assert_eq!(test_clock.timer_count(), 1);
1794
1795 test_clock.advance_time(start_time + DurationNanos::new(10_000), true);
1796
1797 assert!(!test_clock.timer_exists(&name));
1799 assert_eq!(test_clock.timer_count(), 0);
1800 assert!(test_clock.timer_names().is_empty());
1801 }
1802
1803 #[rstest]
1804 fn test_timer_rejects_past_stop_time_when_not_allowed(mut test_clock: TestClock) {
1805 test_clock.set_time(UnixNanos::from(10_000));
1806 let current = test_clock.timestamp_ns();
1807
1808 let result = test_clock.set_timer_ns(
1809 "past_stop",
1810 DurationNanos::new(10_000),
1811 Some(current - DurationNanos::new(500)),
1812 Some(current - DurationNanos::new(100)),
1813 None,
1814 Some(false),
1815 None,
1816 );
1817
1818 let err = result.expect_err("expected stop time validation error");
1819 let err_msg = err.to_string();
1820 assert!(err_msg.contains("stop time"));
1821 assert!(err_msg.contains("in the past"));
1822 }
1823
1824 #[rstest]
1825 fn test_timer_accepts_future_stop_time(mut test_clock: TestClock) {
1826 let current = test_clock.timestamp_ns();
1827
1828 let result = test_clock.set_timer_ns(
1829 "future_stop",
1830 DurationNanos::new(1_000),
1831 Some(current),
1832 Some(current + DurationNanos::new(10_000)),
1833 None,
1834 Some(false),
1835 None,
1836 );
1837
1838 assert!(result.is_ok());
1839 }
1840
1841 #[rstest]
1842 fn test_timer_fire_immediately_at_exact_stop_time(mut test_clock: TestClock) {
1843 let start_time = test_clock.timestamp_ns();
1844 let interval_ns = DurationNanos::new(1000);
1845 let stop_time = start_time + interval_ns; test_clock
1848 .set_timer_ns(
1849 "exact_stop",
1850 interval_ns,
1851 Some(start_time),
1852 Some(stop_time),
1853 None,
1854 None,
1855 Some(true),
1856 )
1857 .unwrap();
1858
1859 let events = test_clock.advance_time(stop_time, true);
1860
1861 assert_eq!(events.len(), 2);
1863 assert_eq!(*events[0].ts_event, *start_time); assert_eq!(*events[1].ts_event, *stop_time); }
1866
1867 #[rstest]
1868 fn test_timer_advance_to_exact_next_time(mut test_clock: TestClock) {
1869 let start_time = test_clock.timestamp_ns();
1870 let interval_ns = DurationNanos::new(1000);
1871
1872 test_clock
1873 .set_timer_ns(
1874 "exact_advance",
1875 interval_ns,
1876 Some(start_time),
1877 None,
1878 None,
1879 None,
1880 Some(false),
1881 )
1882 .unwrap();
1883
1884 let next_time = test_clock.next_time_ns("exact_advance").unwrap();
1886 let events = test_clock.advance_time(next_time, true);
1887
1888 assert_eq!(events.len(), 1);
1889 assert_eq!(*events[0].ts_event, *next_time);
1890 }
1891
1892 #[rstest]
1893 fn test_allow_past_bar_aggregation_use_case(mut test_clock: TestClock) {
1894 test_clock.set_time(UnixNanos::from(100_500)); let bar_start_time = UnixNanos::from(100_000); let interval_ns = DurationNanos::new(1000); let result = test_clock.set_timer_ns(
1904 "bar_timer",
1905 interval_ns,
1906 Some(bar_start_time),
1907 None,
1908 None,
1909 Some(false), Some(false), );
1912
1913 assert!(result.is_ok());
1915 assert_eq!(test_clock.timer_count(), 1);
1916
1917 let next_time = test_clock.next_time_ns("bar_timer").unwrap();
1919 assert_eq!(*next_time, 101_000);
1920 }
1921
1922 #[rstest]
1923 fn test_allow_past_false_rejects_when_next_event_in_past(mut test_clock: TestClock) {
1924 test_clock.set_time(UnixNanos::from(102_000)); let past_start_time = UnixNanos::from(100_000); let interval_ns = DurationNanos::new(1000); let result = test_clock.set_timer_ns(
1933 "past_event_timer",
1934 interval_ns,
1935 Some(past_start_time),
1936 None,
1937 None,
1938 Some(false), Some(false), );
1941
1942 assert!(result.is_err());
1944 assert!(
1945 result
1946 .unwrap_err()
1947 .to_string()
1948 .contains("would be in the past")
1949 );
1950 }
1951
1952 #[rstest]
1953 fn test_allow_past_false_with_fire_immediately_true(mut test_clock: TestClock) {
1954 test_clock.set_time(UnixNanos::from(100_500)); let past_start_time = UnixNanos::from(100_000); let interval_ns = DurationNanos::new(1000);
1958
1959 let result = test_clock.set_timer_ns(
1962 "immediate_past_timer",
1963 interval_ns,
1964 Some(past_start_time),
1965 None,
1966 None,
1967 Some(false), Some(true), );
1970
1971 assert!(result.is_err());
1973 assert!(
1974 result
1975 .unwrap_err()
1976 .to_string()
1977 .contains("would be in the past")
1978 );
1979 }
1980
1981 #[rstest]
1982 fn test_cancel_timer_during_execution(mut test_clock: TestClock) {
1983 let start_time = test_clock.timestamp_ns();
1984
1985 test_clock
1986 .set_timer_ns(
1987 "cancel_test",
1988 DurationNanos::new(1000),
1989 Some(start_time),
1990 None,
1991 None,
1992 None,
1993 None,
1994 )
1995 .unwrap();
1996
1997 assert_eq!(test_clock.timer_count(), 1);
1998
1999 test_clock.cancel_timer("cancel_test");
2001
2002 assert_eq!(test_clock.timer_count(), 0);
2003
2004 let events = test_clock.advance_time(start_time + DurationNanos::new(2000), true);
2006 assert_eq!(events.len(), 0);
2007 }
2008
2009 #[rstest]
2010 fn test_cancelled_timer_queue_entry_is_skipped(mut test_clock: TestClock) {
2011 let start_time = test_clock.timestamp_ns();
2012 test_clock
2013 .set_time_alert_ns(
2014 "cancelled",
2015 start_time + DurationNanos::new(1000),
2016 None,
2017 None,
2018 )
2019 .unwrap();
2020 test_clock
2021 .set_time_alert_ns("active", start_time + DurationNanos::new(2000), None, None)
2022 .unwrap();
2023
2024 test_clock.cancel_timer("cancelled");
2025 assert_eq!(test_clock.timer_count(), 1);
2026 assert_eq!(test_clock.timer_queue.len(), 2);
2027
2028 let events = test_clock.advance_time(start_time + DurationNanos::new(1000), true);
2029 assert!(events.is_empty());
2030 assert_eq!(test_clock.timer_names(), vec!["active"]);
2031
2032 let events = test_clock.advance_time(start_time + DurationNanos::new(2000), true);
2033 assert_eq!(events.len(), 1);
2034 assert_eq!(events[0].name, "active");
2035 }
2036
2037 #[rstest]
2038 fn test_timer_queue_compacts_stale_entries(mut test_clock: TestClock) {
2039 let start_time = test_clock.timestamp_ns();
2040 test_clock
2041 .set_time_alert_ns("active", start_time + DurationNanos::new(1000), None, None)
2042 .unwrap();
2043 test_clock
2044 .set_time_alert_ns(
2045 "cancelled-1",
2046 start_time + DurationNanos::new(2000),
2047 None,
2048 None,
2049 )
2050 .unwrap();
2051 test_clock
2052 .set_time_alert_ns(
2053 "cancelled-2",
2054 start_time + DurationNanos::new(3000),
2055 None,
2056 None,
2057 )
2058 .unwrap();
2059
2060 test_clock.cancel_timer("cancelled-1");
2061 assert_eq!(test_clock.timer_queue.len(), 3);
2062
2063 test_clock.cancel_timer("cancelled-2");
2064 assert_eq!(test_clock.timer_count(), 1);
2065 assert_eq!(test_clock.timer_queue.len(), 1);
2066 }
2067
2068 #[rstest]
2069 fn test_cancel_all_timers(mut test_clock: TestClock) {
2070 test_clock
2072 .set_timer_ns(
2073 "timer1",
2074 DurationNanos::new(1000),
2075 None,
2076 None,
2077 None,
2078 None,
2079 None,
2080 )
2081 .unwrap();
2082 test_clock
2083 .set_timer_ns(
2084 "timer2",
2085 DurationNanos::new(1500),
2086 None,
2087 None,
2088 None,
2089 None,
2090 None,
2091 )
2092 .unwrap();
2093 test_clock
2094 .set_timer_ns(
2095 "timer3",
2096 DurationNanos::new(2000),
2097 None,
2098 None,
2099 None,
2100 None,
2101 None,
2102 )
2103 .unwrap();
2104
2105 assert_eq!(test_clock.timer_count(), 3);
2106
2107 test_clock.cancel_timers();
2109
2110 assert_eq!(test_clock.timer_count(), 0);
2111
2112 let events = test_clock.advance_time(UnixNanos::from(5000), true);
2114 assert_eq!(events.len(), 0);
2115 }
2116
2117 #[rstest]
2118 fn test_clock_reset_clears_timers(mut test_clock: TestClock) {
2119 test_clock
2120 .set_timer_ns(
2121 "reset_test",
2122 DurationNanos::new(1000),
2123 None,
2124 None,
2125 None,
2126 None,
2127 None,
2128 )
2129 .unwrap();
2130
2131 assert_eq!(test_clock.timer_count(), 1);
2132
2133 test_clock.reset();
2135
2136 assert_eq!(test_clock.timer_count(), 0);
2137 assert_eq!(test_clock.timestamp_ns(), UnixNanos::default()); }
2139
2140 #[rstest]
2141 fn test_cancel_default_handler_clears_default(mut test_clock: TestClock) {
2142 test_clock.cancel_default_handler();
2144
2145 let alert_time: UnixNanos = test_clock.timestamp_ns() + DurationNanos::new(1000);
2147 let err = test_clock
2148 .set_time_alert_ns("alert", alert_time, None, None)
2149 .unwrap_err();
2150 assert!(
2151 err.to_string().contains("No callbacks provided"),
2152 "unexpected error: {err}"
2153 );
2154 }
2155
2156 #[rstest]
2157 fn test_cancel_default_handler_is_idempotent_on_empty_registry() {
2158 let mut clock = TestClock::new();
2160 clock.cancel_default_handler();
2161 clock.cancel_default_handler();
2162 }
2163
2164 #[rstest]
2165 fn test_cancel_callbacks_clears_named(mut test_clock: TestClock) {
2166 let alert_time: UnixNanos = test_clock.timestamp_ns() + DurationNanos::new(1000);
2167 let callback = TimeEventCallback::from(TestCallback::default());
2168 test_clock
2169 .set_time_alert_ns("named_alert", alert_time, Some(callback), None)
2170 .unwrap();
2171 test_clock.cancel_timer("named_alert");
2172
2173 test_clock.cancel_default_handler();
2175 test_clock.cancel_callbacks();
2176
2177 let err = test_clock
2178 .set_time_alert_ns("named_alert", alert_time, None, None)
2179 .unwrap_err();
2180 assert!(
2181 err.to_string().contains("No callbacks provided"),
2182 "unexpected error: {err}"
2183 );
2184 }
2185
2186 #[rstest]
2187 fn test_failed_set_time_alert_ns_preserves_existing_timer() {
2188 let mut clock = TestClock::new();
2190 let alert_time: UnixNanos = clock.timestamp_ns() + DurationNanos::new(1000);
2191 let callback = TimeEventCallback::from(TestCallback::default());
2192 clock
2193 .set_time_alert_ns("alert", alert_time, Some(callback), None)
2194 .unwrap();
2195 assert_eq!(clock.next_time_ns("alert"), Some(alert_time));
2196
2197 clock.cancel_callbacks();
2199
2200 let err = clock
2203 .set_time_alert_ns("alert", alert_time + DurationNanos::new(1000), None, None)
2204 .unwrap_err();
2205 assert!(
2206 err.to_string().contains("No callbacks provided"),
2207 "unexpected error: {err}"
2208 );
2209 assert_eq!(clock.timer_count(), 1);
2210 assert_eq!(clock.next_time_ns("alert"), Some(alert_time));
2211 }
2212
2213 #[rstest]
2214 fn test_cancel_default_handler_preserves_named_callbacks(mut test_clock: TestClock) {
2215 let alert_time: UnixNanos = test_clock.timestamp_ns() + DurationNanos::new(1000);
2216 let callback = TimeEventCallback::from(TestCallback::default());
2217 test_clock
2218 .set_time_alert_ns("alert", alert_time, Some(callback), None)
2219 .unwrap();
2220 test_clock.cancel_timer("alert");
2221
2222 test_clock.cancel_default_handler();
2223
2224 test_clock
2226 .set_time_alert_ns("alert", alert_time, None, None)
2227 .unwrap();
2228 }
2229
2230 #[rstest]
2231 fn test_cancel_callbacks_preserves_default_handler(mut test_clock: TestClock) {
2232 test_clock.cancel_callbacks();
2234
2235 let alert_time: UnixNanos = test_clock.timestamp_ns() + DurationNanos::new(1000);
2236 test_clock
2237 .set_time_alert_ns("alert", alert_time, None, None)
2238 .unwrap();
2239 }
2240
2241 #[rstest]
2242 fn test_set_time_alert_default_impl(mut test_clock: TestClock) {
2243 let current_time = test_clock.utc_now();
2244 let alert_time = current_time + jiff::SignedDuration::from_secs(1);
2245
2246 test_clock
2248 .set_time_alert("alert_test", alert_time, None, None)
2249 .unwrap();
2250
2251 assert_eq!(test_clock.timer_count(), 1);
2252 assert_eq!(test_clock.timer_names(), vec!["alert_test"]);
2253
2254 let expected_ns = UnixNanos::from(alert_time);
2256 let next_time = test_clock.next_time_ns("alert_test").unwrap();
2257
2258 let diff = if next_time >= expected_ns {
2260 next_time.as_u64() - expected_ns.as_u64()
2261 } else {
2262 expected_ns.as_u64() - next_time.as_u64()
2263 };
2264
2265 assert!(
2266 diff < 1000,
2267 "Timer should be set within 1 microsecond of expected time"
2268 );
2269 }
2270
2271 #[rstest]
2272 fn test_set_timer_default_impl(mut test_clock: TestClock) {
2273 let current_time = test_clock.utc_now();
2274 let start_time = current_time + jiff::SignedDuration::from_secs(1);
2275 let interval = Duration::from_millis(500);
2276
2277 test_clock
2279 .set_timer(
2280 "timer_test",
2281 interval,
2282 Some(start_time),
2283 None,
2284 None,
2285 None,
2286 None,
2287 )
2288 .unwrap();
2289
2290 assert_eq!(test_clock.timer_count(), 1);
2291 assert_eq!(test_clock.timer_names(), vec!["timer_test"]);
2292
2293 let start_ns = UnixNanos::from(start_time);
2295 let interval_ns = interval.as_nanos() as u64;
2296
2297 let events = test_clock.advance_time(start_ns + DurationNanos::new(interval_ns) * 3, true);
2298 assert_eq!(events.len(), 3); assert_eq!(*events[0].ts_event, *start_ns + interval_ns);
2302 assert_eq!(*events[1].ts_event, *start_ns + interval_ns * 2);
2303 assert_eq!(*events[2].ts_event, *start_ns + interval_ns * 3);
2304 }
2305
2306 #[rstest]
2307 fn test_set_timer_with_stop_time_default_impl(mut test_clock: TestClock) {
2308 let current_time = test_clock.utc_now();
2309 let start_time = current_time + jiff::SignedDuration::from_secs(1);
2310 let stop_time = current_time + jiff::SignedDuration::from_secs(3);
2311 let interval = Duration::from_secs(1);
2312
2313 test_clock
2315 .set_timer(
2316 "timer_with_stop",
2317 interval,
2318 Some(start_time),
2319 Some(stop_time),
2320 None,
2321 None,
2322 None,
2323 )
2324 .unwrap();
2325
2326 assert_eq!(test_clock.timer_count(), 1);
2327
2328 let stop_ns = UnixNanos::from(stop_time);
2330 let events = test_clock.advance_time(stop_ns + DurationNanos::new(1000), true);
2331
2332 assert_eq!(events.len(), 2);
2334
2335 let start_ns = UnixNanos::from(start_time);
2336 let interval_ns = interval.as_nanos() as u64;
2337 assert_eq!(*events[0].ts_event, *start_ns + interval_ns);
2338 assert_eq!(*events[1].ts_event, *start_ns + interval_ns * 2);
2339 }
2340
2341 #[rstest]
2342 fn test_set_timer_fire_immediately_default_impl(mut test_clock: TestClock) {
2343 let current_time = test_clock.utc_now();
2344 let start_time = current_time + jiff::SignedDuration::from_secs(1);
2345 let interval = Duration::from_millis(500);
2346
2347 test_clock
2349 .set_timer(
2350 "immediate_timer",
2351 interval,
2352 Some(start_time),
2353 None,
2354 None,
2355 None,
2356 Some(true),
2357 )
2358 .unwrap();
2359
2360 let start_ns = UnixNanos::from(start_time);
2361 let interval_ns = interval.as_nanos() as u64;
2362
2363 let events = test_clock.advance_time(start_ns + DurationNanos::new(interval_ns), true);
2365
2366 assert_eq!(events.len(), 2);
2368 assert_eq!(*events[0].ts_event, *start_ns); assert_eq!(*events[1].ts_event, *start_ns + interval_ns); }
2371
2372 #[rstest]
2373 fn test_set_time_alert_when_alert_time_equals_current_time(mut test_clock: TestClock) {
2374 let current_time = test_clock.timestamp_ns();
2375
2376 test_clock
2378 .set_time_alert_ns("alert_at_current_time", current_time, None, None)
2379 .unwrap();
2380
2381 assert_eq!(test_clock.timer_count(), 1);
2382
2383 let events = test_clock.advance_time(current_time, true);
2385
2386 assert_eq!(events.len(), 1);
2388 assert_eq!(events[0].name, "alert_at_current_time");
2389 assert_eq!(*events[0].ts_event, *current_time);
2390 }
2391
2392 #[rstest]
2393 fn test_cancel_and_reschedule_same_name(mut test_clock: TestClock) {
2394 let start = test_clock.timestamp_ns();
2395
2396 test_clock
2397 .set_time_alert_ns("timer", start + DurationNanos::new(1000), None, None)
2398 .unwrap();
2399 assert_eq!(test_clock.timer_count(), 1);
2400
2401 test_clock.cancel_timer("timer");
2402 assert_eq!(test_clock.timer_count(), 0);
2403
2404 test_clock
2405 .set_time_alert_ns("timer", start + DurationNanos::new(2000), None, None)
2406 .unwrap();
2407 assert_eq!(test_clock.timer_count(), 1);
2408
2409 let events = test_clock.advance_time(start + DurationNanos::new(1500), true);
2410 assert!(events.is_empty());
2411
2412 let events = test_clock.advance_time(start + DurationNanos::new(2000), true);
2413 assert_eq!(events.len(), 1);
2414 assert_eq!(events[0].ts_event, start + DurationNanos::new(2000));
2415 }
2416
2417 #[rstest]
2418 fn test_multiple_timers_same_timestamp_all_fire(mut test_clock: TestClock) {
2419 let fire_time = test_clock.timestamp_ns() + DurationNanos::new(1000);
2420
2421 for i in 0..5 {
2422 test_clock
2423 .set_time_alert_ns(&format!("timer_{i}"), fire_time, None, None)
2424 .unwrap();
2425 }
2426
2427 assert_eq!(test_clock.timer_count(), 5);
2428
2429 let events = test_clock.advance_time(fire_time, true);
2430 assert_eq!(events.len(), 5);
2431
2432 for event in &events {
2433 assert_eq!(*event.ts_event, *fire_time);
2434 }
2435 }
2436
2437 #[rstest]
2438 fn test_events_ordered_by_timestamp_after_advance() {
2439 let mut clock = TestClock::new();
2440 clock.register_default_handler(TestCallback::default().into());
2441 let start = clock.timestamp_ns();
2442
2443 clock
2444 .set_time_alert_ns("third", start + DurationNanos::new(300), None, None)
2445 .unwrap();
2446 clock
2447 .set_time_alert_ns("first", start + DurationNanos::new(100), None, None)
2448 .unwrap();
2449 clock
2450 .set_time_alert_ns("second", start + DurationNanos::new(200), None, None)
2451 .unwrap();
2452
2453 let events = clock.advance_time(start + DurationNanos::new(400), true);
2454 assert_eq!(events.len(), 3);
2455 assert_eq!(events[0].name, "first");
2456 assert_eq!(events[1].name, "second");
2457 assert_eq!(events[2].name, "third");
2458 }
2459
2460 #[rstest]
2461 fn test_large_interval_does_not_overflow(mut test_clock: TestClock) {
2462 let start = test_clock.timestamp_ns();
2463 let large_interval = DurationNanos::from_days(365);
2464
2465 test_clock
2466 .set_timer_ns(
2467 "large_interval",
2468 large_interval,
2469 Some(start),
2470 None,
2471 None,
2472 None,
2473 None,
2474 )
2475 .unwrap();
2476
2477 let events = test_clock.advance_time(start + large_interval, true);
2478 assert_eq!(events.len(), 1);
2479 assert_eq!(events[0].ts_event, start + large_interval);
2480 }
2481
2482 #[rstest]
2483 fn test_near_zero_interval_fires_correctly(mut test_clock: TestClock) {
2484 let start = test_clock.timestamp_ns();
2485
2486 test_clock
2487 .set_timer_ns(
2488 "tiny",
2489 DurationNanos::new(1),
2490 Some(start),
2491 None,
2492 None,
2493 None,
2494 None,
2495 )
2496 .unwrap();
2497
2498 let events = test_clock.advance_time(start + DurationNanos::new(10), true);
2499 assert_eq!(events.len(), 10);
2500
2501 for i in 1..events.len() {
2502 assert!(events[i].ts_event >= events[i - 1].ts_event);
2503 }
2504 }
2505
2506 #[rstest]
2507 fn test_repeated_advance_to_same_time_no_double_fire(mut test_clock: TestClock) {
2508 let fire_time = test_clock.timestamp_ns() + DurationNanos::new(1000);
2509
2510 test_clock
2511 .set_time_alert_ns("once", fire_time, None, None)
2512 .unwrap();
2513
2514 let events1 = test_clock.advance_time(fire_time, true);
2515 assert_eq!(events1.len(), 1);
2516
2517 let events2 = test_clock.advance_time(fire_time, true);
2518 assert!(events2.is_empty());
2519 }
2520
2521 #[rstest]
2522 fn test_advance_with_no_timers(mut test_clock: TestClock) {
2523 let start = test_clock.timestamp_ns();
2524
2525 let events = test_clock.advance_time(start + DurationNanos::new(1000), true);
2526 assert!(events.is_empty());
2527 assert_eq!(test_clock.timestamp_ns(), start + DurationNanos::new(1000));
2528 }
2529
2530 #[rstest]
2531 fn test_set_time_alert_rejects_unconvertible_datetime(mut test_clock: TestClock) {
2532 let pre_epoch = Timestamp::from_nanosecond(-1).unwrap();
2533
2534 let err = test_clock
2535 .set_time_alert("pre_epoch_alert", pre_epoch, None, None)
2536 .unwrap_err();
2537 assert!(
2538 err.to_string().contains("cannot be negative"),
2539 "unexpected error: {err}"
2540 );
2541
2542 let err = test_clock
2543 .set_time_alert("out_of_range_alert", Timestamp::MAX, None, None)
2544 .unwrap_err();
2545 assert!(
2546 err.to_string().contains("out of range"),
2547 "unexpected error: {err}"
2548 );
2549
2550 assert_eq!(test_clock.timer_count(), 0);
2551 }
2552
2553 #[rstest]
2554 fn test_set_timer_rejects_unconvertible_datetime(mut test_clock: TestClock) {
2555 let pre_epoch = Timestamp::from_nanosecond(-1).unwrap();
2556 let valid_start = test_clock.utc_now() + jiff::SignedDuration::from_secs(1);
2557
2558 let err = test_clock
2559 .set_timer(
2560 "pre_epoch_start",
2561 Duration::from_secs(1),
2562 Some(pre_epoch),
2563 None,
2564 None,
2565 None,
2566 None,
2567 )
2568 .unwrap_err();
2569 assert!(
2570 err.to_string().contains("cannot be negative"),
2571 "unexpected error: {err}"
2572 );
2573
2574 let err = test_clock
2575 .set_timer(
2576 "pre_epoch_stop",
2577 Duration::from_secs(1),
2578 Some(valid_start),
2579 Some(pre_epoch),
2580 None,
2581 None,
2582 None,
2583 )
2584 .unwrap_err();
2585 assert!(
2586 err.to_string().contains("cannot be negative"),
2587 "unexpected error: {err}"
2588 );
2589
2590 assert_eq!(test_clock.timer_count(), 0);
2591 }
2592
2593 #[rstest]
2594 fn test_set_timer_rejects_interval_exceeding_u64_nanos(mut test_clock: TestClock) {
2595 let interval = Duration::from_secs(u64::MAX / NANOSECONDS_IN_SECOND + 1);
2596
2597 let err = test_clock
2598 .set_timer("overflow", interval, None, None, None, None, None)
2599 .unwrap_err();
2600
2601 assert_eq!(err.to_string(), "Interval exceeds u64 nanoseconds");
2602 assert_eq!(test_clock.timer_count(), 0);
2603 }
2604
2605 #[rstest]
2606 fn test_set_timer_ns_rejects_unrepresentable_first_event_without_replacing_timer(
2607 mut test_clock: TestClock,
2608 ) {
2609 test_clock.set_time(UnixNanos::from(1));
2610 test_clock
2611 .set_timer_ns(
2612 "overflow",
2613 DurationNanos::new(1),
2614 None,
2615 None,
2616 None,
2617 None,
2618 None,
2619 )
2620 .unwrap();
2621
2622 let err = test_clock
2623 .set_timer_ns("overflow", DurationNanos::MAX, None, None, None, None, None)
2624 .unwrap_err();
2625
2626 assert_eq!(
2627 err.to_string(),
2628 "Timer 'overflow' first event time exceeds UnixNanos range"
2629 );
2630 assert_eq!(test_clock.timer_count(), 1);
2631 assert_eq!(
2632 test_clock.next_time_ns("overflow"),
2633 Some(UnixNanos::from(2))
2634 );
2635 }
2636
2637 #[rstest]
2638 fn test_clock_api_handlers_reject_invalid_time_inputs() {
2639 let calls = Arc::new(Mutex::new(Vec::new()));
2640 let calls_for_alert = Arc::clone(&calls);
2641 let calls_for_timer = Arc::clone(&calls);
2642
2643 let clock = ClockApi::from_handlers(
2644 || UnixNanos::from(1_700_000_000_000_000_000),
2645 move |name, _, _, _| {
2646 calls_for_alert.lock().push(name.to_string());
2647 Ok(())
2648 },
2649 move |name, _, _, _, _, _, _| {
2650 calls_for_timer.lock().push(name.to_string());
2651 Ok(())
2652 },
2653 Vec::new,
2654 || 0,
2655 |_| false,
2656 |_| None,
2657 |_| {},
2658 || {},
2659 );
2660
2661 let pre_epoch = Timestamp::from_nanosecond(-1).unwrap();
2662 clock
2663 .set_time_alert("alert", pre_epoch, None, None)
2664 .unwrap_err();
2665 clock
2666 .set_timer(
2667 "timer",
2668 Duration::from_secs(1),
2669 Some(pre_epoch),
2670 None,
2671 None,
2672 None,
2673 None,
2674 )
2675 .unwrap_err();
2676 let interval = Duration::from_secs(u64::MAX / NANOSECONDS_IN_SECOND + 1);
2677 let err = clock
2678 .set_timer("overflow", interval, None, None, None, None, None)
2679 .unwrap_err();
2680
2681 assert_eq!(err.to_string(), "Interval exceeds u64 nanoseconds");
2682 assert!(calls.lock().is_empty());
2683 }
2684
2685 #[rstest]
2686 fn test_clock_api_new_uses_native_backing(test_clock: TestClock) {
2687 let clock = RefCell::new(test_clock);
2688 let api = ClockApi::new(&clock);
2689
2690 api.set_timer_ns(
2691 "native-timer",
2692 DurationNanos::new(1_000),
2693 None,
2694 None,
2695 None,
2696 Some(true),
2697 Some(false),
2698 )
2699 .unwrap();
2700
2701 assert_eq!(api.timer_count(), 1);
2702 assert_eq!(api.timer_names(), vec!["native-timer".to_string()]);
2703 assert_eq!(
2704 api.next_time_ns("native-timer"),
2705 Some(UnixNanos::from(1_000))
2706 );
2707 }
2708
2709 #[rstest]
2710 fn test_clock_api_handlers_back_full_surface() {
2711 let alerts = Arc::new(Mutex::new(Vec::new()));
2712 let timers = Arc::new(Mutex::new(Vec::new()));
2713 let cancellations = Arc::new(Mutex::new(Vec::new()));
2714 let cancel_all = Arc::new(Mutex::new(false));
2715
2716 let alerts_for_handler = Arc::clone(&alerts);
2717 let timers_for_handler = Arc::clone(&timers);
2718 let cancellations_for_handler = Arc::clone(&cancellations);
2719 let cancel_all_for_handler = Arc::clone(&cancel_all);
2720
2721 let clock = ClockApi::from_handlers(
2722 || UnixNanos::from(1_700_000_000_123_456_789),
2723 move |name, alert_time_ns, _callback, allow_past| {
2724 alerts_for_handler
2725 .lock()
2726 .push((name.to_string(), alert_time_ns, allow_past));
2727 Ok(())
2728 },
2729 move |name,
2730 interval_ns,
2731 start_time_ns,
2732 stop_time_ns,
2733 _callback,
2734 allow_past,
2735 fire_immediately| {
2736 timers_for_handler.lock().push((
2737 name.to_string(),
2738 interval_ns,
2739 start_time_ns,
2740 stop_time_ns,
2741 allow_past,
2742 fire_immediately,
2743 ));
2744 Ok(())
2745 },
2746 || vec!["alpha".to_string(), "beta".to_string()],
2747 || 2,
2748 |name| name == "alpha",
2749 |name| (name == "alpha").then(|| UnixNanos::from(1_700_000_000_999_000_000)),
2750 move |name| {
2751 cancellations_for_handler.lock().push(name.to_string());
2752 },
2753 move || {
2754 *cancel_all_for_handler.lock() = true;
2755 },
2756 );
2757
2758 let alert_time = Timestamp::from_nanosecond(1_700_000_000_333_000_000).unwrap();
2759 let start_time = Timestamp::from_nanosecond(1_700_000_000_444_000_000).unwrap();
2760 let stop_time = Timestamp::from_nanosecond(1_700_000_001_444_000_000).unwrap();
2761 clock
2762 .set_time_alert("alert", alert_time, None, Some(false))
2763 .unwrap();
2764 clock
2765 .set_time_alert_ns(
2766 "alert-ns",
2767 UnixNanos::from(1_700_000_000_555_000_000),
2768 None,
2769 Some(true),
2770 )
2771 .unwrap();
2772 clock
2773 .set_timer(
2774 "timer",
2775 Duration::from_millis(250),
2776 Some(start_time),
2777 Some(stop_time),
2778 None,
2779 Some(true),
2780 Some(false),
2781 )
2782 .unwrap();
2783 clock
2784 .set_timer_ns(
2785 "timer-ns",
2786 DurationNanos::from_millis(500),
2787 Some(UnixNanos::from(1_700_000_000_666_000_000)),
2788 Some(UnixNanos::from(1_700_000_001_666_000_000)),
2789 None,
2790 Some(false),
2791 Some(true),
2792 )
2793 .unwrap();
2794 clock.cancel_timer("alpha");
2795 clock.cancel_timers();
2796
2797 assert_eq!(
2798 clock.timestamp_ns(),
2799 UnixNanos::from(1_700_000_000_123_456_789)
2800 );
2801 assert_eq!(clock.timestamp_us(), 1_700_000_000_123_456);
2802 assert_eq!(clock.timestamp_ms(), 1_700_000_000_123);
2803 assert_eq!(clock.timestamp(), 1_700_000_000.123_456_7);
2804 assert_eq!(
2805 clock.utc_now(),
2806 Timestamp::from_nanosecond(1_700_000_000_123_456_789).unwrap()
2807 );
2808 assert_eq!(clock.timer_names(), vec!["alpha", "beta"]);
2809 assert_eq!(clock.timer_count(), 2);
2810 assert!(clock.timer_exists("alpha"));
2811 assert!(!clock.timer_exists("gamma"));
2812 assert_eq!(
2813 clock.next_time_ns("alpha"),
2814 Some(UnixNanos::from(1_700_000_000_999_000_000))
2815 );
2816 assert_eq!(
2817 alerts.lock().as_slice(),
2818 &[
2819 (
2820 "alert".to_string(),
2821 UnixNanos::from(1_700_000_000_333_000_000),
2822 Some(false)
2823 ),
2824 (
2825 "alert-ns".to_string(),
2826 UnixNanos::from(1_700_000_000_555_000_000),
2827 Some(true)
2828 )
2829 ]
2830 );
2831 assert_eq!(
2832 timers.lock().as_slice(),
2833 &[
2834 (
2835 "timer".to_string(),
2836 DurationNanos::from_millis(250),
2837 Some(UnixNanos::from(1_700_000_000_444_000_000)),
2838 Some(UnixNanos::from(1_700_000_001_444_000_000)),
2839 Some(true),
2840 Some(false)
2841 ),
2842 (
2843 "timer-ns".to_string(),
2844 DurationNanos::from_millis(500),
2845 Some(UnixNanos::from(1_700_000_000_666_000_000)),
2846 Some(UnixNanos::from(1_700_000_001_666_000_000)),
2847 Some(false),
2848 Some(true)
2849 )
2850 ]
2851 );
2852 assert_eq!(cancellations.lock().as_slice(), &["alpha".to_string()]);
2853 assert!(*cancel_all.lock());
2854 }
2855
2856 proptest! {
2857 #[rstest]
2858 fn prop_test_clock_operations_match_reference(
2859 initial_time_ns in clock_time_strategy(),
2860 operations in prop::collection::vec(clock_operation_strategy(), 1..=50),
2861 ) {
2862 check_clock_operations(initial_time_ns, operations)?;
2863 }
2864
2865 #[rstest]
2866 fn prop_test_clock_max_time_alert(initial_time_ns in clock_time_strategy()) {
2867 check_clock_max_time_alert(initial_time_ns)?;
2868 }
2869 }
2870
2871 #[derive(Clone, Debug)]
2872 enum ClockOperation {
2873 Set {
2874 name_index: usize,
2875 interval_ns: u64,
2876 stop_after_ns: Option<u64>,
2877 fire_immediately: bool,
2878 },
2879 Cancel(usize),
2880 Advance {
2881 delta_ns: u64,
2882 set_time: bool,
2883 },
2884 }
2885
2886 #[derive(Clone, Debug)]
2887 struct TimerModel {
2888 interval: u64,
2889 next: u64,
2890 stop: Option<u64>,
2891 }
2892
2893 fn clock_operation_strategy() -> impl Strategy<Value = ClockOperation> {
2894 prop_oneof![
2895 5 => (
2896 0usize..CLOCK_TIMER_NAMES.len(),
2897 1u64..=15,
2898 prop::option::of(1u64..=60),
2899 prop::bool::ANY,
2900 )
2901 .prop_map(
2902 |(name_index, interval_ns, stop_after_ns, fire_immediately)| {
2903 ClockOperation::Set {
2904 name_index,
2905 interval_ns,
2906 stop_after_ns,
2907 fire_immediately,
2908 }
2909 },
2910 ),
2911 2 => (0usize..CLOCK_TIMER_NAMES.len()).prop_map(ClockOperation::Cancel),
2912 5 => (0u64..=30, prop::bool::ANY)
2913 .prop_map(|(delta_ns, set_time)| ClockOperation::Advance { delta_ns, set_time }),
2914 ]
2915 }
2916
2917 fn clock_time_strategy() -> impl Strategy<Value = u64> {
2918 prop_oneof![
2919 6 => 0u64..=u64::MAX - CLOCK_TIME_HEADROOM,
2920 2 => 0u64..=1_000_000,
2921 1 => Just(1_700_000_000_000_000_000),
2922 1 => Just(u64::MAX - CLOCK_TIME_HEADROOM),
2923 ]
2924 }
2925
2926 fn check_clock_operations(
2927 initial_time_ns: u64,
2928 operations: Vec<ClockOperation>,
2929 ) -> TestCaseResult {
2930 let mut clock = TestClock::new();
2931 clock.register_default_handler(TestCallback::default().into());
2932 clock.set_time(UnixNanos::from(initial_time_ns));
2933
2934 let mut time_ns = initial_time_ns;
2935 let mut timers = BTreeMap::new();
2936
2937 for operation in operations {
2938 match operation {
2939 ClockOperation::Set {
2940 name_index,
2941 interval_ns,
2942 stop_after_ns,
2943 fire_immediately,
2944 } => {
2945 let name = clock_timer_name(name_index);
2946 let stop_time_ns = stop_after_ns.map(|offset| time_ns + offset);
2947 clock
2948 .set_timer_ns(
2949 name.as_str(),
2950 DurationNanos::new(interval_ns),
2951 Some(UnixNanos::from(time_ns)),
2952 stop_time_ns.map(UnixNanos::from),
2953 None,
2954 None,
2955 Some(fire_immediately),
2956 )
2957 .expect("generated timer configuration should be valid");
2958 timers.insert(
2959 name,
2960 TimerModel {
2961 interval: interval_ns,
2962 next: if fire_immediately {
2963 time_ns
2964 } else {
2965 time_ns + interval_ns
2966 },
2967 stop: stop_time_ns,
2968 },
2969 );
2970 }
2971 ClockOperation::Cancel(name_index) => {
2972 let name = clock_timer_name(name_index);
2973 clock.cancel_timer(name.as_str());
2974 timers.remove(&name);
2975 }
2976 ClockOperation::Advance { delta_ns, set_time } => {
2977 let to_time_ns = time_ns + delta_ns;
2978 let actual: Vec<(u64, Ustr, u64)> = clock
2979 .advance_time(UnixNanos::from(to_time_ns), set_time)
2980 .into_iter()
2981 .map(|event| (event.ts_event.as_u64(), event.name, event.ts_init.as_u64()))
2982 .collect();
2983 let expected = advance_clock_timers(&mut timers, to_time_ns);
2984
2985 prop_assert_eq!(actual, expected);
2986
2987 if set_time {
2988 time_ns = to_time_ns;
2989 }
2990 }
2991 }
2992
2993 assert_clock_state(&clock, &timers, time_ns)?;
2994 }
2995
2996 Ok(())
2997 }
2998
2999 fn check_clock_max_time_alert(initial_time_ns: u64) -> TestCaseResult {
3000 let mut clock = TestClock::new();
3001 clock.register_default_handler(TestCallback::default().into());
3002 clock.set_time(UnixNanos::from(initial_time_ns));
3003 let name = Ustr::from("terminal-alert");
3004 clock
3005 .set_time_alert_ns(name.as_str(), UnixNanos::max(), None, None)
3006 .expect("maximum timestamp should be a valid time alert");
3007
3008 let events: Vec<(u64, Ustr, u64)> = clock
3009 .advance_time(UnixNanos::max(), true)
3010 .into_iter()
3011 .map(|event| (event.ts_event.as_u64(), event.name, event.ts_init.as_u64()))
3012 .collect();
3013
3014 prop_assert_eq!(events, vec![(u64::MAX, name, u64::MAX)]);
3015 prop_assert_eq!(clock.timestamp_ns(), UnixNanos::max());
3016 prop_assert_eq!(clock.timer_count(), 0);
3017 prop_assert!(clock.timer_names().is_empty());
3018 prop_assert!(!clock.timer_exists(&name));
3019 prop_assert_eq!(clock.next_time_ns(name.as_str()), None);
3020
3021 Ok(())
3022 }
3023
3024 fn advance_clock_timers(
3025 timers: &mut BTreeMap<Ustr, TimerModel>,
3026 to_time_ns: u64,
3027 ) -> Vec<(u64, Ustr, u64)> {
3028 let mut events = Vec::new();
3029
3030 timers.retain(|name, timer| {
3031 while timer.next <= to_time_ns {
3032 if timer
3033 .stop
3034 .is_some_and(|stop_time_ns| timer.next > stop_time_ns)
3035 {
3036 return false;
3037 }
3038
3039 let event_time_ns = timer.next;
3040 events.push((event_time_ns, *name, event_time_ns));
3041
3042 let Some(following_time_ns) = event_time_ns.checked_add(timer.interval) else {
3043 return false;
3044 };
3045
3046 timer.next = following_time_ns;
3047 if timer.stop == Some(event_time_ns) {
3048 return false;
3049 }
3050 }
3051
3052 true
3053 });
3054
3055 events.sort_by(|a, b| a.0.cmp(&b.0).then_with(|| a.1.cmp(&b.1)));
3056 events
3057 }
3058
3059 fn assert_clock_state(
3060 clock: &TestClock,
3061 timers: &BTreeMap<Ustr, TimerModel>,
3062 time_ns: u64,
3063 ) -> TestCaseResult {
3064 let expected_names: Vec<&str> = timers.keys().map(Ustr::as_str).collect();
3065
3066 prop_assert_eq!(clock.timestamp_ns(), UnixNanos::from(time_ns));
3067 prop_assert_eq!(clock.timer_count(), timers.len());
3068 prop_assert_eq!(clock.timer_names(), expected_names);
3069
3070 for name in CLOCK_TIMER_NAMES.map(Ustr::from) {
3071 let expected = timers.get(&name);
3072 prop_assert_eq!(clock.timer_exists(&name), expected.is_some());
3073 prop_assert_eq!(
3074 clock
3075 .next_time_ns(name.as_str())
3076 .map(|next_time_ns| next_time_ns.as_u64()),
3077 expected.map(|timer| timer.next),
3078 );
3079 }
3080
3081 Ok(())
3082 }
3083
3084 const CLOCK_TIMER_NAMES: [&str; 4] = ["timer-0", "timer-1", "timer-2", "timer-3"];
3085 const CLOCK_TIME_HEADROOM: u64 = 100_000;
3086
3087 fn clock_timer_name(index: usize) -> Ustr {
3088 Ustr::from(CLOCK_TIMER_NAMES[index])
3089 }
3090}