use timely::dataflow::Scope;
use timely::progress::Timestamp;
use timely::dataflow::operators::vec::{Filter, Map};
use differential_dataflow::{AsCollection, VecCollection, Data};
use differential_dataflow::difference::Abelian;
use crate::altneu::AltNeu;
pub trait Differentiate<'scope, T: Timestamp, D: Data, R: Abelian> {
fn differentiate<'inner>(self, child: Scope<'inner, AltNeu<T>>) -> VecCollection<'inner, AltNeu<T>, D, R>;
}
pub trait Integrate<'scope, T: Timestamp, D: Data, R: Abelian> {
fn integrate<'outer>(self, outer: Scope<'outer, T>) -> VecCollection<'outer, T, D, R>;
}
impl<'scope, T, D, R> Differentiate<'scope, T, D, R> for VecCollection<'scope, T, D, R>
where
T: Timestamp,
D: Data,
R: Abelian + 'static,
{
fn differentiate<'inner>(self, child: Scope<'inner, AltNeu<T>>) -> VecCollection<'inner, AltNeu<T>, D, R> {
self.enter(child)
.inner
.flat_map(|(data, time, diff)| {
let mut neg_diff = diff.clone();
neg_diff.negate();
let neu = (data.clone(), AltNeu::neu(time.time.clone()), neg_diff);
let alt = (data, time, diff);
Some(alt).into_iter().chain(Some(neu))
})
.as_collection()
}
}
impl<'scope, T, D, R> Integrate<'scope, T, D, R> for VecCollection<'scope, AltNeu<T>, D, R>
where
T: Timestamp,
D: Data,
R: Abelian + 'static,
{
fn integrate<'outer>(self, outer: Scope<'outer, T>) -> VecCollection<'outer, T, D, R> {
self.inner
.filter(|(_d,t,_r)| !t.neu)
.as_collection()
.leave(outer)
}
}