Skip to main content

nautilus_common/live/
timer.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live timer scheduling and callback dispatch.
17//!
18//! # Scheduling
19//!
20//! Runtime deadlines use the remaining duration to the nominal schedule and are computed before
21//! the timer task is spawned. A deadline already in the past starts immediately. Event timestamps
22//! retain the nominal schedule, and stop times are inclusive.
23//!
24//! # Task lifecycle
25//!
26//! Each [`LiveTimer::start`] creates fresh public schedule and task-state atomics. Because aborting
27//! a runtime task does not join it, task retirement linearizes restart, cancellation, and drop
28//! against event reservation. Retirement either prevents the old task from dispatching or observes
29//! the schedule advanced by its reserved event; fresh atomics keep that task from overwriting its
30//! replacement's schedule.
31//!
32//! # Callback dispatch
33//!
34//! Thread-safe callbacks cross the worker boundary as [`TimeEventMessage`] values. `RustLocal`
35//! callbacks remain in an owner-thread registry and cross the boundary only through tokens and
36//! leases. Senderless Python callbacks run inline and publish their following schedule after the
37//! callback returns.
38
39use 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/// A live timer for use with a `LiveClock`.
75///
76/// `LiveTimer` triggers events at specified intervals in a real-time environment,
77/// using Tokio's async runtime to handle scheduling and execution.
78///
79/// # Threading
80///
81/// The timer runs on the runtime thread that created it and dispatches events across threads as needed.
82#[derive(Debug)]
83pub struct LiveTimer {
84    /// The name of the timer.
85    pub name: Ustr,
86    /// The interval between timer events in nanoseconds.
87    pub interval_ns: NonZeroU64,
88    /// The start time of the timer in UNIX nanoseconds.
89    pub start_time_ns: UnixNanos,
90    /// The optional stop time of the timer in UNIX nanoseconds.
91    pub stop_time_ns: Option<UnixNanos>,
92    /// If the timer should fire immediately at start time.
93    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    /// Creates a new [`LiveTimer`] instance.
104    ///
105    /// # Panics
106    ///
107    /// Panics if:
108    /// - `name` is not a valid string.
109    /// - `fire_immediately` is false and `start_time_ns + interval_ns` overflows `UnixNanos`.
110    #[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    /// Returns the next time in UNIX nanoseconds when the timer will fire.
159    ///
160    /// Provides the scheduled time for the next event based on the current state of the timer.
161    #[must_use]
162    pub fn next_time_ns(&self) -> UnixNanos {
163        UnixNanos::from(self.next_time_ns.load(atomic::Ordering::SeqCst))
164    }
165
166    /// Returns whether the timer is expired.
167    ///
168    /// An expired timer will not trigger any further events.
169    /// A timer that has not been started is not expired.
170    #[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    /// Starts the timer.
180    ///
181    /// Time events will begin triggering at the specified intervals.
182    /// The generated events are handled by the provided callback function.
183    ///
184    /// Starting a timer whose task is still active aborts that task first
185    /// (restart semantics); a previously fired event that is already queued
186    /// still dispatches.
187    ///
188    /// # Panics
189    ///
190    /// Panics if using a Rust callback (`Rust` or `RustLocal`) without a `TimeEventSender`.
191    #[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        // Close the old token before registering its replacement;
212        // any lease acquired before retirement remains valid.
213        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                // A cancel or a final stop-time fire closed the token; a
221                // restart needs a fresh registration or acquire() would
222                // return None on every fire.
223                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        // Get current time
241        let clock = get_atomic_clock_realtime();
242        let now_ns = clock.get_time_ns();
243
244        // Check if the timer's alert time is in the past and adjust if needed
245        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        // Floor the next time to the nearest microsecond which is within the timers accuracy
260        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                // Never fire an event scheduled past the stop time. The event's
277                // `ts_event` is the scheduled `next_time_ns`, so the bound is
278                // enforced on the scheduled time (matching `TestTimer`), not on
279                // the wall-clock read used only for `ts_init`.
280                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; // Timer expired before this event
291                }
292
293                // `timer.tick` is cancellation safe, if the cancel branch completes
294                // first then no tick has been consumed (no event was ready).
295                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                // The event scheduled exactly at the stop time fires (inclusive
301                // boundary), then the timer expires.
302                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                // Reserve this fire and its following schedule together. A
306                // restart either observes the advanced schedule or retires
307                // the task before it can dispatch. Registered callbacks
308                // acquire their lease first so token closure cannot suppress
309                // an event whose schedule already advanced.
310                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; // Timer expired at the stop boundary
355                }
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    /// Cancels the timer.
369    ///
370    /// The timer will not generate a final event.
371    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                // The firing section contains only atomic schedule publication
516                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    /// The callback is retained owner-side so a restarted timer (cancel or
529    /// natural expiry closed the token) can register a fresh token.
530    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, // time_event_sender
776        );
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        // The timer retains an owner-side clone for restart; only after it
873        // drops must no registry copy remain.
874        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        // Second start with the first task still active: the old task is
971        // aborted and the token stays live for the new one.
972        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        // The timer retains an owner-side clone for restart; only after it
1014        // drops must no registry copy remain.
1015        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}