1#![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#[inline]
22pub fn sleep<T: Into<Millis>>(dur: T) -> Sleep {
23 Sleep::new(dur.into())
24}
25
26#[inline]
30pub fn deadline<T: Into<Millis>>(dur: T) -> Deadline {
31 Deadline::new(dur.into())
32}
33
34#[inline]
41pub fn interval<T: Into<Millis>>(period: T) -> Interval {
42 Interval::new(period.into())
43}
44
45#[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#[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#[derive(Debug)]
91#[must_use = "futures do nothing unless you `.await` or poll them"]
92pub struct Sleep {
93 hnd: TimerHandle,
95}
96
97impl Sleep {
98 #[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 #[inline]
110 pub fn is_elapsed(&self) -> bool {
111 self.hnd.is_elapsed()
112 }
113
114 #[inline]
116 pub fn elapse(&self) {
117 self.hnd.elapse();
118 }
119
120 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 pub async fn wait(&self) {
134 poll_fn(|cx| self.hnd.poll_elapsed(cx)).await;
135 }
136
137 #[inline]
138 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#[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 #[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 pub async fn wait(&self) {
191 poll_fn(|cx| self.poll_elapsed(cx)).await;
192 }
193
194 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 #[inline]
220 pub fn is_elapsed(&self) -> bool {
221 self.hnd.as_ref().is_none_or(TimerHandle::is_elapsed)
222 }
223
224 #[inline]
225 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 #[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 if let Poll::Ready(v) = this.value.poll(cx) {
271 return Poll::Ready(Ok(v));
272 }
273
274 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 #[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#[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 #[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 pub async fn tick(&self) {
356 poll_fn(|cx| self.poll_tick(cx)).await;
357 }
358
359 #[inline]
360 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 #[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 #[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 #[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 #[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 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 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 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 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 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}