mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-07 00:00:09 +02:00
time: wake DelayQueue after resetting to expired (#8274)
Resetting an item into the expired stack bypassed the existing Sleep reset path, leaving a pending consumer asleep until an unrelated deadline.
This commit is contained in:
@@ -864,11 +864,19 @@ impl<T> DelayQueue<T> {
|
|||||||
self.slab[*key].expired = false;
|
self.slab[*key].expired = false;
|
||||||
|
|
||||||
self.insert_idx(when, *key);
|
self.insert_idx(when, *key);
|
||||||
|
let inserted_expired = self.slab[*key].expired;
|
||||||
|
|
||||||
let next_deadline = self.next_deadline();
|
let next_deadline = self.next_deadline();
|
||||||
if let (Some(ref mut delay), Some(deadline)) = (&mut self.delay, next_deadline) {
|
match (next_deadline, &mut self.delay) {
|
||||||
// This should awaken us if necessary (ie, if already expired)
|
(None, _) => self.delay = None,
|
||||||
delay.as_mut().reset(deadline);
|
(Some(deadline), Some(delay)) => delay.as_mut().reset(deadline),
|
||||||
|
(Some(deadline), None) => self.delay = Some(Box::pin(sleep_until(deadline))),
|
||||||
|
}
|
||||||
|
|
||||||
|
if inserted_expired {
|
||||||
|
if let Some(waker) = self.waker.take() {
|
||||||
|
waker.wake();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -202,6 +202,27 @@ async fn reset_entry() {
|
|||||||
assert!(entry.is_none())
|
assert!(entry.is_none())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn reset_to_past_wakes_pending_queue() {
|
||||||
|
time::pause();
|
||||||
|
|
||||||
|
let mut queue = task::spawn(DelayQueue::new());
|
||||||
|
let key = queue.insert("foo", ms(10_000));
|
||||||
|
|
||||||
|
assert_pending!(poll!(queue));
|
||||||
|
assert!(!queue.is_woken());
|
||||||
|
|
||||||
|
queue.reset_at(&key, Instant::now() - ms(100));
|
||||||
|
|
||||||
|
assert!(queue.is_woken());
|
||||||
|
|
||||||
|
let entry = assert_ready_some!(poll!(queue));
|
||||||
|
assert_eq!(*entry.get_ref(), "foo");
|
||||||
|
|
||||||
|
let entry = assert_ready!(poll!(queue));
|
||||||
|
assert!(entry.is_none());
|
||||||
|
}
|
||||||
|
|
||||||
// Reproduces tokio-rs/tokio#849.
|
// Reproduces tokio-rs/tokio#849.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn reset_much_later() {
|
async fn reset_much_later() {
|
||||||
|
|||||||
Reference in New Issue
Block a user