Skip to main content

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