mod tests_utils;
use crate::tests_utils::drop_probe::{DropCallback, DropProbe};
use rx_rust::utils::mutable::MutableExt;
use rx_rust::utils::mutable::MutableHelper;
use rx_rust::{
observer::{BoxedObserverExt, Flow, Observer, Termination, boxed_observer::BoxedObserver},
utils::{
mutable::Mutable,
pending_events::EventBatch,
serialized_delivery::{
DeliveryStopped, SerializedDelivery, UpdateOutcome, WeakSerializedDelivery,
},
types::{MaybeSend, Shared},
},
};
type TestError = &'static str;
const ERROR: TestError = "boom";
type TestObserver = BoxedObserver<'static, Value, TestError>;
type TestDelivery = SerializedDelivery<Value, TestError, TestObserver, Resources>;
type WeakTestDelivery = WeakSerializedDelivery<Value, TestError, TestObserver, Resources>;
type TestOutcome<R> = UpdateOutcome<Value, TestError, R>;
#[derive(Debug, Clone, PartialEq, Eq)]
enum Record {
Next(i32),
Termination(Termination<TestError>),
ObserverDropped,
ResourcesDropped,
Note(&'static str),
}
type Log = Shared<Mutable<Vec<Record>>>;
fn new_log() -> Log {
Shared::new(Mutable::new(Vec::new()))
}
fn record(log: &Log, record: Record) {
log.with_mut(|values| values.push(record));
}
fn records(log: &Log) -> Vec<Record> {
log.clone_value()
}
fn values(log: &Log) -> Vec<i32> {
records(log)
.into_iter()
.filter_map(|record| match record {
Record::Next(value) => Some(value),
_ => None,
})
.collect()
}
fn record_on_drop(log: &Log, record_value: Record) -> DropCallback {
let log = log.clone();
Box::new(move || record(&log, record_value))
}
fn note(log: &Log, text: &'static str) -> DropCallback {
record_on_drop(log, Record::Note(text))
}
fn note_and_reenter(handle: &DeliveryHandle, log: &Log, text: &'static str) -> DropCallback {
let handle = handle.clone();
let log = log.clone();
Box::new(move || {
record(&log, Record::Note(text));
let _ = handle.get().send(next(99));
})
}
struct Value {
number: i32,
probe: DropProbe,
}
impl Value {
fn new(number: i32) -> Self {
Self {
number,
probe: DropProbe::new(),
}
}
fn on_drop(mut self, callback: DropCallback) -> Self {
self.probe = self.probe.on_drop(callback);
self
}
}
struct Resources {
model: i32,
probe: DropProbe,
}
impl Resources {
fn new(model: i32, log: &Log) -> Self {
Self {
model,
probe: DropProbe::new().on_drop(record_on_drop(log, Record::ResourcesDropped)),
}
}
fn on_drop(&mut self, callback: DropCallback) {
self.probe.also_on_drop(callback);
}
}
#[derive(Clone)]
struct DeliveryHandle(Shared<Mutable<Option<WeakTestDelivery>>>);
impl DeliveryHandle {
fn new() -> Self {
Self(Shared::new(Mutable::new(None)))
}
fn of(delivery: &TestDelivery) -> Self {
let handle = Self::new();
handle.install(delivery);
handle
}
fn install(&self, delivery: &TestDelivery) {
self.0.replace_value(Some(delivery.downgrade()));
}
fn get(&self) -> TestDelivery {
self.0
.clone_value()
.expect("the delivery is installed before anything can reach it")
.upgrade()
.expect("the tests keep a handle while their callbacks run")
}
}
struct RecordingObserver<FN, FT> {
_probe: DropProbe,
log: Log,
handle: DeliveryHandle,
on_next: FN,
on_termination: Option<FT>,
stop_after: Option<usize>,
received: usize,
}
impl<FN, FT> Observer<Value, TestError> for RecordingObserver<FN, FT>
where
FN: FnMut(&TestDelivery, i32),
FT: FnOnce(&TestDelivery, &Termination<TestError>),
{
fn on_next(&mut self, value: Value) -> Flow {
let number = value.number;
record(&self.log, Record::Next(number));
let delivery = self.handle.get();
(self.on_next)(&delivery, number);
self.received += 1;
match self.stop_after {
Some(stop_after) if self.received >= stop_after => Flow::Stop,
_ => Flow::Continue,
}
}
fn on_termination(mut self, termination: Termination<TestError>) {
record(&self.log, Record::Termination(termination.clone()));
let delivery = self.handle.get();
let hook = self
.on_termination
.take()
.expect("an observer is terminated at most once");
hook(&delivery, &termination);
}
}
type NoNextHook = fn(&TestDelivery, i32);
type NoTerminationHook = fn(&TestDelivery, &Termination<TestError>);
fn builder(log: &Log) -> Builder<NoNextHook, NoTerminationHook> {
Builder {
log: log.clone(),
handle: DeliveryHandle::new(),
model: 0,
on_next: |_delivery, _number| {},
on_termination: |_delivery, _termination| {},
on_observer_drop: None,
stop_after: None,
}
}
struct Builder<FN, FT> {
log: Log,
handle: DeliveryHandle,
model: i32,
on_next: FN,
on_termination: FT,
on_observer_drop: Option<DropCallback>,
stop_after: Option<usize>,
}
impl<FN, FT> Builder<FN, FT> {
fn handle(&self) -> DeliveryHandle {
self.handle.clone()
}
fn model(mut self, model: i32) -> Self {
self.model = model;
self
}
fn on_next<F>(self, on_next: F) -> Builder<F, FT>
where
F: FnMut(&TestDelivery, i32),
{
Builder {
log: self.log,
handle: self.handle,
model: self.model,
on_next,
on_termination: self.on_termination,
on_observer_drop: self.on_observer_drop,
stop_after: self.stop_after,
}
}
fn stop_after(mut self, count: usize) -> Self {
self.stop_after = Some(count);
self
}
fn on_termination<F>(self, on_termination: F) -> Builder<FN, F>
where
F: FnOnce(&TestDelivery, &Termination<TestError>),
{
Builder {
log: self.log,
handle: self.handle,
model: self.model,
on_next: self.on_next,
on_termination,
on_observer_drop: self.on_observer_drop,
stop_after: self.stop_after,
}
}
fn on_observer_drop(mut self, callback: DropCallback) -> Self {
self.on_observer_drop = Some(callback);
self
}
fn build(self) -> TestDelivery
where
FN: FnMut(&TestDelivery, i32) + MaybeSend + 'static,
FT: FnOnce(&TestDelivery, &Termination<TestError>) + MaybeSend + 'static,
{
let handle = self.handle.clone();
let mut probe =
DropProbe::new().on_drop(record_on_drop(&self.log, Record::ObserverDropped));
if let Some(callback) = self.on_observer_drop {
probe.also_on_drop(callback);
}
let observer: TestObserver = RecordingObserver {
_probe: probe,
log: self.log.clone(),
handle: self.handle,
on_next: self.on_next,
on_termination: Some(self.on_termination),
stop_after: self.stop_after,
received: 0,
}
.into_boxed();
let resources = Resources::new(self.model, &self.log);
let delivery = SerializedDelivery::idle(observer, resources);
handle.install(&delivery);
delivery
}
}
fn next(number: i32) -> EventBatch<Value, TestError> {
EventBatch::Next(Value::new(number))
}
fn next_batch(numbers: impl IntoIterator<Item = i32>) -> EventBatch<Value, TestError> {
EventBatch::NextBatch(numbers.into_iter().map(Value::new).collect())
}
fn completed() -> EventBatch<Value, TestError> {
EventBatch::Termination(Termination::Completed)
}
fn errored() -> EventBatch<Value, TestError> {
EventBatch::Termination(Termination::Error(ERROR))
}
fn next_batch_and_completed(
numbers: impl IntoIterator<Item = i32>,
) -> EventBatch<Value, TestError> {
EventBatch::NextBatchAndTermination(
numbers.into_iter().map(Value::new).collect(),
Termination::Completed,
)
}
#[test]
fn send_next_delivers_to_the_parked_observer() {
let log = new_log();
let delivery = builder(&log).build();
assert!(delivery.send(next(1)).is_continue());
assert_eq!(records(&log), [Record::Next(1)]);
assert!(delivery.send(next(2)).is_continue());
assert_eq!(values(&log), [1, 2]);
}
#[test]
fn send_delivers_a_batch_in_order() {
let log = new_log();
let delivery = builder(&log).build();
assert!(delivery.send(next_batch([1, 2, 3])).is_continue());
assert_eq!(values(&log), [1, 2, 3]);
}
#[test]
fn a_termination_is_notified_before_the_resources_are_dropped() {
let log = new_log();
let delivery = builder(&log).build();
assert!(delivery.send(completed()).is_stop());
assert_eq!(
records(&log),
[
Record::Termination(Termination::Completed),
Record::ObserverDropped,
Record::ResourcesDropped,
]
);
}
#[test]
fn a_batch_delivers_its_values_before_its_error_termination() {
let log = new_log();
let delivery = builder(&log).build();
assert!(
delivery
.send(EventBatch::NextBatchAndTermination(
vec![Value::new(1), Value::new(2)],
Termination::Error(ERROR),
))
.is_stop()
);
assert_eq!(
records(&log),
[
Record::Next(1),
Record::Next(2),
Record::Termination(Termination::Error(ERROR)),
Record::ObserverDropped,
Record::ResourcesDropped,
]
);
}
#[test]
fn a_next_and_termination_batch_delivers_its_value_first() {
let log = new_log();
let delivery = builder(&log).build();
assert!(
delivery
.send(EventBatch::NextAndTermination(
Value::new(1),
Termination::Completed,
))
.is_stop()
);
assert_eq!(
records(&log),
[
Record::Next(1),
Record::Termination(Termination::Completed),
Record::ObserverDropped,
Record::ResourcesDropped,
]
);
}
#[test]
fn an_empty_batch_is_accepted_and_leaves_the_delivery_idle() {
let log = new_log();
let delivery = builder(&log).build();
assert!(delivery.send(EventBatch::NextBatch(vec![])).is_continue());
assert_eq!(records(&log), []);
assert!(delivery.send(next(1)).is_continue());
assert_eq!(values(&log), [1]);
}
#[test]
fn an_empty_batch_sent_from_on_next_is_accepted_as_well() {
let log = new_log();
let delivery = builder(&log)
.on_next(|delivery, number| {
if number == 1 {
assert!(delivery.send(EventBatch::NextBatch(vec![])).is_continue());
}
})
.build();
assert!(delivery.send(next_batch([1, 2])).is_continue());
assert_eq!(values(&log), [1, 2]);
}
#[test]
fn events_sent_after_a_termination_are_rejected() {
let log = new_log();
let delivery = builder(&log).build();
assert!(delivery.send(completed()).is_stop());
assert!(delivery.send(next(1)).is_stop());
assert!(delivery.send(errored()).is_stop());
assert!(delivery.send(EventBatch::NextBatch(vec![])).is_stop());
assert_eq!(
records(&log),
[
Record::Termination(Termination::Completed),
Record::ObserverDropped,
Record::ResourcesDropped,
]
);
}
#[test]
fn events_sent_after_stop_are_rejected() {
let log = new_log();
let delivery = builder(&log).build();
delivery.stop();
assert!(delivery.send(next(1)).is_stop());
assert!(delivery.send(completed()).is_stop());
assert_eq!(
records(&log),
[Record::ObserverDropped, Record::ResourcesDropped]
);
}
#[test]
fn rejected_events_are_dropped_outside_the_lock() {
let log = new_log();
let delivery = builder(&log).build();
assert!(delivery.send(completed()).is_stop());
let rejected = Value::new(1).on_drop(note_and_reenter(
&DeliveryHandle::of(&delivery),
&log,
"rejected dropped",
));
assert!(delivery.send(EventBatch::Next(rejected)).is_stop());
assert_eq!(
records(&log)[3..],
[Record::Note("rejected dropped")],
"the rejected value must be dropped, outside the lock"
);
}
#[test]
fn events_sent_from_on_termination_are_rejected() {
let log = new_log();
let delivery = builder(&log)
.on_termination(|delivery, _termination| {
assert!(delivery.send(next(1)).is_stop());
assert_eq!(
delivery.update(|_resources| UpdateOutcome::empty()),
Err(DeliveryStopped)
);
})
.build();
assert!(delivery.send(completed()).is_stop());
assert_eq!(values(&log), []);
}
#[test]
fn stop_drops_the_observer_without_notifying_it_and_is_idempotent() {
let log = new_log();
let delivery = builder(&log).build();
delivery.stop();
assert_eq!(
records(&log),
[Record::ObserverDropped, Record::ResourcesDropped]
);
delivery.stop();
assert_eq!(
records(&log),
[Record::ObserverDropped, Record::ResourcesDropped],
"stopping again must drop nothing more"
);
}
#[test]
fn stop_drops_the_observer_and_the_resources_outside_the_lock() {
let log = new_log();
let builder = builder(&log);
let handle = builder.handle();
let delivery = builder
.on_observer_drop(note_and_reenter(&handle, &log, "observer dropped"))
.build();
delivery
.update(|resources| {
resources.on_drop(note_and_reenter(&handle, &log, "resources dropped"));
UpdateOutcome::empty()
})
.expect("the delivery is idle");
delivery.stop();
assert_eq!(
records(&log),
[
Record::ObserverDropped,
Record::Note("observer dropped"),
Record::ResourcesDropped,
Record::Note("resources dropped"),
]
);
}
#[test]
fn stop_from_on_next_drops_the_values_that_are_still_queued() {
let log = new_log();
let delivery = builder(&log)
.on_next(|delivery, number| {
if number == 2 {
delivery.stop();
}
})
.build();
assert!(
delivery
.send(EventBatch::NextBatch(vec![
Value::new(1),
Value::new(2),
Value::new(3).on_drop(note(&log, "3 dropped")),
]))
.is_stop()
);
assert_eq!(values(&log), [1, 2]);
assert_eq!(
records(&log)[2..],
[
Record::Note("3 dropped"),
Record::ResourcesDropped,
Record::ObserverDropped,
]
);
assert!(delivery.send(next(4)).is_stop());
}
#[test]
fn an_observer_that_stops_drops_the_values_that_are_still_queued() {
let log = new_log();
let delivery = builder(&log).stop_after(2).build();
assert!(
delivery
.send(EventBatch::NextBatch(vec![
Value::new(1),
Value::new(2),
Value::new(3).on_drop(note(&log, "3 dropped")),
]))
.is_stop()
);
assert_eq!(values(&log), [1, 2]);
assert_eq!(
records(&log)[2..],
[
Record::Note("3 dropped"),
Record::ResourcesDropped,
Record::ObserverDropped,
]
);
assert!(delivery.send(next(4)).is_stop());
}
#[test]
fn an_observer_that_stops_is_dropped_instead_of_being_terminated() {
let log = new_log();
let delivery = builder(&log).stop_after(1).build();
assert!(delivery.send(next_batch_and_completed([1, 2])).is_stop());
assert_eq!(
records(&log),
[
Record::Next(1),
Record::ResourcesDropped,
Record::ObserverDropped,
],
"an observer that ended its own stream must not be terminated on top of that"
);
}
#[test]
fn stop_from_on_next_suppresses_the_queued_termination() {
let log = new_log();
let delivery = builder(&log)
.on_next(|delivery, number| {
if number == 1 {
delivery.stop();
}
})
.build();
assert!(delivery.send(next_batch_and_completed([1, 2])).is_stop());
assert_eq!(
records(&log),
[
Record::Next(1),
Record::ResourcesDropped,
Record::ObserverDropped,
],
"the observer must be dropped instead of being terminated"
);
}
#[test]
fn a_value_sent_from_on_next_is_delivered_after_the_queued_ones() {
let log = new_log();
let log_of_observer = log.clone();
let delivery = builder(&log)
.on_next(move |delivery, number| {
if number == 1 {
assert!(delivery.send(next(10)).is_continue());
assert_eq!(values(&log_of_observer), [1]);
}
})
.build();
assert!(delivery.send(next_batch([1, 2, 3])).is_continue());
assert_eq!(values(&log), [1, 2, 3, 10]);
}
#[test]
fn a_termination_sent_from_on_next_is_delivered_after_the_queued_values() {
let log = new_log();
let delivery = builder(&log)
.on_next(|delivery, number| {
if number == 1 {
assert!(delivery.send(completed()).is_stop());
assert!(delivery.send(next(10)).is_stop());
}
})
.build();
assert!(delivery.send(next_batch([1, 2, 3])).is_stop());
assert_eq!(
records(&log),
[
Record::Next(1),
Record::Next(2),
Record::Next(3),
Record::Termination(Termination::Completed),
Record::ObserverDropped,
Record::ResourcesDropped,
]
);
}
#[test]
fn update_without_events_updates_the_model_and_returns_its_result() {
let log = new_log();
let delivery = builder(&log).model(1).build();
assert_eq!(
delivery.update(|resources| {
resources.model += 10;
UpdateOutcome::new(resources.model)
}),
Ok(11)
);
assert_eq!(
delivery.update(|resources| UpdateOutcome::new(resources.model)),
Ok(11)
);
assert_eq!(records(&log), [], "an update alone delivers nothing");
}
#[test]
fn update_runs_while_a_delivery_is_running() {
let log = new_log();
let delivery = builder(&log)
.on_next(|delivery, number| {
assert_eq!(
delivery.update(|resources| {
resources.model += number;
UpdateOutcome::new(resources.model)
}),
Ok(number * (number + 1) / 2)
);
})
.build();
assert!(delivery.send(next_batch([1, 2, 3])).is_continue());
assert_eq!(
delivery.update(|resources| UpdateOutcome::new(resources.model)),
Ok(6)
);
}
#[test]
fn update_after_a_termination_returns_delivery_stopped() {
let log = new_log();
let delivery = builder(&log).build();
assert!(delivery.send(completed()).is_stop());
assert_eq!(
delivery.update(|resources| UpdateOutcome::new(resources.model)),
Err(DeliveryStopped)
);
}
#[test]
fn a_stopped_update_does_not_run_and_is_dropped_outside_the_lock() {
let log = new_log();
let delivery = builder(&log).build();
delivery.stop();
let probe = Value::new(0).on_drop(note_and_reenter(
&DeliveryHandle::of(&delivery),
&log,
"update dropped",
));
let result: Result<(), _> = delivery.update(move |_resources| -> TestOutcome<()> {
drop(probe); panic!("the update must not run once the delivery has stopped");
});
assert_eq!(result, Err(DeliveryStopped));
assert_eq!(
records(&log)[2..],
[Record::Note("update dropped")],
"the update must be dropped, outside the lock"
);
}
#[test]
fn update_delivers_the_events_of_its_outcome() {
let log = new_log();
let delivery = builder(&log).build();
assert_eq!(
delivery.update(|resources| {
resources.model = 1;
UpdateOutcome::new("first").with_next_event(Value::new(1))
}),
Ok("first")
);
assert_eq!(
delivery.update(|_resources| {
UpdateOutcome::empty()
.with_events(EventBatch::NextBatch(vec![Value::new(2), Value::new(3)]))
}),
Ok(())
);
assert_eq!(values(&log), [1, 2, 3]);
}
#[test]
fn update_terminates_before_dropping_what_it_asked_to_drop_outside() {
let log = new_log();
let delivery = builder(&log).build();
let result = delivery.update(|resources| {
resources.model = 1;
UpdateOutcome::empty()
.with_drop_outside(Value::new(0).on_drop(note(&log, "dropped outside")))
.with_next_and_termination_events(Value::new(1), Termination::Error(ERROR))
});
assert_eq!(result, Ok(()));
assert_eq!(
records(&log),
[
Record::Next(1),
Record::Termination(Termination::Error(ERROR)),
Record::ObserverDropped,
Record::ResourcesDropped,
Record::Note("dropped outside"),
]
);
}
#[test]
fn what_an_update_drops_outside_is_dropped_with_the_lock_released() {
let log = new_log();
let delivery = builder(&log).build();
let probe = Value::new(0).on_drop(note_and_reenter(
&DeliveryHandle::of(&delivery),
&log,
"dropped outside",
));
let result = delivery.update(move |_resources| {
UpdateOutcome::empty()
.with_drop_outside(probe)
.with_next_event(Value::new(1))
});
assert_eq!(result, Ok(()));
assert_eq!(
records(&log),
[
Record::Next(1),
Record::Note("dropped outside"),
Record::Next(99),
]
);
}
#[test]
fn update_with_a_termination_event_stops_the_delivery() {
let log = new_log();
let delivery = builder(&log).build();
assert_eq!(
delivery.update(|_resources| {
UpdateOutcome::empty().with_termination_event(Termination::Completed)
}),
Ok(())
);
assert_eq!(
delivery.update(|_resources| UpdateOutcome::empty()),
Err(DeliveryStopped)
);
assert!(delivery.send(next(1)).is_stop());
}
#[test]
fn update_from_on_next_queues_its_events_after_the_pending_ones() {
let log = new_log();
let delivery = builder(&log)
.on_next(|delivery, number| {
if number == 1 {
assert_eq!(
delivery.update(|resources| {
resources.model += 1;
UpdateOutcome::new(resources.model).with_next_event(Value::new(10))
}),
Ok(1)
);
}
})
.build();
assert!(delivery.send(next_batch([1, 2])).is_continue());
assert_eq!(values(&log), [1, 2, 10]);
}
#[test]
fn update_runs_and_returns_even_when_its_events_are_rejected() {
let log = new_log();
let log_of_observer = log.clone();
let delivery = builder(&log)
.on_next(move |delivery, number| {
if number != 1 {
return;
}
assert!(delivery.send(completed()).is_stop());
let dropped = Value::new(10).on_drop(note(&log_of_observer, "10 dropped"));
assert_eq!(
delivery.update(move |resources| {
resources.model += 1;
UpdateOutcome::new(resources.model).with_next_event(dropped)
}),
Ok(1)
);
})
.build();
assert!(delivery.send(next_batch([1, 2])).is_stop());
assert_eq!(values(&log), [1, 2]);
assert!(records(&log).contains(&Record::Note("10 dropped")));
}
#[test]
fn the_arms_of_one_update_share_the_type_of_their_outcome() {
let log = new_log();
let delivery = builder(&log).build();
let update = || {
delivery.update(|resources| {
resources.model += 1;
if resources.model == 1 {
UpdateOutcome::new(resources.model)
.with_drop_outside(Value::new(0).on_drop(note(&log, "dropped outside")))
.with_next_event(Value::new(1))
} else {
UpdateOutcome::new(resources.model)
.without_drop_outside()
.without_events()
}
})
};
assert_eq!(update(), Ok(1));
assert_eq!(update(), Ok(2));
assert_eq!(
records(&log),
[Record::Next(1), Record::Note("dropped outside")]
);
}
#[test]
fn a_clone_shares_one_delivery() {
let log = new_log();
let delivery = builder(&log).build();
let clone = delivery.clone();
assert!(clone.send(next(1)).is_continue());
delivery.stop();
assert!(clone.send(next(2)).is_stop());
assert_eq!(values(&log), [1]);
}
#[test]
fn a_weak_handle_upgrades_while_a_strong_one_is_alive() {
let log = new_log();
let delivery = builder(&log).build();
let weak = delivery.downgrade();
assert!(weak.upgrade().is_some());
delivery.stop();
let upgraded = weak.upgrade().expect("a stopped delivery is still there");
assert!(upgraded.send(next(1)).is_stop());
drop(upgraded);
drop(delivery);
assert!(weak.upgrade().is_none());
}
#[test]
fn dropping_the_last_handle_drops_the_observer_without_notifying_it() {
let log = new_log();
let delivery = builder(&log).build();
let clone = delivery.clone();
drop(delivery);
assert_eq!(records(&log), [], "another handle is still alive");
drop(clone);
assert_eq!(
records(&log),
[Record::ObserverDropped, Record::ResourcesDropped]
);
}
#[cfg(panic = "unwind")]
#[test]
fn a_panic_from_on_next_stops_the_delivery() {
use crate::tests_utils::panic::expect_panic_on_drop;
let log = new_log();
let token = Shared::new(Mutable::new(None));
let token_of_observer = token.clone();
let delivery = builder(&log)
.on_next(move |_delivery, number| {
if number == 1 {
drop(token_of_observer.take_value());
}
})
.build();
expect_panic_on_drop(|panic_on_drop| {
token.replace_value(Some(panic_on_drop));
let _ = delivery.send(next_batch([1, 2, 3]));
});
assert_eq!(
records(&log),
[
Record::Next(1),
Record::ResourcesDropped,
Record::ObserverDropped,
]
);
assert!(delivery.send(next(4)).is_stop());
}
#[cfg(panic = "unwind")]
#[test]
fn a_panic_from_on_termination_drops_the_resources() {
use crate::tests_utils::panic::expect_panic_on_drop;
let log = new_log();
let token = Shared::new(Mutable::new(None));
let token_of_observer = token.clone();
let delivery = builder(&log)
.on_termination(move |_delivery, _termination| {
drop(token_of_observer.take_value());
})
.build();
expect_panic_on_drop(|panic_on_drop| {
token.replace_value(Some(panic_on_drop));
let _ = delivery.send(completed());
});
assert_eq!(
records(&log),
[
Record::Termination(Termination::Completed),
Record::ObserverDropped,
Record::ResourcesDropped,
]
);
assert!(delivery.send(next(1)).is_stop());
}
#[cfg(not(feature = "single-threaded"))]
#[test]
fn concurrent_sends_are_delivered_one_at_a_time() {
use rx_rust::utils::mutable::{MutableBool, MutableBoolHelper};
const THREADS: i32 = 4;
const VALUES_PER_THREAD: i32 = 25;
let log = new_log();
let delivering = Shared::new(MutableBool::new(false));
let delivering_of_observer = delivering.clone();
let delivery = builder(&log)
.on_next(move |_delivery, _value| {
assert!(
delivering_of_observer.change_if_not_equal(true),
"two values must never be delivered at the same time"
);
std::thread::yield_now();
assert!(delivering_of_observer.change_if_not_equal(false));
})
.build();
std::thread::scope(|scope| {
for thread in 0..THREADS {
let delivery = delivery.clone();
scope.spawn(move || {
for value in 0..VALUES_PER_THREAD {
assert!(
delivery
.send(next(thread * VALUES_PER_THREAD + value))
.is_continue()
);
}
});
}
});
let mut delivered = values(&log);
delivered.sort_unstable();
assert_eq!(
delivered,
(0..THREADS * VALUES_PER_THREAD).collect::<Vec<_>>()
);
}