From d822b721b4c88a693dc04d281c02b61f857158d5 Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Mon, 20 Aug 2018 21:47:10 -0700 Subject: [PATCH] Add DelayQueue implementation to tokio-timer (#550) This patch adds a `DelayQueue` to tokio_timer. The `DelayQueue` allows inserting elements as well as specifying a time at which the element should be returned to the user. This allows handling more complex timeout situations. --- src/timer.rs | 6 + tokio-timer/Cargo.toml | 5 + tokio-timer/src/delay_queue.rs | 835 ++++++++++++++++++++++ tokio-timer/src/lib.rs | 39 +- tokio-timer/src/timer/atomic_stack.rs | 122 ++++ tokio-timer/src/timer/entry.rs | 239 +------ tokio-timer/src/timer/mod.rs | 289 ++------ tokio-timer/src/timer/stack.rs | 123 ++++ tokio-timer/src/{timer => wheel}/level.rs | 43 +- tokio-timer/src/wheel/mod.rs | 289 ++++++++ tokio-timer/src/wheel/stack.rs | 26 + tokio-timer/tests/interval.rs | 10 +- tokio-timer/tests/queue.rs | 223 ++++++ tokio-timer/tests/support/mod.rs | 72 +- 14 files changed, 1803 insertions(+), 518 deletions(-) create mode 100644 tokio-timer/src/delay_queue.rs create mode 100644 tokio-timer/src/timer/atomic_stack.rs create mode 100644 tokio-timer/src/timer/stack.rs rename tokio-timer/src/{timer => wheel}/level.rs (85%) create mode 100644 tokio-timer/src/wheel/mod.rs create mode 100644 tokio-timer/src/wheel/stack.rs create mode 100644 tokio-timer/tests/queue.rs diff --git a/src/timer.rs b/src/timer.rs index 2c4c7569a..79fa790d6 100644 --- a/src/timer.rs +++ b/src/timer.rs @@ -14,6 +14,9 @@ //! a specified `Instant` in time. If the future does not complete in time, //! then it is canceled and an error is returned. //! +//! * [`DelayQueue`]: A queue where items are returned once the requested delay +//! has expired. +//! //! These types are sufficient for handling a large number of scenarios //! involving time. //! @@ -79,10 +82,13 @@ //! [Deadline]: struct.Deadline.html //! [Delay]: struct.Delay.html //! [Interval]: struct.Interval.html +//! [`DelayQueue`]: struct.DelayQueue.html pub use tokio_timer::{ + delay_queue, Deadline, DeadlineError, + DelayQueue, Error, Interval, Delay, diff --git a/tokio-timer/Cargo.toml b/tokio-timer/Cargo.toml index 401097cca..32ef2e88b 100644 --- a/tokio-timer/Cargo.toml +++ b/tokio-timer/Cargo.toml @@ -19,5 +19,10 @@ Timer facilities for Tokio futures = "0.1.19" tokio-executor = { version = "0.1.1", path = "../tokio-executor" } +# Backs `DelayQueue` +slab = "0.4.1" + [dev-dependencies] rand = "0.5" +tokio-mock-task = "0.1.0" +tokio = { version = "0.1.7", path = "../" } diff --git a/tokio-timer/src/delay_queue.rs b/tokio-timer/src/delay_queue.rs new file mode 100644 index 000000000..3f30f4475 --- /dev/null +++ b/tokio-timer/src/delay_queue.rs @@ -0,0 +1,835 @@ +//! A queue of delayed elements. +//! +//! See [`DelayQueue`] for more details. +//! +//! [`DelayQueue`]: struct.DelayQueue.html + +use {Error, Delay}; +use clock::now; +use wheel::{self, Wheel}; +use timer::Handle; + +use futures::{Future, Stream, Poll}; +use slab::Slab; + +use std::cmp; +use std::marker::PhantomData; +use std::time::{Duration, Instant}; + +/// A queue of delayed elements. +/// +/// Once an element is inserted into the `DelayQueue`, it is yielded once the +/// specified deadline has been reached. +/// +/// # Usage +/// +/// Elements are inserted into `DelayQueue` using the [`insert`] or +/// [`insert_at`] methods. A deadline is provided with the item and a [`Key`] is +/// returned. The key is used to remove the entry or to change the deadline at +/// which it should be yielded back. +/// +/// Once delays have been configured, the `DelayQueue` is used via its +/// [`Stream`] implementation. [`poll`] is called. If an entry has reached its +/// deadline, it is returned. If not, `Async::NotReady` indicating that the +/// current task will be notified once the deadline has been reached. +/// +/// # `Stream` implementation +/// +/// Items are retrieved from the queue via [`Stream::poll`]. If no delays have +/// expired, no items are returned. In this case, `NotReady` is returned and the +/// current task is registered to be notified once the next item's delay has +/// expired. +/// +/// If no items are in the queue, i.e. `is_empty()` returns `true`, then `poll` +/// returns `Ready(None)`. This indicates that the stream has reached an end. +/// However, if a new item is inserted *after*, `poll` will once again start +/// returning items or `NotReady. +/// +/// Items are returned ordered by their expirations. Items that are configured +/// to expire first will be returned first. There are no ordering guarantees +/// for items configured to expire the same instant. Also note that delays are +/// rounded to the closest millisecond. +/// +/// # Implementation +/// +/// The `DelayQueue` is backed by the same hashed timing wheel implementation as +/// [`Timer`] as such, it offers the same performance benefits. See [`Timer`] +/// for further implementation notes. +/// +/// State associated with each entry is stored in a [`slab`]. This allows +/// amortizing the cost of allocation. Space created for expired entries is +/// reused when inserting new entries. +/// +/// Capacity can be checked using [`capacity`] and allocated preemptively by using +/// the [`reserve`] method. +/// +/// # Usage +/// +/// Using `DelayQueue` to manage cache entries. +/// +/// ```rust +/// #[macro_use] +/// extern crate futures; +/// extern crate tokio; +/// # type CacheKey = String; +/// # type Value = String; +/// use tokio::timer::{delay_queue, DelayQueue, Error}; +/// use futures::{Async, Poll, Stream}; +/// use std::collections::HashMap; +/// use std::time::Duration; +/// +/// struct Cache { +/// entries: HashMap, +/// expirations: DelayQueue, +/// } +/// +/// const TTL_SECS: u64 = 30; +/// +/// impl Cache { +/// fn insert(&mut self, key: CacheKey, value: Value) { +/// let delay = self.expirations +/// .insert(key.clone(), Duration::from_secs(TTL_SECS)); +/// +/// self.entries.insert(key, (value, delay)); +/// } +/// +/// fn get(&self, key: &CacheKey) -> Option<&Value> { +/// self.entries.get(key) +/// .map(|&(ref v, _)| v) +/// } +/// +/// fn remove(&mut self, key: &CacheKey) { +/// if let Some((_, cache_key)) = self.entries.remove(key) { +/// self.expirations.remove(&cache_key); +/// } +/// } +/// +/// fn poll_purge(&mut self) -> Poll<(), Error> { +/// while let Some(entry) = try_ready!(self.expirations.poll()) { +/// self.entries.remove(entry.get_ref()); +/// } +/// +/// Ok(Async::Ready(())) +/// } +/// } +/// # fn main() {} +/// ``` +/// +/// [`insert`]: #method.insert +/// [`insert_at`]: #method.insert_at +/// [`Key`]: struct.Key.html +/// [`Stream`]: https://docs.rs/futures/0.1/futures/stream/trait.Stream.html +/// [`poll`]: #method.poll +/// [`Stream::poll`]: #method.poll +/// [`Timer`]: ../struct.Timer.html +/// [`slab`]: https://docs.rs/slab +/// [`capacity`]: #method.capacity +/// [`reserve`]: #method.reserve +#[derive(Debug)] +pub struct DelayQueue { + /// Handle to the timer driving the `DelayQueue` + handle: Handle, + + /// Stores data associated with entries + slab: Slab>, + + /// Lookup structure tracking all delays in the queue + wheel: Wheel>, + + /// Delays that were inserted when already expired. These cannot be stored + /// in the wheel + expired: Stack, + + /// Delay expiring when the *first* item in the queue expires + delay: Option, + + /// Wheel polling state + poll: wheel::Poll, + + /// Instant at which the timer starts + start: Instant, +} + +/// An entry in `DelayQueue` that has expired and removed. +/// +/// Values are returned by `DelayQueue::poll`. +#[derive(Debug)] +pub struct Expired { + /// The data stored in the queue + data: T, + + /// The expiration time + deadline: Instant, + + /// The key associated with the entry + key: Key, +} + +/// Token to a value stored in a `DelayQueue`. +/// +/// Instances of `Key` are returned by `DelayQueue::insert`. See [`DelayQueue`] +/// documentation for more details. +/// +/// [`DelayQueue`]: struct.DelayQueue.html +#[derive(Debug)] +pub struct Key { + index: usize, +} + +#[derive(Debug)] +struct Stack { + /// Head of the stack + head: Option, + _p: PhantomData, +} + +#[derive(Debug)] +struct Data { + /// The data being stored in the queue and will be returned at the requested + /// instant. + inner: T, + + /// The instant at which the item is returned. + when: u64, + + /// Set to true when stored in the `expired` queue + expired: bool, + + /// Next entry in the stack + next: Option, + + /// Previous entry in the stac + prev: Option, +} + +/// Maximum number of entries the queue can handle +const MAX_ENTRIES: usize = (1 << 30) - 1; + +impl DelayQueue { + /// Create a new, empty, `DelayQueue` + /// + /// The queue will not allocate storage until items are inserted into it. + /// + /// # Examples + /// + /// ```rust + /// # use tokio_timer::DelayQueue; + /// let delay_queue: DelayQueue = DelayQueue::new(); + /// ``` + pub fn new() -> DelayQueue { + DelayQueue::with_capacity(0) + } + + /// Create a new, empty, `DelayQueue` backed by the specified timer. + /// + /// The queue will not allocate storage until items are inserted into it. + /// + /// # Examples + /// + /// ```rust,no_run + /// # use tokio_timer::DelayQueue; + /// use tokio_timer::timer::Handle; + /// + /// let handle = Handle::default(); + /// let deplay_queue: DelayQueue = DelayQueue::with_capacity_and_handle(0, &handle); + /// ``` + pub fn with_capacity_and_handle(capacity: usize, handle: &Handle) -> DelayQueue { + DelayQueue { + handle: handle.clone(), + wheel: Wheel::new(), + slab: Slab::with_capacity(capacity), + expired: Stack::default(), + delay: None, + poll: wheel::Poll::new(0), + start: now(), + } + } + + /// Create a new, empty, `DelayQueue` with the specified capacity. + /// + /// The queue will be able to hold at least `capacity` elements without + /// reallocating. If `capacity` is 0, the queue will not allocate for + /// storage. + /// + /// # Examples + /// + /// ```rust + /// # use tokio_timer::DelayQueue; + /// # use std::time::Duration; + /// let mut delay_queue = DelayQueue::with_capacity(10); + /// + /// // These insertions are done without further allocation + /// for i in 0..10 { + /// delay_queue.insert(i, Duration::from_secs(i)); + /// } + /// + /// // This will make the queue allocate additional storage + /// delay_queue.insert(11, Duration::from_secs(11)); + /// ``` + pub fn with_capacity(capacity: usize) -> DelayQueue { + DelayQueue::with_capacity_and_handle(capacity, &Handle::default()) + } + + /// Insert `value` into the queue set to expire at a specific instant in + /// time. + /// + /// This function is identical to `insert`, but takes an `Instant` instead + /// of a `Duration`. + /// + /// `value` is stored in the queue until `when` is reached. At which point, + /// `value` will be returned from [`poll`]. If `when` has already been + /// reached, then `value` is immediately made available to poll. + /// + /// The return value represents the insertion and is used at an argument to + /// [`remove`] and [`reset`]. Note that [`Key`] is token and is reused once + /// `value` is removed from the queue eitheer by calling [`poll`] after + /// `when` is reached or by calling [`remove`]. At this point, the caller + /// must take care to not use the returned [`Key`] again as it may reference + /// a different item in the queue. + /// + /// See [type] level documentation for more details. + /// + /// # Panics + /// + /// This function panics if `when` is too far in the future. + /// + /// # Examples + /// + /// Basic usage + /// + /// ```rust + /// # extern crate tokio; + /// use tokio::timer::DelayQueue; + /// use std::time::{Instant, Duration}; + /// + /// # fn main() { + /// let mut delay_queue = DelayQueue::new(); + /// let key = delay_queue.insert_at( + /// "foo", Instant::now() + Duration::from_secs(5)); + /// + /// // Remove the entry + /// let item = delay_queue.remove(&key); + /// assert_eq!(*item.get_ref(), "foo"); + /// # } + /// ``` + /// + /// [`poll`]: #method.poll + /// [`remove`]: #method.remove + /// [`reset`]: #method.reset + /// [`Key`]: struct.Key.html + /// [type]: # + pub fn insert_at(&mut self, value: T, when: Instant) -> Key { + assert!(self.slab.len() < MAX_ENTRIES, "max entries exceeded"); + + // Normalize the deadline. Values cannot be set to expire in the past. + let when = self.normalize_deadline(when); + + // Insert the value in the store + let key = self.slab.insert(Data { + inner: value, + when, + expired: false, + next: None, + prev: None, + }); + + self.insert_idx(when, key); + + Key::new(key) + } + + /// Insert `value` into the queue set to expire after the requested duration + /// elapses. + /// + /// This function is identical to `insert_at`, but takes a `Duration` + /// instead of an `Instant`. + /// + /// `value` is stored in the queue until `when` is reached. At which point, + /// `value` will be returned from [`poll`]. If `when` has already been + /// reached, then `value` is immediately made available to poll. + /// + /// The return value represents the insertion and is used at an argument to + /// [`remove`] and [`reset`]. Note that [`Key`] is token and is reused once + /// `value` is removed from the queue eitheer by calling [`poll`] after + /// `when` is reached or by calling [`remove`]. At this point, the caller + /// must take care to not use the returned [`Key`] again as it may reference + /// a different item in the queue. + /// + /// See [type] level documentation for more details. + /// + /// # Panics + /// + /// This function panics if `timeout` is greater than the maximum supported + /// duration. + /// + /// # Examples + /// + /// Basic usage + /// + /// ```rust + /// # extern crate tokio; + /// use tokio::timer::DelayQueue; + /// use std::time::Duration; + /// + /// # fn main() { + /// let mut delay_queue = DelayQueue::new(); + /// let key = delay_queue.insert("foo", Duration::from_secs(5)); + /// + /// // Remove the entry + /// let item = delay_queue.remove(&key); + /// assert_eq!(*item.get_ref(), "foo"); + /// # } + /// ``` + /// + /// [`poll`]: #method.poll + /// [`remove`]: #method.remove + /// [`reset`]: #method.reset + /// [`Key`]: struct.Key.html + /// [type]: # + pub fn insert(&mut self, value: T, timeout: Duration) -> Key { + self.insert_at(value, now() + timeout) + } + + fn insert_idx(&mut self, when: u64, key: usize) { + use self::wheel::{InsertError, Stack}; + + // Register the deadline with the timer wheel + match self.wheel.insert(when, key, &mut self.slab) { + Ok(_) => {} + Err((_, InsertError::Elapsed)) => { + self.slab[key].expired = true; + // The delay is already expired, store it in the expired queue + self.expired.push(key, &mut self.slab); + } + Err((_, err)) => { + panic!("invalid deadline; err={:?}", err) + } + } + } + + /// Remove the item associated with `key` from the queue. + /// + /// There must be an item associated with `key`. The function returns the + /// removed item as well as the `Instant` at which it will the delay will + /// have expired. + /// + /// # Panics + /// + /// The function panics if `key` is not contained by the queue. + /// + /// # Examples + /// + /// Basic usage + /// + /// ```rust + /// # extern crate tokio; + /// use tokio::timer::DelayQueue; + /// use std::time::Duration; + /// + /// # fn main() { + /// let mut delay_queue = DelayQueue::new(); + /// let key = delay_queue.insert("foo", Duration::from_secs(5)); + /// + /// // Remove the entry + /// let item = delay_queue.remove(&key); + /// assert_eq!(*item.get_ref(), "foo"); + /// # } + /// ``` + pub fn remove(&mut self, key: &Key) -> Expired { + use wheel::Stack; + + // Special case the `expired` queue + if self.slab[key.index].expired { + self.expired.remove(&key.index, &mut self.slab); + } else { + self.wheel.remove(&key.index, &mut self.slab); + } + + let data = self.slab.remove(key.index); + + Expired { + key: Key::new(key.index), + data: data.inner, + deadline: self.start + Duration::from_millis(data.when), + } + } + + /// Sets the delay of the item associated with `key` to expire at `when`. + /// + /// This function is identical to `reset` but takes an `Instant` instead of + /// a `Duration`. + /// + /// The item remains in the queue but the delay is set to expire at `when`. + /// If `when` is in the past, then the item is immediately made available to + /// the caller. + /// + /// # Panics + /// + /// This function panics if `when` is too far in the future or if `key` is + /// not contained by the queue. + /// + /// # Examples + /// + /// Basic usage + /// + /// ```rust + /// # extern crate tokio; + /// use tokio::timer::DelayQueue; + /// use std::time::{Duration, Instant}; + /// + /// # fn main() { + /// let mut delay_queue = DelayQueue::new(); + /// let key = delay_queue.insert("foo", Duration::from_secs(5)); + /// + /// // "foo" is scheduled to be returned in 5 seconds + /// + /// delay_queue.reset_at(&key, Instant::now() + Duration::from_secs(10)); + /// + /// // "foo"is now scheduledto be returned in 10 seconds + /// # } + /// ``` + pub fn reset_at(&mut self, key: &Key, when: Instant) { + self.wheel.remove(&key.index, &mut self.slab); + + // Normalize the deadline. Values cannot be set to expire in the past. + let when = self.normalize_deadline(when); + let old = self.start + Duration::from_millis(self.slab[key.index].when); + + + self.slab[key.index].when = when; + + if let Some(ref mut delay) = self.delay { + debug_assert!(old >= delay.deadline()); + + if old == delay.deadline() { + delay.reset(self.start + Duration::from_millis(when)); + } + } + + self.insert_idx(when, key.index); + } + + /// Sets the delay of the item associated with `key` to expire after + /// `timeout`. + /// + /// This function is identical to `reset_at` but takes a `Duration` instead + /// of an `Instant`. + /// + /// The item remains in the queue but the delay is set to expire after + /// `timeout`. If `timeout` is zero, then the item is immediately made + /// available to the caller. + /// + /// # Panics + /// + /// This function panics if `timeout` is greater than the maximum supported + /// duration or if `key` is not contained by the queue. + /// + /// # Examples + /// + /// Basic usage + /// + /// ```rust + /// # extern crate tokio; + /// use tokio::timer::DelayQueue; + /// use std::time::Duration; + /// + /// # fn main() { + /// let mut delay_queue = DelayQueue::new(); + /// let key = delay_queue.insert("foo", Duration::from_secs(5)); + /// + /// // "foo" is scheduled to be returned in 5 seconds + /// + /// delay_queue.reset(&key, Duration::from_secs(10)); + /// + /// // "foo"is now scheduledto be returned in 10 seconds + /// # } + /// ``` + pub fn reset(&mut self, key: &Key, timeout: Duration) { + self.reset_at(key, now() + timeout); + } + + /// Clears the queue, removing all items. + /// + /// After calling `clear`, [`poll`] will return `Ok(Ready(None))`. + /// + /// Note that this method has no effect on the allocated capacity. + /// + /// [`poll`]: #method.poll + /// + /// # Examples + /// + /// ```rust + /// # extern crate tokio; + /// use tokio::timer::DelayQueue; + /// use std::time::Duration; + /// + /// # fn main() { + /// let mut delay_queue = DelayQueue::new(); + /// + /// delay_queue.insert("foo", Duration::from_secs(5)); + /// + /// assert!(!delay_queue.is_empty()); + /// + /// delay_queue.clear(); + /// + /// assert!(delay_queue.is_empty()); + /// # } + /// ``` + pub fn clear(&mut self) { + self.slab.clear(); + self.expired = Stack::default(); + self.wheel = Wheel::new(); + self.delay = None; + } + + /// Returns the number of elements the queue can hold without reallocating. + /// + /// # Examples + /// + /// ```rust + /// # use tokio_timer::DelayQueue; + /// let delay_queue: DelayQueue = DelayQueue::with_capacity(10); + /// assert_eq!(delay_queue.capacity(), 10); + /// ``` + pub fn capacity(&self) -> usize { + self.slab.capacity() + } + + /// Reserve capacity for at least `additional` more items to be queued + /// without allocating. + /// + /// `reserve` does nothing if the queue already has sufficient capacity for + /// `additional` more values. If more capacity is required, a new segment of + /// memory will be allocated and all existing values will be copied into it. + /// As such, if the queue is already very large, a call to `reserve` can end + /// up being expensive. + /// + /// The queue may reserve more than `additional` extra space in order to + /// avoid frequent reallocations. + /// + /// # Panics + /// + /// Panics if the new capacity exceeds the maximum number of entries the + /// queue can contain. + /// + /// # Examples + /// + /// ``` + /// # use tokio_timer::DelayQueue; + /// # use std::time::Duration; + /// let mut delay_queue = DelayQueue::new(); + /// delay_queue.insert("hello", Duration::from_secs(10)); + /// delay_queue.reserve(10); + /// assert!(delay_queue.capacity() >= 11); + /// ``` + pub fn reserve(&mut self, additional: usize) { + self.slab.reserve(additional); + } + + /// Returns `true` if there are no items in the queue. + /// + /// Note that this function returns `false` even if all items have not yet + /// expired and a call to `poll` will return `NotReady`. + /// + /// # Examples + /// + /// ``` + /// # use tokio_timer::DelayQueue; + /// use std::time::Duration; + /// let mut delay_queue = DelayQueue::new(); + /// assert!(delay_queue.is_empty()); + /// + /// delay_queue.insert("hello", Duration::from_secs(5)); + /// assert!(!delay_queue.is_empty()); + /// ``` + pub fn is_empty(&self) -> bool { + self.slab.is_empty() + } + + /// Polls the queue, returning the index of the next slot in the slab that + /// should be returned. + /// + /// A slot should be returned when the associated deadline has been reached. + fn poll_idx(&mut self) -> Poll, Error> { + use self::wheel::Stack; + + let expired = self.expired.pop(&mut self.slab); + + if expired.is_some() { + return Ok(expired.into()); + } + + loop { + if let Some(ref mut delay) = self.delay { + if !delay.is_elapsed() { + try_ready!(delay.poll()); + } + + let now = ::ms(delay.deadline() - self.start, ::Round::Down); + + self.poll = wheel::Poll::new(now); + } + + self.delay = None; + + if let Some(idx) = self.wheel.poll(&mut self.poll, &mut self.slab) { + return Ok(Some(idx).into()); + } + + let deadline = match self.wheel.poll_at() { + Some(poll_at) => { + self.start + Duration::from_millis(poll_at) + } + None => return Ok(None.into()), + }; + + self.delay = Some(self.handle.delay(deadline)); + } + } + + fn normalize_deadline(&self, when: Instant) -> u64 { + let when = if when < self.start { + 0 + } else { + ::ms(when - self.start, ::Round::Up) + }; + + cmp::max(when, self.wheel.elapsed()) + } +} + +impl Stream for DelayQueue { + type Item = Expired; + type Error = Error; + + fn poll(&mut self) -> Poll, Error> { + let item = try_ready!(self.poll_idx()) + .map(|idx| { + let data = self.slab.remove(idx); + debug_assert!(data.next.is_none()); + debug_assert!(data.prev.is_none()); + + Expired { + key: Key::new(idx), + data: data.inner, + deadline: self.start + Duration::from_millis(data.when), + } + }); + + Ok(item.into()) + } +} + +impl wheel::Stack for Stack { + type Owned = usize; + type Borrowed = usize; + type Store = Slab>; + + fn is_empty(&self) -> bool { + self.head.is_none() + } + + fn push(&mut self, item: Self::Owned, store: &mut Self::Store) { + // Ensure the entry is not already in a stack. + debug_assert!(store[item].next.is_none()); + debug_assert!(store[item].prev.is_none()); + + // Remove the old head entry + let old = self.head.take(); + + if let Some(idx) = old { + store[idx].prev = Some(item); + } + + store[item].next = old; + self.head = Some(item) + } + + fn pop(&mut self, store: &mut Self::Store) -> Option { + if let Some(idx) = self.head { + self.head = store[idx].next; + + if let Some(idx) = self.head { + store[idx].prev = None; + } + + store[idx].next = None; + debug_assert!(store[idx].prev.is_none()); + + Some(idx) + } else { + None + } + } + + fn remove(&mut self, item: &Self::Borrowed, store: &mut Self::Store) { + assert!(store.contains(*item)); + + // Ensure that the entry is in fact contained by the stack + debug_assert!({ + // This walks the full linked list even if an entry is found. + let mut next = self.head; + let mut contains = false; + + while let Some(idx) = next { + if idx == *item { + debug_assert!(!contains); + contains = true; + } + + next = store[idx].next; + } + + contains + }); + + if let Some(next) = store[*item].next { + store[next].prev = store[*item].prev; + } + + if let Some(prev) = store[*item].prev { + store[prev].next = store[*item].next; + } else { + self.head = store[*item].next; + } + + store[*item].next = None; + store[*item].prev = None; + } + + fn when(item: &Self::Borrowed, store: &Self::Store) -> u64 { + store[*item].when + } +} + +impl Default for Stack { + fn default() -> Stack { + Stack { + head: None, + _p: PhantomData, + } + } +} + +impl Key { + pub(crate) fn new(index: usize) -> Key { + Key { index } + } +} + +impl Expired { + /// Returns a reference to the inner value. + pub fn get_ref(&self) -> &T { + &self.data + } + + /// Returns a mutable reference to the inner value. + pub fn get_mut(&mut self) -> &mut T { + &mut self.data + } + + /// Consumes `self` and returns the inner value. + pub fn into_inner(self) -> T { + self.data + } +} diff --git a/tokio-timer/src/lib.rs b/tokio-timer/src/lib.rs index 754f9aa88..b07c5c6d7 100644 --- a/tokio-timer/src/lib.rs +++ b/tokio-timer/src/lib.rs @@ -9,6 +9,9 @@ //! * [`Deadline`]: Wraps a future, requiring it to complete before a specified //! instant in time, erroring if the future takes too long. //! +//! * [`DelayQueue`]: A queue where items are returned once the requested delay +//! has expired. +//! //! These three types are backed by a [`Timer`] instance. In order for //! [`Delay`], [`Interval`], and [`Deadline`] to function, the associated //! [`Timer`] instance must be running on some thread. @@ -25,8 +28,10 @@ extern crate tokio_executor; #[macro_use] extern crate futures; +extern crate slab; pub mod clock; +pub mod delay_queue; pub mod timer; mod atomic; @@ -34,16 +39,46 @@ mod deadline; mod delay; mod error; mod interval; - -use std::time::{Duration, Instant}; +mod wheel; pub use self::deadline::{Deadline, DeadlineError}; +#[doc(inline)] +pub use self::delay_queue::DelayQueue; pub use self::delay::Delay; pub use self::error::Error; pub use self::interval::Interval; +#[doc(inline)] pub use self::timer::{with_default, Timer}; +use std::time::{Duration, Instant}; + /// Create a Future that completes in `duration` from now. pub fn sleep(duration: Duration) -> Delay { Delay::new(Instant::now() + duration) } + +// ===== Internal utils ===== + +enum Round { + Up, + Down, +} + +/// Convert a `Duration` to milliseconds, rounding up and saturating at +/// `u64::MAX`. +/// +/// The saturating is fine because `u64::MAX` milliseconds are still many +/// million years. +#[inline] +fn ms(duration: Duration, round: Round) -> u64 { + const NANOS_PER_MILLI: u32 = 1_000_000; + const MILLIS_PER_SEC: u64 = 1_000; + + // Round up. + let millis = match round { + Round::Up => (duration.subsec_nanos() + NANOS_PER_MILLI - 1) / NANOS_PER_MILLI, + Round::Down => duration.subsec_nanos() / NANOS_PER_MILLI, + }; + + duration.as_secs().saturating_mul(MILLIS_PER_SEC).saturating_add(millis as u64) +} diff --git a/tokio-timer/src/timer/atomic_stack.rs b/tokio-timer/src/timer/atomic_stack.rs new file mode 100644 index 000000000..81d817a90 --- /dev/null +++ b/tokio-timer/src/timer/atomic_stack.rs @@ -0,0 +1,122 @@ +use Error; +use super::Entry; + +use std::ptr; +use std::sync::Arc; +use std::sync::atomic::AtomicPtr; +use std::sync::atomic::Ordering::SeqCst; + +/// A stack of `Entry` nodes +#[derive(Debug)] +pub(crate) struct AtomicStack { + /// Stack head + head: AtomicPtr, +} + +/// Entries that were removed from the stack +#[derive(Debug)] +pub(crate) struct AtomicStackEntries { + ptr: *mut Entry, +} + +/// Used to indicate that the timer has shutdown. +const SHUTDOWN: *mut Entry = 1 as *mut _; + +impl AtomicStack { + pub fn new() -> AtomicStack { + AtomicStack { head: AtomicPtr::new(ptr::null_mut()) } + } + + /// Push an entry onto the stack. + /// + /// Returns `true` if the entry was pushed, `false` if the entry is already + /// on the stack, `Err` if the timer is shutdown. + pub fn push(&self, entry: &Arc) -> Result { + // First, set the queued bit on the entry + let queued = entry.queued.fetch_or(true, SeqCst).into(); + + if queued { + // Already queued, nothing more to do + return Ok(false); + } + + let ptr = Arc::into_raw(entry.clone()) as *mut _; + + let mut curr = self.head.load(SeqCst); + + loop { + if curr == SHUTDOWN { + // Don't leak the entry node + let _ = unsafe { Arc::from_raw(ptr) }; + + return Err(Error::shutdown()); + } + + // Update the `next` pointer. This is safe because setting the queued + // bit is a "lock" on this field. + unsafe { + *(entry.next_atomic.get()) = curr; + } + + let actual = self.head.compare_and_swap(curr, ptr, SeqCst); + + if actual == curr { + break; + } + + curr = actual; + } + + Ok(true) + } + + /// Take all entries from the stack + pub fn take(&self) -> AtomicStackEntries { + let ptr = self.head.swap(ptr::null_mut(), SeqCst); + AtomicStackEntries { ptr } + } + + /// Drain all remaining nodes in the stack and prevent any new nodes from + /// being pushed onto the stack. + pub fn shutdown(&self) { + // Shutdown the processing queue + let ptr = self.head.swap(SHUTDOWN, SeqCst); + + // Let the drop fn of `AtomicStackEntries` handle draining the stack + drop(AtomicStackEntries { ptr }); + } +} + +// ===== impl AtomicStackEntries ===== + +impl Iterator for AtomicStackEntries { + type Item = Arc; + + fn next(&mut self) -> Option { + if self.ptr.is_null() { + return None; + } + + // Convert the pointer to an `Arc` + let entry = unsafe { Arc::from_raw(self.ptr) }; + + // Update `self.ptr` to point to the next element of the stack + self.ptr = unsafe { (*entry.next_atomic.get()) }; + + // Unset the queued flag + let res = entry.queued.fetch_and(false, SeqCst); + debug_assert!(res); + + // Return the entry + Some(entry) + } +} + +impl Drop for AtomicStackEntries { + fn drop(&mut self) { + while let Some(entry) = self.next() { + // Flag the entry as errored + entry.error(); + } + } +} diff --git a/tokio-timer/src/timer/entry.rs b/tokio-timer/src/timer/entry.rs index 613eca1bf..7b2bac495 100644 --- a/tokio-timer/src/timer/entry.rs +++ b/tokio-timer/src/timer/entry.rs @@ -8,7 +8,7 @@ use futures::task::AtomicTask; use std::cell::UnsafeCell; use std::ptr; use std::sync::{Arc, Weak}; -use std::sync::atomic::{AtomicBool, AtomicPtr}; +use std::sync::atomic::AtomicBool; use std::sync::atomic::Ordering::SeqCst; use std::time::Instant; use std::u64; @@ -56,12 +56,12 @@ pub(crate) struct Entry { /// True when the entry is queued in the "process" stack. This value /// is set before pushing the value and unset after popping the value. - queued: AtomicBool, + pub(super) queued: AtomicBool, /// Next entry in the "process" linked list. /// /// Represents a strong Arc ref. - next_atomic: UnsafeCell<*mut Entry>, + pub(super) next_atomic: UnsafeCell<*mut Entry>, /// When the entry expires, relative to the `start` of the timer /// (Inner::start). This is only used by the timer. @@ -80,7 +80,7 @@ pub(crate) struct Entry { /// Next entry in the State's linked list. /// /// This is only accessed by the timer - next_stack: UnsafeCell>>, + pub(super) next_stack: UnsafeCell>>, /// Previous entry in the State's linked list. /// @@ -88,25 +88,7 @@ pub(crate) struct Entry { /// entry. /// /// This is a weak reference. - prev_stack: UnsafeCell<*const Entry>, -} - -/// A doubly linked stack -pub(crate) struct Stack { - head: Option>, -} - -/// A stack of `Entry` nodes -#[derive(Debug)] -pub(crate) struct AtomicStack { - /// Stack head - head: AtomicPtr, -} - -/// Entries that were removed from the stack -#[derive(Debug)] -pub(crate) struct AtomicStackEntries { - ptr: *mut Entry, + pub(super) prev_stack: UnsafeCell<*const Entry>, } /// Flag indicating a timer entry has elapsed @@ -115,9 +97,6 @@ const ELAPSED: u64 = 1 << 63; /// Flag indicating a timer entry has reached an error state const ERROR: u64 = u64::MAX; -/// Used to indicate that the timer has shutdown. -const SHUTDOWN: *mut Entry = 1 as *mut _; - // ===== impl Entry ===== impl Entry { @@ -349,211 +328,3 @@ impl Drop for Entry { unsafe impl Send for Entry {} unsafe impl Sync for Entry {} - -// ===== impl Stack ===== - -impl Stack { - pub fn new() -> Stack { - Stack { head: None } - } - - pub fn is_empty(&self) -> bool { - self.head.is_none() - } - - /// Push an entry to the head of the linked list - pub fn push(&mut self, entry: Arc) { - // Get a pointer to the entry to for the prev link - let ptr: *const Entry = &*entry as *const _; - - // Remove the old head entry - let old = self.head.take(); - - unsafe { - // Ensure the entry is not already in a stack. - debug_assert!((*entry.next_stack.get()).is_none()); - debug_assert!((*entry.prev_stack.get()).is_null()); - - if let Some(ref entry) = old.as_ref() { - debug_assert!({ - // The head is not already set to the entry - ptr != &***entry as *const _ - }); - - // Set the previous link on the old head - *entry.prev_stack.get() = ptr; - } - - // Set this entry's next pointer - *entry.next_stack.get() = old; - - } - - // Update the head pointer - self.head = Some(entry); - } - - /// Pop the head of the linked list - pub fn pop(&mut self) -> Option> { - let entry = self.head.take(); - - unsafe { - if let Some(entry) = entry.as_ref() { - self.head = (*entry.next_stack.get()).take(); - - if let Some(entry) = self.head.as_ref() { - *entry.prev_stack.get() = ptr::null(); - } - - *entry.prev_stack.get() = ptr::null(); - } - } - - entry - } - - /// Remove the entry from the linked list - /// - /// The caller must ensure that the entry actually is contained by the list. - pub fn remove(&mut self, entry: &Entry) { - unsafe { - // Ensure that the entry is in fact contained by the stack - debug_assert!({ - // This walks the full linked list even if an entry is found. - let mut next = self.head.as_ref(); - let mut contains = false; - - while let Some(n) = next { - if entry as *const _ == &**n as *const _ { - debug_assert!(!contains); - contains = true; - } - - next = (*n.next_stack.get()).as_ref(); - } - - contains - }); - - // Unlink `entry` from the next node - let next = (*entry.next_stack.get()).take(); - - if let Some(next) = next.as_ref() { - (*next.prev_stack.get()) = *entry.prev_stack.get(); - } - - // Unlink `entry` from the prev node - - if let Some(prev) = (*entry.prev_stack.get()).as_ref() { - *prev.next_stack.get() = next; - } else { - // It is the head - self.head = next; - } - - // Unset the prev pointer - *entry.prev_stack.get() = ptr::null(); - } - } -} - -// ===== impl AtomicStack ===== - -impl AtomicStack { - pub fn new() -> AtomicStack { - AtomicStack { head: AtomicPtr::new(ptr::null_mut()) } - } - - /// Push an entry onto the stack. - /// - /// Returns `true` if the entry was pushed, `false` if the entry is already - /// on the stack, `Err` if the timer is shutdown. - pub fn push(&self, entry: &Arc) -> Result { - // First, set the queued bit on the entry - let queued = entry.queued.fetch_or(true, SeqCst).into(); - - if queued { - // Already queued, nothing more to do - return Ok(false); - } - - let ptr = Arc::into_raw(entry.clone()) as *mut _; - - let mut curr = self.head.load(SeqCst); - - loop { - if curr == SHUTDOWN { - // Don't leak the entry node - let _ = unsafe { Arc::from_raw(ptr) }; - - return Err(Error::shutdown()); - } - - // Update the `next` pointer. This is safe because setting the queued - // bit is a "lock" on this field. - unsafe { - *(entry.next_atomic.get()) = curr; - } - - let actual = self.head.compare_and_swap(curr, ptr, SeqCst); - - if actual == curr { - break; - } - - curr = actual; - } - - Ok(true) - } - - /// Take all entries from the stack - pub fn take(&self) -> AtomicStackEntries { - let ptr = self.head.swap(ptr::null_mut(), SeqCst); - AtomicStackEntries { ptr } - } - - /// Drain all remaining nodes in the stack and prevent any new nodes from - /// being pushed onto the stack. - pub fn shutdown(&self) { - // Shutdown the processing queue - let ptr = self.head.swap(SHUTDOWN, SeqCst); - - // Let the drop fn of `AtomicStackEntries` handle draining the stack - drop(AtomicStackEntries { ptr }); - } -} - -// ===== impl AtomicStackEntries ===== - -impl Iterator for AtomicStackEntries { - type Item = Arc; - - fn next(&mut self) -> Option { - if self.ptr.is_null() { - return None; - } - - // Convert the pointer to an `Arc` - let entry = unsafe { Arc::from_raw(self.ptr) }; - - // Update `self.ptr` to point to the next element of the stack - self.ptr = unsafe { (*entry.next_atomic.get()) }; - - // Unset the queued flag - let res = entry.queued.fetch_and(false, SeqCst); - debug_assert!(res); - - // Return the entry - Some(entry) - } -} - -impl Drop for AtomicStackEntries { - fn drop(&mut self) { - while let Some(entry) = self.next() { - // Flag the entry as errored - entry.error(); - } - } -} diff --git a/tokio-timer/src/timer/mod.rs b/tokio-timer/src/timer/mod.rs index b77fcca89..a15292136 100644 --- a/tokio-timer/src/timer/mod.rs +++ b/tokio-timer/src/timer/mod.rs @@ -31,15 +31,17 @@ // This allows the usage of the old `Now` trait. #![allow(deprecated)] +mod atomic_stack; mod entry; mod handle; -mod level; mod now; mod registration; +mod stack; +use self::atomic_stack::AtomicStack; use self::entry::Entry; +use self::stack::Stack; use self::handle::HandlePriv; -use self::level::{Level, Expiration}; pub use self::handle::{Handle, with_default}; pub use self::now::{Now, SystemNow}; @@ -47,6 +49,7 @@ pub(crate) use self::registration::Registration; use Error; use atomic::AtomicU64; +use wheel; use tokio_executor::park::{Park, Unpark, ParkThread}; @@ -125,20 +128,8 @@ pub struct Timer { /// Shared state inner: Arc, - /// The number of milliseconds elapsed since the timer started. - elapsed: u64, - - /// Timer wheel. - /// - /// Levels: - /// - /// * 1 ms slots / 64 ms range - /// * 64 ms slots / ~ 4 sec range - /// * ~ 4 sec slots / ~ 4 min range - /// * ~ 4 min slots / ~ 4 hr range - /// * ~ 4 hr slots / ~ 12 day range - /// * ~ 12 day slots / ~ 2 yr range - levels: Vec, + /// Timer wheel + wheel: wheel::Wheel, /// Thread parker. The `Timer` park implementation delegates to this. park: T, @@ -166,20 +157,12 @@ pub(crate) struct Inner { num: AtomicUsize, /// Head of the "process" linked list. - process: entry::AtomicStack, + process: AtomicStack, /// Unparks the timer thread. unpark: Box, } -/// Number of levels. Each level has 64 slots. By using 6 levels with 64 slots -/// each, the timer is able to track time up to 2 years into the future with a -/// precision of 1 millisecond. -const NUM_LEVELS: usize = 6; - -/// The maximum duration of a delay -const MAX_DURATION: u64 = 1 << (6 * NUM_LEVELS); - /// Maximum number of timeouts the system can handle concurrently. const MAX_TIMEOUTS: usize = usize::MAX >> 1; @@ -226,14 +209,9 @@ where T: Park, pub fn new_with_now(park: T, mut now: N) -> Self { let unpark = Box::new(park.unpark()); - let levels = (0..NUM_LEVELS) - .map(Level::new) - .collect(); - Timer { inner: Arc::new(Inner::new(now.now(), unpark)), - elapsed: 0, - levels, + wheel: wheel::Wheel::new(), park, now, } @@ -277,102 +255,29 @@ where T: Park, Ok(Turn(())) } - /// Returns the instant at which the next timeout expires. - fn next_expiration(&self) -> Option { - // Check all levels - for level in 0..NUM_LEVELS { - if let Some(expiration) = self.levels[level].next_expiration(self.elapsed) { - // There cannot be any expirations at a higher level that happen - // before this one. - debug_assert!({ - let mut res = true; - - for l2 in (level+1)..NUM_LEVELS { - if let Some(e2) = self.levels[l2].next_expiration(self.elapsed) { - if e2.deadline < expiration.deadline { - res = false; - } - } - } - - res - }); - - return Some(expiration); - } - } - - None - } - /// Converts an `Expiration` to an `Instant`. - fn expiration_instant(&self, expiration: &Expiration) -> Instant { - self.inner.start + Duration::from_millis(expiration.deadline) + fn expiration_instant(&self, when: u64) -> Instant { + self.inner.start + Duration::from_millis(when) } /// Run timer related logic fn process(&mut self) { - let now = ms(self.now.now() - self.inner.start, Round::Down); + let now = ::ms(self.now.now() - self.inner.start, ::Round::Down); + let mut poll = wheel::Poll::new(now); - loop { - let expiration = match self.next_expiration() { - Some(expiration) => expiration, - None => break, - }; + while let Some(entry) = self.wheel.poll(&mut poll, &mut ()) { + let when = entry.when_internal() + .expect("invalid internal entry state"); - if expiration.deadline > now { - // This expiration should not fire on this tick - break; - } + // Fire the entry + entry.fire(when); - // Process the slot, either moving it down a level or firing the - // timeout if currently at the final (boss) level. - self.process_expiration(&expiration); - - self.set_elapsed(expiration.deadline); + // Track that the entry has been fired + entry.set_when_internal(None); } - self.set_elapsed(now); - } - - fn set_elapsed(&mut self, when: u64) { - assert!(self.elapsed <= when, "elapsed={:?}; when={:?}", self.elapsed, when); - - if when > self.elapsed { - self.elapsed = when; - self.inner.elapsed.store(when, SeqCst); - } else { - assert_eq!(self.elapsed, when); - } - } - - fn process_expiration(&mut self, expiration: &Expiration) { - while let Some(entry) = self.pop_entry(expiration) { - if expiration.level == 0 { - let when = entry.when_internal() - .expect("invalid internal entry state"); - - debug_assert_eq!(when, expiration.deadline); - - // Fire the entry - entry.fire(when); - - // Track that the entry has been fired - entry.set_when_internal(None); - } else { - let when = entry.when_internal() - .expect("entry not tracked"); - - let next_level = expiration.level - 1; - - self.levels[next_level] - .add_entry(entry, when); - } - } - } - - fn pop_entry(&mut self, expiration: &Expiration) -> Option> { - self.levels[expiration.level].pop_entry_slot(expiration.slot) + // Update the elapsed cache + self.inner.elapsed.store(self.wheel.elapsed(), SeqCst); } /// Process the entry queue @@ -384,27 +289,24 @@ where T: Park, (None, None) => { // Nothing to do } - (Some(when), None) => { + (Some(_), None) => { // Remove the entry - self.clear_entry(&entry, when); + self.clear_entry(&entry); } (None, Some(when)) => { // Queue the entry self.add_entry(entry, when); } - (Some(curr), Some(next)) => { - self.clear_entry(&entry, curr); + (Some(_), Some(next)) => { + self.clear_entry(&entry); self.add_entry(entry, next); } } } } - fn clear_entry(&mut self, entry: &Arc, when: u64) { - // Get the level at which the entry should be stored - let level = self.level_for(when); - self.levels[level].remove_entry(entry, when); - + fn clear_entry(&mut self, entry: &Arc) { + self.wheel.remove(entry, &mut ()); entry.set_when_internal(None); } @@ -412,48 +314,26 @@ where T: Park, /// /// Returns `None` if the entry was fired. fn add_entry(&mut self, entry: Arc, when: u64) { - if when <= self.elapsed { - // The entry's deadline has elapsed, so fire it and update the - // internal state accordingly. - entry.set_when_internal(None); - entry.fire(when); - - return; - } else if when - self.elapsed > MAX_DURATION { - // The entry's deadline is invalid, so error it and update the - // internal state accordingly. - entry.set_when_internal(None); - entry.error(); - - return; - } - - // Get the level at which the entry should be stored - let level = self.level_for(when); + use wheel::InsertError; entry.set_when_internal(Some(when)); - self.levels[level].add_entry(entry, when); - debug_assert!({ - self.levels[level].next_expiration(self.elapsed) - .map(|e| e.deadline >= self.elapsed) - .unwrap_or(true) - }); + match self.wheel.insert(when, entry, &mut ()) { + Ok(_) => {} + Err((entry, InsertError::Elapsed)) => { + // The entry's deadline has elapsed, so fire it and update the + // internal state accordingly. + entry.set_when_internal(None); + entry.fire(when); + } + Err((entry, InsertError::Invalid)) => { + // The entry's deadline is invalid, so error it and update the + // internal state accordingly. + entry.set_when_internal(None); + entry.error(); + } + } } - - fn level_for(&self, when: u64) -> usize { - level_for(self.elapsed, when) - } -} - -fn level_for(elapsed: u64, when: u64) -> usize { - let masked = elapsed ^ when; - - assert!(masked != 0, "elapsed={}; when={}", elapsed, when); - - let leading_zeros = masked.leading_zeros() as usize; - let significant = 63 - leading_zeros; - significant / 6 } impl Default for Timer { @@ -476,10 +356,10 @@ where T: Park, fn park(&mut self) -> Result<(), Self::Error> { self.process_queue(); - match self.next_expiration() { - Some(expiration) => { + match self.wheel.poll_at() { + Some(when) => { let now = self.now.now(); - let deadline = self.expiration_instant(&expiration); + let deadline = self.expiration_instant(when); if deadline > now { self.park.park_timeout(deadline - now)?; @@ -500,10 +380,10 @@ where T: Park, fn park_timeout(&mut self, duration: Duration) -> Result<(), Self::Error> { self.process_queue(); - match self.next_expiration() { - Some(expiration) => { + match self.wheel.poll_at() { + Some(when) => { let now = self.now.now(); - let deadline = self.expiration_instant(&expiration); + let deadline = self.expiration_instant(when); if deadline > now { self.park.park_timeout(cmp::min(deadline - now, duration))?; @@ -524,9 +404,18 @@ where T: Park, impl Drop for Timer { fn drop(&mut self) { + use std::u64; + // Shutdown the stack of entries to process, preventing any new entries // from being pushed. self.inner.process.shutdown(); + + // Clear the wheel, using u64::MAX allows us to drain everything + let mut poll = wheel::Poll::new(u64::MAX); + + while let Some(entry) = self.wheel.poll(&mut poll, &mut ()) { + entry.error(); + } } } @@ -537,7 +426,7 @@ impl Inner { Inner { num: AtomicUsize::new(0), elapsed: AtomicU64::new(0), - process: entry::AtomicStack::new(), + process: AtomicStack::new(), start, unpark, } @@ -586,7 +475,7 @@ impl Inner { return 0; } - ms(deadline - self.start, Round::Up) + ::ms(deadline - self.start, ::Round::Up) } } @@ -596,59 +485,3 @@ impl fmt::Debug for Inner { .finish() } } - -enum Round { - Up, - Down, -} - -/// Convert a `Duration` to milliseconds, rounding up and saturating at -/// `u64::MAX`. -/// -/// The saturating is fine because `u64::MAX` milliseconds are still many -/// million years. -#[inline] -fn ms(duration: Duration, round: Round) -> u64 { - const NANOS_PER_MILLI: u32 = 1_000_000; - const MILLIS_PER_SEC: u64 = 1_000; - - // Round up. - let millis = match round { - Round::Up => (duration.subsec_nanos() + NANOS_PER_MILLI - 1) / NANOS_PER_MILLI, - Round::Down => duration.subsec_nanos() / NANOS_PER_MILLI, - }; - - duration.as_secs().saturating_mul(MILLIS_PER_SEC).saturating_add(millis as u64) -} - -#[cfg(test)] -mod test { - use super::*; - - #[test] - fn test_level_for() { - for pos in 1..64 { - assert_eq!(0, level_for(0, pos), "level_for({}) -- binary = {:b}", pos, pos); - } - - for level in 1..5 { - for pos in level..64 { - let a = pos * 64_usize.pow(level as u32); - assert_eq!(level, level_for(0, a as u64), - "level_for({}) -- binary = {:b}", a, a); - - if pos > level { - let a = a - 1; - assert_eq!(level, level_for(0, a as u64), - "level_for({}) -- binary = {:b}", a, a); - } - - if pos < 64 { - let a = a + 1; - assert_eq!(level, level_for(0, a as u64), - "level_for({}) -- binary = {:b}", a, a); - } - } - } - } -} diff --git a/tokio-timer/src/timer/stack.rs b/tokio-timer/src/timer/stack.rs new file mode 100644 index 000000000..9b438fa47 --- /dev/null +++ b/tokio-timer/src/timer/stack.rs @@ -0,0 +1,123 @@ +use super::Entry; +use wheel; + +use std::ptr; +use std::sync::Arc; + +/// A doubly linked stack +#[derive(Debug)] +pub(crate) struct Stack { + head: Option>, +} + +impl Default for Stack { + fn default() -> Stack { + Stack { head: None } + } +} + +impl wheel::Stack for Stack { + type Owned = Arc; + type Borrowed = Entry; + type Store = (); + + fn is_empty(&self) -> bool { + self.head.is_none() + } + + fn push(&mut self, entry: Self::Owned, _: &mut Self::Store) { + // Get a pointer to the entry to for the prev link + let ptr: *const Entry = &*entry as *const _; + + // Remove the old head entry + let old = self.head.take(); + + unsafe { + // Ensure the entry is not already in a stack. + debug_assert!((*entry.next_stack.get()).is_none()); + debug_assert!((*entry.prev_stack.get()).is_null()); + + if let Some(ref entry) = old.as_ref() { + debug_assert!({ + // The head is not already set to the entry + ptr != &***entry as *const _ + }); + + // Set the previous link on the old head + *entry.prev_stack.get() = ptr; + } + + // Set this entry's next pointer + *entry.next_stack.get() = old; + + } + + // Update the head pointer + self.head = Some(entry); + } + + /// Pop an item from the stack + fn pop(&mut self, _: &mut ()) -> Option> { + let entry = self.head.take(); + + unsafe { + if let Some(entry) = entry.as_ref() { + self.head = (*entry.next_stack.get()).take(); + + if let Some(entry) = self.head.as_ref() { + *entry.prev_stack.get() = ptr::null(); + } + + *entry.prev_stack.get() = ptr::null(); + } + } + + entry + } + + fn remove(&mut self, entry: &Entry, _: &mut ()) { + unsafe { + // Ensure that the entry is in fact contained by the stack + debug_assert!({ + // This walks the full linked list even if an entry is found. + let mut next = self.head.as_ref(); + let mut contains = false; + + while let Some(n) = next { + if entry as *const _ == &**n as *const _ { + debug_assert!(!contains); + contains = true; + } + + next = (*n.next_stack.get()).as_ref(); + } + + contains + }); + + // Unlink `entry` from the next node + let next = (*entry.next_stack.get()).take(); + + if let Some(next) = next.as_ref() { + (*next.prev_stack.get()) = *entry.prev_stack.get(); + } + + // Unlink `entry` from the prev node + + if let Some(prev) = (*entry.prev_stack.get()).as_ref() { + *prev.next_stack.get() = next; + } else { + // It is the head + self.head = next; + } + + // Unset the prev pointer + *entry.prev_stack.get() = ptr::null(); + } + } + + fn when(item: &Entry, _: &()) -> u64 { + item.when_internal() + .expect("invalid internal state") + } +} diff --git a/tokio-timer/src/timer/level.rs b/tokio-timer/src/wheel/level.rs similarity index 85% rename from tokio-timer/src/timer/level.rs rename to tokio-timer/src/wheel/level.rs index 8be0ba180..4f9aec2da 100644 --- a/tokio-timer/src/timer/level.rs +++ b/tokio-timer/src/wheel/level.rs @@ -1,10 +1,9 @@ -use timer::{entry, Entry}; +use wheel::Stack; use std::fmt; -use std::sync::Arc; /// Wheel for a single level in the timer. This wheel contains 64 slots. -pub(crate) struct Level { +pub(crate) struct Level { level: usize, /// Bit field tracking which slots currently contain entries. @@ -17,12 +16,12 @@ pub(crate) struct Level { occupied: u64, /// Slots - slot: [entry::Stack; LEVEL_MULT], + slot: [T; LEVEL_MULT], } /// Indicates when a slot must be processed next. #[derive(Debug)] -pub struct Expiration { +pub(crate) struct Expiration { /// The level containing the slot. pub level: usize, @@ -38,13 +37,13 @@ pub struct Expiration { /// Being a power of 2 is very important. const LEVEL_MULT: usize = 64; -impl Level { - pub fn new(level: usize) -> Level { +impl Level { + pub fn new(level: usize) -> Level { // Rust's derived implementations for arrays require that the value // contained by the array be `Copy`. So, here we have to manually // initialize every single slot. macro_rules! s { - () => { entry::Stack::new() }; + () => { T::default() }; }; Level { @@ -109,17 +108,17 @@ impl Level { Some(slot) } - pub fn add_entry(&mut self, entry: Arc, when: u64) { + pub fn add_entry(&mut self, when: u64, item: T::Owned, store: &mut T::Store) { let slot = slot_for(when, self.level); - self.slot[slot].push(entry); + self.slot[slot].push(item, store); self.occupied |= occupied_bit(slot); } - pub fn remove_entry(&mut self, entry: &Entry, when: u64) { + pub fn remove_entry(&mut self, when: u64, item: &T::Borrowed, store: &mut T::Store) { let slot = slot_for(when, self.level); - self.slot[slot].remove(entry); + self.slot[slot].remove(item, store); if self.slot[slot].is_empty() { // The bit is currently set @@ -130,8 +129,8 @@ impl Level { } } - pub fn pop_entry_slot(&mut self, slot: usize) -> Option> { - let ret = self.slot[slot].pop(); + pub fn pop_entry_slot(&mut self, slot: usize, store: &mut T::Store) -> Option { + let ret = self.slot[slot].pop(store); if ret.is_some() && self.slot[slot].is_empty() { // The bit is currently set @@ -144,19 +143,7 @@ impl Level { } } -impl Drop for Level { - fn drop(&mut self) { - while let Some(slot) = self.next_occupied_slot(0) { - // This should always have one - let entry = self.pop_entry_slot(slot) - .expect("occupied bit set invalid"); - - entry.error(); - } - } -} - -impl fmt::Debug for Level { +impl fmt::Debug for Level { fn fmt(&self, fmt: &mut fmt::Formatter) -> fmt::Result { fmt.debug_struct("Level") .field("occupied", &self.occupied) @@ -181,6 +168,7 @@ fn slot_for(duration: u64, level: usize) -> usize { ((duration >> (level * 6)) % LEVEL_MULT as u64) as usize } +/* #[cfg(test)] mod test { use super::*; @@ -199,3 +187,4 @@ mod test { } } } +*/ diff --git a/tokio-timer/src/wheel/mod.rs b/tokio-timer/src/wheel/mod.rs new file mode 100644 index 000000000..e9085523a --- /dev/null +++ b/tokio-timer/src/wheel/mod.rs @@ -0,0 +1,289 @@ +mod level; +mod stack; + +pub(crate) use self::stack::Stack; +pub(crate) use self::level::Expiration; +use self::level::Level; + +use std::borrow::Borrow; +use std::usize; + +/// Timing wheel implementation. +/// +/// This type provides the hashed timing wheel implementation that backs `Timer` +/// and `DelayQueue`. +/// +/// The structure is generic over `T: Stack`. This allows handling timeout data +/// being stored on the heap or in a slab. In order to support the latter case, +/// the slab must be passed into each function allowing the implementation to +/// lookup timer entries. +/// +/// See `Timer` documentation for some implementation notes. +#[derive(Debug)] +pub(crate) struct Wheel { + /// The number of milliseconds elapsed since the wheel started. + elapsed: u64, + + /// Timer wheel. + /// + /// Levels: + /// + /// * 1 ms slots / 64 ms range + /// * 64 ms slots / ~ 4 sec range + /// * ~ 4 sec slots / ~ 4 min range + /// * ~ 4 min slots / ~ 4 hr range + /// * ~ 4 hr slots / ~ 12 day range + /// * ~ 12 day slots / ~ 2 yr range + levels: Vec>, +} + +/// Number of levels. Each level has 64 slots. By using 6 levels with 64 slots +/// each, the timer is able to track time up to 2 years into the future with a +/// precision of 1 millisecond. +const NUM_LEVELS: usize = 6; + +/// The maximum duration of a delay +const MAX_DURATION: u64 = 1 << (6 * NUM_LEVELS); + +#[derive(Debug)] +pub(crate) enum InsertError { + Elapsed, + Invalid, +} + +/// Poll expirations from the wheel +#[derive(Debug, Default)] +pub(crate) struct Poll { + now: u64, + expiration: Option, +} + +impl Wheel +where + T: Stack, +{ + /// Create a new timing wheel + pub fn new() -> Wheel { + let levels = (0..NUM_LEVELS) + .map(Level::new) + .collect(); + + Wheel { + elapsed: 0, + levels, + } + } + + /// Return the number of milliseconds that have elapsed since the timing + /// wheele's creation. + pub fn elapsed(&self) -> u64 { + self.elapsed + } + + /// Insert an entry into the timing wheel. + /// + /// # Arguments + /// + /// * `when`: is the instant at which the the entry should be fired. It is + /// represented as the number of milliseconds since the creation + /// of the timing wheel. + /// + /// * `item`: The item to insert into the wheel. + /// + /// * `store`: The slab or `()` when using heap storage. + /// + /// # Return + /// + /// Returns `Ok` when the item is successfully inserted, `Err` otherwise. + /// + /// `Err(Elapsed)` indicates that `when` represents an instant that has + /// already passed. In this case, the caller should fire the timeout + /// immediateely. + /// + /// `Err(Invalid)` indicates an invalid `when` argumeent as been supplied. + pub fn insert(&mut self, when: u64, item: T::Owned, store: &mut T::Store) + -> Result<(), (T::Owned, InsertError)> + { + if when <= self.elapsed { + return Err((item, InsertError::Elapsed)); + } else if when - self.elapsed > MAX_DURATION { + return Err((item, InsertError::Invalid)); + } + + // Get the level at which the entry should be stored + let level = self.level_for(when); + + self.levels[level].add_entry(when, item, store); + + debug_assert!({ + self.levels[level].next_expiration(self.elapsed) + .map(|e| e.deadline >= self.elapsed) + .unwrap_or(true) + }); + + Ok(()) + } + + /// Remove `item` from thee timing wheel. + pub fn remove(&mut self, item: &T::Borrowed, store: &mut T::Store) { + let when = T::when(item, store); + let level = self.level_for(when); + + self.levels[level].remove_entry(when, item, store); + } + + /// Instant at which to poll + pub fn poll_at(&self) -> Option { + self.next_expiration() + .map(|expiration| expiration.deadline) + } + + pub fn poll(&mut self, poll: &mut Poll, store: &mut T::Store) + -> Option + { + loop { + if poll.expiration.is_none() { + poll.expiration = self.next_expiration() + .and_then(|expiration| { + if expiration.deadline > poll.now { + None + } else { + Some(expiration) + } + }); + } + + match poll.expiration { + Some(ref expiration) => { + if let Some(item) = self.poll_expiration(expiration, store) { + return Some(item); + } + + self.set_elapsed(expiration.deadline); + } + None => { + self.set_elapsed(poll.now); + return None; + } + } + + poll.expiration = None; + } + } + + /// Returns the instant at which the next timeout expires. + fn next_expiration(&self) -> Option { + // Check all levels + for level in 0..NUM_LEVELS { + if let Some(expiration) = self.levels[level].next_expiration(self.elapsed) { + // There cannot be any expirations at a higher level that happen + // before this one. + debug_assert!({ + let mut res = true; + + for l2 in (level+1)..NUM_LEVELS { + if let Some(e2) = self.levels[l2].next_expiration(self.elapsed) { + if e2.deadline < expiration.deadline { + res = false; + } + } + } + + res + }); + + return Some(expiration); + } + } + + None + } + + pub fn poll_expiration(&mut self, expiration: &Expiration, store: &mut T::Store) + -> Option + { + while let Some(item) = self.pop_entry(expiration, store) { + if expiration.level == 0 { + debug_assert_eq!(T::when(item.borrow(), store), expiration.deadline); + + return Some(item); + } else { + let when = T::when(item.borrow(), store); + + let next_level = expiration.level - 1; + + self.levels[next_level] + .add_entry(when, item, store); + } + } + + None + } + + fn set_elapsed(&mut self, when: u64) { + assert!(self.elapsed <= when, "elapsed={:?}; when={:?}", self.elapsed, when); + + if when > self.elapsed { + self.elapsed = when; + } + } + + fn pop_entry(&mut self, expiration: &Expiration, store: &mut T::Store) -> Option { + self.levels[expiration.level].pop_entry_slot(expiration.slot, store) + } + + fn level_for(&self, when: u64) -> usize { + level_for(self.elapsed, when) + } +} + +fn level_for(elapsed: u64, when: u64) -> usize { + let masked = elapsed ^ when; + + assert!(masked != 0, "elapsed={}; when={}", elapsed, when); + + let leading_zeros = masked.leading_zeros() as usize; + let significant = 63 - leading_zeros; + significant / 6 +} + +impl Poll { + pub fn new(now: u64) -> Poll { + Poll { + now, + expiration: None, + } + } +} + +#[cfg(test)] +mod test { + use super::*; + + #[test] + fn test_level_for() { + for pos in 1..64 { + assert_eq!(0, level_for(0, pos), "level_for({}) -- binary = {:b}", pos, pos); + } + + for level in 1..5 { + for pos in level..64 { + let a = pos * 64_usize.pow(level as u32); + assert_eq!(level, level_for(0, a as u64), + "level_for({}) -- binary = {:b}", a, a); + + if pos > level { + let a = a - 1; + assert_eq!(level, level_for(0, a as u64), + "level_for({}) -- binary = {:b}", a, a); + } + + if pos < 64 { + let a = a + 1; + assert_eq!(level, level_for(0, a as u64), + "level_for({}) -- binary = {:b}", a, a); + } + } + } + } +} diff --git a/tokio-timer/src/wheel/stack.rs b/tokio-timer/src/wheel/stack.rs new file mode 100644 index 000000000..6e55c38cc --- /dev/null +++ b/tokio-timer/src/wheel/stack.rs @@ -0,0 +1,26 @@ +use std::borrow::Borrow; + +/// Abstracts the stack operations needed to track timeouts. +pub(crate) trait Stack: Default { + /// Type of the item stored in the stack + type Owned: Borrow; + + /// Borrowed item + type Borrowed; + + /// Item storage, this allows a slab to be used instead of just the heap + type Store; + + /// Returns `true` if the stack is empty + fn is_empty(&self) -> bool; + + /// Push an item onto the stack + fn push(&mut self, item: Self::Owned, store: &mut Self::Store); + + /// Pop an item from the stack + fn pop(&mut self, store: &mut Self::Store) -> Option; + + fn remove(&mut self, item: &Self::Borrowed, store: &mut Self::Store); + + fn when(item: &Self::Borrowed, store: &Self::Store) -> u64; +} diff --git a/tokio-timer/tests/interval.rs b/tokio-timer/tests/interval.rs index 60d6d8d62..42630b857 100644 --- a/tokio-timer/tests/interval.rs +++ b/tokio-timer/tests/interval.rs @@ -24,23 +24,23 @@ fn usage() { let start = time.now(); let mut int = Interval::new(start, ms(300)); - assert_ready!(int, Some(start)); + assert_ready_eq!(int, Some(start)); assert_not_ready!(int); advance(timer, ms(100)); assert_not_ready!(int); advance(timer, ms(200)); - assert_ready!(int, Some(start + ms(300))); + assert_ready_eq!(int, Some(start + ms(300))); assert_not_ready!(int); advance(timer, ms(400)); - assert_ready!(int, Some(start + ms(600))); + assert_ready_eq!(int, Some(start + ms(600))); assert_not_ready!(int); advance(timer, ms(500)); - assert_ready!(int, Some(start + ms(900))); - assert_ready!(int, Some(start + ms(1200))); + assert_ready_eq!(int, Some(start + ms(900))); + assert_ready_eq!(int, Some(start + ms(1200))); assert_not_ready!(int); }); } diff --git a/tokio-timer/tests/queue.rs b/tokio-timer/tests/queue.rs new file mode 100644 index 000000000..6b0ffc5ba --- /dev/null +++ b/tokio-timer/tests/queue.rs @@ -0,0 +1,223 @@ +extern crate futures; +extern crate tokio_executor; +extern crate tokio_timer; +extern crate tokio_mock_task; + +#[macro_use] +mod support; +use support::*; + +use tokio_timer::*; +use tokio_mock_task::MockTask; + +use futures::Stream; + +#[test] +fn single_immediate_delay() { + mocked(|_timer, time| { + let mut queue = DelayQueue::new(); + let _key = queue.insert_at("foo", time.now()); + + let entry = assert_ready!(queue).unwrap(); + assert_eq!(*entry.get_ref(), "foo"); + + let entry = assert_ready!(queue); + assert!(entry.is_none()) + }); +} + +#[test] +fn multi_immediate_delays() { + mocked(|_timer, time| { + let mut queue = DelayQueue::new(); + + let _k = queue.insert_at("1", time.now()); + let _k = queue.insert_at("2", time.now()); + let _k = queue.insert_at("3", time.now()); + + let mut res = vec![]; + + while res.len() < 3 { + let entry = assert_ready!(queue).unwrap(); + res.push(entry.into_inner()); + } + + let entry = assert_ready!(queue); + assert!(entry.is_none()); + + res.sort(); + + assert_eq!("1", res[0]); + assert_eq!("2", res[1]); + assert_eq!("3", res[2]); + }); +} + +#[test] +fn single_short_delay() { + mocked(|timer, time| { + let mut queue = DelayQueue::new(); + let _key = queue.insert_at("foo", time.now() + ms(5)); + + let mut task = MockTask::new(); + + task.enter(|| { + assert_not_ready!(queue); + }); + + turn(timer, ms(1)); + + assert!(!task.is_notified()); + + turn(timer, ms(5)); + + assert!(task.is_notified()); + + let entry = assert_ready!(queue).unwrap(); + assert_eq!(*entry.get_ref(), "foo"); + + let entry = assert_ready!(queue); + assert!(entry.is_none()); + }); +} + +#[test] +fn multi_delay_at_start() { + let long = 262_144 + 9 * 4096; + let delays = &[1000, 2, 234, long, 60, 10]; + + mocked(|timer, time| { + let mut queue = DelayQueue::new(); + let mut task = MockTask::new(); + + // Setup the delays + for &i in delays { + let _key = queue.insert_at(i, time.now() + ms(i)); + } + + task.enter(|| { + assert_not_ready!(queue); + }); + + assert!(!task.is_notified()); + + for elapsed in 0..1200 { + turn(timer, ms(1)); + let elapsed = elapsed + 1; + + if delays.contains(&elapsed) { + assert!(task.is_notified()); + + task.enter(|| { + assert_ready!(queue); + assert_not_ready!(queue); + }); + } else { + if task.is_notified() { + let cascade = &[192, 960]; + assert!(cascade.contains(&elapsed), "elapsed={}", elapsed); + + task.enter(|| { + assert_not_ready!(queue, "elapsed={}", elapsed); + }); + } + } + } + }); +} + +#[test] +fn insert_in_past_fires_immediately() { + mocked(|timer, time| { + let mut queue = DelayQueue::new(); + + let now = time.now(); + + turn(timer, ms(10)); + + queue.insert_at("foo", now); + + assert_ready!(queue); + }); +} + +#[test] +fn remove_entry() { + mocked(|timer, time| { + let mut queue = DelayQueue::new(); + let mut task = MockTask::new(); + + let key = queue.insert_at("foo", time.now() + ms(5)); + + task.enter(|| { + assert_not_ready!(queue); + }); + + let entry = queue.remove(&key); + assert_eq!(entry.into_inner(), "foo"); + + turn(timer, ms(10)); + + task.enter(|| { + let entry = assert_ready!(queue); + assert!(entry.is_none()); + }); + }); +} + +#[test] +fn reset_entry() { + mocked(|timer, time| { + let mut queue = DelayQueue::new(); + let mut task = MockTask::new(); + + let now = time.now(); + let key = queue.insert_at("foo", now + ms(5)); + + task.enter(|| { + assert_not_ready!(queue); + }); + + turn(timer, ms(1)); + + queue.reset_at(&key, now + ms(10)); + + task.enter(|| { + assert_not_ready!(queue); + }); + + turn(timer, ms(7)); + + assert!(!task.is_notified()); + + task.enter(|| { + assert_not_ready!(queue); + }); + + turn(timer, ms(3)); + + assert!(task.is_notified()); + + let entry = assert_ready!(queue).unwrap(); + assert_eq!(*entry.get_ref(), "foo"); + + let entry = assert_ready!(queue); + assert!(entry.is_none()) + }); +} + +#[test] +fn remove_expired_item() { + mocked(|timer, time| { + let mut queue = DelayQueue::new(); + + let now = time.now(); + + turn(timer, ms(10)); + + let key = queue.insert_at("foo", now); + + let entry = queue.remove(&key); + assert_eq!(entry.into_inner(), "foo"); + }) +} diff --git a/tokio-timer/tests/support/mod.rs b/tokio-timer/tests/support/mod.rs index a79b0dab0..cc30e166e 100644 --- a/tokio-timer/tests/support/mod.rs +++ b/tokio-timer/tests/support/mod.rs @@ -1,7 +1,8 @@ #![allow(unused_macros, unused_imports, dead_code, deprecated)] use tokio_executor::park::{Park, Unpark}; -use tokio_timer::timer::{Timer, Now}; +use tokio_timer::clock::Now; +use tokio_timer::timer::Timer; use futures::future::{lazy, Future}; @@ -11,18 +12,45 @@ use std::sync::{Arc, Mutex}; use std::time::{Instant, Duration}; macro_rules! assert_ready { - ($f:expr) => { - assert!($f.poll().unwrap().is_ready()); - }; + ($f:expr) => {{ + use ::futures::Async::*; + + match $f.poll().unwrap() { + Ready(v) => v, + NotReady => panic!("NotReady"), + } + }}; + ($f:expr, $($msg:expr),+) => {{ + use ::futures::Async::*; + + match $f.poll().unwrap() { + Ready(v) => v, + NotReady => { + let msg = format!($($msg),+); + panic!("NotReady; {}", msg) + } + } + }} +} + +macro_rules! assert_ready_eq { ($f:expr, $expect:expr) => { assert_eq!($f.poll().unwrap(), ::futures::Async::Ready($expect)); }; } macro_rules! assert_not_ready { - ($f:expr) => { - assert!(!$f.poll().unwrap().is_ready()); - } + ($f:expr) => {{ + let res = $f.poll().unwrap(); + assert!(!res.is_ready(), "actual={:?}", res) + }}; + ($f:expr, $($msg:expr),+) => {{ + let res = $f.poll().unwrap(); + if res.is_ready() { + let msg = format!($($msg),+); + panic!("actual={:?}; {}", res, msg); + } + }}; } macro_rules! assert_elapsed { @@ -40,7 +68,6 @@ pub struct MockTime { #[derive(Debug)] pub struct MockNow { inner: Inner, - _p: PhantomData>, } #[derive(Debug)] @@ -85,12 +112,12 @@ impl IntoTimeout for Duration { } /// Turn the timer state once -pub fn turn(timer: &mut Timer, duration: T) { +pub fn turn(timer: &mut Timer, duration: T) { timer.turn(duration.into_timeout()).unwrap(); } /// Advance the timer the specified amount -pub fn advance(timer: &mut Timer, duration: Duration) { +pub fn advance(timer: &mut Timer, duration: Duration) { let inner = timer.get_park().inner.clone(); let deadline = inner.lock().unwrap().now() + duration; @@ -101,27 +128,29 @@ pub fn advance(timer: &mut Timer, duration: Duration) { } pub fn mocked(f: F) -> R -where F: FnOnce(&mut Timer, &mut MockTime) -> R +where F: FnOnce(&mut Timer, &mut MockTime) -> R { mocked_with_now(Instant::now(), f) } pub fn mocked_with_now(now: Instant, f: F) -> R -where F: FnOnce(&mut Timer, &mut MockTime) -> R +where F: FnOnce(&mut Timer, &mut MockTime) -> R { let mut time = MockTime::new(now); let park = time.mock_park(); - let now = time.mock_now(); - - let mut timer = Timer::new_with_now(park, now); - let handle = timer.handle(); + let now = ::tokio_timer::clock::Clock::new_with_now(time.mock_now()); let mut enter = ::tokio_executor::enter().unwrap(); - ::tokio_timer::with_default(&handle, &mut enter, |_| { - lazy(|| { - Ok::<_, ()>(f(&mut timer, &mut time)) - }).wait().unwrap() + ::tokio_timer::clock::with_default(&now, &mut enter, |enter| { + let mut timer = Timer::new(park); + let handle = timer.handle(); + + ::tokio_timer::with_default(&handle, enter, |_| { + lazy(|| { + Ok::<_, ()>(f(&mut timer, &mut time)) + }).wait().unwrap() + }) }) } @@ -144,7 +173,6 @@ impl MockTime { let inner = self.inner.clone(); MockNow { inner, - _p: PhantomData, } } @@ -218,7 +246,7 @@ impl Unpark for MockUnpark { } impl Now for MockNow { - fn now(&mut self) -> Instant { + fn now(&self) -> Instant { self.inner.lock().unwrap().now() } }