Skip to main content

zeal_sdk/
observables.rs

1//! Observable stream extensions for event processing
2
3use futures_util::Stream;
4use std::pin::Pin;
5use std::task::{Context, Poll};
6
7/// Extension trait for observable streams
8pub trait ZealObservable<T>: Stream<Item = T> + Sized {
9    /// Filter items based on a predicate
10    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
23/// Alias for the main observable extension trait
24pub use ZealObservable as ObservableExt;
25
26/// Stream that filters items
27#[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                    // Continue to next item if predicate failed
51                }
52                Poll::Ready(None) => return Poll::Ready(None),
53                Poll::Pending => return Poll::Pending,
54            }
55        }
56    }
57}