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#[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 pub fn new(task: T) -> Self {
53 Self {
54 task: Some(task),
55 waker: None,
56 }
57 }
58
59 pub fn with_future(future: T) -> Self
61 where
62 T: Future,
63 {
64 Self::new(future)
65 }
66
67 pub fn with_stream(stream: T) -> Self
69 where
70 T: Stream,
71 {
72 Self::new(stream)
73 }
74
75 pub fn take(&mut self) -> Option<T> {
77 let fut = self.task.take();
78 if let Some(waker) = self.waker.take() {
82 waker.wake();
83 }
84 fut
85 }
86
87 pub fn is_some(&self) -> bool {
89 self.task.is_some()
90 }
91
92 pub fn is_none(&self) -> bool {
94 self.task.is_none()
95 }
96
97 pub fn as_ref(&self) -> Option<&T> {
99 self.task.as_ref()
100 }
101
102 pub fn as_mut(&mut self) -> Option<&mut T> {
104 self.task.as_mut()
105 }
106
107 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 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 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 pub fn pinned_as_mut(self: Pin<&mut Self>) -> Option<Pin<&mut T>> {
137 self.project().task.as_pin_mut()
138 }
139
140 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 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}