Skip to main content

ntex_util/time/
mod.rs

1//! Utilities for tracking time.
2#![allow(
3    clippy::cast_possible_truncation,
4    clippy::cast_sign_loss,
5    clippy::cast_possible_wrap
6)]
7use std::{cmp, future::Future, future::poll_fn, pin::Pin, task, task::Poll};
8
9mod types;
10mod wheel;
11
12pub use self::types::{Millis, Seconds};
13pub use self::wheel::{TimerHandle, now, system_time};
14
15/// Waits until `dur` has elapsed.
16///
17/// No work is performed while awaiting the returned [`Sleep`]. Timers have a
18/// granularity of approximately 16 milliseconds and are not suitable for
19/// high-resolution timing. A zero duration still waits for at least one timer
20/// tick.
21#[inline]
22pub fn sleep<T: Into<Millis>>(dur: T) -> Sleep {
23    Sleep::new(dur.into())
24}
25
26/// Waits until `dur` has elapsed.
27///
28/// Unlike [`sleep`], a zero-duration deadline never completes.
29#[inline]
30pub fn deadline<T: Into<Millis>>(dur: T) -> Deadline {
31    Deadline::new(dur.into())
32}
33
34/// Creates an [`Interval`] that ticks every `period`.
35///
36/// The first tick completes immediately, and each later tick completes
37/// `period` after the previous one was observed. An interval will tick
38/// indefinitely. At any time, the [`Interval`] value can be dropped. This
39/// cancels the interval.
40#[inline]
41pub fn interval<T: Into<Millis>>(period: T) -> Interval {
42    Interval::new(period.into())
43}
44
45/// Requires a future to complete before `dur` has elapsed.
46///
47/// If the future completes before the duration has elapsed, then the completed
48/// value is returned. Otherwise, an error is returned and the future is
49/// canceled. A zero duration still represents an active timeout of at least
50/// one timer tick; use [`timeout_checked`] to disable the timeout with zero. If
51/// the future and timer are both ready during the same poll, the future wins.
52#[inline]
53pub fn timeout<T, U>(dur: U, future: T) -> Timeout<T>
54where
55    T: Future,
56    U: Into<Millis>,
57{
58    Timeout::new_with_delay(future, Sleep::new(dur.into()))
59}
60
61/// Requires a future to complete before `dur` has elapsed.
62///
63/// If the future completes before the duration has elapsed, then the completed
64/// value is returned. Otherwise, an error is returned and the future is
65/// canceled. A zero duration disables the timeout.
66#[inline]
67pub fn timeout_checked<T, U>(dur: U, future: T) -> TimeoutChecked<T>
68where
69    T: Future,
70    U: Into<Millis>,
71{
72    TimeoutChecked::new_with_delay(future, dur.into())
73}
74
75/// Future returned by [`sleep`].
76///
77/// # Examples
78///
79/// Wait 100ms and print "100 ms have elapsed".
80///
81/// ```
82/// use ntex::time::{sleep, Millis};
83///
84/// #[ntex::main]
85/// async fn main() {
86///     sleep(Millis(100)).await;
87///     println!("100 ms have elapsed");
88/// }
89/// ```
90#[derive(Debug)]
91#[must_use = "futures do nothing unless you `.await` or poll them"]
92pub struct Sleep {
93    // The link between the `Sleep` instance and the timer that drives it.
94    hnd: TimerHandle,
95}
96
97impl Sleep {
98    /// Creates a new sleep future.
99    ///
100    /// A zero duration is rounded up to 1 ms.
101    #[inline]
102    pub fn new(duration: Millis) -> Sleep {
103        Sleep {
104            hnd: TimerHandle::new(u64::from(cmp::max(duration.0, 1))),
105        }
106    }
107
108    /// Returns `true` if `Sleep` has elapsed.
109    #[inline]
110    pub fn is_elapsed(&self) -> bool {
111        self.hnd.is_elapsed()
112    }
113
114    /// Completes the timer immediately.
115    #[inline]
116    pub fn elapse(&self) {
117        self.hnd.elapse();
118    }
119
120    /// Resets the `Sleep` instance to a new deadline.
121    ///
122    /// Calling this function allows changing the instant at which the `Sleep`
123    /// future completes without having to create new associated state.
124    ///
125    /// This function can be called both before and after the future has
126    /// completed. A zero duration is rounded up to 1 ms.
127    pub fn reset<T: Into<Millis>>(&self, millis: T) {
128        self.hnd.reset(u64::from(cmp::max(millis.into().0, 1)));
129    }
130
131    #[inline]
132    /// Waits until this timer has elapsed.
133    pub async fn wait(&self) {
134        poll_fn(|cx| self.hnd.poll_elapsed(cx)).await;
135    }
136
137    #[inline]
138    /// Polls until this timer has elapsed.
139    pub fn poll_elapsed(&self, cx: &mut task::Context<'_>) -> Poll<()> {
140        self.hnd.poll_elapsed(cx)
141    }
142}
143
144impl Future for Sleep {
145    type Output = ();
146
147    fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
148        self.hnd.poll_elapsed(cx)
149    }
150}
151
152/// Future returned by [`deadline`].
153///
154/// # Examples
155///
156/// Wait 100ms and print "100 ms have elapsed".
157///
158/// ```
159/// use ntex::time::{deadline, Millis};
160///
161/// #[ntex::main]
162/// async fn main() {
163///     deadline(Millis(100)).await;
164///     println!("100 ms have elapsed");
165/// }
166/// ```
167#[derive(Debug)]
168#[must_use = "futures do nothing unless you `.await` or poll them"]
169pub struct Deadline {
170    hnd: Option<TimerHandle>,
171}
172
173impl Deadline {
174    /// Creates a new deadline future.
175    ///
176    /// A zero duration creates a deadline that never completes.
177    #[inline]
178    pub fn new(duration: Millis) -> Deadline {
179        if duration.0 != 0 {
180            Deadline {
181                hnd: Some(TimerHandle::new(u64::from(duration.0))),
182            }
183        } else {
184            Deadline { hnd: None }
185        }
186    }
187
188    #[inline]
189    /// Waits until this deadline has elapsed.
190    pub async fn wait(&self) {
191        poll_fn(|cx| self.poll_elapsed(cx)).await;
192    }
193
194    /// Resets the `Deadline` instance to a new deadline.
195    ///
196    /// Calling this function allows changing the instant at which the `Deadline`
197    /// future completes without having to create new associated state.
198    ///
199    /// This function can be called both before and after the future has
200    /// completed.
201    pub fn reset<T: Into<Millis>>(&mut self, millis: T) {
202        let millis = millis.into();
203        if millis.0 != 0 {
204            if let Some(ref mut hnd) = self.hnd {
205                hnd.reset(u64::from(millis.0));
206            } else {
207                self.hnd = Some(TimerHandle::new(u64::from(millis.0)));
208            }
209        } else {
210            let _ = self.hnd.take();
211        }
212    }
213
214    /// Returns `true` if `Deadline` has elapsed.
215    ///
216    /// A disabled zero-duration deadline is reported as elapsed here even
217    /// though polling it remains pending. Use this method to determine whether
218    /// there is an active future deadline to wait for.
219    #[inline]
220    pub fn is_elapsed(&self) -> bool {
221        self.hnd.as_ref().is_none_or(TimerHandle::is_elapsed)
222    }
223
224    #[inline]
225    /// Polls until this deadline has elapsed.
226    ///
227    /// A zero-duration deadline remains pending until it is reset.
228    pub fn poll_elapsed(&self, cx: &mut task::Context<'_>) -> Poll<()> {
229        self.hnd
230            .as_ref()
231            .map_or(Poll::Pending, |t| t.poll_elapsed(cx))
232    }
233}
234
235impl Future for Deadline {
236    type Output = ();
237
238    fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
239        self.poll_elapsed(cx)
240    }
241}
242
243pin_project_lite::pin_project! {
244    /// Future returned by [`timeout`](timeout).
245    #[must_use = "futures do nothing unless you `.await` or poll them"]
246    #[derive(Debug)]
247    pub struct Timeout<T> {
248        #[pin]
249        value: T,
250        delay: Sleep,
251    }
252}
253
254impl<T> Timeout<T> {
255    pub(crate) fn new_with_delay(value: T, delay: Sleep) -> Timeout<T> {
256        Timeout { value, delay }
257    }
258}
259
260impl<T> Future for Timeout<T>
261where
262    T: Future,
263{
264    type Output = Result<T::Output, ()>;
265
266    fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
267        let this = self.project();
268
269        // First, try polling the future
270        if let Poll::Ready(v) = this.value.poll(cx) {
271            return Poll::Ready(Ok(v));
272        }
273
274        // Now check the timer
275        match this.delay.poll_elapsed(cx) {
276            Poll::Ready(()) => Poll::Ready(Err(())),
277            Poll::Pending => Poll::Pending,
278        }
279    }
280}
281
282pin_project_lite::pin_project! {
283    /// Future returned by [`timeout_checked`](timeout_checked).
284    #[must_use = "futures do nothing unless you `.await` or poll them"]
285    pub struct TimeoutChecked<T> {
286        #[pin]
287        state: TimeoutCheckedState<T>,
288    }
289}
290
291pin_project_lite::pin_project! {
292    #[project = TimeoutCheckedStateProject]
293    enum TimeoutCheckedState<T> {
294        Timeout{ #[pin] fut: Timeout<T> },
295        NoTimeout{ #[pin] fut: T },
296    }
297}
298
299impl<T> TimeoutChecked<T> {
300    pub(crate) fn new_with_delay(value: T, delay: Millis) -> TimeoutChecked<T> {
301        if delay.is_zero() {
302            TimeoutChecked {
303                state: TimeoutCheckedState::NoTimeout { fut: value },
304            }
305        } else {
306            TimeoutChecked {
307                state: TimeoutCheckedState::Timeout {
308                    fut: Timeout::new_with_delay(value, sleep(delay)),
309                },
310            }
311        }
312    }
313}
314
315impl<T> Future for TimeoutChecked<T>
316where
317    T: Future,
318{
319    type Output = Result<T::Output, ()>;
320
321    fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
322        match self.project().state.as_mut().project() {
323            TimeoutCheckedStateProject::Timeout { fut } => fut.poll(cx),
324            TimeoutCheckedStateProject::NoTimeout { fut } => fut.poll(cx).map(Result::Ok),
325        }
326    }
327}
328
329/// An interval returned by [`interval`].
330///
331/// This type allows you to wait on a sequence of instants with a certain
332/// duration between each instant.
333#[must_use = "futures do nothing unless you `.await` or poll them"]
334#[derive(Debug)]
335pub struct Interval {
336    hnd: TimerHandle,
337    period: u32,
338}
339
340impl Interval {
341    /// Creates an interval with the specified period.
342    ///
343    /// The first tick completes immediately. A zero period is rounded up to
344    /// 1 ms.
345    #[inline]
346    pub fn new(period: Millis) -> Interval {
347        Interval {
348            hnd: TimerHandle::new(0),
349            period: cmp::max(period.0, 1),
350        }
351    }
352
353    #[inline]
354    /// Waits for the next interval tick.
355    pub async fn tick(&self) {
356        poll_fn(|cx| self.poll_tick(cx)).await;
357    }
358
359    #[inline]
360    /// Polls for the next interval tick.
361    pub fn poll_tick(&self, cx: &mut task::Context<'_>) -> Poll<()> {
362        if self.hnd.poll_elapsed(cx).is_ready() {
363            self.hnd.reset(u64::from(self.period));
364            Poll::Ready(())
365        } else {
366            Poll::Pending
367        }
368    }
369}
370
371impl crate::Stream for Interval {
372    type Item = ();
373
374    #[inline]
375    fn poll_next(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Option<Self::Item>> {
376        self.poll_tick(cx).map(|()| Some(()))
377    }
378}
379
380#[cfg(test)]
381mod tests {
382    use futures_util::StreamExt;
383    use std::{future::poll_fn, rc::Rc, time};
384
385    use super::*;
386    use crate::future::lazy;
387
388    /// State Under Test: Two calls of `now()` return the same value if they are done within resolution interval.
389    ///
390    /// Expected Behavior: Two back-to-back calls of `now()` return the same value.
391    #[ntex::test]
392    async fn lowres_time_does_not_immediately_change() {
393        sleep(Millis(25)).await;
394
395        assert_eq!(now(), now());
396    }
397
398    /// State Under Test: `now()` updates returned value every ~1ms period.
399    ///
400    /// Expected Behavior: Two calls of `now()` made in subsequent resolution interval return different values
401    /// and second value is greater than the first one at least by a 1ms interval.
402    #[ntex::test]
403    async fn lowres_time_updates_after_resolution_interval() {
404        sleep(Millis(50)).await;
405
406        let first_time = now();
407
408        sleep(Millis(25)).await;
409
410        let second_time = now();
411        assert!(second_time - first_time >= time::Duration::from_millis(25));
412    }
413
414    /// State Under Test: Two calls of `system_time()` return the same value if they are done within the
415    /// resolution interval.
416    ///
417    /// Expected Behavior: Two back-to-back calls of `now()` return the same value.
418    #[ntex::test]
419    async fn system_time_service_time_does_not_immediately_change() {
420        sleep(Seconds(1)).await;
421
422        assert_eq!(system_time(), system_time());
423    }
424
425    /// State Under Test: `system_time()` updates returned value every resolution interval (150ms).
426    ///
427    /// Expected Behavior: Two calls of `system_time()` made in subsequent resolution interval return different values
428    /// and second value is greater than the first one at least by a resolution interval.
429    #[ntex::test]
430    async fn system_time_service_time_updates_after_resolution_interval() {
431        sleep(Millis(100)).await;
432
433        let wait_time = 600;
434
435        let first_time = system_time()
436            .duration_since(time::SystemTime::UNIX_EPOCH)
437            .unwrap();
438
439        sleep(Millis(wait_time)).await;
440
441        let second_time = system_time()
442            .duration_since(time::SystemTime::UNIX_EPOCH)
443            .unwrap();
444
445        assert!(
446            second_time.checked_sub(first_time).unwrap()
447                >= time::Duration::from_millis(u64::from(wait_time))
448        );
449    }
450
451    #[ntex::test]
452    async fn test_sleep_0() {
453        sleep(Seconds(1)).await;
454
455        let first_time = time::Instant::now();
456        sleep(Millis(0)).await;
457        let second_time = time::Instant::now();
458        assert!(second_time - first_time >= time::Duration::from_millis(1));
459
460        let first_time = time::Instant::now();
461        sleep(Millis(1)).await;
462        let second_time = time::Instant::now();
463        assert!(second_time - first_time >= time::Duration::from_millis(1));
464
465        let first_time = time::Instant::now();
466        let fut = sleep(Millis(10000));
467        assert!(!fut.is_elapsed());
468        // a zero delay is rounded up to 1 ms
469        fut.reset(Millis::ZERO);
470        assert!(!fut.is_elapsed());
471        fut.await;
472        let second_time = time::Instant::now();
473        assert!(second_time - first_time >= time::Duration::from_millis(1));
474
475        let fut = sleep(Millis::ZERO);
476        assert!(!fut.is_elapsed());
477        fut.await;
478
479        // the cached `now()` may be refreshed between two calls, measure
480        // with the clock. An elapsed timer completes without waiting for the
481        // timer driver, the bound leaves room for scheduling latency
482        let first_time = time::Instant::now();
483        let fut = Sleep {
484            hnd: TimerHandle::new(0),
485        };
486        assert!(fut.is_elapsed());
487        fut.await;
488        let second_time = time::Instant::now();
489        assert!(second_time - first_time < time::Duration::from_millis(50));
490
491        let first_time = time::Instant::now();
492        let fut = Rc::new(sleep(Millis(10_0000)));
493        let s = fut.clone();
494        ntex::rt::spawn(async move {
495            s.elapse();
496        });
497        poll_fn(|cx| fut.poll_elapsed(cx)).await;
498        assert!(fut.is_elapsed());
499        let second_time = time::Instant::now();
500        assert!(second_time - first_time < time::Duration::from_millis(50));
501    }
502
503    #[ntex::test]
504    async fn test_deadline() {
505        sleep(Seconds(1)).await;
506
507        let first_time = now();
508        let dl = deadline(Millis(1));
509        dl.await;
510        let second_time = now();
511        assert!(second_time - first_time >= time::Duration::from_millis(1));
512        assert!(timeout(Millis(100), deadline(Millis(0))).await.is_err());
513
514        let mut dl = deadline(Millis(1));
515        dl.reset(Millis::ZERO);
516        assert!(lazy(|cx| dl.poll_elapsed(cx)).await.is_pending());
517
518        let mut dl = deadline(Millis(1));
519        dl.reset(Millis(100));
520        let first_time = now();
521        dl.await;
522        let second_time = now();
523        assert!(second_time - first_time >= time::Duration::from_millis(100));
524
525        let mut dl = deadline(Millis(0));
526        assert!(dl.is_elapsed());
527        dl.reset(Millis(1));
528        assert!(lazy(|cx| dl.poll_elapsed(cx)).await.is_pending());
529
530        assert!(format!("{dl:?}").contains("Deadline"));
531    }
532
533    #[ntex::test]
534    async fn test_interval_zero_period() {
535        let int = interval(Millis::ZERO);
536        int.tick().await;
537
538        // a zero period is rounded up to 1 ms instead of ticking on every poll
539        assert!(lazy(|cx| int.poll_tick(cx)).await.is_pending());
540        int.tick().await;
541    }
542
543    #[ntex::test]
544    async fn test_interval() {
545        let mut int = interval(Millis(250));
546
547        // the first tick completes immediately
548        let time = time::Instant::now();
549        int.tick().await;
550        assert!(time.elapsed() < time::Duration::from_millis(50));
551
552        let time = time::Instant::now();
553        int.tick().await;
554        let elapsed = time.elapsed();
555        assert!(
556            elapsed > time::Duration::from_millis(200)
557                && elapsed < time::Duration::from_millis(450),
558            "elapsed: {elapsed:?}"
559        );
560
561        let time = time::Instant::now();
562        int.next().await;
563        let elapsed = time.elapsed();
564        assert!(
565            elapsed > time::Duration::from_millis(200)
566                && elapsed < time::Duration::from_millis(450),
567            "elapsed: {elapsed:?}"
568        );
569    }
570
571    #[ntex::test]
572    async fn test_interval_one_sec() {
573        let int = interval(Millis::ONE_SEC);
574        int.tick().await;
575
576        for _i in 0..3 {
577            let time = time::Instant::now();
578            int.tick().await;
579            let elapsed = time.elapsed();
580            assert!(
581                elapsed > time::Duration::from_secs(1)
582                    && elapsed < time::Duration::from_millis(1300),
583                "elapsed: {elapsed:?}"
584            );
585        }
586    }
587
588    #[ntex::test]
589    async fn test_timeout_checked() {
590        let result = timeout_checked(Millis(200), sleep(Millis(100))).await;
591        assert!(result.is_ok());
592
593        // a future that never completes, a late timer driver would elapse
594        // a competing sleep in the same pass
595        let result = timeout_checked(Millis(5), std::future::pending::<()>()).await;
596        assert!(result.is_err());
597
598        let result = timeout_checked(Millis(0), sleep(Millis(100))).await;
599        assert!(result.is_ok());
600    }
601
602    #[ntex::test]
603    async fn sleep_and_deadline_wait() {
604        let start = now();
605        sleep(Millis(5)).wait().await;
606        deadline(Millis(5)).wait().await;
607        assert!(now() >= start);
608
609        let hnd = Deadline::new(Millis::ZERO);
610        assert!(
611            crate::time::timeout(Millis(20), hnd.wait()).await.is_err(),
612            "zero deadline never completes"
613        );
614    }
615}