Skip to main content

rx_rust/operators/mathematical_aggregate/
reduce.rs

1use crate::utils::types::MaybeSend;
2use crate::{
3    observable::Observable,
4    observable::Subscription,
5    observer::{Flow, Observer, Termination},
6    utils::types::MarkerType,
7};
8use educe::Educe;
9use std::marker::PhantomData;
10
11/// Applies a function to each item emitted by an Observable, sequentially, and emits the final accumulated value.
12/// See <https://reactivex.io/documentation/operators/reduce.html>
13///
14/// # Examples
15/// ```rust
16/// use rx_rust::{
17///     observable::ObservableExt,
18///     observer::Termination,
19///     operators::{
20///         creating::from_iter::FromIter,
21///         mathematical_aggregate::reduce::Reduce,
22///     },
23/// };
24///
25/// let mut values = Vec::new();
26/// let mut terminations = Vec::new();
27///
28/// let observable = Reduce::new(FromIter::new(vec![1, 2, 3]), 0, |acc, value| acc + value);
29/// observable.subscribe_with_callback(
30///     |value| values.push(value),
31///     |termination| terminations.push(termination),
32/// );
33///
34/// assert_eq!(values, vec![6]);
35/// assert_eq!(terminations, vec![Termination::Completed]);
36/// ```
37#[derive(Educe)]
38#[educe(Debug, Clone)]
39pub struct Reduce<T, T1, OE, F> {
40    source: OE,
41    initial_value: T,
42    callback: F,
43    _marker: MarkerType<T1>,
44}
45
46impl<T, T1, OE, F> Reduce<T, T1, OE, F> {
47    pub fn new<'or, E>(source: OE, initial_value: T, callback: F) -> Self
48    where
49        OE: Observable<'or, T1, E>,
50        F: FnMut(T, T1) -> T,
51    {
52        Self {
53            source,
54            initial_value,
55            callback,
56            _marker: PhantomData,
57        }
58    }
59}
60
61impl<'or, T, T1, E, OE, F> Observable<'or, T, E> for Reduce<T, T1, OE, F>
62where
63    T: MaybeSend + 'or,
64    OE: Observable<'or, T1, E>,
65    F: FnMut(T, T1) -> T + MaybeSend + 'or,
66{
67    type D = OE::D;
68
69    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
70        let observer = ReduceObserver {
71            observer,
72            value: Some(self.initial_value),
73            callback: self.callback,
74        };
75        self.source.subscribe(observer)
76    }
77}
78
79struct ReduceObserver<T, OR, F> {
80    observer: OR,
81    value: Option<T>,
82    callback: F,
83}
84
85impl<T, T1, E, OR, F> Observer<T1, E> for ReduceObserver<T, OR, F>
86where
87    OR: Observer<T, E>,
88    F: FnMut(T, T1) -> T,
89{
90    fn on_next(&mut self, value: T1) -> Flow {
91        self.value = Some((self.callback)(self.value.take().unwrap(), value));
92        Flow::Continue
93    }
94
95    fn on_termination(mut self, termination: Termination<E>) {
96        // The final value ends the stream, so a downstream that stopped on it is not completed
97        // on top of that: it has already ended itself.
98        if matches!(termination, Termination::Completed)
99            && self.observer.on_next(self.value.take().unwrap()).is_stop()
100        {
101            return;
102        }
103        self.observer.on_termination(termination)
104    }
105}