Skip to main content

ntex_util/time/
wheel.rs

1//! Hierarchical timer wheel backing ntex timers.
2//!
3//! The design follows the Linux kernel timer wheel (`kernel/time/timer.c`).
4//!
5//! # Clock
6//!
7//! The wheel counts time in *units* of 16 milliseconds (`1 << UNITS`). The
8//! wheel clock `elapsed` is the unit of the last processed expiry, and
9//! `elapsed_time` is the [`Instant`] that corresponds to it. Timer deadlines
10//! are converted to units relative to this pair.
11//!
12//! # Levels
13//!
14//! The wheel has `LVL_DEPTH` (8) levels of `LVL_SIZE` (64) buckets. Each level
15//! is 8 times coarser than the previous one, a timer is placed on the first
16//! level that can represent its delay:
17//!
18//! | Level | Granularity | Range            |
19//! |-------|-------------|------------------|
20//! | 0     | 16 ms       | 0 .. ~1 s        |
21//! | 1     | 128 ms      | ~1 s .. ~8 s     |
22//! | 2     | ~1 s        | ~8 s .. ~64 s    |
23//! | 3     | ~8 s        | ~64 s .. ~8.6 m  |
24//! | 4     | ~65 s       | ~8.6 m .. ~69 m  |
25//! | 5     | ~8.7 m      | ~69 m .. ~9.2 h  |
26//! | 6     | ~70 m       | ~9.2 h .. ~3 d   |
27//! | 7     | ~9.3 h      | ~3 d .. ~24.5 d  |
28//!
29//! Longer delays are clamped to the capacity of the wheel. The expiry is
30//! rounded up to the granularity of the level and the delay is measured from
31//! [`Instant::now()`], so a timer never fires early but may fire up to one
32//! granularity late. Timers are not cascaded to finer
33//! levels, which keeps insertion and removal `O(1)`.
34//!
35//! A bitmap per level tracks occupied buckets, the next expiry is found by
36//! scanning the bitmaps instead of the buckets.
37//!
38//! # Drivers
39//!
40//! Two tasks are spawned lazily on the current thread:
41//!
42//! * [`TimerDriver`] sleeps until the next occupied bucket expires and wakes
43//!   the timers stored in it. The clock advances to the scheduled time of the
44//!   bucket, not to the time the driver woke up, so a late wakeup does not
45//!   delay later timers. Buckets that are overdue are processed at once.
46//! * [`LowresTimerDriver`] invalidates the cached [`now()`] and
47//!   [`system_time()`] values every 150 milliseconds.
48//!
49//! Dropping the timer driver, i.e. when the runtime stops, stops the wheel and
50//! marks all timers as elapsed. Dropping the lowres driver invalidates the
51//! cached time.
52//!
53//! A runtime does not always drop its pending tasks when it stops, the timer
54//! is also reset once the arbiter storage of the stopped system is cleared,
55//! so that the next runtime on the thread starts its own drivers. Drivers of
56//! an older generation exit without touching the state.
57use std::num::NonZeroUsize;
58use std::time::{Duration, Instant, SystemTime};
59use std::{cell::Cell, cmp, future::Future, pin::Pin, rc::Rc, task, task::Poll};
60
61use futures_timer::Delay;
62use slab::Slab;
63
64use crate::task::LocalWaker;
65
66/// Resolution of the wheel clock, a unit is `1 << UNITS` milliseconds.
67const UNITS: u64 = 4;
68
69/// Each level is `LVL_CLK_DIV` times coarser than the previous one.
70const LVL_CLK_SHIFT: u64 = 3;
71const LVL_CLK_DIV: u64 = 1 << LVL_CLK_SHIFT;
72const LVL_CLK_MASK: u64 = LVL_CLK_DIV - 1;
73
74/// Number of buckets per level.
75const LVL_BITS: u64 = 6;
76const LVL_SIZE: u64 = 1 << LVL_BITS;
77const LVL_MASK: u64 = LVL_SIZE - 1;
78
79/// Number of levels.
80const LVL_DEPTH: u64 = 8;
81
82/// Total number of buckets.
83const WHEEL_SIZE: usize = (LVL_SIZE * LVL_DEPTH) as usize;
84
85/// Delays at or above the cutoff are clamped to `WHEEL_TIMEOUT_MAX`.
86const WHEEL_TIMEOUT_CUTOFF: u64 = lvl_start(LVL_DEPTH);
87const WHEEL_TIMEOUT_MAX: u64 = WHEEL_TIMEOUT_CUTOFF - lvl_gran(LVL_DEPTH - 1);
88
89/// Refresh interval of the cached time.
90const LOWRES_RESOLUTION: Duration = Duration::from_millis(150);
91
92/// Shift of the level clock relative to the wheel clock.
93const fn lvl_shift(lvl: u64) -> u64 {
94    lvl * LVL_CLK_SHIFT
95}
96
97/// Granularity of a level in units.
98const fn lvl_gran(lvl: u64) -> u64 {
99    1 << lvl_shift(lvl)
100}
101
102/// Smallest delay in units that is stored on level `lvl`, `lvl` must be at
103/// least 1.
104const fn lvl_start(lvl: u64) -> u64 {
105    (LVL_SIZE - 1) << ((lvl - 1) * LVL_CLK_SHIFT)
106}
107
108const fn to_units(millis: u64) -> u64 {
109    millis >> UNITS
110}
111
112const fn to_millis(units: u64) -> u64 {
113    units << UNITS
114}
115
116const fn as_millis(dur: Duration) -> u64 {
117    dur.as_secs() * 1_000 + (dur.subsec_millis() as u64)
118}
119
120/// Returns a cached approximation of the current instant.
121///
122/// The cached value is refreshed at roughly 150 millisecond intervals.
123#[inline]
124pub fn now() -> Instant {
125    TIMER.with(Timer::now)
126}
127
128/// Returns a cached approximation of the current system time.
129///
130/// The cached value is refreshed at roughly 150 millisecond intervals.
131#[inline]
132pub fn system_time() -> SystemTime {
133    TIMER.with(Timer::system_time)
134}
135
136#[derive(Debug)]
137/// Handle to a timer registered with ntex's local timer wheel.
138///
139/// Dropping the handle cancels the timer. A handle may be reset and reused
140/// after it has elapsed.
141pub struct TimerHandle(NonZeroUsize);
142
143impl TimerHandle {
144    /// Registers a timer that elapses after `millis`.
145    pub fn new(millis: u64) -> Self {
146        TIMER.with(|t| t.add_timer(millis))
147    }
148
149    /// Restarts the timer with a new delay in milliseconds.
150    pub fn reset(&self, millis: u64) {
151        TIMER.with(|t| t.update_timer(self.0.get(), millis));
152    }
153
154    /// Completes the timer immediately and wakes its registered task.
155    pub fn elapse(&self) {
156        TIMER.with(|t| t.remove_timer(self.0.get()));
157    }
158
159    /// Returns `true` if this timer has elapsed.
160    pub fn is_elapsed(&self) -> bool {
161        TIMER.with(|t| t.with_wheel(|w| w.timers[self.0.get()].bucket.is_none()))
162    }
163
164    /// Polls until this timer has elapsed.
165    pub fn poll_elapsed(&self, cx: &mut task::Context<'_>) -> Poll<()> {
166        TIMER.with(|t| {
167            t.with_wheel(|w| {
168                let entry = &w.timers[self.0.get()];
169                if entry.bucket.is_none() {
170                    Poll::Ready(())
171                } else {
172                    entry.task.register(cx.waker());
173                    Poll::Pending
174                }
175            })
176        })
177    }
178}
179
180impl Drop for TimerHandle {
181    fn drop(&mut self) {
182        // the wheel is already destroyed if the handle is dropped by
183        // another thread-local destructor
184        let _ = TIMER.try_with(|t| {
185            t.with_wheel(|w| {
186                w.unlink(self.0.get());
187                w.timers.remove(self.0.get());
188            });
189        });
190    }
191}
192
193bitflags::bitflags! {
194    #[derive(Copy, Clone, Debug, Eq, PartialEq)]
195    struct Flags: u8 {
196        /// The timer driver task is spawned.
197        const DRIVER_STARTED = 0b0000_0001;
198        /// The cached time is populated and its refresh sleep is armed.
199        const LOWRES_TIMER   = 0b0000_1000;
200        /// The lowres driver task is spawned.
201        const LOWRES_DRIVER  = 0b0001_0000;
202        /// A timer was scheduled, `now()` and `system_time()` cache the time.
203        const RUNNING        = 0b0010_0000;
204    }
205}
206
207thread_local! {
208    static TIMER: Rc<Timer> = Rc::new(Timer::new());
209}
210
211/// Per-thread timer state, shared by the timer handles and the driver tasks.
212struct Timer {
213    /// Wheel clock in units, the expiry that was processed last.
214    elapsed: Cell<u64>,
215    /// Instant that corresponds to `elapsed`, the scheduled time of the last
216    /// processed expiry. Set lazily when the wheel is idle.
217    elapsed_time: Cell<Option<Instant>>,
218    /// Expiry of the earliest occupied bucket, `u64::MAX` if the wheel is empty.
219    next_expiry: Cell<u64>,
220    flags: Cell<Flags>,
221    /// Incremented on reset, drivers of an older generation are stale.
222    generation: Cell<u64>,
223    driver: LocalWaker,
224    lowres_time: Cell<Option<Instant>>,
225    lowres_stime: Cell<Option<SystemTime>>,
226    lowres_driver: LocalWaker,
227    /// Taken out of the cell for the duration of an operation.
228    wheel: Cell<Option<Box<Wheel>>>,
229}
230
231/// Timer storage.
232struct Wheel {
233    /// Timer entries indexed by handle, slot 0 is reserved so handles are
234    /// non-zero.
235    timers: Slab<TimerEntry>,
236    /// Handles of the timers stored in each bucket, `lvl * LVL_SIZE + offset`.
237    buckets: Box<[Slab<usize>]>,
238    /// Occupied buckets of each level, bit `n` is set if bucket `n` has timers.
239    occupied: [u64; LVL_DEPTH as usize],
240}
241
242#[derive(Debug)]
243struct TimerEntry {
244    /// Bucket index, `None` once the timer has elapsed.
245    bucket: Option<u16>,
246    /// Key of the entry in the bucket.
247    bucket_entry: usize,
248    task: LocalWaker,
249}
250
251impl Timer {
252    fn new() -> Self {
253        let mut timers = Slab::default();
254        timers.insert(TimerEntry {
255            bucket: None,
256            bucket_entry: 0,
257            task: LocalWaker::new(),
258        });
259
260        Timer {
261            elapsed: Cell::new(0),
262            elapsed_time: Cell::new(None),
263            next_expiry: Cell::new(u64::MAX),
264            flags: Cell::new(Flags::empty()),
265            generation: Cell::new(0),
266            driver: LocalWaker::new(),
267            lowres_time: Cell::new(None),
268            lowres_stime: Cell::new(None),
269            lowres_driver: LocalWaker::new(),
270            wheel: Cell::new(Some(Box::new(Wheel {
271                timers,
272                buckets: (0..WHEEL_SIZE).map(|_| Slab::new()).collect(),
273                occupied: [0; LVL_DEPTH as usize],
274            }))),
275        }
276    }
277
278    fn with_wheel<F, R>(&self, f: F) -> R
279    where
280        F: FnOnce(&mut Wheel) -> R,
281    {
282        let mut wheel = self.wheel.take().unwrap();
283        let result = f(&mut wheel);
284        self.wheel.set(Some(wheel));
285        result
286    }
287
288    fn insert_flags(&self, flags: Flags) -> Flags {
289        let mut f = self.flags.get();
290        f.insert(flags);
291        self.flags.set(f);
292        f
293    }
294
295    /// Returns the cached instant, the cache is populated only while the wheel
296    /// is running.
297    fn now(self: &Rc<Self>) -> Instant {
298        if let Some(cur) = self.lowres_time.get() {
299            cur
300        } else {
301            let now = Instant::now();
302            if self.flags.get().contains(Flags::RUNNING) {
303                self.lowres_time.set(Some(now));
304                self.refresh_lowres();
305            }
306            now
307        }
308    }
309
310    fn system_time(self: &Rc<Self>) -> SystemTime {
311        if let Some(cur) = self.lowres_stime.get() {
312            cur
313        } else {
314            let now = SystemTime::now();
315            if self.flags.get().contains(Flags::RUNNING) {
316                self.lowres_stime.set(Some(now));
317                self.refresh_lowres();
318            }
319            now
320        }
321    }
322
323    /// Arms the invalidation of the cached time.
324    fn refresh_lowres(self: &Rc<Self>) {
325        if self.flags.get().contains(Flags::LOWRES_DRIVER) {
326            self.lowres_driver.wake();
327        } else {
328            LowresTimerDriver::start(self);
329        }
330    }
331
332    /// Instant that corresponds to the wheel clock.
333    fn elapsed_time(&self) -> Instant {
334        if let Some(elapsed_time) = self.elapsed_time.get() {
335            elapsed_time
336        } else {
337            let elapsed_time = Instant::now();
338            self.elapsed_time.set(Some(elapsed_time));
339            elapsed_time
340        }
341    }
342
343    fn add_timer(self: &Rc<Self>, millis: u64) -> TimerHandle {
344        let no = if millis == 0 {
345            // elapsed immediately
346            self.with_wheel(|w| {
347                w.timers.insert(TimerEntry {
348                    bucket: None,
349                    bucket_entry: 0,
350                    task: LocalWaker::new(),
351                })
352            })
353        } else {
354            let (idx, expiry) = self.calc_bucket(millis);
355            let no = self.with_wheel(|w| {
356                let no = w.timers.insert(TimerEntry {
357                    bucket: None,
358                    bucket_entry: 0,
359                    task: LocalWaker::new(),
360                });
361                w.link(no, idx);
362                no
363            });
364            self.update_next_expiry(expiry);
365            no
366        };
367
368        // slot 0 is reserved in `Timer::new()`
369        TimerHandle(NonZeroUsize::new(no).unwrap())
370    }
371
372    /// Moves the timer to a new bucket, a zero delay elapses it and wakes its
373    /// task.
374    fn update_timer(self: &Rc<Self>, hnd: usize, millis: u64) {
375        if millis == 0 {
376            self.remove_timer(hnd);
377        } else {
378            let (idx, expiry) = self.calc_bucket(millis);
379            self.with_wheel(|w| w.relink(hnd, idx));
380            self.update_next_expiry(expiry);
381        }
382    }
383
384    /// Elapses the timer and wakes its task.
385    fn remove_timer(&self, hnd: usize) {
386        self.with_wheel(|w| {
387            if w.unlink(hnd) {
388                w.timers[hnd].task.wake();
389            }
390        });
391    }
392
393    /// Returns the bucket index and the bucket expiry for a delay starting now.
394    fn calc_bucket(self: &Rc<Self>, millis: u64) -> (usize, u64) {
395        self.insert_flags(Flags::RUNNING);
396
397        // The delay is measured from the wheel clock. The cached time is not
398        // used, it goes stale while the thread is blocked and the timer would
399        // fire early by its age
400        let since = Instant::now().saturating_duration_since(self.elapsed_time());
401        let delta = to_units(as_millis(since) + millis);
402        self.calc_wheel_index(self.elapsed.get().wrapping_add(delta), delta)
403    }
404
405    /// Selects the level for a timer that expires at `expires`, `delta` units
406    /// from now.
407    fn calc_wheel_index(&self, expires: u64, delta: u64) -> (usize, u64) {
408        for lvl in 0..LVL_DEPTH {
409            if delta < lvl_start(lvl + 1) {
410                return Self::calc_index(expires, lvl);
411            }
412        }
413        // expire larger delays at the capacity limit of the wheel
414        Self::calc_index(
415            self.elapsed.get().wrapping_add(WHEEL_TIMEOUT_MAX),
416            LVL_DEPTH - 1,
417        )
418    }
419
420    /// Returns the bucket index and the bucket expiry on level `lvl`.
421    fn calc_index(expires: u64, lvl: u64) -> (usize, u64) {
422        // The timer must not fire early. Early expiry can happen because the
423        // timer is armed at the edge of a tick, or because the expiry is
424        // truncated to the level granularity, round up to prevent it.
425        let expires = (expires + lvl_gran(lvl)) >> lvl_shift(lvl);
426        (
427            (lvl * LVL_SIZE + (expires & LVL_MASK)) as usize,
428            expires << lvl_shift(lvl),
429        )
430    }
431
432    /// Wakes the driver if the new bucket expires before the current deadline.
433    fn update_next_expiry(self: &Rc<Self>, expiry: u64) {
434        if expiry < self.next_expiry.get() {
435            self.next_expiry.set(expiry);
436            if self.flags.get().contains(Flags::DRIVER_STARTED) {
437                self.driver.wake();
438            } else {
439                TimerDriver::start(self);
440            }
441        }
442    }
443
444    /// Instant at which the bucket expiring at `expiry` is due.
445    fn expiry_time(&self, expiry: u64) -> Instant {
446        self.elapsed_time()
447            + Duration::from_millis(to_millis(expiry.saturating_sub(self.elapsed.get())))
448    }
449
450    /// Returns the expiry of the earliest occupied bucket.
451    fn next_pending_bucket(&self, wheel: &Wheel) -> Option<u64> {
452        let mut clk = self.elapsed.get();
453        let mut next = u64::MAX;
454
455        for lvl in 0..LVL_DEPTH {
456            let lvl_clk = clk & LVL_CLK_MASK;
457            let occupied = wheel.occupied[lvl as usize];
458
459            if occupied != 0 {
460                // distance to the next occupied bucket, wrapping around the level
461                let pos = u64::from(
462                    occupied
463                        .rotate_right((clk & LVL_MASK) as u32)
464                        .trailing_zeros(),
465                );
466                next = cmp::min(next, (clk + pos) << lvl_shift(lvl));
467
468                // The next level is reached once the clock of this level wraps
469                // to a multiple of `LVL_CLK_DIV`, an earlier bucket here cannot
470                // be preceded by one of the next level.
471                if pos <= ((LVL_CLK_DIV - lvl_clk) & LVL_CLK_MASK) {
472                    break;
473                }
474            }
475
476            // Clock of the next level. A partially elapsed tick of this level
477            // is rounded up, as the next level bucket it belongs to was already
478            // processed.
479            clk >>= LVL_CLK_SHIFT;
480            clk += u64::from(lvl_clk != 0);
481        }
482
483        if next < u64::MAX { Some(next) } else { None }
484    }
485
486    /// Removes `flags`, and `RUNNING` so the cached time is not populated and
487    /// no driver is spawned after the runtime stopped.
488    fn remove_flags(&self, flags: Flags) {
489        let mut f = self.flags.get();
490        f.remove(flags | Flags::RUNNING);
491        self.flags.set(f);
492    }
493
494    /// Invalidates the cached time, called when the lowres driver is dropped.
495    fn stop_lowres(&self) {
496        self.remove_flags(Flags::LOWRES_DRIVER | Flags::LOWRES_TIMER);
497        self.lowres_time.set(None);
498        self.lowres_stime.set(None);
499    }
500
501    /// Stops the drivers of the current generation, called when the system
502    /// that runs them has stopped.
503    fn reset(&self) {
504        self.generation.set(self.generation.get().wrapping_add(1));
505        self.stop_wheel();
506        self.stop_lowres();
507        self.driver.take();
508        self.lowres_driver.take();
509    }
510
511    /// Marks all timers as elapsed and resets the wheel, called when the timer
512    /// driver is dropped. Tasks are not woken, the runtime is stopping.
513    fn stop_wheel(&self) {
514        self.remove_flags(Flags::DRIVER_STARTED);
515
516        // the wheel is in use if a driver is dropped from a timer operation
517        if let Some(mut wheel) = self.wheel.take() {
518            let Wheel {
519                timers,
520                buckets,
521                occupied,
522            } = &mut *wheel;
523            for bucket in buckets.iter_mut() {
524                for no in bucket.drain() {
525                    timers[no].bucket = None;
526                }
527            }
528            *occupied = [0; LVL_DEPTH as usize];
529
530            self.next_expiry.set(u64::MAX);
531            self.elapsed.set(0);
532            self.elapsed_time.set(None);
533            self.wheel.set(Some(wheel));
534        }
535    }
536}
537
538impl Wheel {
539    /// Level and occupied bit of a bucket.
540    fn bucket_bit(idx: usize) -> (usize, u64) {
541        (idx / LVL_SIZE as usize, 1 << (idx % LVL_SIZE as usize))
542    }
543
544    /// Stores the timer in bucket `idx`.
545    fn link(&mut self, hnd: usize, idx: usize) {
546        let entry = &mut self.timers[hnd];
547        entry.bucket = Some(idx as u16);
548        entry.bucket_entry = self.buckets[idx].insert(hnd);
549
550        let (lvl, bit) = Self::bucket_bit(idx);
551        self.occupied[lvl] |= bit;
552    }
553
554    /// Removes the timer from its bucket and marks it as elapsed.
555    ///
556    /// Returns `false` if the timer has already elapsed.
557    fn unlink(&mut self, hnd: usize) -> bool {
558        let entry = &mut self.timers[hnd];
559        if let Some(idx) = entry.bucket.take() {
560            let idx = idx as usize;
561            let bucket = &mut self.buckets[idx];
562            bucket.remove(entry.bucket_entry);
563            if bucket.is_empty() {
564                let (lvl, bit) = Self::bucket_bit(idx);
565                self.occupied[lvl] &= !bit;
566            }
567            true
568        } else {
569            false
570        }
571    }
572
573    /// Moves the timer to bucket `idx`.
574    fn relink(&mut self, hnd: usize, idx: usize) {
575        if self.timers[hnd].bucket != Some(idx as u16) {
576            self.unlink(hnd);
577            self.link(hnd, idx);
578        }
579    }
580
581    /// Wakes the timers of the buckets that expire at `clk`.
582    fn execute_expired_timers(&mut self, mut clk: u64) {
583        for lvl in 0..LVL_DEPTH {
584            let idx = ((clk & LVL_MASK) + lvl * LVL_SIZE) as usize;
585            let bucket = &mut self.buckets[idx];
586            if !bucket.is_empty() {
587                let (lvl, bit) = Self::bucket_bit(idx);
588                self.occupied[lvl] &= !bit;
589                for no in bucket.drain() {
590                    let entry = &mut self.timers[no];
591                    entry.bucket = None;
592                    entry.task.wake();
593                }
594            }
595
596            // The next level expires only when the clock is a multiple of its
597            // granularity
598            if (clk & LVL_CLK_MASK) != 0 {
599                break;
600            }
601            clk >>= LVL_CLK_SHIFT;
602        }
603    }
604}
605
606/// Resets the timer when dropped with the arbiter storage of a stopped system.
607#[derive(Default)]
608struct TimerReset;
609
610impl TimerReset {
611    fn register() {
612        ntex_rt::with_item::<TimerReset, _, _>(|_| ());
613    }
614}
615
616impl Drop for TimerReset {
617    fn drop(&mut self) {
618        // the timer is already destroyed if the thread is exiting
619        let _ = TIMER.try_with(|t| t.reset());
620    }
621}
622
623/// Task that sleeps until the next bucket expires and wakes its timers.
624struct TimerDriver {
625    timer: Rc<Timer>,
626    generation: u64,
627    sleep: Delay,
628    /// Deadline the sleep is armed for.
629    armed: Option<Instant>,
630}
631
632impl TimerDriver {
633    fn start(timer: &Rc<Timer>) {
634        timer.insert_flags(Flags::DRIVER_STARTED);
635        TimerReset::register();
636
637        let deadline = timer.expiry_time(timer.next_expiry.get());
638        crate::spawn(TimerDriver {
639            timer: timer.clone(),
640            generation: timer.generation.get(),
641            sleep: Delay::new(deadline.saturating_duration_since(Instant::now())),
642            armed: Some(deadline),
643        });
644
645        // start lowres driver
646        timer.refresh_lowres();
647    }
648}
649
650impl Drop for TimerDriver {
651    fn drop(&mut self) {
652        if self.timer.generation.get() == self.generation {
653            self.timer.stop_wheel();
654        }
655    }
656}
657
658impl Future for TimerDriver {
659    type Output = ();
660
661    fn poll(mut self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
662        let this = &mut *self;
663        let timer = &this.timer;
664        if timer.generation.get() != this.generation {
665            return Poll::Ready(());
666        }
667        timer.driver.register(cx.waker());
668
669        let now = Instant::now();
670        timer.lowres_time.set(Some(now));
671        if !timer.flags.get().contains(Flags::LOWRES_TIMER) {
672            // arm the invalidation, otherwise the cached time stays stale
673            // until the next wakeup of the driver
674            timer.refresh_lowres();
675        }
676
677        loop {
678            let expiry = timer.next_expiry.get();
679            if expiry == u64::MAX {
680                // the wheel is empty, a new timer wakes the driver
681                return Poll::Pending;
682            }
683
684            let deadline = timer.expiry_time(expiry);
685            if deadline > now {
686                if this.armed != Some(deadline) {
687                    this.armed = Some(deadline);
688                    this.sleep.reset(deadline.saturating_duration_since(now));
689                }
690                if Pin::new(&mut this.sleep).poll(cx).is_pending() {
691                    return Poll::Pending;
692                }
693                if deadline > now {
694                    // the sleep fired before the deadline, re-arm it
695                    this.armed = None;
696                    continue;
697                }
698            }
699
700            // Advance the clock to the scheduled time of the bucket, a late
701            // wakeup must not shift the timers that expire later
702            timer.elapsed.set(expiry);
703            timer.elapsed_time.set(Some(deadline));
704
705            let next = timer.with_wheel(|w| {
706                w.execute_expired_timers(expiry);
707                timer.next_pending_bucket(w)
708            });
709
710            if let Some(next) = next {
711                timer.next_expiry.set(next);
712            } else {
713                timer.next_expiry.set(u64::MAX);
714                timer.elapsed_time.set(None);
715            }
716        }
717    }
718}
719
720/// Task that invalidates the cached time every `LOWRES_RESOLUTION`.
721struct LowresTimerDriver {
722    timer: Rc<Timer>,
723    generation: u64,
724    sleep: Delay,
725}
726
727impl LowresTimerDriver {
728    fn start(timer: &Rc<Timer>) {
729        timer.insert_flags(Flags::LOWRES_DRIVER);
730        TimerReset::register();
731
732        crate::spawn(LowresTimerDriver {
733            timer: timer.clone(),
734            generation: timer.generation.get(),
735            sleep: Delay::new(LOWRES_RESOLUTION),
736        });
737    }
738}
739
740impl Drop for LowresTimerDriver {
741    fn drop(&mut self) {
742        if self.timer.generation.get() == self.generation {
743            self.timer.stop_lowres();
744        }
745    }
746}
747
748impl Future for LowresTimerDriver {
749    type Output = ();
750
751    fn poll(mut self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
752        let this = &mut *self;
753        let timer = &this.timer;
754        if timer.generation.get() != this.generation {
755            return Poll::Ready(());
756        }
757        timer.lowres_driver.register(cx.waker());
758
759        // the cache was populated, invalidate it after `LOWRES_RESOLUTION`
760        let mut flags = timer.flags.get();
761        if !flags.contains(Flags::LOWRES_TIMER) {
762            flags.insert(Flags::LOWRES_TIMER);
763            timer.flags.set(flags);
764            this.sleep.reset(LOWRES_RESOLUTION);
765        }
766
767        if Pin::new(&mut this.sleep).poll(cx).is_ready() {
768            timer.lowres_time.set(None);
769            timer.lowres_stime.set(None);
770            flags.remove(Flags::LOWRES_TIMER);
771            timer.flags.set(flags);
772        }
773        Poll::Pending
774    }
775}
776
777#[cfg(test)]
778mod tests {
779    use super::*;
780    use crate::time::{Millis, interval, sleep};
781
782    /// A handle dropped by a thread-local destructor after the wheel is
783    /// destroyed must not panic.
784    #[test]
785    fn test_drop_handle_after_wheel_destroyed() {
786        use std::cell::RefCell;
787
788        thread_local! {
789            static HOLDER: RefCell<Option<TimerHandle>> = const { RefCell::new(None) };
790        }
791
792        let res = std::thread::spawn(|| {
793            // register the holder destructor first, the wheel is destroyed
794            // before it, destructors run in reverse registration order
795            HOLDER.with(|h| h.borrow_mut().take());
796            ntex::rt::System::build()
797                .build(ntex::rt::DefaultRuntime)
798                .block_on(async {
799                    let hnd = TimerHandle::new(1000);
800                    HOLDER.with(|h| *h.borrow_mut() = Some(hnd));
801                });
802        })
803        .join();
804        assert!(res.is_ok());
805    }
806
807    /// The drivers of a stopped system may never be dropped, the next system
808    /// on the same thread starts its own drivers.
809    #[test]
810    fn test_timer_in_next_system() {
811        let res = std::thread::spawn(|| {
812            for _ in 0..3 {
813                let start = Instant::now();
814                ntex::rt::System::build()
815                    .build(ntex::rt::DefaultRuntime)
816                    .block_on(async {
817                        sleep(Millis(10)).await;
818                        let _ = now();
819                        // pending timer of the stopped system
820                        let _hnd = sleep(Millis(10_000));
821                        crate::spawn(async { sleep(Millis(10_000)).await });
822                    });
823                assert!(start.elapsed() < Duration::from_secs(5));
824            }
825        });
826        let (tx, rx) = std::sync::mpsc::channel();
827        std::thread::spawn(move || tx.send(res.join().is_ok()));
828        assert_eq!(rx.recv_timeout(Duration::from_secs(10)), Ok(true));
829    }
830
831    /// A stale driver does not stop the drivers of a newer generation.
832    #[test]
833    fn test_stale_driver_drop() {
834        let timer = Rc::new(Timer::new());
835        timer.insert_flags(Flags::RUNNING | Flags::DRIVER_STARTED | Flags::LOWRES_DRIVER);
836        timer.reset();
837        assert_eq!(timer.flags.get(), Flags::empty());
838
839        timer.insert_flags(Flags::RUNNING | Flags::DRIVER_STARTED | Flags::LOWRES_DRIVER);
840        drop(TimerDriver {
841            timer: timer.clone(),
842            generation: 0,
843            sleep: Delay::new(Duration::ZERO),
844            armed: None,
845        });
846        drop(LowresTimerDriver {
847            timer: timer.clone(),
848            generation: 0,
849            sleep: Delay::new(Duration::ZERO),
850        });
851        assert_eq!(
852            timer.flags.get(),
853            Flags::RUNNING | Flags::DRIVER_STARTED | Flags::LOWRES_DRIVER
854        );
855    }
856
857    fn entry() -> TimerEntry {
858        TimerEntry {
859            bucket: None,
860            bucket_entry: 0,
861            task: LocalWaker::new(),
862        }
863    }
864
865    /// Bucket expiry is never earlier than the requested delay, and at most
866    /// one level granularity later.
867    #[test]
868    fn test_bucket_expiry_bounds() {
869        let timer = Timer::new();
870        for elapsed in [0, 1, 7, 8, 63, 64, 511, 12_345, 1 << 30] {
871            timer.elapsed.set(elapsed);
872            let mut delta = 0;
873            while delta < WHEEL_TIMEOUT_CUTOFF {
874                let (idx, expiry) = timer.calc_wheel_index(elapsed + delta, delta);
875                let lvl = (idx / LVL_SIZE as usize) as u64;
876                assert!(expiry > elapsed + delta, "{elapsed} {delta} {expiry}");
877                assert!(
878                    expiry <= elapsed + delta + lvl_gran(lvl),
879                    "{elapsed} {delta} {expiry}"
880                );
881                assert_eq!(expiry % lvl_gran(lvl), 0);
882                assert_eq!(
883                    ((expiry >> lvl_shift(lvl)) & LVL_MASK) as usize,
884                    idx % LVL_SIZE as usize
885                );
886                delta = delta * 2 + 1;
887            }
888
889            // clamped to the capacity of the wheel
890            let (_, expiry) =
891                timer.calc_wheel_index(elapsed + u64::from(u32::MAX), u64::from(u32::MAX));
892            assert!(expiry <= elapsed + WHEEL_TIMEOUT_CUTOFF);
893        }
894    }
895
896    /// The next pending bucket is the earliest occupied one, and executing it
897    /// wakes only its timers.
898    #[test]
899    fn test_next_pending_bucket() {
900        let timer = Timer::new();
901        timer.elapsed.set(1000);
902
903        timer.with_wheel(|w| {
904            assert_eq!(timer.next_pending_bucket(w), None);
905
906            let mut expected = Vec::new();
907            for delta in [5000, 30, 700, 100_000] {
908                let (idx, expiry) = timer.calc_wheel_index(1000 + delta, delta);
909                let no = w.timers.insert(entry());
910                w.link(no, idx);
911                expected.push((expiry, no));
912            }
913            expected.sort_unstable();
914
915            for (expiry, no) in expected {
916                assert_eq!(timer.next_pending_bucket(w), Some(expiry));
917                timer.elapsed.set(expiry);
918                w.execute_expired_timers(expiry);
919                assert!(w.timers[no].bucket.is_none());
920            }
921            assert_eq!(timer.next_pending_bucket(w), None);
922            assert_eq!(w.occupied, [0; LVL_DEPTH as usize]);
923        });
924    }
925
926    #[test]
927    fn test_unlink_relink() {
928        let timer = Timer::new();
929        timer.with_wheel(|w| {
930            let a = w.timers.insert(entry());
931            let b = w.timers.insert(entry());
932            w.link(a, 3);
933            w.link(b, 3);
934            assert_eq!(w.occupied[0], 1 << 3);
935
936            assert!(w.unlink(a));
937            assert!(!w.unlink(a));
938            assert_eq!(w.occupied[0], 1 << 3);
939
940            w.relink(b, LVL_SIZE as usize + 5);
941            assert_eq!(w.occupied[0], 0);
942            assert_eq!(w.occupied[1], 1 << 5);
943            assert_eq!(w.timers[b].bucket, Some(LVL_SIZE as u16 + 5));
944
945            assert!(w.unlink(b));
946            assert_eq!(w.occupied, [0; LVL_DEPTH as usize]);
947        });
948    }
949
950    /// `reset(0)` elapses the timer and wakes the task waiting for it.
951    #[ntex::test]
952    async fn test_reset_zero_wakes_task() {
953        let hnd = Rc::new(TimerHandle::new(10_000));
954        let hnd2 = hnd.clone();
955        crate::spawn(async move { hnd2.reset(0) });
956
957        // the sleep wakes the task if `reset(0)` does not
958        let start = Instant::now();
959        crate::future::select(
960            std::future::poll_fn(|cx| hnd.poll_elapsed(cx)),
961            sleep(Millis(500)),
962        )
963        .await;
964        let elapsed = start.elapsed();
965        assert!(elapsed < Duration::from_millis(100), "elapsed: {elapsed:?}");
966    }
967
968    /// Dropping one driver does not reset the state of the other one.
969    #[test]
970    fn test_driver_drop_keeps_other_driver() {
971        let timer = Rc::new(Timer::new());
972        let no = timer.with_wheel(|w| {
973            let no = w.timers.insert(entry());
974            w.link(no, 3);
975            no
976        });
977        timer.next_expiry.set(3);
978        timer.lowres_time.set(Some(Instant::now()));
979        timer.insert_flags(
980            Flags::RUNNING | Flags::DRIVER_STARTED | Flags::LOWRES_DRIVER | Flags::LOWRES_TIMER,
981        );
982
983        drop(LowresTimerDriver {
984            timer: timer.clone(),
985            generation: 0,
986            sleep: Delay::new(Duration::ZERO),
987        });
988        assert_eq!(timer.flags.get(), Flags::DRIVER_STARTED);
989        assert!(timer.lowres_time.get().is_none());
990        // timers are still pending
991        assert_eq!(timer.next_expiry.get(), 3);
992        timer.with_wheel(|w| assert_eq!(w.timers[no].bucket, Some(3)));
993
994        timer.insert_flags(Flags::RUNNING | Flags::LOWRES_DRIVER);
995        drop(TimerDriver {
996            timer: timer.clone(),
997            generation: 0,
998            sleep: Delay::new(Duration::ZERO),
999            armed: None,
1000        });
1001        assert_eq!(timer.flags.get(), Flags::LOWRES_DRIVER);
1002        assert_eq!(timer.next_expiry.get(), u64::MAX);
1003        timer.with_wheel(|w| {
1004            assert!(w.timers[no].bucket.is_none());
1005            assert_eq!(w.occupied, [0; LVL_DEPTH as usize]);
1006        });
1007    }
1008
1009    /// A short timer is measured from the current time, not from the cached
1010    /// time that went stale while the thread was blocked.
1011    #[ntex::test]
1012    async fn test_short_timer_after_blocking() {
1013        let _hnd = sleep(Millis(10_000));
1014        let _ = now();
1015
1016        // the lowres driver cannot run, the cached time goes stale
1017        std::thread::sleep(Duration::from_millis(100));
1018
1019        let start = Instant::now();
1020        sleep(Millis(50)).await;
1021        let elapsed = start.elapsed();
1022        assert!(elapsed >= Duration::from_millis(50), "elapsed: {elapsed:?}");
1023    }
1024
1025    /// The time cached by the wheel driver is invalidated after
1026    /// `LOWRES_RESOLUTION`, like the time cached by `now()`.
1027    #[ntex::test]
1028    async fn test_driver_cached_time_expires() {
1029        // the lowres timer of the cache populated at start fires first
1030        sleep(Millis(400)).await;
1031
1032        // the wheel is idle, the thread waits for a non-timer event
1033        let (tx, rx) = async_channel::bounded(1);
1034        std::thread::spawn(move || {
1035            std::thread::sleep(Duration::from_millis(600));
1036            let _ = tx.send_blocking(());
1037        });
1038        rx.recv().await.unwrap();
1039
1040        let stale = Instant::now().saturating_duration_since(now());
1041        assert!(stale < Duration::from_millis(300), "stale: {stale:?}");
1042    }
1043
1044    /// A late wakeup must not delay the timers that expire later.
1045    #[ntex::test]
1046    async fn test_late_wakeup_does_not_drift() {
1047        let start = Instant::now();
1048        let fut1 = sleep(Millis(50));
1049        let fut2 = sleep(Millis(600));
1050
1051        // block the thread, the driver wakes up ~450ms late for `fut1`.
1052        // `fut2` expires at ~620ms, with drift it would expire ~550ms after
1053        // the late wakeup, i.e. after ~1050ms. The bound leaves room for
1054        // scheduling latency on loaded machines.
1055        std::thread::sleep(Duration::from_millis(500));
1056        fut1.await;
1057        fut2.await;
1058
1059        let elapsed = start.elapsed();
1060        assert!(
1061            elapsed >= Duration::from_millis(600) && elapsed < Duration::from_millis(900),
1062            "elapsed: {elapsed:?}"
1063        );
1064    }
1065
1066    #[ntex::test]
1067    #[allow(unused_variables, clippy::used_underscore_binding)]
1068    async fn test_timer() {
1069        crate::spawn(async {
1070            let s = interval(Millis(25));
1071            loop {
1072                s.tick().await;
1073            }
1074        });
1075        let time = Instant::now();
1076        let fut1 = sleep(Millis(1000));
1077        let fut2 = sleep(Millis(200));
1078
1079        fut2.await;
1080        #[cfg(not(target_os = "macos"))]
1081        {
1082            let _elapsed = time.elapsed();
1083            assert!(
1084                _elapsed > Duration::from_millis(200) && _elapsed < Duration::from_millis(300),
1085                "elapsed: {_elapsed:?}"
1086            );
1087        }
1088
1089        fut1.await;
1090
1091        #[cfg(not(target_os = "macos"))]
1092        {
1093            let _elapsed = time.elapsed();
1094            assert!(
1095                _elapsed > Duration::from_secs(1) && _elapsed < Duration::from_millis(1200), // osx
1096                "elapsed: {_elapsed:?}",
1097            );
1098        }
1099
1100        let time = Instant::now();
1101        sleep(Millis(25)).await;
1102        #[cfg(not(target_os = "macos"))]
1103        {
1104            let _elapsed = time.elapsed();
1105            assert!(
1106                _elapsed > Duration::from_millis(20) && _elapsed < Duration::from_millis(50),
1107                "elapsed: {_elapsed:?}",
1108            );
1109        }
1110    }
1111}