From 18f295c6a336a4c7685a3a89409c3593aa31b6b2 Mon Sep 17 00:00:00 2001 From: Moritz Hoffmann Date: Tue, 29 Sep 2026 10:48:02 +0200 Subject: [PATCH] Build an operator's trace with a caller-supplied factory `arrange_core` and `reduce_with_tactic` construct their trace with `Trace::new` and move it into a `TraceAgent`, after which the only way to reach the trace instance is `TraceAgent::trace_box_unstable`. A caller whose trace shares state with something outside the operator, for example a handle through which other threads observe its batches, has no stable way to obtain that state. `arrange_core_with_trace` and `reduce_with_tactic_and_trace` take a factory with the signature of `Trace::new` and call it once, where the operator builds its trace. The caller can take a handle from the trace before returning it. `arrange_core` and `reduce_with_tactic` delegate with `Trace::new`, so existing callers are unchanged. 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/arrangement.rs | 26 ++++++- differential-dataflow/src/operators/reduce.rs | 27 ++++++- differential-dataflow/tests/trace_factory.rs | 75 +++++++++++++++++++ 4 files changed, 128 insertions(+), 4 deletions(-) create mode 100644 differential-dataflow/tests/trace_factory.rs 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