use std::rc::{Rc, Weak};
use std::cell::RefCell;
use std::default::Default;
use std::collections::VecDeque;
use timely::dataflow::Scope;
use timely::dataflow::operators::generic::source;
use timely::progress::Timestamp;
use timely::dataflow::operators::CapabilitySet;
use lattice::Lattice;
use trace::{Trace, TraceReader, Batch, BatchReader, Cursor};
use trace::wrappers::rc::TraceBox;
use timely::scheduling::Activator;
use super::{TraceWriter, TraceAgentQueueWriter, TraceAgentQueueReader, Arranged};
use super::TraceReplayInstruction;
pub struct TraceAgent<Tr>
where
Tr: TraceReader,
Tr::Time: Lattice+Ord+Clone+'static,
{
trace: Rc<RefCell<TraceBox<Tr>>>,
queues: Weak<RefCell<Vec<TraceAgentQueueWriter<Tr>>>>,
advance: Vec<Tr::Time>,
through: Vec<Tr::Time>,
}
impl<Tr> TraceReader for TraceAgent<Tr>
where
Tr: TraceReader,
Tr::Time: Lattice+Ord+Clone+'static,
{
type Key = Tr::Key;
type Val = Tr::Val;
type Time = Tr::Time;
type R = Tr::R;
type Batch = Tr::Batch;
type Cursor = Tr::Cursor;
fn advance_by(&mut self, frontier: &[Tr::Time]) {
self.trace.borrow_mut().adjust_advance_frontier(&self.advance[..], frontier);
self.advance.clear();
self.advance.extend(frontier.iter().cloned());
}
fn advance_frontier(&mut self) -> &[Tr::Time] {
&self.advance[..]
}
fn distinguish_since(&mut self, frontier: &[Tr::Time]) {
self.trace.borrow_mut().adjust_through_frontier(&self.through[..], frontier);
self.through.clear();
self.through.extend(frontier.iter().cloned());
}
fn distinguish_frontier(&mut self) -> &[Tr::Time] {
&self.through[..]
}
fn cursor_through(&mut self, frontier: &[Tr::Time]) -> Option<(Tr::Cursor, <Tr::Cursor as Cursor<Tr::Key, Tr::Val, Tr::Time, Tr::R>>::Storage)> {
self.trace.borrow_mut().trace.cursor_through(frontier)
}
fn map_batches<F: FnMut(&Self::Batch)>(&mut self, f: F) { self.trace.borrow_mut().trace.map_batches(f) }
}
impl<Tr> TraceAgent<Tr>
where
Tr: TraceReader,
Tr::Time: Timestamp+Lattice,
{
pub fn new(trace: Tr) -> (Self, TraceWriter<Tr>)
where
Tr: Trace,
Tr::Batch: Batch<Tr::Key,Tr::Val,Tr::Time,Tr::R>,
{
let trace = Rc::new(RefCell::new(TraceBox::new(trace)));
let queues = Rc::new(RefCell::new(Vec::new()));
let reader = TraceAgent {
trace: trace.clone(),
queues: Rc::downgrade(&queues),
advance: trace.borrow().advance_frontiers.frontier().to_vec(),
through: trace.borrow().through_frontiers.frontier().to_vec(),
};
let writer = TraceWriter::new(
vec![Default::default()],
Rc::downgrade(&trace),
queues,
);
(reader, writer)
}
pub fn new_listener(&mut self, activator: Activator) -> TraceAgentQueueReader<Tr>
where
Tr::Time: Default
{
let mut new_queue = VecDeque::new();
let mut upper = None;
self.trace
.borrow_mut()
.trace
.map_batches(|batch| {
new_queue.push_back(TraceReplayInstruction::Batch(batch.clone(), Some(Default::default())));
upper = Some(batch.upper().to_vec());
});
if let Some(upper) = upper {
new_queue.push_back(TraceReplayInstruction::Frontier(upper));
}
let reference = Rc::new((activator, RefCell::new(new_queue)));
if let Some(queue) = self.queues.upgrade() {
queue.borrow_mut().push(Rc::downgrade(&reference));
}
reference.0.activate();
reference
}
}
impl<Tr> TraceAgent<Tr>
where
Tr: TraceReader+'static,
Tr::Time: Lattice+Ord+Clone+'static,
{
pub fn import<G>(&mut self, scope: &G) -> Arranged<G, TraceAgent<Tr>>
where
G: Scope<Timestamp=Tr::Time>,
Tr::Time: Timestamp,
{
self.import_named(scope, "ArrangedSource")
}
pub fn import_named<G>(&mut self, scope: &G, name: &str) -> Arranged<G, TraceAgent<Tr>>
where
G: Scope<Timestamp=Tr::Time>,
Tr::Time: Timestamp,
{
self.import_core(scope, name).0
}
pub fn import_core<G>(&mut self, scope: &G, name: &str) -> (Arranged<G, TraceAgent<Tr>>, ShutdownButton<CapabilitySet<Tr::Time>>)
where
G: Scope<Timestamp=Tr::Time>,
Tr::Time: Timestamp,
{
let trace = self.clone();
let mut shutdown_button = None;
let stream = {
let shutdown_button_ref = &mut shutdown_button;
source(scope, name, move |capability, info| {
let capabilities = Rc::new(RefCell::new(Some(CapabilitySet::new())));
let activator = scope.activator_for(&info.address[..]);
let queue = self.new_listener(activator);
let activator = scope.activator_for(&info.address[..]);
*shutdown_button_ref = Some(ShutdownButton::new(capabilities.clone(), activator));
capabilities.borrow_mut().as_mut().unwrap().insert(capability);
move |output| {
let mut capabilities = capabilities.borrow_mut();
if let Some(ref mut capabilities) = *capabilities {
let mut borrow = queue.1.borrow_mut();
for instruction in borrow.drain(..) {
match instruction {
TraceReplayInstruction::Frontier(frontier) => {
capabilities.downgrade(&frontier[..]);
},
TraceReplayInstruction::Batch(batch, hint) => {
if let Some(time) = hint {
let delayed = capabilities.delayed(&time);
output.session(&delayed).give(batch);
}
}
}
}
}
}
})
};
(Arranged { stream, trace }, shutdown_button.unwrap())
}
}
pub struct ShutdownButton<T> {
reference: Rc<RefCell<Option<T>>>,
activator: Activator,
}
impl<T> ShutdownButton<T> {
pub fn new(reference: Rc<RefCell<Option<T>>>, activator: Activator) -> Self {
Self { reference, activator }
}
pub fn press(&mut self) {
*self.reference.borrow_mut() = None;
self.activator.activate();
}
pub fn press_on_drop(self) -> ShutdownDeadmans<T> {
ShutdownDeadmans {
button: self
}
}
}
pub struct ShutdownDeadmans<T> {
button: ShutdownButton<T>,
}
impl<T> Drop for ShutdownDeadmans<T> {
fn drop(&mut self) {
self.button.press();
}
}
impl<Tr> Clone for TraceAgent<Tr>
where
Tr: TraceReader,
Tr::Time: Lattice+Ord+Clone+'static,
{
fn clone(&self) -> Self {
self.trace.borrow_mut().adjust_advance_frontier(&[], &self.advance[..]);
self.trace.borrow_mut().adjust_through_frontier(&[], &self.through[..]);
TraceAgent {
trace: self.trace.clone(),
queues: self.queues.clone(),
advance: self.advance.clone(),
through: self.through.clone(),
}
}
}
impl<Tr> Drop for TraceAgent<Tr>
where
Tr: TraceReader,
Tr::Time: Lattice+Ord+Clone+'static,
{
fn drop(&mut self) {
self.trace.borrow_mut().adjust_advance_frontier(&self.advance[..], &[]);
self.trace.borrow_mut().adjust_through_frontier(&self.through[..], &[]);
}
}