1use std::{
40 num::NonZeroU64,
41 sync::{
42 Arc,
43 atomic::{self, AtomicU8, AtomicU64},
44 },
45};
46
47use nautilus_core::{
48 UUID4, UnixNanos,
49 correctness::{FAILED, check_valid_string_utf8},
50 datetime::floor_to_nearest_microsecond,
51 time::get_atomic_clock_realtime,
52};
53use ustr::Ustr;
54
55use super::dst::{
56 self,
57 task::JoinHandle,
58 time::{Duration, Instant},
59};
60#[cfg(not(all(feature = "simulation", madsim)))]
61use super::runtime::get_runtime;
62use crate::{
63 runner::{
64 TimeEventCallbackLease, TimeEventCallbackToken, TimeEventMessage, TimeEventMessageFactory,
65 TimeEventSender, register_time_event_callback,
66 },
67 timer::{TimeEvent, TimeEventCallback, Timer},
68};
69
70const TASK_ACTIVE: u8 = 0;
71const TASK_FIRING: u8 = 1;
72const TASK_RETIRED: u8 = 2;
73
74#[derive(Debug)]
83pub struct LiveTimer {
84 pub name: Ustr,
86 pub interval_ns: NonZeroU64,
88 pub start_time_ns: UnixNanos,
90 pub stop_time_ns: Option<UnixNanos>,
92 pub fire_immediately: bool,
94 next_time_ns: Arc<AtomicU64>,
95 callback: OwnerCallback,
96 task_handle: Option<JoinHandle<()>>,
97 task_state: Option<Arc<TimerTaskState>>,
98 canceled: bool,
99 sender: Option<Arc<dyn TimeEventSender>>,
100}
101
102impl LiveTimer {
103 #[must_use]
111 pub fn new(
112 name: Ustr,
113 interval_ns: NonZeroU64,
114 start_time_ns: UnixNanos,
115 stop_time_ns: Option<UnixNanos>,
116 callback: TimeEventCallback,
117 fire_immediately: bool,
118 sender: Option<Arc<dyn TimeEventSender>>,
119 ) -> Self {
120 check_valid_string_utf8(name, stringify!(name)).expect(FAILED);
121
122 let next_time_ns = if fire_immediately {
123 start_time_ns.as_u64()
124 } else {
125 (start_time_ns + interval_ns.get()).as_u64()
126 };
127
128 log::trace!("Creating timer '{name}'");
129
130 let owner_callback = if sender.is_some() {
131 if callback.is_local() {
132 OwnerCallback::Registered {
133 token: register_time_event_callback(callback.clone()),
134 callback,
135 }
136 } else {
137 OwnerCallback::Direct(TimeEventMessageFactory::new(&callback))
138 }
139 } else {
140 OwnerCallback::Senderless(callback)
141 };
142
143 Self {
144 name,
145 interval_ns,
146 start_time_ns,
147 stop_time_ns,
148 fire_immediately,
149 next_time_ns: Arc::new(AtomicU64::new(next_time_ns)),
150 callback: owner_callback,
151 task_handle: None,
152 task_state: None,
153 canceled: false,
154 sender,
155 }
156 }
157
158 #[must_use]
162 pub fn next_time_ns(&self) -> UnixNanos {
163 UnixNanos::from(self.next_time_ns.load(atomic::Ordering::SeqCst))
164 }
165
166 #[must_use]
171 pub fn is_expired(&self) -> bool {
172 self.canceled
173 || self
174 .task_handle
175 .as_ref()
176 .is_some_and(JoinHandle::is_finished)
177 }
178
179 #[allow(unused_variables)]
192 pub fn start(&mut self) {
193 if let OwnerCallback::Senderless(callback) = &self.callback {
194 match callback {
195 #[cfg(feature = "python")]
196 TimeEventCallback::Python(_) => {}
197 TimeEventCallback::Rust(_) | TimeEventCallback::RustLocal(_) => {
198 panic!("timer event sender was unset for Rust callback system");
199 }
200 }
201 }
202
203 let event_name = self.name;
204 let stop_time_ns = self.stop_time_ns;
205 let interval_ns = self.interval_ns.get();
206
207 let mut observed_next = self
208 .retire_task()
209 .unwrap_or_else(|| self.next_time_ns.load(atomic::Ordering::SeqCst));
210
211 if let Some(handle) = self.task_handle.take() {
214 self.close_registered_callback();
215 handle.abort();
216 }
217
218 let worker_dispatch = match &mut self.callback {
219 OwnerCallback::Registered { token, callback } => {
220 if token.is_closed() {
224 *token = register_time_event_callback(callback.clone());
225 }
226 WorkerDispatch::Registered(token.clone())
227 }
228 OwnerCallback::Direct(factory) => WorkerDispatch::Direct(factory.clone()),
229 OwnerCallback::Senderless(callback) => match callback {
230 #[cfg(feature = "python")]
231 TimeEventCallback::Python(callback) => {
232 WorkerDispatch::SenderlessPython(callback.clone())
233 }
234 TimeEventCallback::Rust(_) | TimeEventCallback::RustLocal(_) => {
235 unreachable!("senderless Rust callback rejected at start")
236 }
237 },
238 };
239
240 let clock = get_atomic_clock_realtime();
242 let now_ns = clock.get_time_ns();
243
244 let now_raw = now_ns.as_u64();
246
247 if should_adjust_past_due_time(observed_next, now_ns, stop_time_ns) {
248 if observed_next < now_raw {
249 let original = UnixNanos::from(observed_next);
250 log::warn!(
251 "Timer '{event_name}' alert time {} was in the past, adjusted to current time for immediate fire",
252 original.to_rfc3339(),
253 );
254 }
255
256 observed_next = now_raw;
257 }
258
259 let mut next_time_ns = normalize_start_time_ns(observed_next, now_ns, stop_time_ns);
261 let next_time_atomic = Arc::new(AtomicU64::new(next_time_ns.as_u64()));
262 let task_state = Arc::new(TimerTaskState::new(next_time_ns.as_u64()));
263 self.next_time_ns = next_time_atomic.clone();
264 self.task_state = Some(task_state.clone());
265
266 let sender = self.sender.clone();
267 let now_ns = clock.get_time_ns();
268 let start = Instant::now() + timer_start_delay(next_time_ns, now_ns);
269
270 let task = async move {
271 let clock = get_atomic_clock_realtime();
272
273 let mut timer = dst::time::interval_at(start, Duration::from_nanos(interval_ns));
274
275 loop {
276 if !should_fire_scheduled_time(next_time_ns, stop_time_ns) {
281 if let (Some(sender), WorkerDispatch::Registered(token)) =
282 (sender.as_ref(), &worker_dispatch)
283 && let Some(lease) = token.acquire()
284 {
285 token.close();
286 let now_ns = clock.get_time_ns();
287 let event = TimeEvent::new(event_name, UUID4::new(), next_time_ns, now_ns);
288 sender.send(TimeEventMessage::cleanup(event, lease));
289 }
290 break; }
292
293 timer.tick().await;
296 let now_ns = clock.get_time_ns();
297
298 let event = TimeEvent::new(event_name, UUID4::new(), next_time_ns, now_ns);
299
300 let expires_after_fire = expires_after_scheduled_time(next_time_ns, stop_time_ns);
303 let following_next_time_ns = next_time_ns + interval_ns;
304
305 let registered_lease = if let WorkerDispatch::Registered(token) = &worker_dispatch {
311 match task_state.reserve_registered_fire(token, following_next_time_ns.as_u64())
312 {
313 Some(lease) => Some(lease),
314 None => break,
315 }
316 } else {
317 if !task_state.reserve_fire(following_next_time_ns.as_u64()) {
318 break;
319 }
320 None
321 };
322
323 if sender.is_some() {
324 next_time_atomic
325 .store(following_next_time_ns.as_u64(), atomic::Ordering::SeqCst);
326 }
327
328 match (&sender, &worker_dispatch) {
329 (Some(sender), WorkerDispatch::Direct(factory)) => {
330 sender.send(factory.message(event));
331 }
332 (Some(sender), WorkerDispatch::Registered(token)) => {
333 let lease =
334 registered_lease.expect("registered callback lease was not acquired");
335
336 if expires_after_fire {
337 token.close();
338 }
339 sender.send(TimeEventMessage::registered(event, lease));
340 }
341 #[cfg(feature = "python")]
342 (None, WorkerDispatch::SenderlessPython(callback)) => callback.call(event),
343 _ => unreachable!("timer callback dispatch did not match its sender"),
344 }
345
346 if sender.is_none() {
347 next_time_atomic
348 .store(following_next_time_ns.as_u64(), atomic::Ordering::SeqCst);
349 }
350
351 next_time_ns = following_next_time_ns;
352
353 if expires_after_fire {
354 break; }
356 }
357 };
358
359 #[cfg(all(feature = "simulation", madsim))]
360 let handle = dst::task::spawn(task);
361 #[cfg(not(all(feature = "simulation", madsim)))]
362 let handle = get_runtime().spawn(task);
363
364 self.task_handle = Some(handle);
365 self.canceled = false;
366 }
367
368 pub fn cancel(&mut self) {
372 log::trace!("Cancel timer '{}'", self.name);
373
374 self.close_registered_callback();
375
376 self.retire_task();
377
378 if let Some(handle) = self.task_handle.take() {
379 handle.abort();
380 }
381 self.canceled = true;
382 }
383
384 fn close_registered_callback(&self) {
385 if let OwnerCallback::Registered { token, .. } = &self.callback {
386 token.close();
387 }
388 }
389
390 fn retire_task(&mut self) -> Option<u64> {
391 let task_state = self.task_state.take()?;
392 let next_time_ns = task_state.retire();
393 self.next_time_ns
394 .store(next_time_ns, atomic::Ordering::SeqCst);
395 Some(next_time_ns)
396 }
397}
398
399impl Timer for LiveTimer {
400 fn is_expired(&self) -> bool {
401 Self::is_expired(self)
402 }
403
404 fn cancel(&mut self) {
405 Self::cancel(self);
406 }
407}
408
409impl Drop for LiveTimer {
410 fn drop(&mut self) {
411 self.close_registered_callback();
412 self.retire_task();
413
414 if let Some(handle) = self.task_handle.take() {
415 handle.abort();
416 }
417 }
418}
419
420fn should_fire_scheduled_time(next_time_ns: UnixNanos, stop_time_ns: Option<UnixNanos>) -> bool {
421 stop_time_ns.is_none_or(|stop_time_ns| next_time_ns <= stop_time_ns)
422}
423
424fn expires_after_scheduled_time(next_time_ns: UnixNanos, stop_time_ns: Option<UnixNanos>) -> bool {
425 stop_time_ns == Some(next_time_ns)
426}
427
428fn is_stop_boundary(next_time_ns: u64, stop_time_ns: Option<UnixNanos>) -> bool {
429 stop_time_ns == Some(UnixNanos::from(next_time_ns))
430}
431
432fn should_adjust_past_due_time(
433 observed_next: u64,
434 now_ns: UnixNanos,
435 stop_time_ns: Option<UnixNanos>,
436) -> bool {
437 observed_next <= now_ns.as_u64() && !is_stop_boundary(observed_next, stop_time_ns)
438}
439
440fn normalize_start_time_ns(
441 observed_next: u64,
442 now_ns: UnixNanos,
443 stop_time_ns: Option<UnixNanos>,
444) -> UnixNanos {
445 if is_stop_boundary(observed_next, stop_time_ns) {
446 return UnixNanos::from(observed_next);
447 }
448
449 let now_raw = now_ns.as_u64();
450 let start_time_ns = if observed_next <= now_raw {
451 now_raw
452 } else {
453 observed_next
454 };
455
456 UnixNanos::from(floor_to_nearest_microsecond(start_time_ns))
457}
458
459fn timer_start_delay(next_time_ns: UnixNanos, now_ns: UnixNanos) -> Duration {
460 Duration::from_nanos(next_time_ns.saturating_sub(now_ns.as_u64()))
461}
462
463#[derive(Debug)]
464struct TimerTaskState {
465 status: AtomicU8,
466 next_time_ns: AtomicU64,
467}
468
469impl TimerTaskState {
470 fn new(next_time_ns: u64) -> Self {
471 Self {
472 status: AtomicU8::new(TASK_ACTIVE),
473 next_time_ns: AtomicU64::new(next_time_ns),
474 }
475 }
476
477 fn reserve_registered_fire(
478 &self,
479 token: &TimeEventCallbackToken,
480 following_next_time_ns: u64,
481 ) -> Option<TimeEventCallbackLease> {
482 let lease = token.acquire()?;
483 self.reserve_fire(following_next_time_ns).then_some(lease)
484 }
485
486 fn reserve_fire(&self, following_next_time_ns: u64) -> bool {
487 if self
488 .status
489 .compare_exchange(
490 TASK_ACTIVE,
491 TASK_FIRING,
492 atomic::Ordering::SeqCst,
493 atomic::Ordering::SeqCst,
494 )
495 .is_err()
496 {
497 return false;
498 }
499
500 self.next_time_ns
501 .store(following_next_time_ns, atomic::Ordering::SeqCst);
502 self.status.store(TASK_ACTIVE, atomic::Ordering::SeqCst);
503 true
504 }
505
506 fn retire(&self) -> u64 {
507 loop {
508 match self.status.compare_exchange(
509 TASK_ACTIVE,
510 TASK_RETIRED,
511 atomic::Ordering::SeqCst,
512 atomic::Ordering::SeqCst,
513 ) {
514 Ok(_) | Err(TASK_RETIRED) => break,
515 Err(TASK_FIRING) => std::hint::spin_loop(),
517 Err(status) => unreachable!("invalid timer task state {status}"),
518 }
519 }
520
521 self.next_time_ns.load(atomic::Ordering::SeqCst)
522 }
523}
524
525#[derive(Debug)]
526enum OwnerCallback {
527 Direct(TimeEventMessageFactory),
528 Registered {
531 token: TimeEventCallbackToken,
532 callback: TimeEventCallback,
533 },
534 Senderless(TimeEventCallback),
535}
536
537#[derive(Clone, Debug)]
538enum WorkerDispatch {
539 Direct(TimeEventMessageFactory),
540 Registered(TimeEventCallbackToken),
541 #[cfg(feature = "python")]
542 SenderlessPython(Arc<crate::timer::PythonTimeEventCallback>),
543}
544
545#[cfg(test)]
546mod tests {
547 #[cfg(not(all(feature = "simulation", madsim)))]
548 use std::rc::Rc;
549 #[cfg(all(feature = "python", not(all(feature = "simulation", madsim))))]
550 use std::sync::{OnceLock, atomic::AtomicU64};
551 use std::{
552 num::NonZeroU64,
553 sync::{
554 Arc,
555 atomic::{AtomicUsize, Ordering},
556 },
557 };
558 #[cfg(any(feature = "python", not(all(feature = "simulation", madsim))))]
559 use std::{sync::mpsc, time::Duration as StdDuration};
560
561 use nautilus_core::{UnixNanos, time::get_atomic_clock_realtime};
562 #[cfg(not(all(feature = "simulation", madsim)))]
563 use parking_lot::Mutex;
564 #[cfg(feature = "python")]
565 use pyo3::{
566 Python,
567 types::{PyAnyMethods, PyList, PyListMethods},
568 };
569 use rstest::*;
570 use ustr::Ustr;
571
572 use super::LiveTimer;
573 #[cfg(not(all(feature = "simulation", madsim)))]
574 use crate::runner::register_time_event_callback;
575 #[cfg(not(all(feature = "simulation", madsim)))]
576 use crate::testing::wait_until;
577 use crate::{
578 runner::{TimeEventMessage, TimeEventSender},
579 timer::TimeEventCallback,
580 };
581
582 #[cfg(any(feature = "python", not(all(feature = "simulation", madsim))))]
583 #[derive(Debug)]
584 struct ChannelSender {
585 tx: mpsc::Sender<TimeEventMessage>,
586 }
587
588 #[cfg(any(feature = "python", not(all(feature = "simulation", madsim))))]
589 impl TimeEventSender for ChannelSender {
590 fn send(&self, message: TimeEventMessage) {
591 self.tx.send(message).expect("message should send");
592 }
593 }
594
595 #[cfg(not(all(feature = "simulation", madsim)))]
596 #[derive(Debug)]
597 struct PausingChannelSender {
598 tx: mpsc::Sender<TimeEventMessage>,
599 release_rx: Mutex<mpsc::Receiver<()>>,
600 }
601
602 #[cfg(not(all(feature = "simulation", madsim)))]
603 impl TimeEventSender for PausingChannelSender {
604 fn send(&self, message: TimeEventMessage) {
605 self.tx.send(message).expect("message should send");
606 self.release_rx
607 .lock()
608 .recv()
609 .expect("timer send should release");
610 }
611 }
612
613 #[cfg(all(feature = "simulation", madsim))]
614 #[derive(Debug)]
615 struct CountingSender {
616 count: Arc<AtomicUsize>,
617 }
618
619 #[cfg(all(feature = "simulation", madsim))]
620 impl TimeEventSender for CountingSender {
621 fn send(&self, _message: TimeEventMessage) {
622 self.count.fetch_add(1, Ordering::Relaxed);
623 }
624 }
625
626 #[rstest]
627 #[case::unbounded(100, None, true, false)]
628 #[case::past_stop(110, Some(100), false, false)]
629 #[case::before_stop(90, Some(100), true, false)]
630 #[case::at_stop(100, Some(100), true, true)]
631 fn test_live_timer_stop_bound(
632 #[case] next_time_ns: u64,
633 #[case] stop_time_ns: Option<u64>,
634 #[case] should_fire: bool,
635 #[case] expires: bool,
636 ) {
637 let next_time_ns = UnixNanos::from(next_time_ns);
638 let stop_time_ns = stop_time_ns.map(UnixNanos::from);
639
640 assert_eq!(
641 super::should_fire_scheduled_time(next_time_ns, stop_time_ns),
642 should_fire
643 );
644 assert_eq!(
645 super::expires_after_scheduled_time(next_time_ns, stop_time_ns),
646 expires
647 );
648 }
649
650 #[rstest]
651 #[case::stop_boundary(100, 110, 100, false)]
652 #[case::before_stop(90, 110, 120, true)]
653 fn test_live_timer_past_due_adjustment(
654 #[case] observed_next: u64,
655 #[case] now: u64,
656 #[case] stop_time_ns: u64,
657 #[case] expected: bool,
658 ) {
659 assert_eq!(
660 super::should_adjust_past_due_time(
661 observed_next,
662 UnixNanos::from(now),
663 Some(UnixNanos::from(stop_time_ns)),
664 ),
665 expected
666 );
667 }
668
669 #[rstest]
670 #[case::past_due(1_234_567, 2_345_678, None, 2_345_000)]
671 #[case::future(3_456_789, 2_345_678, None, 3_456_000)]
672 #[case::stop_boundary(1_234_567, 2_345_678, Some(1_234_567), 1_234_567)]
673 fn test_live_timer_start_time_normalization(
674 #[case] observed_next: u64,
675 #[case] now: u64,
676 #[case] stop_time_ns: Option<u64>,
677 #[case] expected: u64,
678 ) {
679 assert_eq!(
680 super::normalize_start_time_ns(
681 observed_next,
682 UnixNanos::from(now),
683 stop_time_ns.map(UnixNanos::from),
684 ),
685 UnixNanos::from(expected)
686 );
687 }
688
689 #[rstest]
690 #[case::full(12_000_000, 10_000_000, 2_000_000)]
691 #[case::sub_millisecond(10_500_000, 10_000_000, 500_000)]
692 fn test_live_timer_start_delay(
693 #[case] next_time_ns: u64,
694 #[case] now: u64,
695 #[case] expected_ns: u64,
696 ) {
697 assert_eq!(
698 super::timer_start_delay(UnixNanos::from(next_time_ns), UnixNanos::from(now)),
699 tokio::time::Duration::from_nanos(expected_ns)
700 );
701 }
702
703 #[rstest]
704 fn test_timer_task_retirement_prevents_a_late_fire() {
705 let state = super::TimerTaskState::new(100);
706
707 let restart_time_ns = state.retire();
708 let reserved = state.reserve_fire(200);
709
710 assert_eq!(restart_time_ns, 100);
711 assert!(!reserved);
712 assert_eq!(state.next_time_ns.load(Ordering::SeqCst), 100);
713 }
714
715 #[rstest]
716 fn test_timer_task_retirement_preserves_a_reserved_fire() {
717 let state = super::TimerTaskState::new(100);
718
719 let reserved = state.reserve_fire(200);
720 let restart_time_ns = state.retire();
721
722 assert!(reserved);
723 assert_eq!(restart_time_ns, 200);
724 assert_eq!(state.next_time_ns.load(Ordering::SeqCst), 200);
725 }
726
727 #[cfg(not(all(feature = "simulation", madsim)))]
728 #[rstest]
729 fn test_closed_registered_callback_does_not_reserve_fire() {
730 let state = super::TimerTaskState::new(100);
731 let callback = TimeEventCallback::RustLocal(Rc::new(|_| {}));
732 let token = register_time_event_callback(callback);
733 token.close();
734
735 let lease = state.reserve_registered_fire(&token, 200);
736 let restart_time_ns = state.retire();
737
738 assert!(lease.is_none());
739 assert_eq!(restart_time_ns, 100);
740 assert_eq!(state.next_time_ns.load(Ordering::SeqCst), 100);
741 }
742
743 #[rstest]
744 #[case::immediate(true, 100)]
745 #[case::after_interval(false, 1_100)]
746 fn test_live_timer_fire_immediately(
747 #[case] fire_immediately: bool,
748 #[case] expected_next_time_ns: u64,
749 ) {
750 let timer = LiveTimer::new(
751 Ustr::from("TEST_TIMER"),
752 NonZeroU64::new(1000).unwrap(),
753 UnixNanos::from(100),
754 None,
755 TimeEventCallback::from(|_| {}),
756 fire_immediately,
757 None,
758 );
759
760 assert_eq!(timer.fire_immediately, fire_immediately);
761 assert_eq!(timer.next_time_ns(), UnixNanos::from(expected_next_time_ns));
762 }
763
764 #[rstest]
765 #[should_panic(expected = "timer event sender was unset for Rust callback system")]
766 fn test_live_timer_start_panics_on_senderless_rust_callback() {
767 let now = get_atomic_clock_realtime().get_time_ns();
768 let mut timer = LiveTimer::new(
769 Ustr::from("SENDERLESS_RUST"),
770 NonZeroU64::new(1_000_000).unwrap(),
771 now,
772 None,
773 TimeEventCallback::from(|_| {}),
774 false,
775 None, );
777
778 timer.start();
779 }
780
781 #[cfg(not(all(feature = "simulation", madsim)))]
782 #[rstest]
783 fn test_live_timer_uses_global_runtime() {
784 let (tx, rx) = mpsc::channel();
785 let sender = Arc::new(ChannelSender { tx });
786 let now = get_atomic_clock_realtime().get_time_ns();
787 let mut timer = LiveTimer::new(
788 Ustr::from("LIVE_TIMER"),
789 NonZeroU64::new(1_000_000).unwrap(),
790 now,
791 Some(now),
792 TimeEventCallback::from(|_| {}),
793 true,
794 Some(sender),
795 );
796
797 timer.start();
798 let message = rx
799 .recv_timeout(StdDuration::from_secs(1))
800 .expect("timer message should arrive on the global runtime");
801 wait_until(|| timer.is_expired(), StdDuration::from_secs(1));
802
803 assert_eq!(message.event().ts_event, now);
804 assert!(timer.is_expired());
805 }
806
807 #[cfg(not(all(feature = "simulation", madsim)))]
808 #[rstest]
809 fn test_live_timer_dispatches_rust_local_callback_on_owner_thread() {
810 let (tx, rx) = mpsc::channel();
811 let sender = Arc::new(ChannelSender { tx });
812 let count = Rc::new(std::cell::Cell::new(0));
813 let callback_count = count.clone();
814 let callback: Rc<dyn Fn(crate::timer::TimeEvent)> =
815 Rc::new(move |_| callback_count.set(callback_count.get() + 1));
816 let now = get_atomic_clock_realtime().get_time_ns();
817 let mut timer = LiveTimer::new(
818 Ustr::from("LOCAL_TIMER"),
819 NonZeroU64::new(1_000_000).unwrap(),
820 now,
821 Some(now),
822 TimeEventCallback::RustLocal(callback),
823 true,
824 Some(sender),
825 );
826
827 timer.start();
828 let message = rx
829 .recv_timeout(StdDuration::from_secs(1))
830 .expect("registered timer message should arrive");
831
832 assert!(message.dispatch());
833 assert_eq!(count.get(), 1);
834 }
835
836 #[cfg(not(all(feature = "simulation", madsim)))]
837 #[rstest]
838 fn test_live_timer_cancel_preserves_queued_rust_local_callback_lease() {
839 let (tx, rx) = mpsc::channel();
840 let (release_tx, release_rx) = mpsc::channel();
841 let sender = Arc::new(PausingChannelSender {
842 tx,
843 release_rx: Mutex::new(release_rx),
844 });
845 let count = Rc::new(std::cell::Cell::new(0));
846 let callback_count = count.clone();
847 let callback: Rc<dyn Fn(crate::timer::TimeEvent)> =
848 Rc::new(move |_| callback_count.set(callback_count.get() + 1));
849 let callback_weak = Rc::downgrade(&callback);
850 let now = get_atomic_clock_realtime().get_time_ns();
851 let mut timer = LiveTimer::new(
852 Ustr::from("CANCEL_QUEUED"),
853 NonZeroU64::new(1_000_000).unwrap(),
854 now,
855 None,
856 TimeEventCallback::RustLocal(callback),
857 true,
858 Some(sender),
859 );
860
861 timer.start();
862 let message = rx
863 .recv_timeout(StdDuration::from_secs(1))
864 .expect("registered timer message should arrive");
865 timer.cancel();
866 release_tx.send(()).expect("timer send should release");
867
868 assert!(callback_weak.upgrade().is_some());
869 assert!(message.dispatch());
870 assert_eq!(count.get(), 1);
871
872 drop(timer);
875 assert!(callback_weak.upgrade().is_none());
876 }
877
878 #[cfg(not(all(feature = "simulation", madsim)))]
879 #[rstest]
880 fn test_live_timer_cancel_preserves_queued_direct_callback() {
881 let (tx, rx) = mpsc::channel();
882 let sender = Arc::new(ChannelSender { tx });
883 let count = Arc::new(AtomicUsize::new(0));
884 let callback_count = count.clone();
885 let now = get_atomic_clock_realtime().get_time_ns();
886 let mut timer = LiveTimer::new(
887 Ustr::from("CANCEL_QUEUED_DIRECT"),
888 NonZeroU64::new(1_000_000).unwrap(),
889 now,
890 None,
891 TimeEventCallback::from(move |_| {
892 callback_count.fetch_add(1, Ordering::Relaxed);
893 }),
894 true,
895 Some(sender),
896 );
897
898 timer.start();
899 let message = rx
900 .recv_timeout(StdDuration::from_secs(1))
901 .expect("direct timer message should arrive");
902 timer.cancel();
903
904 assert!(message.dispatch());
905 assert_eq!(count.load(Ordering::Relaxed), 1);
906 }
907
908 #[cfg(not(all(feature = "simulation", madsim)))]
909 #[rstest]
910 fn test_live_timer_restart_after_cancel_re_registers_rust_local_callback() {
911 let (tx, rx) = mpsc::channel();
912 let sender = Arc::new(ChannelSender { tx });
913 let count = Rc::new(std::cell::Cell::new(0));
914 let callback_count = count.clone();
915 let callback: Rc<dyn Fn(crate::timer::TimeEvent)> =
916 Rc::new(move |_| callback_count.set(callback_count.get() + 1));
917 let now = get_atomic_clock_realtime().get_time_ns();
918 let mut timer = LiveTimer::new(
919 Ustr::from("RESTART_TIMER"),
920 NonZeroU64::new(1_000_000).unwrap(),
921 now,
922 None,
923 TimeEventCallback::RustLocal(callback),
924 true,
925 Some(sender),
926 );
927
928 timer.start();
929 let first = rx
930 .recv_timeout(StdDuration::from_secs(1))
931 .expect("first registered timer message should arrive");
932 timer.cancel();
933 assert!(first.dispatch());
934
935 timer.start();
936 let second = rx
937 .recv_timeout(StdDuration::from_secs(1))
938 .expect("restarted timer should re-register and fire");
939 timer.cancel();
940
941 assert!(second.dispatch());
942 assert_eq!(count.get(), 2);
943 }
944
945 #[cfg(not(all(feature = "simulation", madsim)))]
946 #[rstest]
947 fn test_live_timer_start_while_active_restarts_and_keeps_dispatching() {
948 let (tx, rx) = mpsc::channel();
949 let sender = Arc::new(ChannelSender { tx });
950 let count = Rc::new(std::cell::Cell::new(0));
951 let callback_count = count.clone();
952 let callback: Rc<dyn Fn(crate::timer::TimeEvent)> =
953 Rc::new(move |_| callback_count.set(callback_count.get() + 1));
954 let now = get_atomic_clock_realtime().get_time_ns();
955 let mut timer = LiveTimer::new(
956 Ustr::from("DOUBLE_START_TIMER"),
957 NonZeroU64::new(1_000_000).unwrap(),
958 now,
959 None,
960 TimeEventCallback::RustLocal(callback),
961 true,
962 Some(sender),
963 );
964
965 timer.start();
966 let first = rx
967 .recv_timeout(StdDuration::from_secs(1))
968 .expect("first task message should arrive");
969
970 timer.start();
973 let second = rx
974 .recv_timeout(StdDuration::from_secs(1))
975 .expect("restarted task should keep dispatching");
976 timer.cancel();
977
978 assert!(first.dispatch());
979 assert!(second.dispatch());
980 assert_eq!(count.get(), 2);
981 }
982
983 #[cfg(not(all(feature = "simulation", madsim)))]
984 #[rstest]
985 fn test_live_timer_stop_before_first_fire_sends_cleanup_message() {
986 let (tx, rx) = mpsc::channel();
987 let sender = Arc::new(ChannelSender { tx });
988 let count = Rc::new(std::cell::Cell::new(0));
989 let callback_count = count.clone();
990 let callback: Rc<dyn Fn(crate::timer::TimeEvent)> =
991 Rc::new(move |_| callback_count.set(callback_count.get() + 1));
992 let callback_weak = Rc::downgrade(&callback);
993 let now = get_atomic_clock_realtime().get_time_ns();
994 let mut timer = LiveTimer::new(
995 Ustr::from("CLEANUP_TIMER"),
996 NonZeroU64::new(1_000_000).unwrap(),
997 now,
998 Some(now),
999 TimeEventCallback::RustLocal(callback),
1000 false,
1001 Some(sender),
1002 );
1003
1004 timer.start();
1005 let cleanup = rx
1006 .recv_timeout(StdDuration::from_secs(1))
1007 .expect("cleanup message should arrive");
1008
1009 assert!(callback_weak.upgrade().is_some());
1010 assert!(!cleanup.dispatch());
1011 assert_eq!(count.get(), 0);
1012
1013 drop(timer);
1016 assert!(callback_weak.upgrade().is_none());
1017 }
1018
1019 #[cfg(all(feature = "simulation", madsim))]
1020 #[madsim::test]
1021 async fn test_live_timer_uses_dst_runtime() {
1022 let count = Arc::new(AtomicUsize::new(0));
1023 let sender = Arc::new(CountingSender {
1024 count: count.clone(),
1025 });
1026 let now = get_atomic_clock_realtime().get_time_ns();
1027 let mut timer = LiveTimer::new(
1028 Ustr::from("DST_TIMER"),
1029 NonZeroU64::new(1_000_000).unwrap(),
1030 now,
1031 Some(now),
1032 TimeEventCallback::from(|_| {}),
1033 true,
1034 Some(sender),
1035 );
1036
1037 timer.start();
1038 crate::live::dst::time::sleep(crate::live::dst::time::Duration::from_millis(2)).await;
1039 crate::live::dst::task::yield_now().await;
1040
1041 assert_eq!(count.load(Ordering::Relaxed), 1);
1042 assert!(timer.is_expired());
1043 }
1044
1045 #[cfg(feature = "python")]
1046 #[rstest]
1047 fn test_live_timer_with_sender_defers_python_callback_to_handler() {
1048 Python::initialize();
1049
1050 Python::attach(|py| {
1051 let py_list = PyList::empty(py);
1052 let py_append = py_list
1053 .getattr("append")
1054 .expect("append should exist")
1055 .unbind();
1056 let callback = TimeEventCallback::from(py_append);
1057 let (tx, rx) = mpsc::channel();
1058 let sender = Arc::new(ChannelSender { tx });
1059 let now = get_atomic_clock_realtime().get_time_ns();
1060
1061 let mut timer = LiveTimer::new(
1062 Ustr::from("PY_TIMER"),
1063 NonZeroU64::new(1_000_000).unwrap(),
1064 now,
1065 None,
1066 callback,
1067 true,
1068 Some(sender),
1069 );
1070
1071 timer.start();
1072 let message = rx
1073 .recv_timeout(StdDuration::from_secs(1))
1074 .expect("timer message should arrive without acquiring the GIL on the worker");
1075 timer.cancel();
1076
1077 assert_eq!(py_list.len(), 0);
1078 assert!(message.dispatch());
1079 assert_eq!(py_list.len(), 1);
1080 });
1081 }
1082
1083 #[cfg(all(feature = "python", not(all(feature = "simulation", madsim))))]
1084 #[rstest]
1085 fn test_senderless_callback_observes_current_schedule() {
1086 Python::initialize();
1087
1088 Python::attach(|py| {
1089 let schedule = Arc::new(OnceLock::<Arc<AtomicU64>>::new());
1090 let callback_schedule = schedule.clone();
1091 let (tx, rx) = mpsc::channel();
1092 let callback = pyo3::types::PyCFunction::new_closure(
1093 py,
1094 None,
1095 None,
1096 move |_args: &pyo3::Bound<'_, pyo3::types::PyTuple>,
1097 _kwargs: Option<&pyo3::Bound<'_, pyo3::types::PyDict>>|
1098 -> pyo3::PyResult<()> {
1099 let next_time_ns = callback_schedule
1100 .get()
1101 .expect("timer schedule should be available")
1102 .load(Ordering::SeqCst);
1103 tx.send(next_time_ns)
1104 .expect("observed schedule should send");
1105 Ok(())
1106 },
1107 )
1108 .expect("callback should create")
1109 .into_any()
1110 .unbind();
1111 let now = get_atomic_clock_realtime().get_time_ns();
1112 let interval_ns = 10_000_000;
1113 let mut timer = LiveTimer::new(
1114 Ustr::from("SENDERLESS_SCHEDULE"),
1115 NonZeroU64::new(interval_ns).unwrap(),
1116 now,
1117 Some(now),
1118 TimeEventCallback::from(callback),
1119 true,
1120 None,
1121 );
1122
1123 timer.start();
1124 let expected_time_ns = timer.next_time_ns().as_u64();
1125 schedule
1126 .set(timer.next_time_ns.clone())
1127 .expect("timer schedule should set once");
1128 let observed_time_ns = py
1129 .detach(move || rx.recv_timeout(StdDuration::from_secs(1)))
1130 .expect("senderless callback should observe the schedule");
1131 wait_until(
1132 || timer.next_time_ns().as_u64() == expected_time_ns + interval_ns,
1133 StdDuration::from_secs(1),
1134 );
1135 timer.cancel();
1136
1137 assert_eq!(observed_time_ns, expected_time_ns);
1138 });
1139 }
1140}