Skip to main content

pollable_map/futures/
ordered.rs

1use alloc::collections::VecDeque;
2use core::future::Future;
3use core::pin::Pin;
4use core::task::{Context, Poll, Waker};
5use futures::Stream;
6
7/// An unbounded queue of futures imposed a FIFO order while polling one future at a time
8/// and returning the output to stream before popping the next future in queue to be polled.
9#[pin_project::pin_project]
10pub struct OrderedFutureSet<F> {
11    queue: VecDeque<F>,
12    #[pin]
13    current_future: Option<F>,
14    waker: Option<Waker>,
15}
16
17impl<F> Default for OrderedFutureSet<F> {
18    fn default() -> Self {
19        Self {
20            queue: VecDeque::new(),
21            current_future: None,
22            waker: None,
23        }
24    }
25}
26
27impl<F> OrderedFutureSet<F> {
28    /// Constructs a new, empty [`OrderedFutureSet`]
29    pub fn new() -> Self {
30        Self::default()
31    }
32
33    /// Push a future to the back of the queue
34    pub fn push(&mut self, fut: F) {
35        self.queue.push_back(fut);
36        if let Some(waker) = self.waker.take() {
37            waker.wake();
38        }
39    }
40
41    /// Push a future to the back of a pinned queue.
42    pub fn push_pinned(self: Pin<&mut Self>, fut: F) {
43        let this = self.project();
44        this.queue.push_back(fut);
45        if let Some(waker) = this.waker.take() {
46            waker.wake();
47        }
48    }
49
50    /// Remove a future from the front of the queue
51    pub fn pop_front(&mut self) -> Option<F> {
52        let fut = self.queue.pop_front();
53        if let Some(waker) = self.waker.take() {
54            waker.wake();
55        }
56        fut
57    }
58
59    /// Remove a future from the front of a pinned queue.
60    pub fn pop_front_pinned(self: Pin<&mut Self>) -> Option<F> {
61        let this = self.project();
62        let fut = this.queue.pop_front();
63        if let Some(waker) = this.waker.take() {
64            waker.wake();
65        }
66        fut
67    }
68
69    /// Remove a future from the back of the queue
70    pub fn pop_back(&mut self) -> Option<F> {
71        let fut = self.queue.pop_back();
72        if let Some(waker) = self.waker.take() {
73            waker.wake();
74        }
75        fut
76    }
77
78    /// Remove a future from the back of a pinned queue.
79    pub fn pop_back_pinned(self: Pin<&mut Self>) -> Option<F> {
80        let this = self.project();
81        let fut = this.queue.pop_back();
82        if let Some(waker) = this.waker.take() {
83            waker.wake();
84        }
85        fut
86    }
87}
88
89impl<F> FromIterator<F> for OrderedFutureSet<F> {
90    fn from_iter<T: IntoIterator<Item = F>>(iter: T) -> Self {
91        let mut ordered = Self::new();
92        for fut in iter {
93            ordered.push(fut);
94        }
95        ordered
96    }
97}
98
99impl<F> Stream for OrderedFutureSet<F>
100where
101    F: Future,
102{
103    type Item = F::Output;
104    fn poll_next(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Option<Self::Item>> {
105        let mut this = self.project();
106
107        if this.current_future.as_ref().get_ref().is_none() {
108            let Some(fut) = this.queue.pop_front() else {
109                this.waker.replace(cx.waker().clone());
110                return Poll::Pending;
111            };
112            this.current_future.set(Some(fut));
113        }
114
115        let fut = this
116            .current_future
117            .as_mut()
118            .as_pin_mut()
119            .expect("current future was initialized");
120
121        match fut.poll(cx) {
122            Poll::Ready(output) => {
123                this.current_future.set(None);
124                cx.waker().wake_by_ref();
125                Poll::Ready(Some(output))
126            }
127            Poll::Pending => {
128                this.waker.replace(cx.waker().clone());
129                Poll::Pending
130            }
131        }
132    }
133
134    fn size_hint(&self) -> (usize, Option<usize>) {
135        (
136            self.queue.len() + usize::from(self.current_future.is_some()),
137            None,
138        )
139    }
140}
141
142#[cfg(test)]
143mod tests {
144    use crate::futures::ordered::OrderedFutureSet;
145    use alloc::vec;
146    use alloc::vec::Vec;
147    use futures::StreamExt;
148
149    #[test]
150    fn fifo_futures() {
151        futures::executor::block_on(async move {
152            let mut fifo = OrderedFutureSet::new();
153            fifo.push(futures::future::ready(1));
154            fifo.push(futures::future::ready(2));
155            fifo.push(futures::future::ready(4));
156            fifo.push(futures::future::ready(3));
157
158            let items = fifo.take(4).collect::<Vec<u8>>().await;
159
160            assert_eq!(items, vec![1, 2, 4, 3]);
161        });
162    }
163
164    #[test]
165    fn remove_front_entry() {
166        futures::executor::block_on(async move {
167            let mut fifo = OrderedFutureSet::new();
168            fifo.push(futures::future::ready(1));
169            fifo.push(futures::future::ready(2));
170            fifo.push(futures::future::ready(4));
171            fifo.push(futures::future::ready(3));
172
173            let front_fut = fifo.pop_front();
174            // TODO: Write a `Ready` future that supports `Eq` and `PartialEq` for tests
175            //       to use `assert_eq(front_fut, Some(futures::future::ready(1)));`
176            assert!(front_fut.is_some());
177
178            let items = fifo.take(3).collect::<Vec<u8>>().await;
179
180            assert_eq!(items, vec![2, 4, 3]);
181        })
182    }
183
184    #[test]
185    fn remove_back_entry() {
186        futures::executor::block_on(async move {
187            let mut fifo = OrderedFutureSet::new();
188            fifo.push(futures::future::ready(1));
189            fifo.push(futures::future::ready(2));
190            fifo.push(futures::future::ready(4));
191            fifo.push(futures::future::ready(3));
192
193            let front_fut = fifo.pop_back();
194            // TODO: Write a `Ready` future that supports `Eq` and `PartialEq` for tests
195            //       to use `assert_eq(front_fut, Some(futures::future::ready(3)));`
196            assert!(front_fut.is_some());
197
198            let items = fifo.take(3).collect::<Vec<u8>>().await;
199
200            assert_eq!(items, vec![1, 2, 4]);
201        })
202    }
203
204    #[test]
205    fn supports_unboxed_async_futures() {
206        async fn value(value: u8) -> u8 {
207            value
208        }
209
210        futures::executor::block_on(async move {
211            let mut fifo = OrderedFutureSet::new();
212            fifo.push(value(1));
213            fifo.push(value(2));
214            futures::pin_mut!(fifo);
215
216            assert_eq!(fifo.as_mut().next().await, Some(1));
217            assert_eq!(fifo.as_mut().next().await, Some(2));
218
219            fifo.as_mut().push_pinned(value(3));
220            assert_eq!(fifo.as_mut().next().await, Some(3));
221        });
222    }
223}