Skip to main content

rx_rust/operators/combining/
merge.rs

1use crate::utils::serialized_delivery::UpdateOutcome;
2use crate::utils::subscribe_with_context::{
3    self, SubscriptionContext, subscribe_with_context_owning_source,
4};
5use crate::utils::types::MaybeSend;
6use crate::{
7    disposable::Disposable,
8    observable::{Observable, Subscription},
9    observer::{Flow, Observer, Termination},
10};
11use educe::Educe;
12
13/// Combines multiple Observables into a single Observable that emits all of their emissions.
14/// See <https://reactivex.io/documentation/operators/merge.html>
15///
16/// # Examples
17/// ```rust
18/// use rx_rust::{
19///     observable::ObservableExt,
20///     observer::Termination,
21///     operators::{
22///         combining::merge::Merge,
23///         creating::from_iter::FromIter,
24///     },
25/// };
26///
27/// let mut values = Vec::new();
28/// let mut terminations = Vec::new();
29///
30/// let observable = Merge::new(
31///     FromIter::new(vec![1, 3]),
32///     FromIter::new(vec![2, 4]),
33/// );
34/// observable.subscribe_with_callback(
35///     |value| values.push(value),
36///     |termination| terminations.push(termination),
37/// );
38///
39/// assert_eq!(values, vec![1, 3, 2, 4]);
40/// assert_eq!(terminations, vec![Termination::Completed]);
41/// ```
42#[derive(Educe)]
43#[educe(Debug, Clone)]
44pub struct Merge<OE1, OE2> {
45    source_1: OE1,
46    source_2: OE2,
47}
48
49impl<OE1, OE2> Merge<OE1, OE2> {
50    pub fn new<'or, T, E>(source_1: OE1, source_2: OE2) -> Self
51    where
52        OE1: Observable<'or, T, E>,
53        OE2: Observable<'or, T, E>,
54    {
55        Self { source_1, source_2 }
56    }
57}
58
59impl<'or, T, E, OE1, OE2> Observable<'or, T, E> for Merge<OE1, OE2>
60where
61    T: MaybeSend + 'or,
62    E: MaybeSend + 'or,
63    OE1: Observable<'or, T, E>,
64    OE1::D: MaybeSend + 'or,
65    OE2: Observable<'or, T, E>,
66    OE2::D: MaybeSend + 'or,
67{
68    type D = subscribe_with_context::OwningDisposal<'or>;
69
70    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
71        let model = Model {
72            one_is_completed: false,
73        };
74        subscribe_with_context_owning_source(observer, model, |context| {
75            let subscription_1 = self.source_1.subscribe(MergeObserver(context.clone()));
76            let subscription_2 = self.source_2.subscribe(MergeObserver(context));
77            subscription_1.preceded_by_bound(subscription_2)
78        })
79    }
80}
81
82struct Model {
83    one_is_completed: bool,
84}
85
86struct MergeObserver<T, E, OR, D: Disposable>(SubscriptionContext<T, E, OR, Model, D>);
87
88impl<T, E, OR, D> Observer<T, E> for MergeObserver<T, E, OR, D>
89where
90    OR: Observer<T, E>,
91    D: Disposable,
92{
93    fn on_next(&mut self, value: T) -> Flow {
94        self.0.send_next(value)
95    }
96
97    fn on_termination(self, termination: Termination<E>) {
98        match termination {
99            completion @ Termination::Completed => {
100                let _ = self.0.update(|model| {
101                    if model.one_is_completed {
102                        UpdateOutcome::empty().with_termination_event(completion)
103                    } else {
104                        model.one_is_completed = true;
105                        UpdateOutcome::empty().without_events()
106                    }
107                });
108            }
109            error @ Termination::Error(_) => {
110                self.0.send_termination(error);
111            }
112        }
113    }
114}