diff --git a/tokio-timer/src/delay_queue.rs b/tokio-timer/src/delay_queue.rs index 826b70499..ee28020ad 100644 --- a/tokio-timer/src/delay_queue.rs +++ b/tokio-timer/src/delay_queue.rs @@ -338,6 +338,16 @@ impl DelayQueue { self.insert_idx(when, key); + // Set a new delay if the current's deadline is later than the one of the new item + let should_set_delay = if let Some(ref delay) = self.delay { + let current_exp = self.normalize_deadline(delay.deadline()); + current_exp > when + } else { false }; + + if should_set_delay { + self.delay = Some(self.handle.delay(self.start + Duration::from_millis(when))); + } + Key::new(key) } diff --git a/tokio-timer/tests/queue.rs b/tokio-timer/tests/queue.rs index 4ecfc4b35..17d876070 100644 --- a/tokio-timer/tests/queue.rs +++ b/tokio-timer/tests/queue.rs @@ -279,3 +279,35 @@ fn remove_expired_item() { assert_eq!(entry.into_inner(), "foo"); }) } + +#[test] +fn expires_before_last_insert() { + mocked(|timer, time| { + let mut queue = DelayQueue::new(); + let mut task = MockTask::new(); + + + let epoch = time.now(); + + queue.insert_at("foo", epoch + ms(10_000)); + + // Delay should be set to 8.192s here. + task.enter(|| { + assert_not_ready!(queue); + }); + + // Delay should be set to the delay of the new item here + queue.insert_at("bar", epoch + ms(600)); + + task.enter(|| { + assert_not_ready!(queue); + }); + + advance(timer, ms(600)); + + assert!(task.is_notified()); + let entry = assert_ready!(queue).unwrap().into_inner(); + assert_eq!(entry, "bar"); + + }) +}