Skip to main content

rx_rust/operators/utility/
do_after_disposal.rs

1use crate::utils::types::MaybeSend;
2use crate::{
3    disposable::{callback_disposal::CallbackDisposal, chain_disposal::ChainDisposal},
4    observable::Observable,
5    observable::Subscription,
6    observer::Observer,
7};
8use educe::Educe;
9
10/// Invokes a callback when the subscription is disposed.
11/// See <https://reactivex.io/documentation/operators/do.html>
12///
13/// # Examples
14/// ```rust
15/// use rx_rust::{
16///     observable::ObservableExt,
17///     observer::Termination,
18///     operators::{
19///         creating::from_iter::FromIter,
20///         utility::do_after_disposal::DoAfterDisposal,
21///     },
22/// };
23/// use std::sync::{Arc, Mutex};
24///
25/// let disposed = Arc::new(Mutex::new(false));
26/// let disposed_observer = Arc::clone(&disposed);
27///
28/// let subscription = DoAfterDisposal::new(
29///     FromIter::new(vec![1, 2]),
30///     move || *disposed_observer.lock().unwrap() = true,
31/// )
32/// .subscribe_with_callback(
33///     |_value| {},
34///     |_termination| {},
35/// );
36///
37/// drop(subscription);
38///
39/// assert!(*disposed.lock().unwrap());
40/// ```
41#[derive(Educe)]
42#[educe(Debug, Clone)]
43pub struct DoAfterDisposal<OE, F> {
44    source: OE,
45    callback: F,
46}
47
48impl<OE, F> DoAfterDisposal<OE, F> {
49    pub fn new<'or, T, E>(source: OE, callback: F) -> Self
50    where
51        OE: Observable<'or, T, E>,
52        F: FnOnce(),
53    {
54        Self { source, callback }
55    }
56}
57
58impl<'or, T, E, OE, F> Observable<'or, T, E> for DoAfterDisposal<OE, F>
59where
60    T: 'or,
61    E: 'or,
62    OE: Observable<'or, T, E>,
63    F: FnOnce(),
64{
65    type D = ChainDisposal<OE::D, CallbackDisposal<F>>;
66
67    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
68        self.source
69            .subscribe(observer)
70            .then(CallbackDisposal::new(self.callback))
71    }
72}