use crate::prelude::*;
use crate::scheduler::Scheduler;
pub trait SubscribeOn {
fn subscribe_on<SD>(self, scheduler: SD) -> SubscribeOnOP<Self, SD>
where
Self: Sized,
{
SubscribeOnOP {
source: self,
scheduler,
}
}
}
pub struct SubscribeOnOP<S, SD> {
source: S,
scheduler: SD,
}
impl<S> SubscribeOn for S where S: RawSubscribable {}
impl<S, SD> RawSubscribable for SubscribeOnOP<S, SD>
where
S: RawSubscribable + Send + Sync + 'static,
SD: Scheduler,
{
type Item = S::Item;
type Err = S::Err;
fn raw_subscribe(
self,
subscribe: impl RxFn(
RxValue<&'_ Self::Item, &'_ Self::Err>,
) -> RxReturn<Self::Err>
+ Send
+ Sync
+ 'static,
) -> Box<dyn Subscription + Send + Sync> {
let source = self.source;
self
.scheduler
.schedule(move |_: Option<()>| source.raw_subscribe(subscribe), None)
}
}
#[test]
fn new_thread() {
use crate::ops::{Merge, SubscribeOn};
use crate::prelude::*;
use crate::scheduler::new_thread;
use std::sync::{Arc, Mutex};
use std::thread;
let res = Arc::new(Mutex::new(vec![]));
let c_res = res.clone();
let a = observable::from_range(1..5).subscribe_on(new_thread());
let b = observable::from_range(5..10);
let thread = Arc::new(Mutex::new(vec![]));
let c_thread = thread.clone();
a.merge(b).subscribe(move |v| {
res.lock().unwrap().push(*v);
let handle = thread::current();
thread.lock().unwrap().push(handle.id());
});
assert_eq!(*c_res.lock().unwrap(), (1..10).collect::<Vec<_>>());
let first = c_thread.lock().unwrap()[0];
let second = c_thread.lock().unwrap()[4];
assert_ne!(first, second);
let mut thread_list = vec![first; 4];
thread_list.append(&mut vec![second; 5]);
assert_eq!(*c_thread.lock().unwrap(), thread_list);
}