Skip to main content

rx_rust/operators/others/
hook_on_next.rs

1use crate::utils::types::MaybeSend;
2use crate::{
3    observable::{Observable, Subscription},
4    observer::{Flow, Observer, Termination},
5};
6use educe::Educe;
7
8/// Invokes a callback for each item emitted by the source Observable.
9///
10/// # Examples
11/// ```rust
12/// use rx_rust::{
13///     observable::ObservableExt,
14///     observer::Termination,
15///     operators::{
16///         creating::from_iter::FromIter,
17///         others::hook_on_next::HookOnNext,
18///     },
19/// };
20///
21/// let mut values = Vec::new();
22/// let mut terminations = Vec::new();
23///
24/// let observable = HookOnNext::new(FromIter::new(vec![1, 2]), |observer, value| {
25///     observer.on_next(value * 10)
26/// });
27/// observable.subscribe_with_callback(
28///     |value| values.push(value),
29///     |termination| terminations.push(termination),
30/// );
31///
32/// assert_eq!(values, vec![10, 20]);
33/// assert_eq!(terminations, vec![Termination::Completed]);
34/// ```
35#[derive(Educe)]
36#[educe(Debug, Clone)]
37pub struct HookOnNext<OE, F> {
38    source: OE,
39    callback: F,
40}
41
42impl<OE, F> HookOnNext<OE, F> {
43    pub fn new<'or, T, E>(source: OE, callback: F) -> Self
44    where
45        OE: Observable<'or, T, E>,
46        F: FnMut(&mut dyn Observer<T, E>, T) -> Flow,
47    {
48        Self { source, callback }
49    }
50}
51
52impl<'or, T, E, OE, F> Observable<'or, T, E> for HookOnNext<OE, F>
53where
54    OE: Observable<'or, T, E>,
55    F: FnMut(&mut dyn Observer<T, E>, T) -> Flow + MaybeSend + 'or,
56{
57    type D = OE::D;
58
59    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
60        let observer = HookOnNextObserver {
61            observer,
62            callback: self.callback,
63        };
64        self.source.subscribe(observer)
65    }
66}
67
68struct HookOnNextObserver<OR, F> {
69    observer: OR,
70    callback: F,
71}
72
73impl<T, E, OR, F> Observer<T, E> for HookOnNextObserver<OR, F>
74where
75    OR: Observer<T, E>,
76    F: FnMut(&mut dyn Observer<T, E>, T) -> Flow,
77{
78    fn on_next(&mut self, value: T) -> Flow {
79        (self.callback)(&mut self.observer, value)
80    }
81
82    fn on_termination(self, termination: Termination<E>) {
83        self.observer.on_termination(termination);
84    }
85}