From f58a7799babba106be2baebb0427ae1a75c7f240 Mon Sep 17 00:00:00 2001 From: Moritz Hoffmann Date: Tue, 29 Sep 2026 11:40:14 +0200 Subject: [PATCH] Let arranging operators return a caller-chosen agent `arrange_core`, `reduce_with_tactic`, and `arrange_from_upsert` build their trace and hard-wire the reader they return to `TraceAgent`, which hides the trace inside its `TraceBox`. A trace that shares state with something outside the operator, for example a handle through which other threads observe its batches, then has no stable way to hand that state to its caller: `TraceAgent::trace_box_unstable` is the only path. The `Agent` trait abstracts the reader an operator returns: `Agent::new` takes the trace by value and returns the reader and a `TraceWriter`. `TraceAgent` implements it by delegating to `TraceAgent::new`. The new `arrange_core_with_agent`, `reduce_with_tactic_and_agent`, and `arrange_from_upsert_with_agent` return `Arranged` for a caller's agent `A`, which can keep the trace's shared state before wrapping a `TraceAgent`. The existing functions delegate with `TraceAgent`, so their callers are unchanged. The writer stays concrete. `TraceWriter::new` and `TraceBox::new` are public, so any agent can build one. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_015tLhSbZdXrTSK2KwSocT59 --- differential-dataflow/CHANGELOG.md | 4 + .../src/operators/arrange/agent.rs | 22 +++ .../src/operators/arrange/arrangement.rs | 31 +++- .../src/operators/arrange/mod.rs | 2 +- .../src/operators/arrange/upsert.rs | 45 +++-- differential-dataflow/src/operators/reduce.rs | 20 ++- differential-dataflow/tests/generic_agent.rs | 161 ++++++++++++++++++ 7 files changed, 261 insertions(+), 24 deletions(-) create mode 100644 differential-dataflow/tests/generic_agent.rs diff --git a/differential-dataflow/CHANGELOG.md b/differential-dataflow/CHANGELOG.md index 256bb80d2..47419adf6 100644 --- a/differential-dataflow/CHANGELOG.md +++ b/differential-dataflow/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- `Agent` trait over the shared reader an arranging operator returns, implemented by `TraceAgent`, and `arrange_core_with_agent`, `reduce_with_tactic_and_agent`, and `arrange_from_upsert_with_agent`, which return a caller-chosen agent. `Agent::new` receives the trace by value, so an agent can keep state the trace shares before `TraceAgent` takes ownership of it. + ## [0.25.1](https://github.com/TimelyDataflow/differential-dataflow/compare/differential-dataflow-v0.25.0...differential-dataflow-v0.25.1) - 2026-07-15 ### Other diff --git a/differential-dataflow/src/operators/arrange/agent.rs b/differential-dataflow/src/operators/arrange/agent.rs index 8850feecc..8ce484c09 100644 --- a/differential-dataflow/src/operators/arrange/agent.rs +++ b/differential-dataflow/src/operators/arrange/agent.rs @@ -68,6 +68,28 @@ impl TraceReader for TraceAgent { fn map_spans)>(&self, f: F) { self.trace.borrow().trace.map_spans(f) } } +/// A shared reader of a trace, constructed by the operator that maintains the trace. +/// +/// The arranging operators (`arrange_core_with_agent`, `reduce_with_tactic_and_agent`, +/// `arrange_from_upsert_with_agent`) build their trace, hand it to `Agent::new`, keep the writer, +/// and return the agent in the resulting `Arranged`. The operator also reads through its own copy of +/// the agent, so an implementation must honour the `TraceReader` contract of the trace it shares. +pub trait Agent: TraceReader + Sized { + /// The trace the agent shares. + type Trace: Trace = None; // fabricate a data-parallel operator that holds capabilities and consults its input frontier. let reader_ref = &mut reader; @@ -346,21 +363,21 @@ where let mut batcher = batcher(logger.clone(), info.global_id); // Capabilities for the lower envelope of updates in `batcher`. - let mut capabilities = Antichain::>::new(); + let mut capabilities = Antichain::>::new(); let activator = Some(scope.activator_for(std::rc::Rc::clone(&info.address))); - let mut empty_trace = Tr::new(info.clone(), logger.clone(), activator); + let mut empty_trace = A::Trace::new(info.clone(), logger.clone(), activator); // If there is default exertion logic set, install it. if let Some(exert_logic) = scope.worker().config().get::("differential/default_exert_logic").cloned() { empty_trace.set_exert_logic(exert_logic); } - let (reader_local, mut writer) = TraceAgent::new(empty_trace, info, logger); + let (reader_local, mut writer) = A::new(empty_trace, info, logger); *reader_ref = Some(reader_local); // Initialize to the minimal input frontier. - let mut prev_frontier = Antichain::from_elem(Tr::Time::minimum()); + let mut prev_frontier = Antichain::from_elem(::minimum()); move |(input, frontier), output| { @@ -418,7 +435,7 @@ where let description = Description::new( prev_frontier.clone(), frontier.frontier().to_owned(), - Antichain::from_elem(Tr::Time::minimum()), + Antichain::from_elem(::minimum()), ); let (extracted, retained) = batcher.extract(frontier.frontier()); diff --git a/differential-dataflow/src/operators/arrange/mod.rs b/differential-dataflow/src/operators/arrange/mod.rs index 7652f6ec6..0ea0113d9 100644 --- a/differential-dataflow/src/operators/arrange/mod.rs +++ b/differential-dataflow/src/operators/arrange/mod.rs @@ -69,6 +69,6 @@ pub mod arrangement; pub mod upsert; pub use self::writer::TraceWriter; -pub use self::agent::{TraceAgent, ShutdownButton}; +pub use self::agent::{Agent, TraceAgent, ShutdownButton}; pub use self::arrangement::Arranged; diff --git a/differential-dataflow/src/operators/arrange/upsert.rs b/differential-dataflow/src/operators/arrange/upsert.rs index b3cb5f802..f0821b2b4 100644 --- a/differential-dataflow/src/operators/arrange/upsert.rs +++ b/differential-dataflow/src/operators/arrange/upsert.rs @@ -108,12 +108,12 @@ use timely::progress::{Antichain, Timestamp}; use timely::dataflow::operators::Capability; use crate::operators::arrange::arrangement::Arranged; -use crate::trace::{self, BatchCursor, BatchDiff, Builder, Cursor, Description, Navigable, Trace, TraceReader}; +use crate::trace::{self, BatchCursor, BatchDiff, Builder, Cursor, Description, Navigable, Trace}; use crate::{ExchangeData, Hashable}; use crate::trace::implementations::containers::BatchContainer; -use super::TraceAgent; +use super::{Agent, TraceAgent}; /// Arrange data from a stream of keyed upserts. /// @@ -140,14 +140,35 @@ where >, Bu: Builder)>, Output: Into>, { - let mut reader: Option> = None; + arrange_from_upsert_with_agent::, K, V>(stream, name) +} + +/// Arranges data from a stream of keyed upserts like [`arrange_from_upsert`], sharing the trace +/// through the agent `A`. +pub fn arrange_from_upsert_with_agent<'scope, Bu, A, K, V>( + stream: Stream<'scope, A::Time, Vec<(K, Option, A::Time)>>, + name: &str, +) -> Arranged<'scope, A> +where + K: ExchangeData+Hashable+std::hash::Hash, + V: ExchangeData, + A: Agent + Clone + 'static, + A::Trace: 'static, + for<'a> BatchCursor: Cursor< + Key<'a> = &'a K, + Val<'a> = &'a V, + Diff=isize, + >, + Bu: Builder)>, Output: Into>, +{ + let mut reader: Option = None; // fabricate a data-parallel operator that holds capabilities and consults its input frontier. let stream = { let reader = &mut reader; - let exchange = Exchange::new(move |update: &(K,Option,Tr::Time)| (update.0).hashed().into()); + let exchange = Exchange::new(move |update: &(K,Option,A::Time)| (update.0).hashed().into()); let scope = stream.scope(); stream.unary_frontier(exchange, name, move |_capability, info| { @@ -156,24 +177,24 @@ where let logger = scope.worker().logger_for::("differential/arrange").map(Into::into); // Tracks the lower envelope of times in `priority_queue`. - let mut capabilities = Antichain::>::new(); + let mut capabilities = Antichain::>::new(); // Form the trace we will both use internally and publish. let activator = Some(scope.activator_for(std::rc::Rc::clone(&info.address))); - let mut empty_trace = Tr::new(info.clone(), logger.clone(), activator); + let mut empty_trace = A::Trace::new(info.clone(), logger.clone(), activator); if let Some(exert_logic) = scope.worker().config().get::("differential/default_exert_logic").cloned() { empty_trace.set_exert_logic(exert_logic); } - let (mut reader_local, mut writer) = TraceAgent::new(empty_trace, info, logger); + let (mut reader_local, mut writer) = A::new(empty_trace, info, logger); // Capture the reader outside the builder scope. *reader = Some(reader_local.clone()); // Tracks the input frontier, used to populate the lower bound of new batches. - let mut prev_frontier = Antichain::from_elem(Tr::Time::minimum()); + let mut prev_frontier = Antichain::from_elem(::minimum()); // For stashing input upserts, ordered increasing by time (`BinaryHeap` is a max-heap). - let mut priority_queue = BinaryHeap::)>>::new(); + let mut priority_queue = BinaryHeap::)>>::new(); let mut updates = Vec::new(); move |(input, frontier), output| { @@ -234,7 +255,7 @@ where let batches = reader_local.batches_through(Antichain::new().borrow()).unwrap(); let (mut trace_cursor, trace_storage) = crate::trace::cursor::cursor_list(batches); let mut builder = Bu::default(); - let mut key_con = as Cursor>::KeyContainer::with_capacity(1); + let mut key_con = as Cursor>::KeyContainer::with_capacity(1); for (key, mut list) in to_process { key_con.clear(); key_con.push_ref(&key); @@ -248,7 +269,7 @@ where // Determine the prior value associated with the key. while let Some(val) = trace_cursor.get_val(&trace_storage) { let mut count = 0; - trace_cursor.map_times(&trace_storage, |_time, diff| count += as Cursor>::owned_diff(diff)); + trace_cursor.map_times(&trace_storage, |_time, diff| count += as Cursor>::owned_diff(diff)); assert!(count == 0 || count == 1); if count == 1 { assert!(prev_value.is_none()); @@ -277,7 +298,7 @@ where updates.sort(); builder.push(&mut updates); } - let description = Description::new(prev_frontier.clone(), upper.clone(), Antichain::from_elem(Tr::Time::minimum())); + let description = Description::new(prev_frontier.clone(), upper.clone(), Antichain::from_elem(::minimum())); let batch = crate::trace::Span::new(description, builder.done().map(Into::into)); prev_frontier.clone_from(&upper); diff --git a/differential-dataflow/src/operators/reduce.rs b/differential-dataflow/src/operators/reduce.rs index e93f54826..e8b4c0b41 100644 --- a/differential-dataflow/src/operators/reduce.rs +++ b/differential-dataflow/src/operators/reduce.rs @@ -15,7 +15,7 @@ use timely::dataflow::operators::Operator; use timely::dataflow::operators::CapabilitySet; use timely::dataflow::channels::pact::Pipeline; -use crate::operators::arrange::{Arranged, TraceAgent}; +use crate::operators::arrange::{Agent, Arranged, TraceAgent}; use crate::trace::{Span, ExertionLogic, Trace, TraceReader}; /// Sort and deduplicate a list. Shared by the cursor and proxy tactics, which @@ -75,11 +75,23 @@ pub use crate::operators::cursor::reduce::reduce_trace; /// `TraceReader` of its input and `Trace` of its output, never `Navigable`: it extracts batches via /// `spans_through`, and building cursors over them (if that is how the reduce proceeds) is the /// tactic's concern. -pub fn reduce_with_tactic<'scope, Tr1, Tr2, T>(trace: Arranged<'scope, Tr1>, name: &str, mut tactic: T) -> Arranged<'scope, TraceAgent> +pub fn reduce_with_tactic<'scope, Tr1, Tr2, T>(trace: Arranged<'scope, Tr1>, name: &str, tactic: T) -> Arranged<'scope, TraceAgent> where Tr1: TraceReader + 'static, Tr2: Trace