1use futures_util::Stream;
4use std::pin::Pin;
5use std::task::{Context, Poll};
6
7pub trait ZealObservable<T>: Stream<Item = T> + Sized {
9 fn filter<F>(self, predicate: F) -> FilterStream<Self, F>
11 where
12 F: FnMut(&T) -> bool,
13 {
14 FilterStream {
15 stream: self,
16 predicate,
17 }
18 }
19}
20
21impl<S, T> ZealObservable<T> for S where S: Stream<Item = T> {}
22
23pub use ZealObservable as ObservableExt;
25
26#[pin_project::pin_project]
28pub struct FilterStream<S, F> {
29 #[pin]
30 stream: S,
31 predicate: F,
32}
33
34impl<S, F, T> Stream for FilterStream<S, F>
35where
36 S: Stream<Item = T>,
37 F: FnMut(&T) -> bool,
38{
39 type Item = T;
40
41 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
42 let mut this = self.project();
43
44 loop {
45 match this.stream.as_mut().poll_next(cx) {
46 Poll::Ready(Some(item)) => {
47 if (this.predicate)(&item) {
48 return Poll::Ready(Some(item));
49 }
50 }
52 Poll::Ready(None) => return Poll::Ready(None),
53 Poll::Pending => return Poll::Pending,
54 }
55 }
56 }
57}