Skip to main content

pollable_map/
optional.rs

1#[cfg(all(feature = "std", feature = "timeout"))]
2pub mod timeout;
3
4use core::future::Future;
5use core::pin::Pin;
6use core::task::{Context, Poll, Waker};
7use futures::future::FusedFuture;
8use futures::stream::FusedStream;
9use futures::Stream;
10use pin_project::pin_project;
11
12/// A reusable future or stream based on `Option`.
13///
14/// By default, `Optional` will be empty, similar to `Option::None`, which would return [`Poll::Pending`] when polled,
15/// but if a [`Future`] or [`Stream`] is supplied either upon construction via [`Optional::new`] or
16/// is set via [`Optional::replace`], it would then be polled once [`Optional`]
17/// is polled. Once the future is polled to completion, the results will be returned, with
18/// [`Optional`] being empty.
19#[pin_project]
20pub struct Optional<T> {
21    #[pin]
22    task: Option<T>,
23    waker: Option<Waker>,
24}
25
26impl<T> Default for Optional<T> {
27    fn default() -> Self {
28        Self {
29            task: None,
30            waker: None,
31        }
32    }
33}
34
35impl<T> From<Option<T>> for Optional<T> {
36    fn from(task: Option<T>) -> Self {
37        Self { task, waker: None }
38    }
39}
40
41impl<T> From<T> for Optional<T> {
42    fn from(fut: T) -> Self {
43        Self {
44            task: Some(fut),
45            waker: None,
46        }
47    }
48}
49
50impl<T> Optional<T> {
51    /// Construct a new [`Optional`] with an existing [`Future`] or [`Stream`].
52    pub fn new(task: T) -> Self {
53        Self {
54            task: Some(task),
55            waker: None,
56        }
57    }
58
59    /// Construct a new [`Optional`] with an existing [`Future`].
60    pub fn with_future(future: T) -> Self
61    where
62        T: Future,
63    {
64        Self::new(future)
65    }
66
67    /// Construct a new [`Optional`] with an existing [`Stream`].
68    pub fn with_stream(stream: T) -> Self
69    where
70        T: Stream,
71    {
72        Self::new(stream)
73    }
74
75    /// Takes the future or stream out, leaving the [`Optional`] empty.
76    pub fn take(&mut self) -> Option<T> {
77        let fut = self.task.take();
78        // Note: Although we dont have to wake the task, we do it so the task is aware that
79        // the future or stream is no longer valid since it has been taken and returned via
80        // this function.
81        if let Some(waker) = self.waker.take() {
82            waker.wake();
83        }
84        fut
85    }
86
87    /// Returns true if the future or stream still exist.
88    pub fn is_some(&self) -> bool {
89        self.task.is_some()
90    }
91
92    /// Returns false if the future or stream doesn't exist or has been completed.
93    pub fn is_none(&self) -> bool {
94        self.task.is_none()
95    }
96
97    /// Returns reference of the future or stream.
98    pub fn as_ref(&self) -> Option<&T> {
99        self.task.as_ref()
100    }
101
102    /// Returns mutable reference of the future or stream.
103    pub fn as_mut(&mut self) -> Option<&mut T> {
104        self.task.as_mut()
105    }
106
107    /// Replaces the current the future or stream with a new one, returning the previous value if present.
108    pub fn replace(&mut self, task: T) -> Option<T> {
109        let fut = self.task.replace(task);
110        if let Some(waker) = self.waker.take() {
111            waker.wake();
112        }
113        fut
114    }
115
116    /// Replaces the current future or stream in place without moving the previous value.
117    pub fn set(self: Pin<&mut Self>, task: T) {
118        let mut this = self.project();
119
120        this.task.set(Some(task));
121
122        if let Some(waker) = this.waker.take() {
123            waker.wake();
124        }
125    }
126
127    /// Returns a constructed `Option<Pin<&mut T>>`.
128    pub fn as_pin_mut(&mut self) -> Option<Pin<&mut T>>
129    where
130        T: Unpin,
131    {
132        self.task.as_mut().map(Pin::new)
133    }
134
135    /// Return a constructed `Option<Pin<&mut T>>`.
136    pub fn pinned_as_mut(self: Pin<&mut Self>) -> Option<Pin<&mut T>> {
137        self.project().task.as_pin_mut()
138    }
139
140    /// Returns a constructed `Option<Pin<&T>>`.
141    pub fn as_pin_ref(&self) -> Option<Pin<&T>>
142    where
143        T: Unpin,
144    {
145        self.task.as_ref().map(Pin::new)
146    }
147
148    /// Return a constructed `Option<Pin<&T>>`.
149    pub fn pinned_as_ref(self: Pin<&Self>) -> Option<Pin<&T>> {
150        self.project_ref().task.as_pin_ref()
151    }
152}
153
154impl<F> Future for Optional<F>
155where
156    F: Future,
157{
158    type Output = F::Output;
159
160    fn poll(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Self::Output> {
161        let mut this = self.project();
162        let Some(future) = this.task.as_mut().as_pin_mut() else {
163            this.waker.replace(cx.waker().clone());
164            return Poll::Pending;
165        };
166
167        match future.poll(cx) {
168            Poll::Ready(output) => {
169                this.task.set(None);
170                Poll::Ready(output)
171            }
172            Poll::Pending => {
173                this.waker.replace(cx.waker().clone());
174                Poll::Pending
175            }
176        }
177    }
178}
179
180impl<F: Future> FusedFuture for Optional<F>
181where
182    F: Future,
183{
184    fn is_terminated(&self) -> bool {
185        self.task.is_none()
186    }
187}
188
189impl<S> Stream for Optional<S>
190where
191    S: Stream,
192{
193    type Item = S::Item;
194
195    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
196        let mut this = self.project();
197        let Some(stream) = this.task.as_mut().as_pin_mut() else {
198            this.waker.replace(cx.waker().clone());
199            return Poll::Pending;
200        };
201
202        match stream.poll_next(cx) {
203            Poll::Ready(Some(output)) => Poll::Ready(Some(output)),
204            Poll::Ready(None) => {
205                this.task.set(None);
206                Poll::Ready(None)
207            }
208            Poll::Pending => {
209                this.waker.replace(cx.waker().clone());
210                Poll::Pending
211            }
212        }
213    }
214
215    fn size_hint(&self) -> (usize, Option<usize>) {
216        match self.task.as_ref() {
217            Some(st) => st.size_hint(),
218            None => (0, Some(0)),
219        }
220    }
221}
222
223impl<S> FusedStream for Optional<S>
224where
225    S: Stream,
226{
227    fn is_terminated(&self) -> bool {
228        self.task.is_none()
229    }
230}
231
232#[cfg(test)]
233mod test {
234    use super::*;
235    use futures::StreamExt;
236
237    #[test]
238    fn test_optional_future() {
239        let mut future = Optional::new(futures::future::ready(0));
240        assert!(future.is_some());
241        let waker = futures::task::noop_waker_ref();
242
243        let val = Pin::new(&mut future).poll(&mut Context::from_waker(waker));
244        assert_eq!(val, Poll::Ready(0));
245        assert!(future.is_none());
246    }
247
248    #[test]
249    fn reusable_optional_future() {
250        let mut future = Optional::new(futures::future::ready(0));
251        assert!(future.is_some());
252        let waker = futures::task::noop_waker_ref();
253
254        let val = Pin::new(&mut future).poll(&mut Context::from_waker(waker));
255        assert_eq!(val, Poll::Ready(0));
256        assert!(future.is_none());
257
258        future.replace(futures::future::ready(1));
259        assert!(future.is_some());
260
261        let val = Pin::new(&mut future).poll(&mut Context::from_waker(waker));
262        assert_eq!(val, Poll::Ready(1));
263        assert!(future.is_none());
264    }
265
266    #[test]
267    fn reusable_pinned_optional_future() {
268        async fn set_value(value: i32) -> i32 {
269            value
270        }
271
272        let future = Optional::new(set_value(0));
273        futures::pin_mut!(future);
274        assert!(future.is_some());
275        let waker = futures::task::noop_waker_ref();
276
277        let value = future.as_mut().poll(&mut Context::from_waker(waker));
278        assert_eq!(value, Poll::Ready(0));
279        assert!(future.is_none());
280
281        future.as_mut().set(set_value(1));
282        assert!(future.is_some());
283
284        let value = future.as_mut().poll(&mut Context::from_waker(waker));
285        assert_eq!(value, Poll::Ready(1));
286        assert!(future.is_none());
287    }
288
289    #[test]
290    fn convert_future_to_optional_future() {
291        let fut = futures::future::ready(0);
292
293        let mut future = Optional::from(fut);
294        assert!(future.is_some());
295        let waker = futures::task::noop_waker_ref();
296
297        let val = Pin::new(&mut future).poll(&mut Context::from_waker(waker));
298        assert_eq!(val, Poll::Ready(0));
299        assert!(future.is_none());
300    }
301
302    #[test]
303    fn test_optional_stream() {
304        let mut stream = Optional::new(futures::stream::once(async { 0 }).boxed());
305        assert!(stream.is_some());
306        let waker = futures::task::noop_waker_ref();
307
308        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
309        assert_eq!(val, Poll::Ready(Some(0)));
310        assert!(stream.is_some());
311
312        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
313        assert_eq!(val, Poll::Ready(None));
314        assert!(stream.is_none());
315    }
316
317    #[test]
318    fn reusable_optional_stream() {
319        let mut stream = Optional::new(futures::stream::once(async { 0 }).boxed());
320        assert!(stream.is_some());
321        let waker = futures::task::noop_waker_ref();
322
323        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
324        assert_eq!(val, Poll::Ready(Some(0)));
325        assert!(stream.is_some());
326
327        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
328        assert_eq!(val, Poll::Ready(None));
329        assert!(stream.is_none());
330
331        stream.replace(futures::stream::once(async { 1 }).boxed());
332        assert!(stream.is_some());
333
334        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
335        assert_eq!(val, Poll::Ready(Some(1)));
336        assert!(stream.is_some());
337
338        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
339        assert_eq!(val, Poll::Ready(None));
340        assert!(stream.is_none());
341    }
342
343    #[test]
344    fn reusable_pinned_optional_stream() {
345        async fn set_val(value: i32) -> i32 {
346            value
347        }
348
349        let stream = Optional::new(futures::stream::once(set_val(0)));
350        futures::pin_mut!(stream);
351        assert!(stream.is_some());
352        let waker = futures::task::noop_waker_ref();
353
354        let val = stream.as_mut().poll_next(&mut Context::from_waker(waker));
355        assert_eq!(val, Poll::Ready(Some(0)));
356        assert!(stream.is_some());
357
358        let val = stream.as_mut().poll_next(&mut Context::from_waker(waker));
359        assert_eq!(val, Poll::Ready(None));
360        assert!(stream.is_none());
361
362        stream.as_mut().set(futures::stream::once(set_val(1)));
363        assert!(stream.is_some());
364
365        let val = stream.as_mut().poll_next(&mut Context::from_waker(waker));
366        assert_eq!(val, Poll::Ready(Some(1)));
367        assert!(stream.is_some());
368
369        let val = stream.as_mut().poll_next(&mut Context::from_waker(waker));
370        assert_eq!(val, Poll::Ready(None));
371        assert!(stream.is_none());
372    }
373
374    #[test]
375    fn convert_stream_to_optional_stream() {
376        let st = futures::stream::once(async { 0 }).boxed();
377
378        let mut stream = Optional::from(st);
379
380        assert!(stream.is_some());
381        let waker = futures::task::noop_waker_ref();
382
383        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
384        assert_eq!(val, Poll::Ready(Some(0)));
385        assert!(stream.is_some());
386
387        let val = Pin::new(&mut stream).poll_next(&mut Context::from_waker(waker));
388        assert_eq!(val, Poll::Ready(None));
389        assert!(stream.is_none());
390    }
391
392    #[test]
393    fn pinned_accessors_support_not_unpin() {
394        let optional = Optional::new(async { 42 });
395        futures::pin_mut!(optional);
396
397        assert!(optional.as_ref().pinned_as_ref().is_some());
398        assert!(optional.as_mut().pinned_as_mut().is_some());
399    }
400}