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}