use crate::core::scheduler::ffmpeg_scheduler::{FfmpegScheduler, Running};
use crate::error::Error;
use crossbeam_channel::Receiver;
pub(crate) struct OwnedRunIter<T> {
rx: Option<Receiver<T>>,
scheduler: Option<FfmpegScheduler<Running>>,
terminated: bool,
map_terminal_error: fn(Error) -> Error,
}
impl<T> OwnedRunIter<T> {
pub(crate) fn new(
rx: Receiver<T>,
scheduler: FfmpegScheduler<Running>,
map_terminal_error: fn(Error) -> Error,
) -> Self {
Self {
rx: Some(rx),
scheduler: Some(scheduler),
terminated: false,
map_terminal_error,
}
}
fn finish(&mut self) -> Option<Result<T, Error>> {
self.terminated = true;
self.rx = None;
match self.scheduler.take() {
Some(scheduler) => match scheduler.wait() {
Ok(()) => None,
Err(e) => Some(Err((self.map_terminal_error)(e))),
},
None => None,
}
}
}
impl<T> Iterator for OwnedRunIter<T> {
type Item = Result<T, Error>;
fn next(&mut self) -> Option<Self::Item> {
if self.terminated {
return None;
}
let recv = match self.rx.as_ref() {
Some(rx) => rx.recv(),
None => return self.finish(),
};
match recv {
Ok(item) => Some(Ok(item)),
Err(_) => self.finish(),
}
}
}
impl<T> std::iter::FusedIterator for OwnedRunIter<T> {}
impl<T> Drop for OwnedRunIter<T> {
fn drop(&mut self) {
self.rx = None;
if let Some(scheduler) = self.scheduler.take() {
scheduler.abort();
}
}
}