mod tests_utils;
use crate::tests_utils::DURATION_10_MS;
use crate::tests_utils::checker::State;
use crate::tests_utils::test_channel::{ChannelState, test_channel};
use crate::tests_utils::test_runtime::block_on;
use rx_rust::disposable::Disposable;
use rx_rust::disposable::callback_disposal::CallbackDisposal;
use rx_rust::observable::Subscription;
use rx_rust::operators::creating::empty::Empty;
use rx_rust::operators::creating::throw::Throw;
use rx_rust::scheduler::Scheduler;
use rx_rust::subject::behavior_subject::BehaviorSubject;
use rx_rust::utils::mutable::Mutable;
use rx_rust::utils::mutable::MutableExt;
use rx_rust::utils::mutable::MutableHelper;
use rx_rust::utils::types::Shared;
use rx_rust::{
observable::{Observable, ObservableExt},
observer::{Observer, Termination, boxed_observer::BoxedObserver},
operators::{creating::create::Create, transforming::group_by::GroupBy},
subject::publish_subject::PublishSubject,
};
use std::convert::Infallible;
use tests_utils::{checker::Checker, test_struct::TestStruct};
#[test]
fn test_completed() {
let mut subject = PublishSubject::default();
let (checker, observer) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| value % 2);
let mut index = -1;
let _subscription = observable
.flat_map(|v| {
index += 1;
let i = index;
v.map(move |v| 10 * i + v)
})
.subscribe(observer);
assert!(checker.values().is_empty());
assert_eq!(checker.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Active);
subject.on_termination(Termination::<Infallible>::Completed);
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Completed);
}
#[test]
fn test_error() {
let mut subject = PublishSubject::default();
let (checker, observer) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| value % 2);
let mut index = -1;
let _subscription = observable
.flat_map(|v| {
index += 1;
let i = index;
v.map(move |v| 10 * i + v)
})
.subscribe(observer);
assert!(checker.values().is_empty());
assert_eq!(checker.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Active);
subject.on_termination(Termination::Error("error"));
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Error("error"));
}
#[test]
fn test_unsubscribe() {
let mut subject = PublishSubject::default();
let (checker_1, observer_1) = Checker::new();
let (checker_2, observer_2) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| value % 2);
let observable_1 = observable;
let observable_2 = observable_1.clone();
let mut index = -1;
let subscription_1 = observable_1
.flat_map(|v| {
index += 1;
let i = index;
v.map(move |v| 10 * i + v)
})
.subscribe(observer_1);
let mut index = -1;
let _subscription_2 = observable_2
.flat_map(|v| {
index += 1;
let i = index;
v.map(move |v| 10 * i + v)
})
.subscribe(observer_2);
assert!(checker_1.values().is_empty());
assert_eq!(checker_1.state(), State::Active);
assert!(checker_2.values().is_empty());
assert_eq!(checker_2.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker_1.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_1.state(), State::Active);
assert_eq!(checker_2.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_2.state(), State::Active);
subscription_1.dispose();
assert_eq!(checker_1.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_1.state(), State::Dropped);
assert_eq!(checker_2.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_2.state(), State::Active);
assert!(subject.on_next(111).is_continue());
assert_eq!(checker_1.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_1.state(), State::Dropped);
assert_eq!(checker_2.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19, 121]);
assert_eq!(checker_2.state(), State::Active);
subject.on_termination(Termination::Error("error"));
assert_eq!(checker_1.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_1.state(), State::Dropped);
assert_eq!(checker_2.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19, 121]);
assert_eq!(checker_2.state(), State::Error("error"));
}
#[test]
fn test_ref() {
let value_1 = 111;
let value_2 = 222;
let value_3 = 333;
let error = 444;
let mut subject: PublishSubject<'_, &i32, &i32> = PublishSubject::default();
let (checker, observer) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| *value % 2);
let mut index = -1;
let _subscription = observable
.flat_map(|v| {
index += 1;
let i = index;
v.map(move |v| (i, v))
})
.subscribe(observer);
assert!(checker.values().is_empty());
assert_eq!(checker.state(), State::Active);
assert!(subject.on_next(&value_1).is_continue());
assert_eq!(checker.values(), [(0, &value_1)]);
assert_eq!(checker.state(), State::Active);
assert!(subject.on_next(&value_2).is_continue());
assert_eq!(checker.values(), [(0, &value_1), (1, &value_2)]);
assert_eq!(checker.state(), State::Active);
assert!(subject.on_next(&value_3).is_continue());
assert_eq!(
checker.values(),
[(0, &value_1), (1, &value_2), (0, &value_3)]
);
assert_eq!(checker.state(), State::Active);
subject.on_termination(Termination::Error(&error));
assert_eq!(
checker.values(),
[(0, &value_1), (1, &value_2), (0, &value_3)]
);
assert_eq!(checker.state(), State::Error(&error));
}
#[test]
fn test_async() {
block_on(|runtime| async move {
let subject = PublishSubject::default();
let (checker, observer) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| value % 2);
let subscription = runtime
.spawn(async {
let mut index = -1;
observable
.flat_map(move |v| {
index += 1;
let i = index;
v.map(move |v| 10 * i + v)
})
.subscribe(observer)
})
.await
.unwrap();
assert!(checker.values().is_empty());
assert_eq!(checker.state(), State::Active);
let mut subject_cloned = subject.clone();
runtime
.spawn(async move {
for i in 0..=9 {
assert!(subject_cloned.on_next(i).is_continue());
}
})
.await
.unwrap();
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Active);
runtime
.spawn(async { subscription.dispose() })
.await
.unwrap();
runtime.sleep(DURATION_10_MS).await;
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Dropped);
let subject_cloned = subject.clone();
runtime
.spawn(async move {
subject_cloned.on_termination(Termination::Error("error"));
})
.await
.unwrap();
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Dropped);
});
}
#[test]
fn test_subscribe_by_different_observer() {
let mut subject = PublishSubject::default();
let (checker_1, observer_1) = Checker::new();
let (checker_2, observer_2) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| value % 2);
let observable_1 = observable;
let observable_2 = observable_1.clone();
let mut index_1 = -1;
let mut index_2 = -1;
let _subscription_1 = observable_1
.flat_map(|v| {
index_1 += 1;
let i = index_1;
v.map(move |v| 10 * i + v)
})
.subscribe(observer_1);
let (on_next, on_termination) = observer_2.into_callbacks();
let _subscription_2 = observable_2
.flat_map(|v| {
index_2 += 1;
let i = index_2;
v.map(move |v| 10 * i + v)
})
.subscribe_with_callback(on_next, on_termination);
assert!(checker_1.values().is_empty());
assert_eq!(checker_1.state(), State::Active);
assert!(checker_2.values().is_empty());
assert_eq!(checker_2.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker_1.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_1.state(), State::Active);
assert_eq!(checker_2.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_2.state(), State::Active);
subject.on_termination(Termination::Error("error"));
assert_eq!(checker_1.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_1.state(), State::Error("error"));
assert_eq!(checker_2.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker_2.state(), State::Error("error"));
}
#[test]
fn test_unsub_on_next_by_take() {
let mut index = -1;
let (mut sender, observable, channel_checker) = test_channel::<'_, _, Infallible>();
let (checker, observer) = Checker::new();
let observable = observable.group_by(|value| value % 2);
let _subscription = observable
.flat_map(|v| {
index += 1;
let i = index;
v.map(move |v| 10 * i + v)
})
.take(1)
.subscribe(observer);
assert!(checker.values().is_empty());
assert_eq!(checker.state(), State::Active);
assert_eq!(channel_checker.state(), ChannelState::Subscribed);
assert!(sender.on_next(111).is_stop());
assert_eq!(checker.values(), [111]);
assert_eq!(checker.state(), State::Completed);
assert_eq!(channel_checker.state(), ChannelState::Unsubscribed);
}
#[test]
fn test_multiple_operation() {
let mut subject = PublishSubject::default();
let (checker, observer) = Checker::new();
let observable = subject.clone();
let mut index = -1;
let observable = observable.group_by(|value| value % 4).group_by(|_| {
index += 1;
index % 2
});
let mut index_1 = -1;
let mut index_2 = -1;
let _subscription = observable
.flat_map(|v| {
index_1 += 1;
let i = index_1;
v.map(move |v| (i, v))
})
.flat_map(|v| {
index_2 += 1;
let i0 = v.0;
let i = index_2;
v.1.map(move |v| 100 * i0 + 10 * i + v)
})
.subscribe(observer);
assert!(checker.values().is_empty());
assert_eq!(checker.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker.values(), [0, 111, 22, 133, 4, 115, 26, 137, 8, 119]);
assert_eq!(checker.state(), State::Active);
subject.on_termination(Termination::<Infallible>::Completed);
assert_eq!(checker.values(), [0, 111, 22, 133, 4, 115, 26, 137, 8, 119]);
assert_eq!(checker.state(), State::Completed);
}
#[test]
fn test_without_convenient_api() {
let mut subject = PublishSubject::default();
let (checker, observer) = Checker::new();
let observable = subject.clone();
let observable = GroupBy::new(observable, |value| value % 2);
let mut index = -1;
let _subscription = observable
.flat_map(|v| {
index += 1;
let i = index;
v.map(move |v| 10 * i + v)
})
.subscribe(observer);
assert!(checker.values().is_empty());
assert_eq!(checker.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Active);
subject.on_termination(Termination::Error("error"));
assert_eq!(checker.values(), [0, 11, 2, 13, 4, 15, 6, 17, 8, 19]);
assert_eq!(checker.state(), State::Error("error"));
}
#[test]
fn test_revert_completed() {
let mut subject: PublishSubject<'_, i32, &'static str> = PublishSubject::default();
let (checker_1, observer_1) = Checker::new();
let (checker_2, observer_2) = Checker::new();
let (checker_3, observer_3) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| value % 2);
let observable_1 = observable.clone().merge_all();
let observable_2 = observable.clone().concat_all();
let observable_3 = observable.switch();
let _subscription_1 = observable_1.subscribe(observer_1);
let _subscription_2 = observable_2.subscribe(observer_2);
let _subscription_3 = observable_3.subscribe(observer_3);
assert!(checker_1.values().is_empty());
assert_eq!(checker_1.state(), State::Active);
assert!(checker_2.values().is_empty());
assert_eq!(checker_2.state(), State::Active);
assert!(checker_3.values().is_empty());
assert_eq!(checker_3.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker_1.values(), [0, 1, 2, 3, 4, 5, 6, 7, 8, 9]);
assert_eq!(checker_1.state(), State::Active);
assert_eq!(checker_2.values(), [0, 2, 4, 6, 8]);
assert_eq!(checker_2.state(), State::Active);
assert_eq!(checker_3.values(), [0, 1, 3, 5, 7, 9]);
assert_eq!(checker_3.state(), State::Active);
subject.on_termination(Termination::Completed);
assert_eq!(checker_1.values(), [0, 1, 2, 3, 4, 5, 6, 7, 8, 9]);
assert_eq!(checker_1.state(), State::Completed);
assert_eq!(checker_2.values(), [0, 2, 4, 6, 8, 1, 3, 5, 7, 9]);
assert_eq!(checker_2.state(), State::Completed);
assert_eq!(checker_3.values(), [0, 1, 3, 5, 7, 9]);
assert_eq!(checker_3.state(), State::Completed);
}
#[test]
fn test_revert_error() {
let mut subject: PublishSubject<'_, i32, &'static str> = PublishSubject::default();
let (checker_1, observer_1) = Checker::new();
let (checker_2, observer_2) = Checker::new();
let (checker_3, observer_3) = Checker::new();
let observable = subject.clone();
let observable = observable.group_by(|value| value % 2);
let observable_1 = observable.clone().merge_all();
let observable_2 = observable.clone().concat_all();
let observable_3 = observable.switch();
let _subscription_1 = observable_1.subscribe(observer_1);
let _subscription_2 = observable_2.subscribe(observer_2);
let _subscription_3 = observable_3.subscribe(observer_3);
assert!(checker_1.values().is_empty());
assert_eq!(checker_1.state(), State::Active);
assert!(checker_2.values().is_empty());
assert_eq!(checker_2.state(), State::Active);
assert!(checker_3.values().is_empty());
assert_eq!(checker_3.state(), State::Active);
for i in 0..=9 {
assert!(subject.on_next(i).is_continue());
}
assert_eq!(checker_1.values(), [0, 1, 2, 3, 4, 5, 6, 7, 8, 9]);
assert_eq!(checker_1.state(), State::Active);
assert_eq!(checker_2.values(), [0, 2, 4, 6, 8]);
assert_eq!(checker_2.state(), State::Active);
assert_eq!(checker_3.values(), [0, 1, 3, 5, 7, 9]);
assert_eq!(checker_3.state(), State::Active);
subject.on_termination(Termination::Error("error"));
assert_eq!(checker_1.values(), [0, 1, 2, 3, 4, 5, 6, 7, 8, 9]);
assert_eq!(checker_1.state(), State::Error("error"));
assert_eq!(checker_2.values(), [0, 2, 4, 6, 8]);
assert_eq!(checker_2.state(), State::Error("error"));
assert_eq!(checker_3.values(), [0, 1, 3, 5, 7, 9]);
assert_eq!(checker_3.state(), State::Error("error"));
}
#[test]
fn test_next_on_sub() {
let mut subject = BehaviorSubject::new(111);
let (termination_checker, termination_observer) = Checker::<Infallible, _>::new();
let checker_sub_vec = Shared::new(Mutable::new(Vec::new()));
let observable = subject.clone().group_by(|value| value % 2);
let checker_sub_vec_cloned = checker_sub_vec.clone();
let _subscription = observable.subscribe_with_callback(
move |group| {
let (checker, observer) = Checker::new();
let sub = group.subscribe(observer);
checker_sub_vec_cloned.with_mut(|values| values.push((checker, sub)));
},
|termination| {
termination_observer.on_termination(termination);
},
);
assert_eq!(checker_sub_vec.with_ref(Vec::len), 1);
checker_sub_vec.with_ref(|checker_sub_vec| {
for (index, (checker, _)) in checker_sub_vec.iter().enumerate() {
match index {
0 => {
assert_eq!(checker.values(), [111]);
assert_eq!(checker.state(), State::Active);
}
_ => panic!(),
}
}
});
assert_eq!(termination_checker.state(), State::Active);
assert!(subject.on_next(222).is_continue());
assert_eq!(checker_sub_vec.with_ref(Vec::len), 2);
checker_sub_vec.with_ref(|checker_sub_vec| {
for (index, (checker, _)) in checker_sub_vec.iter().enumerate() {
match index {
0 => assert_eq!(checker.values(), [111]),
1 => assert_eq!(checker.values(), [222]),
_ => panic!(),
}
}
});
assert!(subject.on_next(333).is_continue());
assert_eq!(checker_sub_vec.with_ref(Vec::len), 2);
checker_sub_vec.with_ref(|checker_sub_vec| {
for (index, (checker, _)) in checker_sub_vec.iter().enumerate() {
match index {
0 => assert_eq!(checker.values(), [111, 333]),
1 => assert_eq!(checker.values(), [222]),
_ => panic!(),
}
}
});
subject.on_termination(Termination::<Infallible>::Completed);
checker_sub_vec.with_ref(|checker_sub_vec| {
for (index, (checker, _)) in checker_sub_vec.iter().enumerate() {
match index {
0 => {
assert_eq!(checker.values(), [111, 333]);
assert_eq!(checker.state(), State::Completed);
}
1 => {
assert_eq!(checker.values(), [222]);
assert_eq!(checker.state(), State::Completed);
}
_ => panic!(),
}
}
});
assert_eq!(termination_checker.state(), State::Completed);
}
#[test]
fn test_complete_on_sub() {
let (termination_checker, termination_observer) = Checker::<Infallible, _>::new();
let group_vec = Shared::new(Mutable::new(Vec::new()));
let observable = Empty.group_by(|_value| 0);
let group_vec_cloned = group_vec.clone();
let _subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|termination| {
termination_observer.on_termination(termination);
},
);
assert_eq!(group_vec.with_ref(Vec::len), 0);
assert_eq!(termination_checker.state(), State::Completed);
}
#[test]
fn test_error_on_sub() {
let (termination_checker, termination_observer) = Checker::<Infallible, _>::new();
let group_vec = Shared::new(Mutable::new(Vec::new()));
let observable = Throw::new("error").group_by(|_value| 0);
let group_vec_cloned = group_vec.clone();
let _subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|termination| {
termination_observer.on_termination(termination);
},
);
assert_eq!(group_vec.with_ref(Vec::len), 0);
assert_eq!(termination_checker.state(), State::Error("error"));
}
#[test]
fn test_next_on_unsub() {
let (termination_checker, termination_observer) = Checker::<Infallible, _>::new();
let group_vec = Shared::new(Mutable::new(Vec::new()));
let observable = Create::new(|observer: BoxedObserver<'_, i32, Infallible>| {
Subscription::new(CallbackDisposal::new(move || {
let mut observer = observer;
assert!(observer.on_next(111).is_continue());
}))
});
let observable = observable.group_by(|value| value % 2);
let group_vec_cloned = group_vec.clone();
let subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|termination| {
termination_observer.on_termination(termination);
},
);
assert_eq!(group_vec.with_ref(Vec::len), 0);
assert_eq!(termination_checker.state(), State::Active);
subscription.dispose();
assert_eq!(group_vec.with_ref(Vec::len), 1);
assert_eq!(termination_checker.state(), State::Dropped);
}
#[test]
fn test_complete_on_unsub() {
let (termination_checker, termination_observer) = Checker::<Infallible, _>::new();
let group_vec = Shared::new(Mutable::new(Vec::new()));
let observable = Create::new(|observer: BoxedObserver<'_, i32, Infallible>| {
Subscription::new(CallbackDisposal::new(move || {
observer.on_termination(Termination::Completed);
}))
});
let observable = observable.group_by(|value| value % 2);
let group_vec_cloned = group_vec.clone();
let subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|termination| {
termination_observer.on_termination(termination);
},
);
assert_eq!(group_vec.with_ref(Vec::len), 0);
assert_eq!(termination_checker.state(), State::Active);
subscription.dispose();
assert_eq!(group_vec.with_ref(Vec::len), 0);
assert_eq!(termination_checker.state(), State::Completed);
}
#[test]
fn test_error_on_unsub() {
let (termination_checker, termination_observer) = Checker::<Infallible, _>::new();
let group_vec = Shared::new(Mutable::new(Vec::new()));
let observable = Create::new(|observer: BoxedObserver<'_, i32, &str>| {
Subscription::new(CallbackDisposal::new(move || {
observer.on_termination(Termination::Error("error"));
}))
});
let observable = observable.group_by(|value| value % 2);
let group_vec_cloned = group_vec.clone();
let subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|termination| {
termination_observer.on_termination(termination);
},
);
assert_eq!(group_vec.with_ref(Vec::len), 0);
assert_eq!(termination_checker.state(), State::Active);
subscription.dispose();
assert_eq!(group_vec.with_ref(Vec::len), 0);
assert_eq!(termination_checker.state(), State::Error("error"));
}
#[test]
fn test_subscribe_groups_late_with_buffered_values() {
let (mut sender, observable, _channel_checker) = test_channel::<'_, i32, Infallible>();
let observable = observable.group_by(|value| value % 2);
let group_vec = Shared::new(Mutable::new(Vec::new()));
let group_vec_cloned = group_vec.clone();
let _subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|_termination| {},
);
assert!(sender.on_next(1).is_continue());
assert!(sender.on_next(2).is_continue());
assert!(sender.on_next(3).is_continue());
assert!(sender.on_next(4).is_continue());
assert_eq!(group_vec.with_ref(Vec::len), 2);
let mut groups = group_vec.take_value();
let group_even = groups.pop().unwrap();
let group_odd = groups.pop().unwrap();
let (checker_odd, observer_odd) = Checker::new();
let _sub_odd = group_odd.subscribe(observer_odd);
assert_eq!(checker_odd.values(), [1, 3]);
assert_eq!(checker_odd.state(), State::Active);
let (checker_even, observer_even) = Checker::new();
let _sub_even = group_even.subscribe(observer_even);
assert_eq!(checker_even.values(), [2, 4]);
assert_eq!(checker_even.state(), State::Active);
assert!(sender.on_next(5).is_continue());
assert!(sender.on_next(6).is_continue());
assert_eq!(checker_odd.values(), [1, 3, 5]);
assert_eq!(checker_even.values(), [2, 4, 6]);
assert_eq!(group_vec.with_ref(Vec::len), 0);
sender.on_termination(Termination::Completed);
assert_eq!(checker_odd.state(), State::Completed);
assert_eq!(checker_even.state(), State::Completed);
}
#[test]
fn test_subscribe_group_after_termination() {
let (mut sender, observable, _channel_checker) = test_channel::<'_, i32, Infallible>();
let (termination_checker, termination_observer) = Checker::<Infallible, _>::new();
let observable = observable.group_by(|value| value % 2);
let group_vec = Shared::new(Mutable::new(Vec::new()));
let group_vec_cloned = group_vec.clone();
let _subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|termination| termination_observer.on_termination(termination),
);
assert!(sender.on_next(111).is_continue());
assert_eq!(group_vec.with_ref(Vec::len), 1);
sender.on_termination(Termination::Completed);
assert_eq!(termination_checker.state(), State::Completed);
let group = group_vec.with_mut(Vec::pop).unwrap();
let (checker, observer) = Checker::new();
let _sub = group.subscribe(observer);
assert_eq!(checker.values(), [111]);
assert_eq!(checker.state(), State::Completed);
}
#[test]
fn test_subscribe_group_after_unsubscribe() {
let (mut sender, observable, _channel_checker) = test_channel::<'_, i32, Infallible>();
let observable = observable.group_by(|value| value % 2);
let group_vec = Shared::new(Mutable::new(Vec::new()));
let group_vec_cloned = group_vec.clone();
let subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|_termination| {},
);
assert!(sender.on_next(111).is_continue());
assert_eq!(group_vec.with_ref(Vec::len), 1);
drop(subscription);
let group = group_vec.with_mut(Vec::pop).unwrap();
let (checker, observer) = Checker::new();
let _sub = group.subscribe(observer);
assert_eq!(checker.values(), []);
assert_eq!(checker.state(), State::Dropped);
}
#[test]
fn test_values_of_ended_group_are_discarded() {
let (mut sender, observable, _channel_checker) = test_channel::<'_, i32, Infallible>();
let observable = observable.group_by(|value| value % 2);
let group_vec = Shared::new(Mutable::new(Vec::new()));
let group_vec_cloned = group_vec.clone();
let _subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|_termination| {},
);
assert!(sender.on_next(1).is_continue());
assert!(sender.on_next(2).is_continue());
let mut groups = group_vec.take_value();
let group_even = groups.pop().unwrap();
let group_odd = groups.pop().unwrap();
let (checker_odd, observer_odd) = Checker::new();
let sub_odd = group_odd.subscribe(observer_odd);
let (checker_even, observer_even) = Checker::new();
let _sub_even = group_even.subscribe(observer_even);
assert_eq!(checker_odd.values(), [1]);
assert_eq!(checker_even.values(), [2]);
drop(sub_odd);
assert_eq!(checker_odd.state(), State::Dropped);
assert!(sender.on_next(3).is_continue());
assert!(sender.on_next(4).is_continue());
assert_eq!(checker_odd.values(), [1]);
assert_eq!(checker_even.values(), [2, 4]);
assert_eq!(group_vec.with_ref(Vec::len), 0);
sender.on_termination(Termination::Completed);
assert_eq!(checker_even.state(), State::Completed);
}
#[test]
fn test_unsub_on_inner_termination_still_terminates_other_groups() {
let (mut sender, observable, _channel_checker) = test_channel::<'_, i32, Infallible>();
let (outer_termination_checker, outer_termination_observer) = Checker::<Infallible, _>::new();
let observable = observable.group_by(|value| value % 2);
let group_vec = Shared::new(Mutable::new(Vec::new()));
let group_vec_cloned = group_vec.clone();
let subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|termination| outer_termination_observer.on_termination(termination),
);
let outer_subscription = Shared::new(Mutable::new(Some(subscription)));
assert!(sender.on_next(1).is_continue());
assert!(sender.on_next(2).is_continue());
let mut groups = group_vec.take_value();
let group_even = groups.pop().unwrap();
let group_odd = groups.pop().unwrap();
let (checker_odd, observer_odd) = Checker::new();
let (on_next_odd, on_termination_odd) = observer_odd.into_callbacks();
let outer_subscription_cloned = outer_subscription.clone();
let _sub_odd = group_odd.subscribe_with_callback(on_next_odd, move |termination| {
drop(outer_subscription_cloned.take_value());
on_termination_odd(termination);
});
let (checker_even, observer_even) = Checker::new();
let (on_next_even, on_termination_even) = observer_even.into_callbacks();
let outer_subscription_cloned = outer_subscription.clone();
let _sub_even = group_even.subscribe_with_callback(on_next_even, move |termination| {
drop(outer_subscription_cloned.take_value());
on_termination_even(termination);
});
assert_eq!(checker_odd.values(), [1]);
assert_eq!(checker_even.values(), [2]);
sender.on_termination(Termination::Completed);
assert_eq!(checker_odd.values(), [1]);
assert_eq!(checker_odd.state(), State::Completed);
assert_eq!(checker_even.values(), [2]);
assert_eq!(checker_even.state(), State::Completed);
assert_eq!(outer_termination_checker.state(), State::Completed);
}
#[test]
fn test_dropping_unsubscribed_group_releases_buffered_values() {
use crate::tests_utils::drop_probe::{DropCount, DropProbe};
let (mut sender, observable, _channel_checker) = test_channel::<'_, DropProbe, Infallible>();
let observable = observable.group_by(|_value| 0);
let group_vec = Shared::new(Mutable::new(Vec::new()));
let group_vec_cloned = group_vec.clone();
let _outer_subscription = observable.subscribe_with_callback(
move |group| group_vec_cloned.with_mut(|values| values.push(group)),
|_termination| {},
);
let drops = DropCount::new();
assert!(sender.on_next(drops.probe()).is_continue());
assert_eq!(group_vec.with_ref(Vec::len), 1);
assert_eq!(drops.get(), 0);
let group = group_vec.with_mut(Vec::pop).unwrap();
drop(group);
assert_eq!(drops.get(), 1);
}
#[test]
fn test_dropping_ignored_value_does_not_poison_group_context() {
use crate::tests_utils::panic::{PanicOnDrop, expect_panic_on_drop};
let (mut sender, observable, channel_checker) =
test_channel::<'_, Option<PanicOnDrop>, Infallible>();
let observable = observable.group_by(|_value| 0);
let inner_subscriptions = Shared::new(Mutable::new(Vec::new()));
let inner_subscriptions_cloned = inner_subscriptions.clone();
let outer_subscription = observable.subscribe_with_callback(
move |group| {
let sub = group.subscribe_with_callback(|_value| {}, |_termination| {});
inner_subscriptions_cloned.with_mut(|values| values.push(sub));
},
|_termination| {},
);
assert!(sender.on_next(None).is_continue());
assert_eq!(inner_subscriptions.with_ref(Vec::len), 1);
let inner_subscription = inner_subscriptions.with_mut(Vec::pop).unwrap();
drop(inner_subscription);
expect_panic_on_drop(|value| sender.on_next(Some(value)));
drop(outer_subscription);
assert_eq!(channel_checker.state(), ChannelState::Unsubscribed);
}
#[test]
fn test_lifetime_sub() {
let life_marker = TestStruct;
let _subscription;
{
let observable = Create::new(|mut observer| {
assert!(observer.on_next(1).is_continue());
observer.on_termination(Termination::<String>::Completed);
Subscription::new(CallbackDisposal::new(|| {
life_marker.consume_ref();
}))
});
let observable = observable.group_by(|value| value.to_string());
let (_, observer) = Checker::new();
_subscription = observable.subscribe(observer);
}
}
#[test]
fn test_lifetime_or() {
let life_marker_2 = TestStruct;
let mut life_marker_1 = None;
{
let observable = Create::new(|observer| {
life_marker_1 = Some(observer);
Subscription::default()
});
let observable = observable.group_by(|v: &i32| *v).merge_all().map(|_| None);
let (_, mut observer) = Checker::<_, Infallible>::new();
assert!(observer.on_next(Some(&life_marker_2)).is_continue());
let _subscription = observable.subscribe(observer);
}
}
#[test]
fn test_fn() {
let mut s = TestStruct;
let subject: PublishSubject<'_, i32, &str> = PublishSubject::default();
let observable = subject.clone();
let observable = observable.group_by(|value| {
s.consume_mut();
value.to_string()
});
let _ = observable.subscribe_with_callback(|_| {}, |_| {});
}
#[test]
fn test_clone() {
let observable = Create::new(|mut observer| {
assert!(observer.on_next(TestStruct).is_continue());
observer.on_termination(Termination::Error(TestStruct));
Subscription::default()
});
let observable = observable.group_by(|_value| 0);
_ = observable.clone(); }
#[test]
fn test_type_inference_with_subscribe() {
let subject: PublishSubject<'_, i32, String> = PublishSubject::default();
let observable = subject.group_by(|value| value.to_string());
let observable = observable.filter(|_| true);
let (_, observer) = Checker::new();
let _ = observable.subscribe(observer);
}
#[test]
fn test_type_inference_without_subscribe() {
let subject: PublishSubject<'_, i32, String> = PublishSubject::default();
let observable = subject.group_by(|value| value.to_string());
observable.filter(|_| true);
}