rx_rust/operators/mathematical_aggregate/
reduce.rs1use 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#[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 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}