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, query_system_time, system_time};
14
15#[inline]
22pub fn sleep<T: Into<Millis>>(dur: T) -> Sleep {
23 Sleep::new(dur.into())
24}
25
26#[inline]
31pub fn deadline<T: Into<Millis>>(dur: T) -> Deadline {
32 Deadline::new(dur.into())
33}
34
35#[inline]
40pub fn interval<T: Into<Millis>>(period: T) -> Interval {
41 Interval::new(period.into())
42}
43
44#[inline]
50pub fn timeout<T, U>(dur: U, future: T) -> Timeout<T>
51where
52 T: Future,
53 U: Into<Millis>,
54{
55 Timeout::new_with_delay(future, Sleep::new(dur.into()))
56}
57
58#[inline]
64pub fn timeout_checked<T, U>(dur: U, future: T) -> TimeoutChecked<T>
65where
66 T: Future,
67 U: Into<Millis>,
68{
69 TimeoutChecked::new_with_delay(future, dur.into())
70}
71
72#[derive(Debug)]
88#[must_use = "futures do nothing unless you `.await` or poll them"]
89pub struct Sleep {
90 hnd: TimerHandle,
92}
93
94impl Sleep {
95 #[inline]
97 pub fn new(duration: Millis) -> Sleep {
98 Sleep {
99 hnd: TimerHandle::new(u64::from(cmp::max(duration.0, 1))),
100 }
101 }
102
103 #[inline]
105 pub fn is_elapsed(&self) -> bool {
106 self.hnd.is_elapsed()
107 }
108
109 #[inline]
111 pub fn elapse(&self) {
112 self.hnd.elapse();
113 }
114
115 pub fn reset<T: Into<Millis>>(&self, millis: T) {
123 self.hnd.reset(u64::from(millis.into().0));
124 }
125
126 #[inline]
127 pub async fn wait(&self) {
129 poll_fn(|cx| self.hnd.poll_elapsed(cx)).await;
130 }
131
132 #[inline]
133 pub fn poll_elapsed(&self, cx: &mut task::Context<'_>) -> Poll<()> {
134 self.hnd.poll_elapsed(cx)
135 }
136}
137
138impl Future for Sleep {
139 type Output = ();
140
141 fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
142 self.hnd.poll_elapsed(cx)
143 }
144}
145
146#[derive(Debug)]
162#[must_use = "futures do nothing unless you `.await` or poll them"]
163pub struct Deadline {
164 hnd: Option<TimerHandle>,
165}
166
167impl Deadline {
168 #[inline]
170 pub fn new(duration: Millis) -> Deadline {
171 if duration.0 != 0 {
172 Deadline {
173 hnd: Some(TimerHandle::new(u64::from(duration.0))),
174 }
175 } else {
176 Deadline { hnd: None }
177 }
178 }
179
180 #[inline]
181 pub async fn wait(&self) {
183 poll_fn(|cx| self.poll_elapsed(cx)).await;
184 }
185
186 pub fn reset<T: Into<Millis>>(&mut self, millis: T) {
194 let millis = millis.into();
195 if millis.0 != 0 {
196 if let Some(ref mut hnd) = self.hnd {
197 hnd.reset(u64::from(millis.0));
198 } else {
199 self.hnd = Some(TimerHandle::new(u64::from(millis.0)));
200 }
201 } else {
202 let _ = self.hnd.take();
203 }
204 }
205
206 #[inline]
208 pub fn is_elapsed(&self) -> bool {
209 self.hnd.as_ref().is_none_or(TimerHandle::is_elapsed)
210 }
211
212 #[inline]
213 pub fn poll_elapsed(&self, cx: &mut task::Context<'_>) -> Poll<()> {
214 self.hnd
215 .as_ref()
216 .map_or(Poll::Pending, |t| t.poll_elapsed(cx))
217 }
218}
219
220impl Future for Deadline {
221 type Output = ();
222
223 fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
224 self.poll_elapsed(cx)
225 }
226}
227
228pin_project_lite::pin_project! {
229 #[must_use = "futures do nothing unless you `.await` or poll them"]
231 #[derive(Debug)]
232 pub struct Timeout<T> {
233 #[pin]
234 value: T,
235 delay: Sleep,
236 }
237}
238
239impl<T> Timeout<T> {
240 pub(crate) fn new_with_delay(value: T, delay: Sleep) -> Timeout<T> {
241 Timeout { value, delay }
242 }
243}
244
245impl<T> Future for Timeout<T>
246where
247 T: Future,
248{
249 type Output = Result<T::Output, ()>;
250
251 fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
252 let this = self.project();
253
254 if let Poll::Ready(v) = this.value.poll(cx) {
256 return Poll::Ready(Ok(v));
257 }
258
259 match this.delay.poll_elapsed(cx) {
261 Poll::Ready(()) => Poll::Ready(Err(())),
262 Poll::Pending => Poll::Pending,
263 }
264 }
265}
266
267pin_project_lite::pin_project! {
268 #[must_use = "futures do nothing unless you `.await` or poll them"]
270 pub struct TimeoutChecked<T> {
271 #[pin]
272 state: TimeoutCheckedState<T>,
273 }
274}
275
276pin_project_lite::pin_project! {
277 #[project = TimeoutCheckedStateProject]
278 enum TimeoutCheckedState<T> {
279 Timeout{ #[pin] fut: Timeout<T> },
280 NoTimeout{ #[pin] fut: T },
281 }
282}
283
284impl<T> TimeoutChecked<T> {
285 pub(crate) fn new_with_delay(value: T, delay: Millis) -> TimeoutChecked<T> {
286 if delay.is_zero() {
287 TimeoutChecked {
288 state: TimeoutCheckedState::NoTimeout { fut: value },
289 }
290 } else {
291 TimeoutChecked {
292 state: TimeoutCheckedState::Timeout {
293 fut: Timeout::new_with_delay(value, sleep(delay)),
294 },
295 }
296 }
297 }
298}
299
300impl<T> Future for TimeoutChecked<T>
301where
302 T: Future,
303{
304 type Output = Result<T::Output, ()>;
305
306 fn poll(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Self::Output> {
307 match self.project().state.as_mut().project() {
308 TimeoutCheckedStateProject::Timeout { fut } => fut.poll(cx),
309 TimeoutCheckedStateProject::NoTimeout { fut } => fut.poll(cx).map(Result::Ok),
310 }
311 }
312}
313
314#[must_use = "futures do nothing unless you `.await` or poll them"]
319#[derive(Debug)]
320pub struct Interval {
321 hnd: TimerHandle,
322 period: u32,
323}
324
325impl Interval {
326 #[inline]
328 pub fn new(period: Millis) -> Interval {
329 Interval {
330 hnd: TimerHandle::new(u64::from(period.0)),
331 period: period.0,
332 }
333 }
334
335 #[inline]
336 pub async fn tick(&self) {
337 poll_fn(|cx| self.poll_tick(cx)).await;
338 }
339
340 #[inline]
341 pub fn poll_tick(&self, cx: &mut task::Context<'_>) -> Poll<()> {
342 if self.hnd.poll_elapsed(cx).is_ready() {
343 self.hnd.reset(u64::from(self.period));
344 Poll::Ready(())
345 } else {
346 Poll::Pending
347 }
348 }
349}
350
351impl crate::Stream for Interval {
352 type Item = ();
353
354 #[inline]
355 fn poll_next(self: Pin<&mut Self>, cx: &mut task::Context<'_>) -> Poll<Option<Self::Item>> {
356 self.poll_tick(cx).map(|()| Some(()))
357 }
358}
359
360#[cfg(test)]
361mod tests {
362 use futures_util::StreamExt;
363 use std::{future::poll_fn, rc::Rc, time};
364
365 use super::*;
366 use crate::future::lazy;
367
368 #[ntex::test]
372 async fn lowres_time_does_not_immediately_change() {
373 sleep(Millis(25)).await;
374
375 assert_eq!(now(), now());
376 }
377
378 #[ntex::test]
383 async fn lowres_time_updates_after_resolution_interval() {
384 sleep(Millis(50)).await;
385
386 let first_time = now();
387
388 sleep(Millis(25)).await;
389
390 let second_time = now();
391 assert!(second_time - first_time >= time::Duration::from_millis(25));
392 }
393
394 #[ntex::test]
398 async fn system_time_service_time_does_not_immediately_change() {
399 sleep(Seconds(1)).await;
400
401 assert_eq!(system_time(), system_time());
402 assert_eq!(system_time(), query_system_time());
403 }
404
405 #[ntex::test]
410 async fn system_time_service_time_updates_after_resolution_interval() {
411 sleep(Millis(100)).await;
412
413 let wait_time = 300;
414
415 let first_time = system_time()
416 .duration_since(time::SystemTime::UNIX_EPOCH)
417 .unwrap();
418
419 sleep(Millis(wait_time)).await;
420
421 let second_time = system_time()
422 .duration_since(time::SystemTime::UNIX_EPOCH)
423 .unwrap();
424
425 assert!(
426 second_time.checked_sub(first_time).unwrap()
427 >= time::Duration::from_millis(u64::from(wait_time))
428 );
429 }
430
431 #[ntex::test]
432 async fn test_sleep_0() {
433 sleep(Seconds(1)).await;
434
435 let first_time = now();
436 sleep(Millis(0)).await;
437 let second_time = now();
438 assert!(second_time - first_time >= time::Duration::from_millis(1));
439
440 let first_time = now();
441 sleep(Millis(1)).await;
442 let second_time = now();
443 assert!(second_time - first_time >= time::Duration::from_millis(1));
444
445 let first_time = now();
446 let fut = sleep(Millis(10000));
447 assert!(!fut.is_elapsed());
448 fut.reset(Millis::ZERO);
449 fut.await;
450 let second_time = now();
451 assert!(second_time - first_time < time::Duration::from_millis(1));
452
453 let first_time = now();
454 let fut = Sleep {
455 hnd: TimerHandle::new(0),
456 };
457 assert!(fut.is_elapsed());
458 fut.await;
459 let second_time = now();
460 assert!(second_time - first_time < time::Duration::from_millis(1));
461
462 let first_time = now();
463 let fut = Rc::new(sleep(Millis(10_0000)));
464 let s = fut.clone();
465 ntex::rt::spawn(async move {
466 s.elapse();
467 });
468 poll_fn(|cx| fut.poll_elapsed(cx)).await;
469 assert!(fut.is_elapsed());
470 let second_time = now();
471 assert!(second_time - first_time < time::Duration::from_millis(1));
472 }
473
474 #[ntex::test]
475 async fn test_deadline() {
476 sleep(Seconds(1)).await;
477
478 let first_time = now();
479 let dl = deadline(Millis(1));
480 dl.await;
481 let second_time = now();
482 assert!(second_time - first_time >= time::Duration::from_millis(1));
483 assert!(timeout(Millis(100), deadline(Millis(0))).await.is_err());
484
485 let mut dl = deadline(Millis(1));
486 dl.reset(Millis::ZERO);
487 assert!(lazy(|cx| dl.poll_elapsed(cx)).await.is_pending());
488
489 let mut dl = deadline(Millis(1));
490 dl.reset(Millis(100));
491 let first_time = now();
492 dl.await;
493 let second_time = now();
494 assert!(second_time - first_time >= time::Duration::from_millis(100));
495
496 let mut dl = deadline(Millis(0));
497 assert!(dl.is_elapsed());
498 dl.reset(Millis(1));
499 assert!(lazy(|cx| dl.poll_elapsed(cx)).await.is_pending());
500
501 assert!(format!("{dl:?}").contains("Deadline"));
502 }
503
504 #[ntex::test]
505 async fn test_interval() {
506 let mut int = interval(Millis(250));
507
508 let time = time::Instant::now();
509 int.tick().await;
510 let elapsed = time.elapsed();
511 assert!(
512 elapsed > time::Duration::from_millis(200)
513 && elapsed < time::Duration::from_millis(450),
514 "elapsed: {elapsed:?}"
515 );
516
517 let time = time::Instant::now();
518 int.next().await;
519 let elapsed = time.elapsed();
520 assert!(
521 elapsed > time::Duration::from_millis(200)
522 && elapsed < time::Duration::from_millis(450),
523 "elapsed: {elapsed:?}"
524 );
525 }
526
527 #[ntex::test]
528 async fn test_interval_one_sec() {
529 let int = interval(Millis::ONE_SEC);
530
531 for _i in 0..3 {
532 let time = time::Instant::now();
533 int.tick().await;
534 let elapsed = time.elapsed();
535 assert!(
536 elapsed > time::Duration::from_secs(1)
537 && elapsed < time::Duration::from_millis(1300),
538 "elapsed: {elapsed:?}"
539 );
540 }
541 }
542
543 #[ntex::test]
544 async fn test_timeout_checked() {
545 let result = timeout_checked(Millis(200), sleep(Millis(100))).await;
546 assert!(result.is_ok());
547
548 let result = timeout_checked(Millis(5), sleep(Millis(100))).await;
549 assert!(result.is_err());
550
551 let result = timeout_checked(Millis(0), sleep(Millis(100))).await;
552 assert!(result.is_ok());
553 }
554}