Skip to main content

rx_rust/operators/others/
hook_on_termination.rs

1use crate::utils::types::MaybeSend;
2use crate::{
3    observable::{Observable, Subscription},
4    observer::{Flow, Observer, Termination, boxed_observer::BoxedObserver},
5};
6use educe::Educe;
7
8/// Invokes a callback when the source Observable terminates.
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_termination::HookOnTermination,
18///     },
19/// };
20/// use rx_rust::observer::Observer;
21/// use std::cell::Cell;
22/// use std::rc::Rc;
23///
24/// let mut values = Vec::new();
25/// let mut terminations = Vec::new();
26///
27/// let observable =
28///     HookOnTermination::new(FromIter::new(vec![1]), move |observer, termination| {
29///         // Do whatever you want here
30///         observer.on_termination(termination);
31///     });
32/// observable.subscribe_with_callback(
33///     |value| values.push(value),
34///     |termination| terminations.push(termination),
35/// );
36///
37/// assert_eq!(values, vec![1]);
38/// assert_eq!(terminations, vec![Termination::Completed]);
39/// ```
40#[derive(Educe)]
41#[educe(Debug, Clone)]
42pub struct HookOnTermination<OE, F> {
43    source: OE,
44    callback: F,
45}
46
47impl<OE, F> HookOnTermination<OE, F> {
48    pub fn new<'or, T, E>(source: OE, callback: F) -> Self
49    where
50        OE: Observable<'or, T, E>,
51        F: FnOnce(BoxedObserver<'or, T, E>, Termination<E>),
52    {
53        Self { source, callback }
54    }
55}
56
57impl<'or, T, E, OE, F> Observable<'or, T, E> for HookOnTermination<OE, F>
58where
59    T: 'or,
60    E: 'or,
61    OE: Observable<'or, T, E>,
62    F: FnOnce(BoxedObserver<'or, T, E>, Termination<E>) + MaybeSend + 'or,
63{
64    type D = OE::D;
65
66    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
67        let observer = HookOnTerminationObserver {
68            observer: BoxedObserver::new(observer),
69            callback: self.callback,
70        };
71        self.source.subscribe(observer)
72    }
73}
74
75struct HookOnTerminationObserver<OR, F> {
76    observer: OR,
77    callback: F,
78}
79
80impl<T, E, OR, F> Observer<T, E> for HookOnTerminationObserver<OR, F>
81where
82    OR: Observer<T, E>,
83    F: FnOnce(OR, Termination<E>),
84{
85    fn on_next(&mut self, value: T) -> Flow {
86        self.observer.on_next(value)
87    }
88
89    fn on_termination(self, termination: Termination<E>) {
90        (self.callback)(self.observer, termination);
91    }
92}