Skip to main content

rx_rust/operators/error_handling/
map_err.rs

1use crate::utils::types::MaybeSend;
2use crate::{
3    observable::{Observable, Subscription},
4    observer::{Flow, Observer, Termination},
5    utils::types::MarkerType,
6};
7use educe::Educe;
8use std::marker::PhantomData;
9
10/// Transforms an Observable's error while leaving its items unchanged.
11///
12/// # Examples
13/// ```rust
14/// use rx_rust::{
15///     observable::ObservableExt,
16///     observer::Termination,
17///     operators::creating::throw::Throw,
18/// };
19///
20/// let mut terminations = Vec::new();
21///
22/// Throw::new("boom")
23///     .with_item_type::<i32>()
24///     .map_err(|error| error.len())
25///     .subscribe_with_callback(
26///         |_| -> () { unreachable!() },
27///         |termination| terminations.push(termination),
28///     );
29///
30/// assert_eq!(terminations, vec![Termination::Error(4)]);
31/// ```
32#[derive(Educe)]
33#[educe(Debug, Clone)]
34pub struct MapErr<E, OE, F> {
35    source: OE,
36    callback: F,
37    _marker: MarkerType<E>,
38}
39
40impl<E, OE, F> MapErr<E, OE, F> {
41    pub fn new<'or, T, E1>(source: OE, callback: F) -> Self
42    where
43        OE: Observable<'or, T, E>,
44        F: FnOnce(E) -> E1,
45    {
46        Self {
47            source,
48            callback,
49            _marker: PhantomData,
50        }
51    }
52}
53
54impl<'or, T, E, E1, OE, F> Observable<'or, T, E1> for MapErr<E, OE, F>
55where
56    OE: Observable<'or, T, E>,
57    F: FnOnce(E) -> E1 + MaybeSend + 'or,
58{
59    type D = OE::D;
60
61    fn subscribe(self, observer: impl Observer<T, E1> + MaybeSend + 'or) -> Subscription<Self::D> {
62        self.source.subscribe(MapErrObserver {
63            observer,
64            callback: self.callback,
65        })
66    }
67}
68
69struct MapErrObserver<OR, F> {
70    observer: OR,
71    callback: F,
72}
73
74impl<T, E, E1, OR, F> Observer<T, E> for MapErrObserver<OR, F>
75where
76    OR: Observer<T, E1>,
77    F: FnOnce(E) -> E1,
78{
79    fn on_next(&mut self, value: T) -> Flow {
80        self.observer.on_next(value)
81    }
82
83    fn on_termination(self, termination: Termination<E>) {
84        match termination {
85            Termination::Completed => self.observer.on_termination(Termination::Completed),
86            Termination::Error(error) => self
87                .observer
88                .on_termination(Termination::Error((self.callback)(error))),
89        }
90    }
91}