diff --git a/tokio/src/time/delay_queue.rs b/tokio/src/time/delay_queue.rs index f6007d740..1790ada8f 100644 --- a/tokio/src/time/delay_queue.rs +++ b/tokio/src/time/delay_queue.rs @@ -326,7 +326,12 @@ impl DelayQueue { }; if should_set_delay { - self.delay = Some(delay_until(self.start + Duration::from_millis(when))); + let delay_time = self.start + Duration::from_millis(when); + if let Some(ref mut delay) = &mut self.delay { + delay.reset(delay_time); + } else { + self.delay = Some(delay_until(delay_time)); + } } Key::new(key) diff --git a/tokio/tests/time_delay_queue.rs b/tokio/tests/time_delay_queue.rs index 2239fab4c..32e812e03 100644 --- a/tokio/tests/time_delay_queue.rs +++ b/tokio/tests/time_delay_queue.rs @@ -384,6 +384,32 @@ async fn reset_first_expiring_item_to_expire_later() { assert_eq!(entry, "two"); } +#[tokio::test] +async fn insert_before_first_after_poll() { + time::pause(); + + let mut queue = task::spawn(DelayQueue::new()); + + let now = Instant::now(); + + let _one = queue.insert_at("one", now + ms(200)); + + assert_pending!(poll!(queue)); + + let _two = queue.insert_at("two", now + ms(100)); + + delay_for(ms(99)).await; + + assert!(!queue.is_woken()); + + delay_for(ms(1)).await; + + assert!(queue.is_woken()); + + let entry = assert_ready_ok!(poll!(queue)).into_inner(); + assert_eq!(entry, "two"); +} + fn ms(n: u64) -> Duration { Duration::from_millis(n) }