diff --git a/differential-dataflow/CHANGELOG.md b/differential-dataflow/CHANGELOG.md index 256bb80d2..b619cb213 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 + +- `arrange_core_with_trace` and `reduce_with_tactic_and_trace`, which build the operator's trace with a caller-supplied factory in place of `Trace::new`, so a caller can keep a handle into the trace 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/arrangement.rs b/differential-dataflow/src/operators/arrange/arrangement.rs index 457b760d9..83af43648 100644 --- a/differential-dataflow/src/operators/arrange/arrangement.rs +++ b/differential-dataflow/src/operators/arrange/arrangement.rs @@ -20,7 +20,8 @@ use timely::dataflow::operators::{Enter, vec::Map}; use timely::order::PartialOrder; use timely::dataflow::{Scope, Stream}; -use timely::dataflow::operators::generic::Operator; +use timely::dataflow::operators::generic::{Operator, OperatorInfo}; +use timely::scheduling::Activator; use timely::dataflow::channels::pact::{ParallelizationContract, Pipeline}; use timely::progress::Timestamp; use timely::progress::Antichain; @@ -310,6 +311,27 @@ pub fn arrange_core<'scope, P, C, Ba, Tr>( name: &str, batcher: impl FnOnce(Option, usize) -> Ba, ) -> Arranged<'scope, TraceAgent> +where + C: Container + Clone + 'static, + P: ParallelizationContract, + Ba: Batcher> + 'static, + Tr: Trace+'static, +{ + arrange_core_with_trace(stream, pact, name, batcher, Tr::new) +} + +/// Arranges a stream of updates like [`arrange_core`], building the trace with `trace`. +/// +/// `trace` receives the arguments [`Trace::new`] would, and is called once, on the operator's +/// construction. A caller that needs to reach the trace after [`TraceAgent`] takes ownership of it, +/// for example through a handle the trace shares, can obtain that handle here. +pub fn arrange_core_with_trace<'scope, P, C, Ba, Tr>( + stream: Stream<'scope, Tr::Time, C>, + pact: P, + name: &str, + batcher: impl FnOnce(Option, usize) -> Ba, + trace: impl FnOnce(OperatorInfo, Option, Option) -> Tr, +) -> Arranged<'scope, TraceAgent> where C: Container + Clone + 'static, P: ParallelizationContract, @@ -349,7 +371,7 @@ where 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 = trace(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); diff --git a/differential-dataflow/src/operators/reduce.rs b/differential-dataflow/src/operators/reduce.rs index e93f54826..24d0ef438 100644 --- a/differential-dataflow/src/operators/reduce.rs +++ b/differential-dataflow/src/operators/reduce.rs @@ -14,8 +14,11 @@ use timely::progress::Timestamp; use timely::dataflow::operators::Operator; use timely::dataflow::operators::CapabilitySet; use timely::dataflow::channels::pact::Pipeline; +use timely::dataflow::operators::generic::OperatorInfo; +use timely::scheduling::Activator; use crate::operators::arrange::{Arranged, TraceAgent}; +use crate::logging::Logger; use crate::trace::{Span, ExertionLogic, Trace, TraceReader}; /// Sort and deduplicate a list. Shared by the cursor and proxy tactics, which @@ -75,7 +78,27 @@ 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