rx_rust/operators/error_handling/
catch.rs1use 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#[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}