use crate::dataflow::{Scope, StreamCore};
use crate::dataflow::channels::pushers::CounterCore as PushCounter;
use crate::dataflow::channels::pushers::buffer::BufferCore as PushBuffer;
use crate::dataflow::operators::generic::builder_raw::OperatorBuilder;
use crate::progress::Timestamp;
use super::EventCore;
use super::event::EventIteratorCore;
use crate::Container;
pub trait Replay<T: Timestamp, C> : Sized {
fn replay_into<S: Scope<Timestamp=T>>(self, scope: &mut S) -> StreamCore<S, C> {
self.replay_core(scope, Some(std::time::Duration::new(0, 0)))
}
fn replay_core<S: Scope<Timestamp=T>>(self, scope: &mut S, period: Option<std::time::Duration>) -> StreamCore<S, C>;
}
impl<T: Timestamp, C: Container, I> Replay<T, C> for I
where I : IntoIterator,
<I as IntoIterator>::Item: EventIteratorCore<T, C>+'static {
fn replay_core<S: Scope<Timestamp=T>>(self, scope: &mut S, period: Option<std::time::Duration>) -> StreamCore<S, C>{
let mut builder = OperatorBuilder::new("Replay".to_owned(), scope.clone());
let address = builder.operator_info().address;
let activator = scope.activator_for(&address[..]);
let (targets, stream) = builder.new_output();
let mut output = PushBuffer::new(PushCounter::new(targets));
let mut event_streams = self.into_iter().collect::<Vec<_>>();
let mut started = false;
let mut allocation: C = Default::default();
builder.build(
move |progress| {
if !started {
progress.internals[0].update(S::Timestamp::minimum(), (event_streams.len() as i64) - 1);
started = true;
}
for event_stream in event_streams.iter_mut() {
while let Some(event) = event_stream.next() {
match event {
EventCore::Progress(vec) => {
progress.internals[0].extend(vec.iter().cloned());
},
EventCore::Messages(ref time, data) => {
allocation.clone_from(data);
output.session(time).give_container(&mut allocation);
}
}
}
}
if let Some(delay) = period {
activator.activate_after(delay);
}
output.cease();
output.inner().produced().borrow_mut().drain_into(&mut progress.produceds[0]);
false
}
);
stream
}
}