Skip to main content

rx_rust/operators/mathematical_aggregate/
sum.rs

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