use std::sync::{Arc, atomic::AtomicBool, atomic::Ordering};
use timely::{Data, dataflow::{Scope, Stream}, progress::Timestamp};
use timely::dataflow::channels::pushers::{Counter as PushCounter, buffer::Buffer as PushBuffer};
use timely::dataflow::operators::generic::builder_raw::OperatorBuilder;
use timely::progress::frontier::MutableAntichain;
use timely::dataflow::operators::capture::event::{Event, EventIterator};
pub trait ReplayWithShutdown<T: Timestamp, D: Data> {
fn replay_with_shutdown_into<S: Scope<Timestamp=T>>(self, scope: &mut S, is_running: Arc<AtomicBool>) -> Stream<S, D>;
}
impl<T: Timestamp, D: Data, I> ReplayWithShutdown<T, D> for I
where I : IntoIterator,
<I as IntoIterator>::Item: EventIterator<T, D>+'static {
fn replay_with_shutdown_into<S: Scope<Timestamp=T>>(self, scope: &mut S, is_running: Arc<AtomicBool>) -> Stream<S, D> {
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 antichain = MutableAntichain::new();
builder.build(
move |_frontier| { },
move |_consumed, internal, produced| {
if !started {
internal[0].update(Default::default(), (event_streams.len() as i64) - 1);
antichain.update_iter(Some((Default::default(), (event_streams.len() as i64) - 1)).into_iter());
started = true;
}
if is_running.load(Ordering::Acquire) {
for event_stream in event_streams.iter_mut() {
while let Some(event) = event_stream.next() {
match *event {
Event::Progress(ref vec) => {
antichain.update_iter(vec.iter().cloned());
internal[0].extend(vec.iter().cloned());
},
Event::Messages(ref time, ref data) => {
output.session(time).give_iterator(data.iter().cloned());
}
}
}
}
activator.activate();
output.cease();
output.inner().produced().borrow_mut().drain_into(&mut produced[0]);
} else {
while !antichain.is_empty() {
let elements = antichain.frontier().iter().map(|t| (t.clone(), -1)).collect::<Vec<_>>();
for (t, c) in elements.iter() {
internal[0].update(t.clone(), *c);
}
antichain.update_iter(elements);
}
}
false
}
);
stream
}
}