Skip to main content

pollable_map/futures/
set.rs

1use 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    /// Creates an empty [`FutureSet`]
28    pub fn new() -> Self {
29        Self {
30            id: 0,
31            map: FutureMap::default(),
32        }
33    }
34
35    /// Insert a future into the set of futures.
36    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    /// An iterator visiting all futures in arbitrary order.
42    pub fn iter(&self) -> impl Iterator<Item = &S> {
43        self.map.iter().map(|(_, st)| st)
44    }
45
46    /// An iterator visiting all futures pinned valued in arbitrary order
47    pub fn iter_pin(&mut self) -> impl Iterator<Item = Pin<&mut S>> {
48        self.map.iter_pin().map(|(_, st)| st)
49    }
50
51    /// Clears the set.
52    pub fn clear(&mut self) {
53        self.map.clear();
54    }
55
56    /// Returns the number of futures in the set.
57    pub fn len(&self) -> usize {
58        self.map.len()
59    }
60
61    /// Return `true` map contains no elements.
62    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    /// An iterator visiting all futures mutably in arbitrary order.
72    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}