use crate::utils::types::MaybeSend;
use crate::{
observable::{Observable, Subscription},
observer::{Flow, Observer, Termination},
scheduler::Scheduler,
};
use educe::Educe;
use futures::Stream;
#[derive(Educe)]
#[educe(Debug, Clone)]
pub struct FromTryStream<SM, S> {
stream: SM,
scheduler: S,
}
impl<SM, S> FromTryStream<SM, S> {
pub fn new(stream: SM, scheduler: S) -> Self {
Self { stream, scheduler }
}
}
impl<T, E, SM, S> Observable<'static, T, E> for FromTryStream<SM, S>
where
SM: Stream<Item = Result<T, E>> + MaybeSend + 'static,
S: Scheduler,
{
type D = S::D;
fn subscribe(
self,
observer: impl Observer<T, E> + MaybeSend + 'static,
) -> Subscription<Self::D> {
let mut observer = Some(observer);
self.scheduler
.schedule_stream(self.stream, move |result| match result {
Some(Ok(value)) => {
let flow = match observer.as_mut() {
Some(observer) => observer.on_next(value),
None => Flow::Stop,
};
if flow.is_stop() {
drop(observer.take());
}
flow.is_continue()
}
Some(Err(error)) => {
if let Some(observer) = observer.take() {
observer.on_termination(Termination::Error(error))
}
false
}
None => {
if let Some(observer) = observer.take() {
observer.on_termination(Termination::Completed)
}
false
}
})
}
}