From 95ce57958c4a5cb8b88ebda19ca7b97fe736fc7f Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Tue, 1 Sep 2026 22:31:40 -0400 Subject: [PATCH 1/6] leave_dynamic: a capability per element of a multi-element stamp MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both the row and the corgi `leave_dynamic` read `cap.time()`, which now panics ("expected a singleton stamp") when a message carries timestamps from two epochs — late iterations of one alongside early ones of the next, which any steady-state run of an iterative program produces. Hold a capability per stamp element, each truncated exactly as the records are, and open the session on the set. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_012k2GSwxmvD2LvckkoXi6GK --- .../src/columnar/collection/operators.rs | 19 ++++++++++----- differential-dataflow/src/dynamic/mod.rs | 23 +++++++++++-------- interactive/src/backend/corgi.rs | 22 ++++++++++-------- 3 files changed, 40 insertions(+), 24 deletions(-) diff --git a/differential-dataflow/src/columnar/collection/operators.rs b/differential-dataflow/src/columnar/collection/operators.rs index ff7b3f21a..fe7db42e4 100644 --- a/differential-dataflow/src/columnar/collection/operators.rs +++ b/differential-dataflow/src/columnar/collection/operators.rs @@ -116,12 +116,19 @@ where move |_frontier| { let mut output = output.activate(); op_input.for_each(|cap, data| { - // Truncate the capability's timestamp. - let mut new_time = cap.time().clone(); - let mut vec = std::mem::take(&mut new_time.inner).into_inner(); - vec.truncate(level - 1); - new_time.inner = PointStamp::new(vec); - let new_cap = cap.delayed(&new_time, 0); + // A message may carry several timestamps (a multi-element stamp): hold a + // capability for each, truncated exactly as the updates are. + let new_cap: timely::dataflow::operators::CapabilitySet<_> = cap + .stamp() + .iter() + .map(|t| { + let mut new_time = t.clone(); + let mut vec = std::mem::take(&mut new_time.inner).into_inner(); + vec.truncate(level - 1); + new_time.inner = PointStamp::new(vec); + cap.delayed(&new_time, 0) + }) + .collect(); // Push updates with truncated times into the builder. // The builder's form call on flush sorts and consolidates, // handling the duplicate times that truncation can produce. diff --git a/differential-dataflow/src/dynamic/mod.rs b/differential-dataflow/src/dynamic/mod.rs index 0a615ffbf..1020b33af 100644 --- a/differential-dataflow/src/dynamic/mod.rs +++ b/differential-dataflow/src/dynamic/mod.rs @@ -16,6 +16,7 @@ pub mod pointstamp; use timely::order::Product; use timely::progress::Timestamp; use timely::dataflow::operators::generic::{OutputBuilder, builder_rc::OperatorBuilder}; +use timely::dataflow::operators::CapabilitySet; use timely::dataflow::channels::pact::Pipeline; use timely::progress::Antichain; @@ -47,17 +48,21 @@ where builder.build(move |_capability| move |_frontier| { let mut output = output.activate(); input.for_each(|cap, data| { - let mut new_time = cap.time().clone(); - let mut vec = std::mem::take(&mut new_time.inner).into_inner(); - vec.truncate(level - 1); - new_time.inner = PointStamp::new(vec); - let new_cap = cap.delayed(&new_time, 0); - for (_data, time, _diff) in data.iter_mut() { - let mut vec = std::mem::take(&mut time.inner).into_inner(); + // A message may carry several timestamps (a multi-element stamp, e.g. late + // iterations of one epoch alongside early ones of the next): hold a capability + // for each, truncated exactly as the records are. + let truncate = |time: &Product>| { + let mut new_time = time.clone(); + let mut vec = std::mem::take(&mut new_time.inner).into_inner(); vec.truncate(level - 1); - time.inner = PointStamp::new(vec); + new_time.inner = PointStamp::new(vec); + new_time + }; + let caps: CapabilitySet<_> = cap.stamp().iter().map(|t| cap.delayed(&truncate(t), 0)).collect(); + for (_data, time, _diff) in data.iter_mut() { + *time = truncate(time); } - output.session(&new_cap).give_container(data); + output.session(&caps).give_container(data); }); }); diff --git a/interactive/src/backend/corgi.rs b/interactive/src/backend/corgi.rs index 4422b235c..3a1cb3eff 100644 --- a/interactive/src/backend/corgi.rs +++ b/interactive/src/backend/corgi.rs @@ -367,17 +367,21 @@ impl Backend for CorgiBackend { builder.build(move |_capability| move |_frontier| { let mut output = output.activate(); input.for_each(|cap, data| { - let mut new_time = cap.time().clone(); - let mut v = std::mem::take(&mut new_time.inner).into_inner(); - v.truncate(level - 1); - new_time.inner = PointStamp::new(v); - let new_cap = cap.delayed(&new_time, 0); - for t in data.times.iter_mut() { - let mut v = std::mem::take(&mut t.inner).into_inner(); + // A message may carry several timestamps (a multi-element stamp): hold a + // capability for each, truncated exactly as the records are. + let truncate = |t: &Time| { + let mut new_time = t.clone(); + let mut v = std::mem::take(&mut new_time.inner).into_inner(); v.truncate(level - 1); - t.inner = PointStamp::new(v); + new_time.inner = PointStamp::new(v); + new_time + }; + let caps: timely::dataflow::operators::CapabilitySet