Skip to main content

rx_rust/operators/error_handling/
catch.rs

1use crate::delegate_disposal;
2use crate::disposable::{
3    Disposable, chain_disposal::ChainDisposal, shared_disposal::SharedDisposal,
4};
5use crate::observable::Subscription;
6use crate::utils::types::MaybeSend;
7use crate::{
8    observable::Observable,
9    observer::{Flow, Observer, Termination},
10    utils::types::MarkerType,
11};
12use educe::Educe;
13use std::marker::PhantomData;
14
15/// Catches errors on the observable to be handled by returning a new observable or throwing an error.
16/// See <https://reactivex.io/documentation/operators/catch.html>
17///
18/// # Examples
19/// ```rust
20/// use rx_rust::{
21///     observable::ObservableExt,
22///     observer::Termination,
23///     operators::{
24///         creating::{just::Just, throw::Throw},
25///         error_handling::catch::Catch,
26///     },
27/// };
28///
29/// let mut values = Vec::new();
30/// let mut terminations = Vec::new();
31///
32/// let observable = Catch::new(Throw::new("boom").with_item_type(), |error| Just::new(error));
33/// observable.subscribe_with_callback(
34///     |value| values.push(value),
35///     |termination| terminations.push(termination),
36/// );
37///
38/// assert_eq!(values, vec!["boom"]);
39/// assert_eq!(terminations, vec![Termination::Completed]);
40/// ```
41#[derive(Educe)]
42#[educe(Debug, Clone)]
43pub struct Catch<E0, OE, F> {
44    source: OE,
45    callback: F,
46    _marker: MarkerType<E0>,
47}
48
49impl<E0, OE, F> Catch<E0, OE, F> {
50    pub fn new<'or, T, E, OE1>(source: OE, callback: F) -> Self
51    where
52        OE: Observable<'or, T, E0>,
53        OE1: Observable<'or, T, E>,
54        F: FnOnce(E0) -> OE1,
55    {
56        Self {
57            source,
58            callback,
59            _marker: PhantomData,
60        }
61    }
62}
63
64delegate_disposal!(
65    Disposal<D, D1>,
66    ChainDisposal<SharedDisposal<Subscription<D1>>, D>,
67    where D: Disposable, D1: Disposable
68);
69
70impl<'or, T, E0, E, OE, OE1, F> Observable<'or, T, E> for Catch<E0, OE, F>
71where
72    E: 'or,
73    OE: Observable<'or, T, E0>,
74    OE1: Observable<'or, T, E>,
75    OE1::D: MaybeSend + 'or,
76    F: FnOnce(E0) -> OE1 + MaybeSend + 'or,
77{
78    type D = Disposal<OE::D, OE1::D>;
79
80    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
81        let shared_disposal = SharedDisposal::default();
82        let observer = CatchObserver {
83            observer,
84            callback: self.callback,
85            shared_disposal: shared_disposal.clone(),
86            _marker: PhantomData,
87        };
88        self.source
89            .subscribe(observer)
90            .preceded_by(shared_disposal)
91            .map_into()
92    }
93}
94
95struct CatchObserver<E, OR, F, D: Disposable> {
96    observer: OR,
97    callback: F,
98    shared_disposal: SharedDisposal<Subscription<D>>,
99    _marker: MarkerType<E>,
100}
101
102impl<'or, T, E0, E, OR, OE1, F> Observer<T, E0> for CatchObserver<E, OR, F, OE1::D>
103where
104    OR: Observer<T, E> + MaybeSend + 'or,
105    OE1: Observable<'or, T, E>,
106    F: FnOnce(E0) -> OE1,
107{
108    fn on_next(&mut self, value: T) -> Flow {
109        self.observer.on_next(value)
110    }
111
112    fn on_termination(self, termination: Termination<E0>) {
113        match termination {
114            Termination::Completed => self.observer.on_termination(Termination::Completed),
115            Termination::Error(error) => {
116                self.shared_disposal.replace(|| {
117                    let observable = (self.callback)(error);
118                    observable.subscribe(self.observer)
119                });
120            }
121        }
122    }
123}