diff --git a/diagnostics/src/logging.rs b/diagnostics/src/logging.rs index 979ddece8..18bc36d74 100644 --- a/diagnostics/src/logging.rs +++ b/diagnostics/src/logging.rs @@ -448,46 +448,46 @@ fn construct_timely<'scope>( let mut chs = ch_act.session(&cap); let mut els = el_act.session(&cap); let mut msgs = msg_act.session(&cap); - let ts = *cap.time(); - - for (event_time, event) in data.drain(..) { - match event { - TimelyEvent::Operates(e) => { - ops.give(((e.id, e.name.clone(), e.addr.clone()), ts, 1i64)); - state.operators.insert(e.id, e); - } - TimelyEvent::Shutdown(e) => { - if let Some(op) = state.operators.remove(&e.id) { - ops.give(((op.id, op.name, op.addr), ts, -1i64)); + if let Some(&ts) = cap.least() { + for (event_time, event) in data.drain(..) { + match event { + TimelyEvent::Operates(e) => { + ops.give(((e.id, e.name.clone(), e.addr.clone()), ts, 1i64)); + state.operators.insert(e.id, e); } - } - TimelyEvent::Channels(e) => { - chs.give(( - (e.id, e.scope_addr.clone(), e.source, e.target), - ts, - 1i64, - )); - } - TimelyEvent::Schedule(e) => match e.start_stop { - StartStop::Start => { - state.schedule_starts.insert(e.id, event_time); + TimelyEvent::Shutdown(e) => { + if let Some(op) = state.operators.remove(&e.id) { + ops.give(((op.id, op.name, op.addr), ts, -1i64)); + } } - StartStop::Stop => { - if let Some(start) = state.schedule_starts.remove(&e.id) { - let elapsed_ns = - event_time.saturating_sub(start).as_nanos() as i64; - if elapsed_ns > 0 { - els.give((e.id, ts, elapsed_ns)); + TimelyEvent::Channels(e) => { + chs.give(( + (e.id, e.scope_addr.clone(), e.source, e.target), + ts, + 1i64, + )); + } + TimelyEvent::Schedule(e) => match e.start_stop { + StartStop::Start => { + state.schedule_starts.insert(e.id, event_time); + } + StartStop::Stop => { + if let Some(start) = state.schedule_starts.remove(&e.id) { + let elapsed_ns = + event_time.saturating_sub(start).as_nanos() as i64; + if elapsed_ns > 0 { + els.give((e.id, ts, elapsed_ns)); + } } } + }, + TimelyEvent::Messages(e) => { + if e.is_send { + msgs.give((e.channel, ts, e.record_count as i64)); + } } - }, - TimelyEvent::Messages(e) => { - if e.is_send { - msgs.give((e.channel, ts, e.record_count as i64)); - } + _ => {} } - _ => {} } } }); @@ -587,40 +587,40 @@ fn construct_differential<'scope>( let mut b_sz = bs_act.session(&cap); let mut b_cap = bc_act.session(&cap); let mut b_alloc = ba_act.session(&cap); - let ts = *cap.time(); - - for (_event_time, event) in data.drain(..) { - match event { - DifferentialEvent::Batch(e) => { - bat.give((e.operator, ts, 1i64)); - rec.give((e.operator, ts, e.length as i64)); - } - DifferentialEvent::Merge(e) => { - if let Some(complete) = e.complete { + if let Some(&ts) = cap.least() { + for (_event_time, event) in data.drain(..) { + match event { + DifferentialEvent::Batch(e) => { + bat.give((e.operator, ts, 1i64)); + rec.give((e.operator, ts, e.length as i64)); + } + DifferentialEvent::Merge(e) => { + if let Some(complete) = e.complete { + bat.give((e.operator, ts, -1i64)); + let diff = complete as i64 - (e.length1 + e.length2) as i64; + if diff != 0 { + rec.give((e.operator, ts, diff)); + } + } + } + DifferentialEvent::Drop(e) => { bat.give((e.operator, ts, -1i64)); - let diff = complete as i64 - (e.length1 + e.length2) as i64; + let diff = -(e.length as i64); if diff != 0 { rec.give((e.operator, ts, diff)); } } - } - DifferentialEvent::Drop(e) => { - bat.give((e.operator, ts, -1i64)); - let diff = -(e.length as i64); - if diff != 0 { - rec.give((e.operator, ts, diff)); + DifferentialEvent::TraceShare(e) => { + shr.give((e.operator, ts, e.diff as i64)); } + DifferentialEvent::Batcher(e) => { + b_rec.give((e.operator, ts, e.records_diff as i64)); + b_sz.give((e.operator, ts, e.size_diff as i64)); + b_cap.give((e.operator, ts, e.capacity_diff as i64)); + b_alloc.give((e.operator, ts, e.allocations_diff as i64)); + } + _ => {} } - DifferentialEvent::TraceShare(e) => { - shr.give((e.operator, ts, e.diff as i64)); - } - DifferentialEvent::Batcher(e) => { - b_rec.give((e.operator, ts, e.records_diff as i64)); - b_sz.give((e.operator, ts, e.size_diff as i64)); - b_cap.give((e.operator, ts, e.capacity_diff as i64)); - b_alloc.give((e.operator, ts, e.allocations_diff as i64)); - } - _ => {} } } }); diff --git a/differential-dataflow/src/capture.rs b/differential-dataflow/src/capture.rs index ba1522aca..4b194c091 100644 --- a/differential-dataflow/src/capture.rs +++ b/differential-dataflow/src/capture.rs @@ -227,7 +227,7 @@ pub mod source { use std::rc::Rc; use std::marker::{Send, Sync}; use std::sync::Arc; - use timely::dataflow::{Scope, Stream, operators::{Capability, CapabilitySet}}; + use timely::dataflow::{Scope, Stream, operators::CapabilitySet}; use timely::dataflow::operators::generic::OutputBuilder; use timely::progress::Timestamp; use timely::scheduling::SyncActivator; @@ -462,17 +462,17 @@ pub mod source { // If the frontier changes we need a capability to express that. // Any capability should work; the downstream listener doesn't care. - let mut capability: Option> = None; + let mut capability: Option> = None; // Drain all relevant update counts in to the mutable antichain tracking its frontier. counts.for_each(|cap, counts| { updates_frontier.update_iter(counts.iter().cloned()); - capability = Some(cap.retain(0)); + capability = Some(cap.retain_stamp(0)); }); // Drain all progress statements into the queue out of which we will work. input.for_each(|cap, progress| { progress_queue.extend(progress.iter().map(|x| (x.1).clone())); - capability = Some(cap.retain(0)); + capability = Some(cap.retain_stamp(0)); }); // Extract and act on actionable progress messages. diff --git a/differential-dataflow/src/collection.rs b/differential-dataflow/src/collection.rs index 3500aa9cb..e726993a0 100644 --- a/differential-dataflow/src/collection.rs +++ b/differential-dataflow/src/collection.rs @@ -128,15 +128,15 @@ impl<'scope, T: Timestamp, C: Container> Collection<'scope, T, C> { /// scope.new_collection_from(1 .. 10).1 /// .map_in_place(|x| *x *= 2) /// .filter(|x| x % 2 == 1) - /// .inspect_container(|event| println!("event: {:?}", event)); + /// .inspect_core(|event| println!("event: {:?}", event)); /// }); /// ``` - pub fn inspect_container(self, func: F) -> Self + pub fn inspect_core(self, func: F) -> Self where - F: FnMut(Result<(&T, &C), &[T]>)+'static, + F: FnMut(Result<(&timely::progress::Stamp, &C), &[T]>)+'static, { self.inner - .inspect_container(func) + .inspect_core(func) .as_collection() } /// Attaches a timely dataflow probe to the output of a Collection. @@ -562,17 +562,25 @@ pub mod vec { /// ordered, they should have the same order or compare equal once `func` is applied to them (this /// is because we advance the timely capability with the same logic, and it must remain `less_equal` /// to all of the data timestamps). - pub fn delay(self, func: F) -> Collection<'scope, T, D, R> + pub fn delay(self, mut func: F) -> Collection<'scope, T, D, R> where - T: Hash, - F: FnMut(&T) -> T + Clone + 'static, + F: FnMut(&T) -> T + 'static, { - let mut func1 = func.clone(); - let mut func2 = func.clone(); - + use timely::dataflow::channels::pact::Pipeline; + use timely::dataflow::operators::CapabilitySet; self.inner - .delay_batch(move |x| func1(x)) - .map_in_place(move |x| x.1 = func2(&x.1)) + .unary(Pipeline, "Delay", move |_, _| move |input, output| { + input.for_each_stamp(|time, data| { + // Each element of the stamp is delayed by `func`, and each update's time likewise; + // by monotonicity the delayed capabilities still cover the delayed updates. + let caps: CapabilitySet = time.stamp().iter().map(|t| time.delayed(&func(t), output.output_index())).collect(); + let mut session = output.session(&caps); + for data in data { + for (_, t, _) in data.iter_mut() { *t = func(t); } + session.give_container(data); + } + }); + }) .as_collection() } @@ -609,8 +617,8 @@ pub mod vec { } /// Applies a supplied function to each batch of updates. /// - /// This method is analogous to `inspect`, but operates on batches and reveals the timestamp of the - /// timely dataflow capability associated with the batch of updates. The observed batching depends + /// This method is analogous to `inspect`, but operates on batches and reveals the stamp under which + /// the batch of updates travels: the timestamps of the timely dataflow capabilities associated with it. The observed batching depends /// on how the system executes, and may vary run to run. /// /// # Examples @@ -627,10 +635,10 @@ pub mod vec { /// ``` pub fn inspect_batch(self, mut func: F) -> Collection<'scope, T, D, R> where - F: FnMut(&T, &[(D, T, R)])+'static, + F: FnMut(&timely::progress::Stamp, &[(D, T, R)])+'static, { self.inner - .inspect_batch(move |time, data| func(time, data)) + .inspect_core(move |event| if let Ok((stamp, data)) = event { func(stamp, data) }) .as_collection() } diff --git a/differential-dataflow/src/operators/arrange/arrangement.rs b/differential-dataflow/src/operators/arrange/arrangement.rs index a7d0d817e..d688d594f 100644 --- a/differential-dataflow/src/operators/arrange/arrangement.rs +++ b/differential-dataflow/src/operators/arrange/arrangement.rs @@ -343,7 +343,7 @@ where let mut reader: Option> = None; - // fabricate a data-parallel operator using the `unary_notify` pattern. + // fabricate a data-parallel operator that holds capabilities and consults its input frontier. let reader_ref = &mut reader; let scope = stream.scope(); diff --git a/differential-dataflow/src/operators/arrange/upsert.rs b/differential-dataflow/src/operators/arrange/upsert.rs index 466ca4c5f..b3cb5f802 100644 --- a/differential-dataflow/src/operators/arrange/upsert.rs +++ b/differential-dataflow/src/operators/arrange/upsert.rs @@ -142,7 +142,7 @@ where { let mut reader: Option> = None; - // fabricate a data-parallel operator using the `unary_notify` pattern. + // fabricate a data-parallel operator that holds capabilities and consults its input frontier. let stream = { let reader = &mut reader; @@ -180,7 +180,7 @@ where // Stash capabilities and associated data (ordered by time). input.for_each(|cap, data| { - capabilities.insert(cap.retain(0)); + if let Some(cap) = cap.retain_least(0) { capabilities.insert(cap); } for (key, val, time) in data.drain(..) { priority_queue.push(std::cmp::Reverse((time, key, val))) } diff --git a/differential-dataflow/src/operators/reduce.rs b/differential-dataflow/src/operators/reduce.rs index 6a5bd0d56..e93f54826 100644 --- a/differential-dataflow/src/operators/reduce.rs +++ b/differential-dataflow/src/operators/reduce.rs @@ -83,7 +83,7 @@ where { let mut result_trace = None; - // fabricate a data-parallel operator using the `unary_notify` pattern. + // fabricate a data-parallel operator that holds capabilities and consults its input frontier. let stream = { let mut source_trace = trace.trace; diff --git a/differential-dataflow/tests/dynamic.rs b/differential-dataflow/tests/dynamic.rs index 05e896bb7..5ed3c9dec 100644 --- a/differential-dataflow/tests/dynamic.rs +++ b/differential-dataflow/tests/dynamic.rs @@ -48,7 +48,7 @@ fn leave_dynamic_on_a_message_over_two_epochs() { let mut output = output.activate(); input.for_each(|cap, data| { let Some(root) = root.as_ref() else { return }; - let epoch = cap.time().outer; + let epoch = cap.stamp().iter().map(|t| t.outer).min().unwrap(); let t1 = Product::new(epoch, PointStamp::new([3].into_iter().collect())); let t2 = Product::new(epoch + 1, PointStamp::new([0].into_iter().collect())); let caps: CapabilitySet