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