pollable_map/futures/
ordered.rs1use alloc::collections::VecDeque;
2use core::future::Future;
3use core::pin::Pin;
4use core::task::{Context, Poll, Waker};
5use futures::Stream;
6
7#[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 pub fn new() -> Self {
30 Self::default()
31 }
32
33 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 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 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 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 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 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 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 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}