Skip to main content

rx_rust/operators/transforming/
switch_map.rs

1use super::map::Map;
2use crate::operators::combining::switch::Switch;
3use crate::utils::subscribe_with_context;
4use crate::utils::types::MaybeSend;
5use crate::{
6    observable::Observable, observable::Subscription, observer::Observer, utils::types::MarkerType,
7};
8use educe::Educe;
9use std::marker::PhantomData;
10
11/// Projects each source value to an Observable which is merged in the output Observable, emitting values only from the most recently projected Observable.
12/// See <https://reactivex.io/documentation/operators/switch.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::switch_map::SwitchMap,
22///     },
23/// };
24///
25/// let mut values = Vec::new();
26/// let mut terminations = Vec::new();
27///
28/// let observable = SwitchMap::new(FromIter::new(vec![1, 2]), |value| {
29///     FromIter::new(vec![value, value + 10])
30/// });
31/// observable.subscribe_with_callback(
32///     |value| values.push(value),
33///     |termination| terminations.push(termination),
34/// );
35///
36/// assert_eq!(values, vec![1, 11, 2, 12]);
37/// assert_eq!(terminations, vec![Termination::Completed]);
38/// ```
39#[derive(Educe)]
40#[educe(Debug, Clone)]
41pub struct SwitchMap<T0, OE, OE1, F> {
42    source: OE,
43    callback: F,
44    _marker: MarkerType<(T0, OE1)>,
45}
46
47impl<T0, OE, OE1, F> SwitchMap<T0, OE, OE1, F> {
48    pub fn new<'or, T, E>(source: OE, callback: F) -> Self
49    where
50        OE: Observable<'or, T0, E>,
51        OE1: Observable<'or, T, E>,
52        F: FnMut(T0) -> OE1,
53    {
54        Self {
55            source,
56            callback,
57            _marker: PhantomData,
58        }
59    }
60}
61
62impl<'or, T0, T, E, OE, OE1, F> Observable<'or, T, E> for SwitchMap<T0, OE, OE1, F>
63where
64    T: MaybeSend + 'or,
65    E: MaybeSend + 'or,
66    OE: Observable<'or, T0, E>,
67    OE::D: MaybeSend + 'or,
68    OE1: Observable<'or, T, E>,
69    OE1::D: MaybeSend + 'or,
70    F: FnMut(T0) -> OE1 + MaybeSend + 'or,
71{
72    type D = subscribe_with_context::OwningDisposal<'or>;
73
74    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
75        let observable = Map::new(self.source, self.callback);
76        let observable = Switch::new(observable);
77        observable.subscribe(observer)
78    }
79}