use crate::communication::Push;
use crate::Container;
pub mod pushers;
pub mod pullers;
pub mod pact;
pub type BundleCore<T, D> = crate::communication::Message<Message<T, D>>;
pub type Bundle<T, D> = BundleCore<T, Vec<D>>;
#[derive(Clone, Abomonation, Serialize, Deserialize)]
pub struct Message<T, D> {
pub time: T,
pub data: D,
pub from: usize,
pub seq: usize,
}
impl<T, D> Message<T, D> {
#[deprecated = "Use timely::buffer::default_capacity instead"]
pub fn default_length() -> usize {
crate::container::buffer::default_capacity::<D>()
}
}
impl<T, D: Container> Message<T, D> {
pub fn new(time: T, data: D, from: usize, seq: usize) -> Self {
Message { time, data, from, seq }
}
#[inline]
pub fn push_at<P: Push<BundleCore<T, D>>>(buffer: &mut D, time: T, pusher: &mut P) {
let data = ::std::mem::take(buffer);
let message = Message::new(time, data, 0, 0);
let mut bundle = Some(BundleCore::from_typed(message));
pusher.push(&mut bundle);
if let Some(message) = bundle {
if let Some(message) = message.if_typed() {
*buffer = message.data;
buffer.clear();
}
}
}
}