From 6cf1a5b6b8686e5bde107d072d77199aaefcb2ec Mon Sep 17 00:00:00 2001 From: Christofer Nolander Date: Thu, 26 Mar 2020 20:54:56 +0100 Subject: [PATCH] time: fix DelayQueue rewriting delay on insert after Poll::Ready (#2285) When the queue was polled and yielded an index from the wheel, the delay until the next item was never updated. As a result, when one item was yielded from `poll_idx` the following insert erronously updated the delay to the instant of the inserted item. Fixes: #1700 --- tokio/src/time/delay_queue.rs | 11 ++++++----- tokio/tests/time_delay_queue.rs | 33 +++++++++++++++++++++++++++++++++ 2 files changed, 39 insertions(+), 5 deletions(-) diff --git a/tokio/src/time/delay_queue.rs b/tokio/src/time/delay_queue.rs index 1790ada8f..821c0c27c 100644 --- a/tokio/src/time/delay_queue.rs +++ b/tokio/src/time/delay_queue.rs @@ -721,15 +721,16 @@ impl DelayQueue { self.poll = wheel::Poll::new(now); } - self.delay = None; + // We poll the wheel to get the next value out before finding the next deadline. + let wheel_idx = self.wheel.poll(&mut self.poll, &mut self.slab); - if let Some(idx) = self.wheel.poll(&mut self.poll, &mut self.slab) { + self.delay = self.next_deadline().map(delay_until); + + if let Some(idx) = wheel_idx { return Poll::Ready(Some(Ok(idx))); } - if let Some(deadline) = self.next_deadline() { - self.delay = Some(delay_until(deadline)); - } else { + if self.delay.is_none() { return Poll::Ready(None); } } diff --git a/tokio/tests/time_delay_queue.rs b/tokio/tests/time_delay_queue.rs index 32e812e03..214b9ebee 100644 --- a/tokio/tests/time_delay_queue.rs +++ b/tokio/tests/time_delay_queue.rs @@ -410,6 +410,39 @@ async fn insert_before_first_after_poll() { assert_eq!(entry, "two"); } +#[tokio::test] +async fn insert_after_ready_poll() { + time::pause(); + + let mut queue = task::spawn(DelayQueue::new()); + + let now = Instant::now(); + + queue.insert_at("1", now + ms(100)); + queue.insert_at("2", now + ms(100)); + queue.insert_at("3", now + ms(100)); + + assert_pending!(poll!(queue)); + + delay_for(ms(100)).await; + + assert!(queue.is_woken()); + + let mut res = vec![]; + + while res.len() < 3 { + let entry = assert_ready_ok!(poll!(queue)); + res.push(entry.into_inner()); + queue.insert_at("foo", now + ms(500)); + } + + res.sort(); + + assert_eq!("1", res[0]); + assert_eq!("2", res[1]); + assert_eq!("3", res[2]); +} + fn ms(n: u64) -> Duration { Duration::from_millis(n) }