Skip to main content

rx_rust/operators/transforming/
map.rs

1use crate::utils::types::MaybeSend;
2use crate::{
3    observable::Observable,
4    observable::Subscription,
5    observer::{Flow, Observer, Termination},
6    utils::types::MarkerType,
7};
8use educe::Educe;
9use std::marker::PhantomData;
10
11/// Transforms items emitted by an Observable by applying a function to each item.
12/// See <https://reactivex.io/documentation/operators/map.html>
13///
14/// # Examples
15/// ```rust
16/// use rx_rust::{
17///     observable::ObservableExt,
18///     observer::Termination,
19///     operators::{
20///         creating::from_iter::FromIter,
21///         transforming::map::Map,
22///     },
23/// };
24///
25/// let mut values = Vec::new();
26/// let mut terminations = Vec::new();
27///
28/// let observable = Map::new(FromIter::new(vec![1, 2]), |value| value * 10);
29/// observable.subscribe_with_callback(
30///     |value| values.push(value),
31///     |termination| terminations.push(termination),
32/// );
33///
34/// assert_eq!(values, vec![10, 20]);
35/// assert_eq!(terminations, vec![Termination::Completed]);
36/// ```
37#[derive(Educe)]
38#[educe(Debug, Clone)]
39pub struct Map<T0, OE, F> {
40    source: OE,
41    callback: F,
42    _marker: MarkerType<T0>,
43}
44
45impl<T0, OE, F> Map<T0, OE, F> {
46    pub fn new<'or, T, E>(source: OE, callback: F) -> Self
47    where
48        OE: Observable<'or, T0, E>,
49        F: FnMut(T0) -> T,
50    {
51        Self {
52            source,
53            callback,
54            _marker: PhantomData,
55        }
56    }
57}
58
59impl<'or, T0, T, E, OE, F> Observable<'or, T, E> for Map<T0, OE, F>
60where
61    OE: Observable<'or, T0, E>,
62    F: FnMut(T0) -> T + MaybeSend + 'or,
63{
64    type D = OE::D;
65
66    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
67        let observer = MapObserver {
68            observer,
69            callback: self.callback,
70        };
71        self.source.subscribe(observer)
72    }
73}
74
75struct MapObserver<OR, F> {
76    observer: OR,
77    callback: F,
78}
79
80impl<T0, T, E, OR, F> Observer<T0, E> for MapObserver<OR, F>
81where
82    OR: Observer<T, E>,
83    F: FnMut(T0) -> T,
84{
85    fn on_next(&mut self, value: T0) -> Flow {
86        self.observer.on_next((self.callback)(value))
87    }
88
89    fn on_termination(self, termination: Termination<E>) {
90        self.observer.on_termination(termination)
91    }
92}