pollable_map/futures/
set.rs1use core::future::Future;
2use core::pin::Pin;
3
4use super::FutureMap;
5use core::task::{Context, Poll};
6use futures::stream::FusedStream;
7use futures::{Stream, StreamExt};
8
9pub struct FutureSet<S> {
10 id: i64,
11 map: FutureMap<i64, S>,
12}
13
14impl<S> Default for FutureSet<S>
15where
16 S: Future,
17{
18 fn default() -> Self {
19 Self::new()
20 }
21}
22
23impl<S> FutureSet<S>
24where
25 S: Future,
26{
27 pub fn new() -> Self {
29 Self {
30 id: 0,
31 map: FutureMap::default(),
32 }
33 }
34
35 pub fn insert(&mut self, fut: S) -> bool {
37 self.id = self.id.wrapping_add(1);
38 self.map.insert(self.id, fut)
39 }
40
41 pub fn iter(&self) -> impl Iterator<Item = &S> {
43 self.map.iter().map(|(_, st)| st)
44 }
45
46 pub fn iter_pin(&mut self) -> impl Iterator<Item = Pin<&mut S>> {
48 self.map.iter_pin().map(|(_, st)| st)
49 }
50
51 pub fn clear(&mut self) {
53 self.map.clear();
54 }
55
56 pub fn len(&self) -> usize {
58 self.map.len()
59 }
60
61 pub fn is_empty(&self) -> bool {
63 self.map.is_empty()
64 }
65}
66
67impl<S> FutureSet<S>
68where
69 S: Future + Unpin,
70{
71 pub fn iter_mut(&mut self) -> impl Iterator<Item = &mut S> {
73 self.map.iter_mut().map(|(_, st)| st)
74 }
75}
76
77impl<S> FromIterator<S> for FutureSet<S>
78where
79 S: Future,
80{
81 fn from_iter<I: IntoIterator<Item = S>>(iter: I) -> Self {
82 let mut maps = Self::new();
83 for st in iter {
84 maps.insert(st);
85 }
86 maps
87 }
88}
89
90impl<S> Stream for FutureSet<S>
91where
92 S: Future,
93{
94 type Item = S::Output;
95
96 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
97 self.map
98 .poll_next_unpin(cx)
99 .map(|output| output.map(|(_, item)| item))
100 }
101
102 fn size_hint(&self) -> (usize, Option<usize>) {
103 self.map.size_hint()
104 }
105}
106
107impl<S> FusedStream for FutureSet<S>
108where
109 S: Future,
110{
111 fn is_terminated(&self) -> bool {
112 self.map.is_terminated()
113 }
114}
115
116#[cfg(test)]
117mod test {
118 use crate::futures::set::FutureSet;
119 use futures::StreamExt;
120
121 #[test]
122 fn valid_future_set() {
123 let mut list = FutureSet::new();
124 assert!(list.insert(futures::future::ready(0)));
125 assert!(list.insert(futures::future::ready(1)));
126
127 futures::executor::block_on(async move {
128 let val = list.next().await;
129 assert_eq!(val, Some(0));
130 let val = list.next().await;
131 assert_eq!(val, Some(1));
132 });
133 }
134
135 #[test]
136 fn supports_unboxed_async_future() {
137 let mut set = FutureSet::new();
138 assert!(set.insert(async { 42 }));
139
140 futures::executor::block_on(async move {
141 assert_eq!(set.next().await, Some(42));
142 });
143 }
144}