use crate::dataflow::channels::{Bundle, BundleCore, Message};
use crate::progress::Timestamp;
use crate::dataflow::operators::Capability;
use crate::communication::Push;
use crate::{Container, Data};
#[derive(Debug)]
pub struct BufferCore<T, D: Container, P: Push<BundleCore<T, D>>> {
time: Option<T>,
buffer: D,
pusher: P,
}
pub type Buffer<T, D, P> = BufferCore<T, Vec<D>, P>;
impl<T, C: Container, P: Push<BundleCore<T, C>>> BufferCore<T, C, P> where T: Eq+Clone {
pub fn new(pusher: P) -> Self {
Self {
time: None,
buffer: Default::default(),
pusher,
}
}
pub fn session(&mut self, time: &T) -> Session<T, C, P> {
if let Some(true) = self.time.as_ref().map(|x| x != time) { self.flush(); }
self.time = Some(time.clone());
Session { buffer: self }
}
pub fn autoflush_session(&mut self, cap: Capability<T>) -> AutoflushSessionCore<T, C, P> where T: Timestamp {
if let Some(true) = self.time.as_ref().map(|x| x != cap.time()) { self.flush(); }
self.time = Some(cap.time().clone());
AutoflushSessionCore {
buffer: self,
_capability: cap,
}
}
pub fn inner(&mut self) -> &mut P { &mut self.pusher }
pub fn cease(&mut self) {
self.flush();
self.pusher.push(&mut None);
}
fn flush(&mut self) {
if !self.buffer.is_empty() {
let time = self.time.as_ref().unwrap().clone();
Message::push_at(&mut self.buffer, time, &mut self.pusher);
}
}
fn give_container(&mut self, vector: &mut C) {
if !vector.is_empty() {
self.flush();
let time = self.time.as_ref().expect("Buffer::give_container(): time is None.").clone();
Message::push_at(vector, time, &mut self.pusher);
}
}
}
impl<T, D: Data, P: Push<Bundle<T, D>>> Buffer<T, D, P> where T: Eq+Clone {
#[inline]
fn give(&mut self, data: D) {
if self.buffer.capacity() < crate::container::buffer::default_capacity::<D>() {
let to_reserve = crate::container::buffer::default_capacity::<D>() - self.buffer.capacity();
self.buffer.reserve(to_reserve);
}
self.buffer.push(data);
if self.buffer.len() == self.buffer.capacity() {
self.flush();
}
}
fn give_vec(&mut self, vector: &mut Vec<D>) {
self.flush();
let time = self.time.as_ref().expect("Buffer::give_vec(): time is None.").clone();
Message::push_at(vector, time, &mut self.pusher);
}
}
pub struct Session<'a, T, C: Container, P: Push<BundleCore<T, C>>+'a> where T: Eq+Clone+'a, C: 'a {
buffer: &'a mut BufferCore<T, C, P>,
}
impl<'a, T, C: Container, P: Push<BundleCore<T, C>>+'a> Session<'a, T, C, P> where T: Eq+Clone+'a, C: 'a {
pub fn give_container(&mut self, container: &mut C) {
self.buffer.give_container(container)
}
}
impl<'a, T, D: Data, P: Push<BundleCore<T, Vec<D>>>+'a> Session<'a, T, Vec<D>, P> where T: Eq+Clone+'a, D: 'a {
#[inline]
pub fn give(&mut self, data: D) {
self.buffer.give(data);
}
#[inline]
pub fn give_iterator<I: Iterator<Item=D>>(&mut self, iter: I) {
for item in iter {
self.give(item);
}
}
#[inline]
pub fn give_vec(&mut self, message: &mut Vec<D>) {
if !message.is_empty() {
self.buffer.give_vec(message);
}
}
}
pub struct AutoflushSessionCore<'a, T: Timestamp, C: Container, P: Push<BundleCore<T, C>>+'a> where
T: Eq+Clone+'a, C: 'a {
buffer: &'a mut BufferCore<T, C, P>,
_capability: Capability<T>,
}
pub type AutoflushSession<'a, T, D, P> = AutoflushSessionCore<'a, T, Vec<D>, P>;
impl<'a, T: Timestamp, D: Data, P: Push<BundleCore<T, Vec<D>>>+'a> AutoflushSessionCore<'a, T, Vec<D>, P> where T: Eq+Clone+'a, D: 'a {
#[inline]
pub fn give(&mut self, data: D) {
self.buffer.give(data);
}
#[inline]
pub fn give_iterator<I: Iterator<Item=D>>(&mut self, iter: I) {
for item in iter {
self.give(item);
}
}
#[inline]
pub fn give_content(&mut self, message: &mut Vec<D>) {
if !message.is_empty() {
self.buffer.give_vec(message);
}
}
}
impl<'a, T: Timestamp, C: Container, P: Push<BundleCore<T, C>>+'a> Drop for AutoflushSessionCore<'a, T, C, P> where T: Eq+Clone+'a, C: 'a {
fn drop(&mut self) {
self.buffer.cease();
}
}