Skip to main content

mz_compute/
render.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Renders a plan into a timely/differential dataflow computation.
11//!
12//! ## Error handling
13//!
14//! Timely and differential have no idioms for computations that can error. The
15//! philosophy is, reasonably, to define the semantics of the computation such
16//! that errors are unnecessary: e.g., by using wrap-around semantics for
17//! integer overflow.
18//!
19//! Unfortunately, SQL semantics are not nearly so elegant, and require errors
20//! in myriad cases. The classic example is a division by zero, but invalid
21//! input for casts, overflowing integer operations, and dozens of other
22//! functions need the ability to produce errors ar runtime.
23//!
24//! At the moment, only *scalar* expression evaluation can fail, so only
25//! operators that evaluate scalar expressions can fail. At the time of writing,
26//! that includes map, filter, reduce, and join operators. Constants are a bit
27//! of a special case: they can be either a constant vector of rows *or* a
28//! constant, singular error.
29//!
30//! The approach taken is to build two parallel trees of computation: one for
31//! the rows that have been successfully evaluated (the "oks tree"), and one for
32//! the errors that have been generated (the "errs tree"). For example:
33//!
34//! ```text
35//!    oks1  errs1       oks2  errs2
36//!      |     |           |     |
37//!      |     |           |     |
38//!   project  |           |     |
39//!      |     |           |     |
40//!      |     |           |     |
41//!     map    |           |     |
42//!      |\    |           |     |
43//!      | \   |           |     |
44//!      |  \  |           |     |
45//!      |   \ |           |     |
46//!      |    \|           |     |
47//!   project  +           +     +
48//!      |     |          /     /
49//!      |     |         /     /
50//!    join ------------+     /
51//!      |     |             /
52//!      |     | +----------+
53//!      |     |/
54//!     oks   errs
55//! ```
56//!
57//! The project operation cannot fail, so errors from errs1 are propagated
58//! directly. Map operators are fallible and so can inject additional errors
59//! into the stream. Join operators combine the errors from each of their
60//! inputs.
61//!
62//! The semantics of the error stream are minimal. From the perspective of SQL,
63//! a dataflow is considered to be in an error state if there is at least one
64//! element in the final errs collection. The error value returned to the user
65//! is selected arbitrarily; SQL only makes provisions to return one error to
66//! the user at a time. There are plans to make the err collection accessible to
67//! end users, so they can see all errors at once.
68//!
69//! To make errors transient, simply ensure that the operator can retract any
70//! produced errors when corrected data arrives. To make errors permanent, write
71//! the operator such that it never retracts the errors it produced. Future work
72//! will likely want to introduce some sort of sort order for errors, so that
73//! permanent errors are returned to the user ahead of transient errors—probably
74//! by introducing a new error type a la:
75//!
76//! ```no_run
77//! # struct EvalError;
78//! # struct SourceError;
79//! enum DataflowError {
80//!     Transient(EvalError),
81//!     Permanent(SourceError),
82//! }
83//! ```
84//!
85//! If the error stream is empty, the oks stream must be correct. If the error
86//! stream is non-empty, then there are no semantics for the oks stream. This is
87//! sufficient to support SQL in its current form, but is likely to be
88//! unsatisfactory long term. We suspect that we can continue to imbue the oks
89//! stream with semantics if we are very careful in describing what data should
90//! and should not be produced upon encountering an error. Roughly speaking, the
91//! oks stream could represent the correct result of the computation where all
92//! rows that caused an error have been pruned from the stream. There are
93//! strange and confusing questions here around foreign keys, though: what if
94//! the optimizer proves that a particular key must exist in a collection, but
95//! the key gets pruned away because its row participated in a scalar expression
96//! evaluation that errored?
97//!
98//! In the meantime, it is probably wise for operators to keep the oks stream
99//! roughly "as correct as possible" even when errors are present in the errs
100//! stream. This reduces the amount of recomputation that must be performed
101//! if/when the errors are retracted.
102
103use std::any::Any;
104use std::cell::RefCell;
105use std::collections::{BTreeMap, BTreeSet};
106use std::convert::Infallible;
107use std::future::Future;
108use std::pin::Pin;
109use std::rc::{Rc, Weak};
110use std::sync::Arc;
111use std::task::Poll;
112
113use ::columnar::{Columnar as ColumnarData, Index as ColumnarIndex, Push as ColumnarPush};
114use differential_dataflow::dynamic::pointstamp::PointStamp;
115use differential_dataflow::lattice::Lattice;
116use differential_dataflow::operators::arrange::Arranged;
117use differential_dataflow::operators::arrange::ShutdownButton;
118use differential_dataflow::operators::iterate::Variable;
119use differential_dataflow::trace::cursor::{BatchCursor, BatchDiff, BatchKey, BatchVal};
120use differential_dataflow::trace::{BatchReader, Cursor, Navigable, TraceReader};
121use differential_dataflow::{AsCollection, Collection, Data, VecCollection};
122use futures::FutureExt;
123use futures::channel::oneshot;
124use itertools::Itertools;
125use mz_compute_types::dataflows::{DataflowDescription, IndexDesc};
126use mz_compute_types::dyncfgs::{
127    COMPUTE_APPLY_COLUMN_DEMANDS, COMPUTE_LOGICAL_BACKPRESSURE_INFLIGHT_SLACK,
128    COMPUTE_LOGICAL_BACKPRESSURE_MAX_RETAINED_CAPABILITIES, ENABLE_COMPUTE_LOGICAL_BACKPRESSURE,
129    ENABLE_COMPUTE_TEMPORAL_BUCKETING, ENABLE_ERROR_DISTINCT, SUBSCRIBE_SNAPSHOT_OPTIMIZATION,
130    TEMPORAL_BUCKETING_SUMMARY,
131};
132use mz_compute_types::plan::render_plan::{
133    self, BindStage, LetBind, LetFreePlan, RecBind, RenderPlan,
134};
135use mz_compute_types::plan::scalar::LirScalarExpr;
136use mz_compute_types::plan::{ArrangementStrategy, LirId};
137use mz_expr::{EvalError, Id, LocalId, permutation_for_arrangement};
138use mz_persist_client::operators::shard_source::{ErrorHandler, SnapshotMode};
139use mz_repr::explain::DummyHumanizer;
140use mz_repr::fixed_length::ExtendDatums;
141use mz_repr::{Datum, DatumVec, Diff, GlobalId, ReprRelationType, Row, RowArena, SharedRow};
142use mz_storage_operators::persist_source;
143use mz_storage_types::controller::CollectionMetadata;
144use mz_timely_util::columnar::Column;
145use mz_timely_util::columnar::builder::ColumnBuilder;
146use mz_timely_util::columnation::ColumnationChunker;
147use mz_timely_util::operator::StreamExt;
148use mz_timely_util::probe::{Handle as MzProbeHandle, ProbeNotify};
149use mz_timely_util::scope_label::ScopeExt;
150use timely::PartialOrder;
151use timely::container::CapacityContainerBuilder;
152use timely::dataflow::channels::pact::Pipeline;
153use timely::dataflow::operators::core::to_stream::ToStreamBuilder;
154use timely::dataflow::operators::vec::Filter;
155use timely::dataflow::operators::vec::ToStream;
156use timely::dataflow::operators::{Capability, Operator, Probe, probe};
157use timely::dataflow::{Scope, Stream, StreamVec};
158use timely::order::{Product, TotalOrder};
159use timely::progress::timestamp::Refines;
160use timely::progress::{Antichain, Timestamp};
161use timely::scheduling::ActivateOnDrop;
162use timely::worker::Worker as TimelyWorker;
163
164use crate::arrangement::manager::TraceBundle;
165use crate::compute_state::ComputeState;
166use crate::extensions::arrange::{KeyCollection, MzArrange};
167use crate::extensions::reduce::MzReduce;
168use crate::extensions::temporal_bucket::TemporalBucketing;
169use crate::logging::compute::{
170    ComputeEvent, DataflowGlobal, LirMapping, LirMetadata, LogDataflowErrors, OperatorHydration,
171};
172use crate::render::columnar::{
173    ColCollection, RecTimestamp, columnar_consolidate, columnar_leave_dynamic, columnar_negate,
174    concat_many, flat_map_datums,
175};
176use crate::render::context::{ArrangementFlavor, Context};
177use crate::render::errors::DataflowErrorSer;
178use crate::typedefs::{ErrBatcher, ErrBuilder, ErrSpine, MzTimestamp};
179use mz_row_spine::{DatumSeq, RowRowBatcher, RowRowBuilder};
180use mz_timely_util::columnar::consolidate::ConsolidatingColumnBuilder;
181
182pub(crate) mod columnar;
183pub mod context;
184pub(crate) mod errors;
185mod flat_map;
186mod join;
187mod reduce;
188pub mod sinks;
189mod threshold;
190mod top_k;
191
192pub use context::CollectionBundle;
193pub use join::LinearJoinSpec;
194
195/// Guard that presses a differential [`ShutdownButton`] when dropped.
196///
197/// Dropping this guard releases the imported trace's capabilities.
198struct PressOnDrop<T>(ShutdownButton<T>);
199
200impl<T> Drop for PressOnDrop<T> {
201    fn drop(&mut self) {
202        self.0.press();
203    }
204}
205
206/// Assemble the "compute"  side of a dataflow, i.e. all but the sources.
207///
208/// This method imports sources from provided assets, and then builds the remaining
209/// dataflow using "compute-local" assets like shared arrangements, and producing
210/// both arrangements and sinks.
211pub fn build_compute_dataflow(
212    timely_worker: &mut TimelyWorker,
213    compute_state: &mut ComputeState,
214    dataflow: DataflowDescription<RenderPlan, CollectionMetadata>,
215    start_signal: StartSignal,
216    until: Antichain<mz_repr::Timestamp>,
217    dataflow_expiration: Antichain<mz_repr::Timestamp>,
218) {
219    // Mutually recursive view definitions require special handling.
220    let recursive = dataflow
221        .objects_to_build
222        .iter()
223        .any(|object| object.plan.is_recursive());
224
225    // Determine indexes to export, and their dependencies.
226    let indexes = dataflow
227        .index_exports
228        .iter()
229        .map(|(idx_id, (idx, _typ))| (*idx_id, dataflow.depends_on(idx.on_id), idx.as_lir()))
230        .collect::<Vec<_>>();
231
232    // Determine sinks to export, and their dependencies.
233    let sinks = dataflow
234        .sink_exports
235        .iter()
236        .map(|(sink_id, sink)| (*sink_id, dataflow.depends_on(sink.from), sink.clone()))
237        .collect::<Vec<_>>();
238
239    let worker_logging = timely_worker.logger_for("timely").map(Into::into);
240    let apply_demands = COMPUTE_APPLY_COLUMN_DEMANDS.get(&compute_state.worker_config);
241    let subscribe_snapshot_optimization =
242        SUBSCRIBE_SNAPSHOT_OPTIMIZATION.get(&compute_state.worker_config);
243
244    let name = format!("Dataflow: {}", dataflow.debug_name);
245    let input_name = format!("InputRegion: {}", dataflow.debug_name);
246    let build_name = format!("BuildRegion: {}", dataflow.debug_name);
247
248    timely_worker.dataflow_core(&name, worker_logging, Box::new(()), |_, scope| {
249        let scope = scope.with_label();
250
251        // The scope.clone() occurs to allow import in the region.
252        // We build a region here to establish a pattern of a scope inside the dataflow,
253        // so that other similar uses (e.g. with iterative scopes) do not require weird
254        // alternate type signatures.
255        let mut imported_sources = Vec::new();
256        let mut tokens: BTreeMap<_, Rc<dyn Any>> = BTreeMap::new();
257        let output_probe = MzProbeHandle::default();
258
259        scope.clone().region_named(&input_name, |region| {
260            // Import declared sources into the rendering context.
261            for (source_id, import) in dataflow.source_imports.iter() {
262                region.region_named(&format!("Source({:?})", source_id), |inner| {
263                    let mut read_schema = None;
264                    let mut mfp = import.desc.arguments.operators.clone().map(|mut ops| {
265                        // If enabled, we read from Persist with a `RelationDesc` that
266                        // omits uneeded columns.
267                        if apply_demands {
268                            let demands = ops.demand();
269                            let new_desc = import
270                                .desc
271                                .storage_metadata
272                                .relation_desc
273                                .apply_demand(&demands);
274                            let new_arity = demands.len();
275                            let remap: BTreeMap<_, _> = demands
276                                .into_iter()
277                                .enumerate()
278                                .map(|(new, old)| (old, new))
279                                .collect();
280                            ops.permute_fn(|old_idx| remap[&old_idx], new_arity);
281                            read_schema = Some(new_desc);
282                        }
283
284                        mz_expr::MfpPlan::create_from(ops)
285                            .expect("Linear operators should always be valid")
286                    });
287
288                    let snapshot_mode = if import.with_snapshot || !subscribe_snapshot_optimization
289                    {
290                        SnapshotMode::Include
291                    } else {
292                        compute_state.metrics.inc_subscribe_snapshot_optimization();
293                        SnapshotMode::Exclude
294                    };
295                    let suppress_early_progress_as_of = dataflow.as_of.clone();
296
297                    // Note: For correctness, we require that sources only emit times advanced by
298                    // `dataflow.as_of`. `persist_source` is documented to provide this guarantee.
299                    let (mut ok_stream, err_stream, token) = persist_source::persist_source::<
300                        DataflowErrorSer,
301                        ConsolidatingColumnBuilder<Row, mz_repr::Timestamp, Diff>,
302                    >(
303                        inner,
304                        *source_id,
305                        Arc::clone(&compute_state.persist_clients),
306                        &compute_state.txns_ctx,
307                        import.desc.storage_metadata.clone(),
308                        read_schema,
309                        dataflow.as_of.clone(),
310                        snapshot_mode,
311                        until.clone(),
312                        mfp.as_mut(),
313                        compute_state.dataflow_max_inflight_bytes(),
314                        start_signal.clone().into_send_future(),
315                        ErrorHandler::Halt("compute_import"),
316                    );
317
318                    // If `mfp` is non-identity, we need to apply what remains.
319                    // For the moment, assert that it is either trivial or `None`.
320                    assert!(mfp.map(|x| x.is_identity()).unwrap_or(true));
321
322                    // To avoid a memory spike during arrangement hydration (database-issues#6368), need to
323                    // ensure that the first frontier we report into the dataflow is beyond the
324                    // `as_of`.
325                    if let Some(as_of) = suppress_early_progress_as_of {
326                        ok_stream = suppress_early_progress(ok_stream, as_of);
327                    }
328
329                    if ENABLE_COMPUTE_LOGICAL_BACKPRESSURE.get(&compute_state.worker_config) {
330                        // Apply logical backpressure to the source.
331                        let limit = COMPUTE_LOGICAL_BACKPRESSURE_MAX_RETAINED_CAPABILITIES
332                            .get(&compute_state.worker_config);
333                        let slack = COMPUTE_LOGICAL_BACKPRESSURE_INFLIGHT_SLACK
334                            .get(&compute_state.worker_config)
335                            .as_millis()
336                            .try_into()
337                            .expect("must fit");
338
339                        let stream = ok_stream.limit_progress(
340                            output_probe.clone(),
341                            slack,
342                            limit,
343                            import.upper.clone(),
344                            name.clone(),
345                        );
346                        ok_stream = stream;
347                    }
348
349                    // Attach a probe reporting the input frontier.
350                    let input_probe =
351                        compute_state.input_probe_for(*source_id, dataflow.export_ids());
352                    ok_stream = ok_stream.probe_with(&input_probe);
353
354                    let (oks, errs) = (
355                        ok_stream
356                            .as_collection()
357                            .leave_region(region)
358                            .leave_region(scope),
359                        err_stream
360                            .as_collection()
361                            .leave_region(region)
362                            .leave_region(scope),
363                    );
364
365                    imported_sources.push((mz_expr::Id::Global(*source_id), (oks, errs)));
366
367                    // Associate returned tokens with the source identifier.
368                    tokens.insert(*source_id, Rc::new(token));
369                });
370            }
371        });
372
373        // If there exists a recursive expression, we'll need to use a non-region scope,
374        // in order to support additional timestamp coordinates for iteration.
375        if recursive {
376            scope.clone().iterative::<PointStamp<u64>, _, _>(|region| {
377                let mut context = Context::for_dataflow_in(
378                    &dataflow,
379                    region.clone(),
380                    compute_state,
381                    until,
382                    dataflow_expiration,
383                );
384
385                for (id, (oks, errs)) in imported_sources.into_iter() {
386                    let bundle = crate::render::CollectionBundle::from_edge(
387                        oks.enter(region),
388                        errs.enter(region),
389                    );
390                    // Associate collection bundle with the source identifier.
391                    context.insert_id(id, bundle);
392                }
393
394                // Import declared indexes into the rendering context.
395                for (idx_id, idx) in &dataflow.index_imports {
396                    let input_probe = compute_state.input_probe_for(*idx_id, dataflow.export_ids());
397                    let snapshot_mode = if idx.with_snapshot || !subscribe_snapshot_optimization {
398                        SnapshotMode::Include
399                    } else {
400                        compute_state.metrics.inc_subscribe_snapshot_optimization();
401                        SnapshotMode::Exclude
402                    };
403                    context.import_index(
404                        scope,
405                        compute_state,
406                        &mut tokens,
407                        input_probe,
408                        *idx_id,
409                        &idx.desc.as_lir(),
410                        &idx.typ,
411                        snapshot_mode,
412                        start_signal.clone(),
413                    );
414                }
415
416                // Build declared objects.
417                for object in dataflow.objects_to_build {
418                    let bundle = context.scope.clone().region_named(
419                        &format!("BuildingObject({:?})", object.id),
420                        |region| {
421                            let depends = object.plan.depends();
422                            let in_let = object.plan.is_recursive();
423                            context
424                                .enter_region(region, Some(&depends))
425                                .render_recursive_plan(
426                                    object.id,
427                                    0,
428                                    object.plan,
429                                    // recursive plans _must_ have bodies in a let
430                                    BindingInfo::Body { in_let },
431                                )
432                                .leave_region(context.scope)
433                        },
434                    );
435                    let global_id = object.id;
436
437                    context.log_dataflow_global_id(
438                        *bundle
439                            .scope()
440                            .addr()
441                            .first()
442                            .expect("Dataflow root id must exist"),
443                        global_id,
444                    );
445                    context.insert_id(Id::Global(object.id), bundle);
446                }
447
448                // Export declared indexes.
449                for (idx_id, dependencies, idx) in indexes {
450                    context.export_index_iterative(
451                        scope,
452                        compute_state,
453                        &tokens,
454                        dependencies,
455                        idx_id,
456                        &idx,
457                        &output_probe,
458                    );
459                }
460
461                // Export declared sinks.
462                for (sink_id, dependencies, sink) in sinks {
463                    context.export_sink(
464                        compute_state,
465                        &tokens,
466                        dependencies,
467                        sink_id,
468                        &sink,
469                        start_signal.clone(),
470                        &output_probe,
471                        scope,
472                    );
473                }
474            });
475        } else {
476            scope.clone().region_named(&build_name, |region| {
477                let mut context = Context::for_dataflow_in(
478                    &dataflow,
479                    region.clone(),
480                    compute_state,
481                    until,
482                    dataflow_expiration,
483                );
484
485                for (id, (oks, errs)) in imported_sources.into_iter() {
486                    let bundle = crate::render::CollectionBundle::from_edge(
487                        oks.enter_region(region),
488                        errs.enter_region(region),
489                    );
490                    // Associate collection bundle with the source identifier.
491                    context.insert_id(id, bundle);
492                }
493
494                // Import declared indexes into the rendering context.
495                for (idx_id, idx) in &dataflow.index_imports {
496                    let input_probe = compute_state.input_probe_for(*idx_id, dataflow.export_ids());
497                    let snapshot_mode = if idx.with_snapshot || !subscribe_snapshot_optimization {
498                        SnapshotMode::Include
499                    } else {
500                        compute_state.metrics.inc_subscribe_snapshot_optimization();
501                        SnapshotMode::Exclude
502                    };
503                    context.import_index(
504                        scope,
505                        compute_state,
506                        &mut tokens,
507                        input_probe,
508                        *idx_id,
509                        &idx.desc.as_lir(),
510                        &idx.typ,
511                        snapshot_mode,
512                        start_signal.clone(),
513                    );
514                }
515
516                // Build declared objects.
517                for object in dataflow.objects_to_build {
518                    let bundle = context.scope.clone().region_named(
519                        &format!("BuildingObject({:?})", object.id),
520                        |region| {
521                            let depends = object.plan.depends();
522                            context
523                                .enter_region(region, Some(&depends))
524                                .render_plan(object.id, object.plan)
525                                .leave_region(context.scope)
526                        },
527                    );
528                    let global_id = object.id;
529                    context.log_dataflow_global_id(
530                        *bundle
531                            .scope()
532                            .addr()
533                            .first()
534                            .expect("Dataflow root id must exist"),
535                        global_id,
536                    );
537                    context.insert_id(Id::Global(object.id), bundle);
538                }
539
540                // Export declared indexes.
541                for (idx_id, dependencies, idx) in indexes {
542                    context.export_index(
543                        compute_state,
544                        &tokens,
545                        dependencies,
546                        idx_id,
547                        &idx,
548                        &output_probe,
549                    );
550                }
551
552                // Export declared sinks.
553                for (sink_id, dependencies, sink) in sinks {
554                    context.export_sink(
555                        compute_state,
556                        &tokens,
557                        dependencies,
558                        sink_id,
559                        &sink,
560                        start_signal.clone(),
561                        &output_probe,
562                        scope,
563                    );
564                }
565            });
566        }
567    });
568}
569
570// This implementation block allows child timestamps to vary from parent timestamps,
571// but requires the parent timestamp to be `repr::Timestamp`.
572impl<'g, T> Context<'g, T>
573where
574    T: Refines<mz_repr::Timestamp> + RenderTimestamp,
575{
576    /// Import the collection from the arrangement, discarding batches from the snapshot.
577    /// (This does not guarantee that no records from the snapshot are included; the assumption is
578    /// that we'll filter those out later if necessary.)
579    fn import_filtered_index_collection<
580        'outer,
581        Tr: TraceReader<Time = mz_repr::Timestamp, Batch: Navigable> + Clone,
582        V: Data,
583    >(
584        &self,
585        arranged: Arranged<'outer, Tr>,
586        start_signal: StartSignal,
587        mut logic: impl FnMut(BatchKey<'_, Tr>, BatchVal<'_, Tr>) -> V + 'static,
588    ) -> VecCollection<'g, T, V, BatchDiff<Tr>>
589    where
590        // This is implied by the fact that the outer timestamp = mz_repr::Timestamp, but it's essential
591        // for our batch-level filtering to be safe, so we document it here regardless.
592        mz_repr::Timestamp: TotalOrder,
593        BatchCursor<Tr>: Cursor<Time = mz_repr::Timestamp>,
594    {
595        let oks = arranged.stream.with_start_signal(start_signal).filter({
596            let as_of = self.as_of_frontier.clone();
597            move |b| !<Antichain<mz_repr::Timestamp> as PartialOrder>::less_equal(b.upper(), &as_of)
598        });
599        Arranged::<'outer, Tr>::flat_map_batches(oks, move |a, b| [logic(a, b)]).enter(self.scope)
600    }
601
602    /// Extracts a filtered index's contents onto the columnar collection edge, discarding
603    /// batches that the snapshot covers.
604    ///
605    /// `logic` packs each `(key, val)` pair into the row buffer it is handed. The buffer is
606    /// reused across records and pushed borrowed, so a pair that carries many times costs one
607    /// pack rather than an owned [`Row`] per time.
608    ///
609    /// The `TotalOrder` bound is what makes discarding a batch by its upper safe.
610    fn import_filtered_index_edge<'outer, Tr>(
611        &self,
612        arranged: Arranged<'outer, Tr>,
613        start_signal: StartSignal,
614        mut logic: impl FnMut(BatchKey<'_, Tr>, BatchVal<'_, Tr>, &mut Row) + 'static,
615    ) -> ColCollection<'g, T>
616    where
617        Tr: TraceReader<Time = mz_repr::Timestamp, Batch: Navigable> + Clone,
618        mz_repr::Timestamp: TotalOrder,
619        BatchCursor<Tr>: Cursor<Time = mz_repr::Timestamp, Diff = Diff>,
620    {
621        let as_of = self.as_of_frontier.clone();
622        arranged
623            .stream
624            .with_start_signal(start_signal)
625            .unary::<ColumnBuilder<(Row, mz_repr::Timestamp, Diff)>, _, _, _>(
626                Pipeline,
627                "IndexToColumnar",
628                move |_cap, _info| {
629                    let mut row_buf = Row::default();
630                    move |input, output| {
631                        input.for_each(|time, data| {
632                            let mut session = output.session_with_builder(&time);
633                            for batch in data.iter() {
634                                if <Antichain<mz_repr::Timestamp> as PartialOrder>::less_equal(
635                                    batch.upper(),
636                                    &as_of,
637                                ) {
638                                    continue;
639                                }
640                                let mut cursor = batch.cursor();
641                                while let Some(key) = cursor.get_key(batch) {
642                                    while let Some(val) = cursor.get_val(batch) {
643                                        logic(key, val, &mut row_buf);
644                                        cursor.map_times(batch, |t, d| {
645                                            let t = <BatchCursor<Tr> as Cursor>::owned_time(t);
646                                            let d = <BatchCursor<Tr> as Cursor>::owned_diff(d);
647                                            session.give((&*row_buf, &t, &d));
648                                        });
649                                        cursor.step_val(batch);
650                                    }
651                                    cursor.step_key(batch);
652                                }
653                            }
654                        });
655                    }
656                },
657            )
658            .as_collection()
659            .enter(self.scope)
660    }
661
662    pub(crate) fn import_index<'outer>(
663        &mut self,
664        outer: Scope<'outer, mz_repr::Timestamp>,
665        compute_state: &mut ComputeState,
666        tokens: &mut BTreeMap<GlobalId, Rc<dyn std::any::Any>>,
667        input_probe: probe::Handle<mz_repr::Timestamp>,
668        idx_id: GlobalId,
669        idx: &IndexDesc<LirScalarExpr>,
670        typ: &ReprRelationType,
671        snapshot_mode: SnapshotMode,
672        start_signal: StartSignal,
673    ) {
674        if let Some(traces) = compute_state.traces.get_mut(&idx_id) {
675            assert!(
676                PartialOrder::less_equal(&traces.compaction_frontier(), &self.as_of_frontier),
677                "Index {idx_id} has been allowed to compact beyond the dataflow as_of"
678            );
679
680            let token = traces.to_drop().clone();
681
682            let (mut oks, ok_button) = traces.oks_mut().import_frontier_core(
683                outer,
684                &format!("Index({}, {:?})", idx.on_id, idx.key),
685                self.as_of_frontier.clone(),
686                self.until.clone(),
687            );
688
689            oks.stream = oks.stream.probe_with(&input_probe);
690
691            let (err_arranged, err_button) = traces.errs_mut().import_frontier_core(
692                outer,
693                &format!("ErrIndex({}, {:?})", idx.on_id, idx.key),
694                self.as_of_frontier.clone(),
695                self.until.clone(),
696            );
697
698            let bundle = match snapshot_mode {
699                SnapshotMode::Include => {
700                    let ok_arranged = oks
701                        .enter(self.scope)
702                        .with_start_signal(start_signal.clone());
703                    let err_arranged = err_arranged
704                        .enter(self.scope)
705                        .with_start_signal(start_signal);
706                    CollectionBundle::from_expressions(
707                        idx.key.clone(),
708                        ArrangementFlavor::Trace(idx_id, ok_arranged, err_arranged),
709                    )
710                }
711                SnapshotMode::Exclude => {
712                    // When we import an index without a snapshot, we have two balancing considerations:
713                    // - It's easy to filter out irrelevant batches from the stream, but hard to filter them out from an arrangement.
714                    //   (The `TraceFrontier` wrapper allows us to set an "until" frontier, but not a lower.)
715                    // - We do not actually need to reference the arrangement in this dataflow, since all operators that use the arrangement
716                    //   (joins, reduces, etc.) also require the snapshot data.
717                    // So: when the snapshot is excluded, we import only the (filtered) collection itself and ignore the arrangement.
718                    let oks = {
719                        let mut datums = DatumVec::new();
720                        let (permutation, _thinning) =
721                            permutation_for_arrangement(&idx.key, typ.arity());
722                        self.import_filtered_index_edge(
723                            oks,
724                            start_signal.clone(),
725                            move |k: DatumSeq, v: DatumSeq, row: &mut Row| {
726                                let temp_storage = RowArena::new();
727                                let mut datums_borrow = datums.borrow();
728                                k.extend_datums(&temp_storage, &mut datums_borrow, None);
729                                v.extend_datums(&temp_storage, &mut datums_borrow, None);
730                                row.packer()
731                                    .extend(permutation.iter().map(|i| datums_borrow[*i]));
732                            },
733                        )
734                    };
735                    let errs = self.import_filtered_index_collection(
736                        err_arranged,
737                        start_signal,
738                        |e, _| e.clone(),
739                    );
740                    CollectionBundle::from_edge(oks, errs)
741                }
742            };
743            self.update_id(Id::Global(idx.on_id), bundle);
744            tokens.insert(
745                idx_id,
746                Rc::new((PressOnDrop(ok_button), PressOnDrop(err_button), token)),
747            );
748        } else {
749            panic!(
750                "import of index {} failed while building dataflow {}",
751                idx_id, self.dataflow_id
752            );
753        }
754    }
755}
756
757// This implementation block requires the scopes have the same timestamp as the trace manager.
758// That makes some sense, because we are hoping to deposit an arrangement in the trace manager.
759impl<'g> Context<'g, mz_repr::Timestamp> {
760    pub(crate) fn export_index(
761        &self,
762        compute_state: &mut ComputeState,
763        tokens: &BTreeMap<GlobalId, Rc<dyn std::any::Any>>,
764        dependency_ids: BTreeSet<GlobalId>,
765        idx_id: GlobalId,
766        idx: &IndexDesc<LirScalarExpr>,
767        output_probe: &MzProbeHandle<mz_repr::Timestamp>,
768    ) {
769        // put together tokens that belong to the export
770        let mut needed_tokens = Vec::new();
771        for dep_id in dependency_ids {
772            if let Some(token) = tokens.get(&dep_id) {
773                needed_tokens.push(Rc::clone(token));
774            }
775        }
776        let bundle = self.lookup_id(Id::Global(idx_id)).unwrap_or_else(|| {
777            panic!(
778                "Arrangement alarmingly absent! id: {:?}",
779                Id::Global(idx_id)
780            )
781        });
782
783        let key = &idx.key;
784        match bundle.arrangement(key) {
785            Some(ArrangementFlavor::Local(mut oks, mut errs)) => {
786                // NOTE: Do not give an exported arrangement a second reader that holds a trace
787                // handle, such as a `reduce`. Such a reader pins the shared spine's physical
788                // frontier at its own lagging progress, and `ArrangementManager::maintenance` can
789                // then no longer advance it, so batches pile up in `Spine::pending`. A cursor is
790                // only checked for straddling over pending batches, so an importing dataflow's
791                // `cursor_through` eventually panics with `upper` straddles batch. Watching
792                // `errs.stream` in `output_probe` does not help, and neither does discarding the
793                // reader's output. Stream-level readers like `as_collection` are unaffected. This is
794                // why error multiplicity is not collapsed here, leaving multiplicity that crosses
795                // an index boundary unbounded. TODO(CPU-209): bound it without a trace reader.
796
797                // Ensure that the frontier does not advance past the expiration time, if set.
798                // Otherwise, we might write down incorrect data.
799                if let Some(&expiration) = self.dataflow_expiration.as_option() {
800                    oks.stream = oks.stream.expire_stream_at(
801                        &format!("{}_export_index_oks", self.debug_name),
802                        expiration,
803                    );
804                    errs.stream = errs.stream.expire_stream_at(
805                        &format!("{}_export_index_errs", self.debug_name),
806                        expiration,
807                    );
808                }
809
810                oks.stream = oks.stream.probe_notify_with(vec![output_probe.clone()]);
811
812                // Attach logging of dataflow errors.
813                if let Some(logger) = compute_state.compute_logger.clone() {
814                    errs.stream = errs.stream.log_dataflow_errors(logger, idx_id);
815                }
816
817                compute_state.traces.set(
818                    idx_id,
819                    TraceBundle::new(oks.trace, errs.trace).with_drop(needed_tokens),
820                );
821            }
822            Some(ArrangementFlavor::Trace(gid, _, _)) => {
823                // Duplicate of existing arrangement with id `gid`, so
824                // just create another handle to that arrangement.
825                let trace = compute_state.traces.get(&gid).unwrap().clone();
826                compute_state.traces.set(idx_id, trace);
827            }
828            None => {
829                println!("collection available: {:?}", bundle.collection.is_none());
830                println!(
831                    "keys available: {:?}",
832                    bundle.arranged.keys().collect::<Vec<_>>()
833                );
834                panic!(
835                    "Arrangement alarmingly absent! id: {:?}, keys: {:?}",
836                    Id::Global(idx_id),
837                    key
838                );
839            }
840        };
841    }
842}
843
844// This implementation block requires the scopes have the same timestamp as the trace manager.
845// That makes some sense, because we are hoping to deposit an arrangement in the trace manager.
846impl<'g, T> Context<'g, T>
847where
848    T: RenderTimestamp,
849{
850    pub(crate) fn export_index_iterative<'outer>(
851        &self,
852        outer: Scope<'outer, mz_repr::Timestamp>,
853        compute_state: &mut ComputeState,
854        tokens: &BTreeMap<GlobalId, Rc<dyn std::any::Any>>,
855        dependency_ids: BTreeSet<GlobalId>,
856        idx_id: GlobalId,
857        idx: &IndexDesc<LirScalarExpr>,
858        output_probe: &MzProbeHandle<mz_repr::Timestamp>,
859    ) {
860        // put together tokens that belong to the export
861        let mut needed_tokens = Vec::new();
862        for dep_id in dependency_ids {
863            if let Some(token) = tokens.get(&dep_id) {
864                needed_tokens.push(Rc::clone(token));
865            }
866        }
867        let bundle = self.lookup_id(Id::Global(idx_id)).unwrap_or_else(|| {
868            panic!(
869                "Arrangement alarmingly absent! id: {:?}",
870                Id::Global(idx_id)
871            )
872        });
873
874        let key = &idx.key;
875        match bundle.arrangement(key) {
876            Some(ArrangementFlavor::Local(oks, errs)) => {
877                // TODO: The following as_collection/leave/arrange sequence could be optimized.
878                //   * Combine as_collection and leave into a single function.
879                //   * Use columnar to extract columns from the batches to implement leave.
880                let mut oks = oks
881                    .as_collection(|k, v| (k.to_row(), v.to_row()))
882                    .leave(outer)
883                    .mz_arrange::<
884                        ColumnationChunker<_>,
885                        RowRowBatcher<_, _>,
886                        RowRowBuilder<_, _>,
887                        _,
888                    >(
889                        "Arrange export iterative",
890                    );
891
892                let mut errs = errs
893                    .as_collection(|k, v| (k.clone(), v.clone()))
894                    .leave(outer)
895                    .mz_arrange::<ColumnationChunker<_>, ErrBatcher<_, _>, ErrBuilder<_, _>, _>(
896                        "Arrange export iterative err",
897                    );
898
899                // Ensure that the frontier does not advance past the expiration time, if set.
900                // Otherwise, we might write down incorrect data.
901                if let Some(&expiration) = self.dataflow_expiration.as_option() {
902                    oks.stream = oks.stream.expire_stream_at(
903                        &format!("{}_export_index_iterative_oks", self.debug_name),
904                        expiration,
905                    );
906                    errs.stream = errs.stream.expire_stream_at(
907                        &format!("{}_export_index_iterative_err", self.debug_name),
908                        expiration,
909                    );
910                }
911
912                oks.stream = oks.stream.probe_notify_with(vec![output_probe.clone()]);
913
914                // Attach logging of dataflow errors.
915                if let Some(logger) = compute_state.compute_logger.clone() {
916                    errs.stream = errs.stream.log_dataflow_errors(logger, idx_id);
917                }
918
919                compute_state.traces.set(
920                    idx_id,
921                    TraceBundle::new(oks.trace, errs.trace).with_drop(needed_tokens),
922                );
923            }
924            Some(ArrangementFlavor::Trace(gid, _, _)) => {
925                // Duplicate of existing arrangement with id `gid`, so
926                // just create another handle to that arrangement.
927                let trace = compute_state.traces.get(&gid).unwrap().clone();
928                compute_state.traces.set(idx_id, trace);
929            }
930            None => {
931                println!("collection available: {:?}", bundle.collection.is_none());
932                println!(
933                    "keys available: {:?}",
934                    bundle.arranged.keys().collect::<Vec<_>>()
935                );
936                panic!(
937                    "Arrangement alarmingly absent! id: {:?}, keys: {:?}",
938                    Id::Global(idx_id),
939                    key,
940                );
941            }
942        };
943    }
944}
945
946/// Information about bindings, tracked in `render_recursive_plan` and
947/// `render_plan`, to be passed to `render_letfree_plan`.
948///
949/// `render_letfree_plan` uses these to produce nice output (e.g., `With ...
950/// Returning ...`) for local bindings in the `mz_lir_mapping` output.
951enum BindingInfo {
952    Body { in_let: bool },
953    Let { id: LocalId, last: bool },
954    LetRec { id: LocalId, last: bool },
955}
956
957impl<'scope> Context<'scope, Product<mz_repr::Timestamp, PointStamp<u64>>> {
958    /// Renders a plan to a differential dataflow, producing the collection of results.
959    ///
960    /// This method allows for `plan` to contain [`RecBind`]s, and is planned
961    /// in the context of `level` pre-existing iteration coordinates.
962    ///
963    /// This method recursively descends [`RecBind`] values, establishing nested scopes for each
964    /// and establishing the appropriate recursive dependencies among the bound variables.
965    /// Once all [`RecBind`]s have been rendered it calls in to `render_plan` which will error if
966    /// further [`RecBind`]s are found.
967    ///
968    /// The method requires that all variables conclude with a physical representation that
969    /// contains a collection (i.e. a non-arrangement), and it will panic otherwise.
970    fn render_recursive_plan(
971        &mut self,
972        object_id: GlobalId,
973        level: usize,
974        plan: RenderPlan,
975        binding: BindingInfo,
976    ) -> CollectionBundle<'scope, Product<mz_repr::Timestamp, PointStamp<u64>>> {
977        for BindStage { lets, recs } in plan.binds {
978            // Render the let bindings in order.
979            let mut let_iter = lets.into_iter().peekable();
980            while let Some(LetBind { id, value }) = let_iter.next() {
981                let bundle =
982                    self.scope
983                        .clone()
984                        .region_named(&format!("Binding({:?})", id), |region| {
985                            let depends = value.depends();
986                            let last = let_iter.peek().is_none();
987                            let binding = BindingInfo::Let { id, last };
988                            self.enter_region(region, Some(&depends))
989                                .render_letfree_plan(object_id, value, binding)
990                                .leave_region(self.scope)
991                        });
992                let bundle = self.distinct_binding_errs(bundle);
993                self.insert_id(Id::Local(id), bundle);
994            }
995
996            let rec_ids: Vec<_> = recs.iter().map(|r| r.id).collect();
997
998            // A binding's `Variable` serves the `Get`s rendered before the rec
999            // loop binds the real value, which are exactly the values of
1000            // `recs[0..=i]`. A binding no such value reads has no use for a
1001            // `Variable`-backed bundle, and installing one would build a
1002            // re-encode that repacks the whole collection once per iteration
1003            // with nothing to consume it.
1004            let mut variable_read = BTreeSet::new();
1005            let mut read_so_far = BTreeSet::new();
1006            for rec in recs.iter() {
1007                read_so_far.extend(rec.value.depends());
1008                if read_so_far.contains(&Id::Local(rec.id)) {
1009                    variable_read.insert(rec.id);
1010                }
1011            }
1012
1013            // Define variables for rec bindings.
1014            // It is important that we only use the `Variable` until the object is bound.
1015            // At that point, all subsequent uses should have access to the object itself.
1016            let mut variables = BTreeMap::new();
1017            for id in rec_ids.iter() {
1018                use differential_dataflow::dynamic::feedback_summary;
1019                let inner = feedback_summary::<u64>(level + 1, 1);
1020                let (oks_v, oks_collection) =
1021                    Variable::new(self.scope, Product::new(Default::default(), inner.clone()));
1022                let (err_v, err_collection) =
1023                    Variable::new(self.scope, Product::new(Default::default(), inner));
1024
1025                if variable_read.contains(id) {
1026                    self.insert_id(
1027                        Id::Local(*id),
1028                        CollectionBundle::from_edge(oks_collection, err_collection),
1029                    );
1030                }
1031                variables.insert(Id::Local(*id), (oks_v, err_v));
1032            }
1033            // The bound value is kept so the extraction below reuses it rather than
1034            // reaching back into a bundle that forward reads have since collapsed.
1035            let mut bound_oks = BTreeMap::new();
1036            let mut rec_iter = recs.into_iter().peekable();
1037            while let Some(RecBind { id, value, limit }) = rec_iter.next() {
1038                let last = rec_iter.peek().is_none();
1039                let binding = BindingInfo::LetRec { id, last };
1040                let bundle = self.render_recursive_plan(object_id, level + 1, value, binding);
1041                // We need to ensure that the raw collection exists, but do not have enough information
1042                // here to cause that to happen.
1043                let (oks, mut err) = bundle.collection.clone().unwrap();
1044                bound_oks.insert(id, oks.clone());
1045                // Collapses what forward reads see. `err_v` below feeds reads rendered before this
1046                // binding and is collapsed separately; without this, a `Get` in a later rec binding
1047                // or in the body resolves to the bundle stored here and compounds level over level,
1048                // which is exactly what the collapse prevents for non-recursive bindings.
1049                let bundle = self.distinct_binding_errs(bundle);
1050                self.insert_id(Id::Local(id), bundle);
1051                let (oks_v, err_v) = variables.remove(&Id::Local(id)).unwrap();
1052
1053                // Set oks variable to `oks` but consolidated to ensure iteration ceases at fixed point.
1054                let mut oks = columnar_consolidate(oks, "LetRecConsolidation");
1055
1056                if let Some(limit) = limit {
1057                    // We swallow the results of the `max_iter`th iteration, because
1058                    // these results would go into the `max_iter + 1`th iteration.
1059                    //
1060                    // The split consults each update's own time. A container's capability
1061                    // is only a lower bound on the times it carries, so deciding per
1062                    // container by capability could retain updates from an iteration
1063                    // beyond the one the capability names.
1064                    //
1065                    // The predicate reads the time out of the column into `time`, whose
1066                    // coordinate vector it reuses across records.
1067                    let mut time = RecTimestamp::default();
1068                    let (over_limit, in_limit) = oks
1069                        .inner
1070                        .partition_by::<ColumnBuilder<(Row, RecTimestamp, Diff)>, _>(
1071                            "LetRecLimit",
1072                            move |(_data, reference, _diff)| {
1073                                time.copy_from(*reference);
1074                                // The iteration number, or zero if absent, since trailing zero
1075                                // coordinates are truncated.
1076                                let iteration_index = time.inner.get(level).copied().unwrap_or(0);
1077                                // The pointstamp starts counting from 0, so we need to add 1.
1078                                iteration_index + 1 >= limit.max_iters.into()
1079                            },
1080                        );
1081                    oks = Collection::new(in_limit);
1082                    if !limit.return_at_limit {
1083                        // One error per record that passed the limit, replacing the record.
1084                        // `flat_map_datums` reads the record without decoding it, and the ok
1085                        // side stays empty, so the builder it names never ships a container.
1086                        let (_, over_limit) = flat_map_datums::<
1087                            _,
1088                            CapacityContainerBuilder<Vec<(Row, RecTimestamp, Diff)>>,
1089                            _,
1090                        >(
1091                            Collection::new(over_limit),
1092                            "LetRecLimitExceeded",
1093                            0,
1094                            move |_datums, time, diff, _ok_session, err_session| {
1095                                err_session.give((
1096                                    DataflowErrorSer::from(EvalError::LetRecLimitExceeded(
1097                                        format!("{}", limit.max_iters.get()).into(),
1098                                    )),
1099                                    time,
1100                                    diff,
1101                                ));
1102                                1
1103                            },
1104                        );
1105                        err = err.concat(over_limit.as_collection());
1106                    }
1107                }
1108
1109                // Set err variable to the distinct elements of `err`.
1110                // Distinctness is important, as we otherwise might add the same error each iteration,
1111                // say if the limit of `oks` has an error. This would result in non-terminating rather
1112                // than a clean report of the error. The trade-off is that we lose information about
1113                // multiplicities of errors, but .. this seems to be the better call.
1114                let err: KeyCollection<_, _, _> = err.into();
1115                let errs = err
1116                    .mz_arrange::<
1117                        ColumnationChunker<_>,
1118                        ErrBatcher<_, _>,
1119                        ErrBuilder<_, _>,
1120                        ErrSpine<_, _>,
1121                    >("Arrange recursive err")
1122                    .mz_reduce_abelian::<_, ErrBuilder<_, _>, ErrSpine<_, _>, _>(
1123                        "Distinct recursive err",
1124                        move |_k, _s, t| t.push(((), Diff::ONE)),
1125                    )
1126                    .as_collection(|k, _| k.clone());
1127
1128                oks_v.set(oks);
1129                err_v.set(errs);
1130            }
1131            // Now extract each of the rec bindings into the outer scope.
1132            for id in rec_ids.into_iter() {
1133                let bundle = self.remove_id(Id::Local(id)).unwrap();
1134                let (_, err) = bundle.collection.unwrap();
1135                let oks = bound_oks
1136                    .remove(&id)
1137                    .expect("rec binding bound while rendering above");
1138                self.insert_id(
1139                    Id::Local(id),
1140                    CollectionBundle::from_edge(
1141                        columnar_leave_dynamic(oks, level + 1),
1142                        err.leave_dynamic(level + 1),
1143                    ),
1144                );
1145            }
1146        }
1147
1148        self.render_letfree_plan(object_id, plan.body, binding)
1149    }
1150}
1151
1152impl<'scope, T: RenderTimestamp + MaybeBucketByTime> Context<'scope, T> {
1153    /// Renders a non-recursive plan to a differential dataflow, producing the collection of
1154    /// results.
1155    ///
1156    /// The return type reflects the uncertainty about the data representation, perhaps
1157    /// as a stream of data, perhaps as an arrangement, perhaps as a stream of batches.
1158    ///
1159    /// # Panics
1160    ///
1161    /// Panics if the given plan contains any [`RecBind`]s. Recursive plans must be rendered using
1162    /// `render_recursive_plan` instead.
1163    fn render_plan(
1164        &mut self,
1165        object_id: GlobalId,
1166        plan: RenderPlan,
1167    ) -> CollectionBundle<'scope, T> {
1168        let mut in_let = false;
1169        for BindStage { lets, recs } in plan.binds {
1170            assert!(recs.is_empty());
1171
1172            let mut let_iter = lets.into_iter().peekable();
1173            while let Some(LetBind { id, value }) = let_iter.next() {
1174                // if we encounter a single let, the body is in a let
1175                in_let = true;
1176                let bundle =
1177                    self.scope
1178                        .clone()
1179                        .region_named(&format!("Binding({:?})", id), |region| {
1180                            let depends = value.depends();
1181                            let last = let_iter.peek().is_none();
1182                            let binding = BindingInfo::Let { id, last };
1183                            self.enter_region(region, Some(&depends))
1184                                .render_letfree_plan(object_id, value, binding)
1185                                .leave_region(self.scope)
1186                        });
1187                let bundle = self.distinct_binding_errs(bundle);
1188                self.insert_id(Id::Local(id), bundle);
1189            }
1190        }
1191
1192        self.scope.clone().region_named("Main Body", |region| {
1193            let depends = plan.body.depends();
1194            self.enter_region(region, Some(&depends))
1195                .render_letfree_plan(object_id, plan.body, BindingInfo::Body { in_let })
1196                .leave_region(self.scope)
1197        })
1198    }
1199
1200    /// Collapses a binding's error multiplicities.
1201    ///
1202    /// Applied to every binding, not only the multiply-read ones. Gating on the reference count
1203    /// would be a pure optimization, since collapsing a binding one `Get` reads is harmless, and
1204    /// there is almost nothing to gate: `NormalizeLets` inlines single-use bindings, so the ones
1205    /// reaching rendering are shared. See [`CollectionBundle::distinct_errs`] for why the collapse
1206    /// is needed at all, and why a binding's definition is the place for it rather than the
1207    /// multi-input operators where the duplicate copies happen to meet again.
1208    fn distinct_binding_errs(
1209        &self,
1210        bundle: CollectionBundle<'scope, T>,
1211    ) -> CollectionBundle<'scope, T> {
1212        if ENABLE_ERROR_DISTINCT.get(&self.config_set) {
1213            bundle.distinct_errs()
1214        } else {
1215            bundle
1216        }
1217    }
1218
1219    /// Renders a let-free plan to a differential dataflow, producing the collection of results.
1220    fn render_letfree_plan(
1221        &self,
1222        object_id: GlobalId,
1223        plan: LetFreePlan,
1224        binding: BindingInfo,
1225    ) -> CollectionBundle<'scope, T> {
1226        let (mut nodes, root_id, topological_order) = plan.destruct();
1227
1228        // Rendered collections by their `LirId`.
1229        let mut collections = BTreeMap::new();
1230
1231        // Mappings to send along.
1232        // To save overhead, we'll only compute mappings when we need to,
1233        // which means things get gated behind options. Unfortunately, that means we
1234        // have several `Option<...>` types that are _all_ `Some` or `None` together,
1235        // but there's no convenient way to express the invariant.
1236        let should_compute_lir_metadata = self.compute_logger.is_some();
1237        let mut lir_mapping_metadata = if should_compute_lir_metadata {
1238            Some(Vec::with_capacity(nodes.len()))
1239        } else {
1240            None
1241        };
1242
1243        let mut topo_iter = topological_order.into_iter().peekable();
1244        while let Some(lir_id) = topo_iter.next() {
1245            let node = nodes.remove(&lir_id).unwrap();
1246
1247            // TODO(mgree) need ExprHumanizer in DataflowDescription to get nice column names
1248            // ActiveComputeState can't have a catalog reference, so we'll need to capture the names
1249            // in some other structure and have that structure impl ExprHumanizer
1250            let metadata = if should_compute_lir_metadata {
1251                let operator = node.expr.humanize(&DummyHumanizer);
1252
1253                // mark the last operator in topo order with any binding decoration
1254                let operator = if topo_iter.peek().is_none() {
1255                    match &binding {
1256                        BindingInfo::Body { in_let: true } => format!("Returning {operator}"),
1257                        BindingInfo::Body { in_let: false } => operator,
1258                        BindingInfo::Let { id, last: true } => {
1259                            format!("With {id} = {operator}")
1260                        }
1261                        BindingInfo::Let { id, last: false } => {
1262                            format!("{id} = {operator}")
1263                        }
1264                        BindingInfo::LetRec { id, last: true } => {
1265                            format!("With Recursive {id} = {operator}")
1266                        }
1267                        BindingInfo::LetRec { id, last: false } => {
1268                            format!("{id} = {operator}")
1269                        }
1270                    }
1271                } else {
1272                    operator
1273                };
1274
1275                let operator_id_start = self.scope.worker().peek_identifier();
1276                Some((operator, operator_id_start))
1277            } else {
1278                None
1279            };
1280
1281            let mut bundle = self.render_plan_expr(node.expr, &collections);
1282
1283            if let Some((operator, operator_id_start)) = metadata {
1284                let operator_id_end = self.scope.worker().peek_identifier();
1285                let operator_span = (operator_id_start, operator_id_end);
1286
1287                if let Some(lir_mapping_metadata) = &mut lir_mapping_metadata {
1288                    lir_mapping_metadata.push((
1289                        lir_id,
1290                        LirMetadata::new(operator, node.parent, node.nesting, operator_span),
1291                    ))
1292                }
1293            }
1294
1295            self.log_operator_hydration(&mut bundle, lir_id);
1296
1297            collections.insert(lir_id, bundle);
1298        }
1299
1300        if let Some(lir_mapping_metadata) = lir_mapping_metadata {
1301            self.log_lir_mapping(object_id, lir_mapping_metadata);
1302        }
1303
1304        collections
1305            .remove(&root_id)
1306            .expect("LetFreePlan invariant (1)")
1307    }
1308
1309    /// Renders a [`render_plan::Expr`], producing the collection of results.
1310    ///
1311    /// # Panics
1312    ///
1313    /// Panics if any of the expr's inputs is not found in `collections`.
1314    /// Callers must ensure that input nodes have been rendered previously.
1315    fn render_plan_expr(
1316        &self,
1317        expr: render_plan::Expr,
1318        collections: &BTreeMap<LirId, CollectionBundle<'scope, T>>,
1319    ) -> CollectionBundle<'scope, T> {
1320        use render_plan::Expr::*;
1321
1322        let expect_input = |id| {
1323            collections
1324                .get(&id)
1325                .cloned()
1326                .unwrap_or_else(|| panic!("missing input collection: {id}"))
1327        };
1328
1329        match expr {
1330            Constant { rows } => {
1331                // Produce both rows and errs to avoid conditional dataflow construction.
1332                let (rows, errs) = match rows {
1333                    Ok(rows) => (rows, Vec::new()),
1334                    Err(e) => (Vec::new(), vec![e]),
1335                };
1336
1337                // We should advance times in constant collections to start from `as_of`.
1338                let as_of_frontier = self.as_of_frontier.clone();
1339                let until = self.until.clone();
1340                // Advancing times to `as_of` can collapse distinct times onto one, so
1341                // rows the planner left distinct can become duplicates. The
1342                // `ConsolidatingColumnBuilder` folds those within the batch.
1343                let ok_collection = rows
1344                    .into_iter()
1345                    .filter_map(move |(row, mut time, diff)| {
1346                        time.advance_by(as_of_frontier.borrow());
1347                        if !until.less_equal(&time) {
1348                            Some((
1349                                row.0,
1350                                <T as Refines<mz_repr::Timestamp>>::to_inner(time),
1351                                diff,
1352                            ))
1353                        } else {
1354                            None
1355                        }
1356                    })
1357                    .to_stream_with_builder::<_, ConsolidatingColumnBuilder<Row, T, Diff>>(
1358                        self.scope,
1359                    )
1360                    .as_collection();
1361
1362                let mut error_time: mz_repr::Timestamp = Timestamp::minimum();
1363                error_time.advance_by(self.as_of_frontier.borrow());
1364                let err_collection = errs
1365                    .into_iter()
1366                    .map(move |e| {
1367                        (
1368                            DataflowErrorSer::from(e),
1369                            <T as Refines<mz_repr::Timestamp>>::to_inner(error_time),
1370                            Diff::ONE,
1371                        )
1372                    })
1373                    .to_stream(self.scope)
1374                    .as_collection();
1375
1376                CollectionBundle::from_edge(ok_collection, err_collection)
1377            }
1378            Get { id, keys, plan } => {
1379                // Recover the collection from `self` and then apply `mfp` to it.
1380                // If `mfp` happens to be trivial, we can just return the collection.
1381                let mut collection = self
1382                    .lookup_id(id)
1383                    .unwrap_or_else(|| panic!("Get({:?}) not found at render time", id));
1384                match plan {
1385                    mz_compute_types::plan::GetPlan::PassArrangements => {
1386                        // Assert that each of `keys` are present in `collection`.
1387                        assert!(
1388                            keys.arranged
1389                                .iter()
1390                                .all(|(key, _, _)| collection.arranged.contains_key(key))
1391                        );
1392                        assert!(keys.raw <= collection.collection.is_some());
1393                        // Retain only those keys we want to import.
1394                        collection.arranged.retain(|key, _value| {
1395                            keys.arranged.iter().any(|(key2, _, _)| key2 == key)
1396                        });
1397                        collection
1398                    }
1399                    mz_compute_types::plan::GetPlan::Arrangement(key, row, mfp) => {
1400                        let (oks, errs) = collection.as_collection_core(
1401                            mfp,
1402                            Some((key, row)),
1403                            self.until.clone(),
1404                        );
1405                        CollectionBundle::from_edge(oks, errs)
1406                    }
1407                    mz_compute_types::plan::GetPlan::Collection(mfp) => {
1408                        let (oks, errs) =
1409                            collection.as_collection_core(mfp, None, self.until.clone());
1410                        CollectionBundle::from_edge(oks, errs)
1411                    }
1412                }
1413            }
1414            Mfp {
1415                input,
1416                mfp,
1417                input_key_val,
1418            } => {
1419                let input = expect_input(input);
1420                // If `mfp` is non-trivial, we should apply it and produce a collection.
1421                if mfp.is_identity() {
1422                    input
1423                } else {
1424                    let (oks, errs) =
1425                        input.as_collection_core(mfp, input_key_val, self.until.clone());
1426                    CollectionBundle::from_edge(oks, errs)
1427                }
1428            }
1429            FlatMap {
1430                input_key,
1431                input,
1432                exprs,
1433                func,
1434                mfp_after: mfp,
1435            } => {
1436                let input = expect_input(input);
1437                self.render_flat_map(input_key, input, exprs, func, mfp)
1438            }
1439            Join { inputs, plan } => {
1440                let inputs = inputs.into_iter().map(expect_input).collect();
1441                match plan {
1442                    mz_compute_types::plan::join::JoinPlan::Linear(linear_plan) => {
1443                        self.render_join(inputs, linear_plan)
1444                    }
1445                    mz_compute_types::plan::join::JoinPlan::Delta(delta_plan) => {
1446                        self.render_delta_join(inputs, delta_plan)
1447                    }
1448                }
1449            }
1450            Reduce {
1451                input_key,
1452                input,
1453                key_val_plan,
1454                plan,
1455                mfp_after,
1456                temporal_bucketing_strategy,
1457            } => {
1458                let input = expect_input(input);
1459                let mfp_option = (!mfp_after.is_identity()).then_some(mfp_after);
1460                self.render_reduce(
1461                    input_key,
1462                    input,
1463                    key_val_plan,
1464                    plan,
1465                    mfp_option,
1466                    temporal_bucketing_strategy,
1467                )
1468            }
1469            TopK {
1470                input,
1471                top_k_plan,
1472                temporal_bucketing_strategy,
1473            } => {
1474                let input = expect_input(input);
1475                self.render_topk(input, top_k_plan, temporal_bucketing_strategy)
1476            }
1477            Negate { input } => {
1478                let input = expect_input(input);
1479                let (oks, errs) = input
1480                    .collection
1481                    .clone()
1482                    .expect("Negate input must be an unarranged collection");
1483                CollectionBundle::from_edge(columnar_negate(oks), errs)
1484            }
1485            Threshold {
1486                input,
1487                threshold_plan,
1488            } => {
1489                let input = expect_input(input);
1490                self.render_threshold(input, threshold_plan)
1491            }
1492            Union {
1493                inputs,
1494                consolidate_output,
1495                temporal_bucketing_strategies,
1496            } => {
1497                let mut oks = Vec::new();
1498                let mut errs = Vec::new();
1499                for (input, strategy) in inputs.into_iter().zip_eq(temporal_bucketing_strategies) {
1500                    let (os, es) = expect_input(input)
1501                        .collection
1502                        .clone()
1503                        .expect("Union input must be an unarranged collection");
1504                    // Apply per-input temporal bucketing. No-op for `Direct`.
1505                    // Only consolidating Unions carry non-`Direct` strategies;
1506                    // see the `Union` arm of `lower_mir_expr_stack_safe`.
1507                    let os = if matches!(strategy, ArrangementStrategy::TemporalBucketing)
1508                        && ENABLE_COMPUTE_TEMPORAL_BUCKETING.get(&self.config_set)
1509                    {
1510                        let summary: mz_repr::Timestamp = TEMPORAL_BUCKETING_SUMMARY
1511                            .get(&self.config_set)
1512                            .try_into()
1513                            .expect("must fit");
1514                        T::maybe_apply_temporal_bucketing(
1515                            os.inner,
1516                            self.as_of_frontier.clone(),
1517                            summary,
1518                        )
1519                    } else {
1520                        os
1521                    };
1522                    oks.push(os);
1523                    errs.push(es);
1524                }
1525                let oks = concat_many(self.scope, oks);
1526                let oks = if consolidate_output {
1527                    columnar_consolidate(oks, "UnionConsolidation")
1528                } else {
1529                    oks
1530                };
1531                let errs = differential_dataflow::collection::concatenate(self.scope, errs);
1532                CollectionBundle::from_edge(oks, errs)
1533            }
1534            ArrangeBy {
1535                input_key,
1536                input,
1537                input_mfp,
1538                forms: keys,
1539                strategy,
1540            } => {
1541                let input = expect_input(input);
1542                input.ensure_collections(
1543                    keys,
1544                    input_key,
1545                    input_mfp,
1546                    self.as_of_frontier.clone(),
1547                    self.until.clone(),
1548                    &self.config_set,
1549                    strategy,
1550                )
1551            }
1552        }
1553    }
1554
1555    fn log_dataflow_global_id(&self, dataflow_index: usize, global_id: GlobalId) {
1556        if let Some(logger) = &self.compute_logger {
1557            logger.log(&ComputeEvent::DataflowGlobal(DataflowGlobal {
1558                dataflow_index,
1559                global_id,
1560            }));
1561        }
1562    }
1563
1564    fn log_lir_mapping(&self, global_id: GlobalId, mapping: Vec<(LirId, LirMetadata)>) {
1565        if let Some(logger) = &self.compute_logger {
1566            logger.log(&ComputeEvent::LirMapping(LirMapping { global_id, mapping }));
1567        }
1568    }
1569
1570    fn log_operator_hydration(&self, bundle: &mut CollectionBundle<'scope, T>, lir_id: LirId) {
1571        // A `CollectionBundle` can contain more than one collection, which makes it not obvious to
1572        // which we should attach the logging operator.
1573        //
1574        // We could attach to each collection and track the lower bound of output frontiers.
1575        // However, that would be of limited use because we expect all collections to hydrate at
1576        // roughly the same time: The `ArrangeBy` operator is not fueled, so as soon as it sees the
1577        // frontier of the unarranged collection advance, it will perform all work necessary to
1578        // also advance its own frontier. We don't expect significant delays between frontier
1579        // advancements of the unarranged and arranged collections, so attaching the logging
1580        // operator to any one of them should produce accurate results.
1581        //
1582        // If the `CollectionBundle` contains both unarranged and arranged representations it is
1583        // beneficial to attach the logging operator to one of the arranged representation to avoid
1584        // unnecessary cloning of data. The unarranged collection feeds into the arrangements, so
1585        // if we attached the logging operator to it, we would introduce a fork in its output
1586        // stream, which would necessitate that all output data is cloned. In contrast, we can hope
1587        // that the output streams of the arrangements don't yet feed into anything else, so
1588        // attaching a (pass-through) logging operator does not introduce a fork.
1589
1590        match bundle.arranged.values_mut().next() {
1591            Some(arrangement) => {
1592                use ArrangementFlavor::*;
1593
1594                match arrangement {
1595                    Local(a, _) => {
1596                        a.stream = self.log_operator_hydration_inner(a.stream.clone(), lir_id);
1597                    }
1598                    Trace(_, a, _) => {
1599                        a.stream = self.log_operator_hydration_inner(a.stream.clone(), lir_id);
1600                    }
1601                }
1602            }
1603            None => {
1604                let (oks, _) = bundle
1605                    .collection
1606                    .as_mut()
1607                    .expect("CollectionBundle invariant");
1608                let stream = self.log_operator_hydration_inner(oks.inner.clone(), lir_id);
1609                *oks = stream.as_collection();
1610            }
1611        }
1612    }
1613
1614    fn log_operator_hydration_inner<D>(
1615        &self,
1616        stream: Stream<'scope, T, D>,
1617        lir_id: LirId,
1618    ) -> Stream<'scope, T, D>
1619    where
1620        D: timely::Container + Clone + 'static,
1621    {
1622        let Some(logger) = self.compute_logger.clone() else {
1623            return stream.clone(); // hydration logging disabled
1624        };
1625
1626        let export_ids = self.export_ids.clone();
1627
1628        // Convert the dataflow as-of into a frontier we can compare with input frontiers.
1629        //
1630        // We (somewhat arbitrarily) define operators in iterative scopes to be hydrated when their
1631        // frontier advances to an outer time that's greater than the `as_of`. Comparing
1632        // `refine(as_of) < input_frontier` would find the moment when the first iteration was
1633        // complete, which is not what we want. We want `refine(as_of + 1) <= input_frontier`
1634        // instead.
1635        let mut hydration_frontier = Antichain::new();
1636        for time in self.as_of_frontier.iter() {
1637            if let Some(time) = time.try_step_forward() {
1638                hydration_frontier.insert(Refines::to_inner(time));
1639            }
1640        }
1641
1642        let name = format!("LogOperatorHydration ({lir_id})");
1643        stream.unary_frontier(Pipeline, &name, |_cap, _info| {
1644            let mut hydrated = false;
1645
1646            for &export_id in &export_ids {
1647                logger.log(&ComputeEvent::OperatorHydration(OperatorHydration {
1648                    export_id,
1649                    lir_id,
1650                    hydrated,
1651                }));
1652            }
1653
1654            move |(input, frontier), output| {
1655                // Pass through inputs.
1656                input.for_each(|cap, data| {
1657                    output.session(&cap).give_container(data);
1658                });
1659
1660                if hydrated {
1661                    return;
1662                }
1663
1664                if PartialOrder::less_equal(&hydration_frontier.borrow(), &frontier.frontier()) {
1665                    hydrated = true;
1666
1667                    for &export_id in &export_ids {
1668                        logger.log(&ComputeEvent::OperatorHydration(OperatorHydration {
1669                            export_id,
1670                            lir_id,
1671                            hydrated,
1672                        }));
1673                    }
1674                }
1675            }
1676        })
1677    }
1678}
1679
1680#[allow(dead_code)] // Some of the methods on this trait are unused, but useful to have.
1681/// A timestamp type that can be used for operations within MZ's dataflow layer.
1682pub trait RenderTimestamp: MzTimestamp + Default + Refines<mz_repr::Timestamp> {
1683    /// The system timestamp component of the timestamp.
1684    ///
1685    /// This is useful for manipulating the system time, as when delaying
1686    /// updates for subsequent cancellation, as with monotonic reduction.
1687    fn system_time(&mut self) -> &mut mz_repr::Timestamp;
1688    /// Effects a system delay in terms of the timestamp summary.
1689    fn system_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary;
1690    /// The event timestamp component of the timestamp.
1691    fn event_time(&self) -> mz_repr::Timestamp;
1692    /// The event timestamp component of the timestamp, as a mutable reference.
1693    fn event_time_mut(&mut self) -> &mut mz_repr::Timestamp;
1694    /// Effects an event delay in terms of the timestamp summary.
1695    fn event_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary;
1696    /// Steps the timestamp back so that logical compaction to the output will
1697    /// not conflate `self` with any historical times.
1698    fn step_back(&self) -> Self;
1699}
1700
1701/// Apply temporal bucketing to a stream when the timestamp type supports it.
1702///
1703/// Sibling to [`RenderTimestamp`]: bucketing is an arrangement-time concern, not a
1704/// general property of a render timestamp, so the dispatch lives in its own trait.
1705/// Total-ordered timestamps perform real bucketing; partially-ordered timestamps
1706/// (e.g. `Product<…>` in iterative scopes) implement this as a no-op.
1707pub trait MaybeBucketByTime: Timestamp + ColumnarData {
1708    /// Buckets a columnar dataflow edge, keeping it columnar.
1709    fn maybe_apply_temporal_bucketing<'scope, D>(
1710        stream: Stream<'scope, Self, Column<(D, Self, Diff)>>,
1711        as_of: Antichain<mz_repr::Timestamp>,
1712        summary: mz_repr::Timestamp,
1713    ) -> Collection<'scope, Self, Column<(D, Self, Diff)>>
1714    where
1715        D: differential_dataflow::ExchangeData
1716            + crate::typedefs::MzData
1717            + differential_dataflow::Hashable
1718            + ColumnarData,
1719        for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1720        for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1721        for<'a> ::columnar::Ref<'a, Diff>: Ord,
1722        for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1723            ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>;
1724
1725    /// Buckets a `Vec` stream, keeping it `Vec`.
1726    ///
1727    /// For a consumer that re-encodes what it reads, where a `Vec` hands it moved
1728    /// allocations rather than copied bytes. The reduce key-value path is the one
1729    /// such caller, since its bucketed output feeds an arrangement.
1730    fn maybe_apply_temporal_bucketing_vec<'scope, D>(
1731        stream: StreamVec<'scope, Self, (D, Self, Diff)>,
1732        as_of: Antichain<mz_repr::Timestamp>,
1733        summary: mz_repr::Timestamp,
1734    ) -> VecCollection<'scope, Self, D, Diff>
1735    where
1736        D: differential_dataflow::ExchangeData
1737            + crate::typedefs::MzData
1738            + differential_dataflow::Hashable
1739            + ColumnarData,
1740        for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1741        for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1742        for<'a> ::columnar::Ref<'a, Diff>: Ord,
1743        for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1744            ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>;
1745}
1746
1747impl RenderTimestamp for mz_repr::Timestamp {
1748    fn system_time(&mut self) -> &mut mz_repr::Timestamp {
1749        self
1750    }
1751    fn system_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1752        delay
1753    }
1754    fn event_time(&self) -> mz_repr::Timestamp {
1755        *self
1756    }
1757    fn event_time_mut(&mut self) -> &mut mz_repr::Timestamp {
1758        self
1759    }
1760    fn event_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1761        delay
1762    }
1763    fn step_back(&self) -> Self {
1764        self.saturating_sub(1)
1765    }
1766}
1767
1768impl MaybeBucketByTime for mz_repr::Timestamp {
1769    fn maybe_apply_temporal_bucketing<'scope, D>(
1770        stream: Stream<'scope, Self, Column<(D, Self, Diff)>>,
1771        as_of: Antichain<mz_repr::Timestamp>,
1772        summary: mz_repr::Timestamp,
1773    ) -> Collection<'scope, Self, Column<(D, Self, Diff)>>
1774    where
1775        D: differential_dataflow::ExchangeData
1776            + crate::typedefs::MzData
1777            + differential_dataflow::Hashable
1778            + ColumnarData,
1779        for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1780        for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1781        for<'a> ::columnar::Ref<'a, Diff>: Ord,
1782        for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1783            ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1784    {
1785        stream.bucket(as_of, summary).as_collection()
1786    }
1787
1788    fn maybe_apply_temporal_bucketing_vec<'scope, D>(
1789        stream: StreamVec<'scope, Self, (D, Self, Diff)>,
1790        as_of: Antichain<mz_repr::Timestamp>,
1791        summary: mz_repr::Timestamp,
1792    ) -> VecCollection<'scope, Self, D, Diff>
1793    where
1794        D: differential_dataflow::ExchangeData
1795            + crate::typedefs::MzData
1796            + differential_dataflow::Hashable
1797            + ColumnarData,
1798        for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1799        for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1800        for<'a> ::columnar::Ref<'a, Diff>: Ord,
1801        for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1802            ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1803    {
1804        stream.bucket(as_of, summary).as_collection()
1805    }
1806}
1807
1808impl RenderTimestamp for Product<mz_repr::Timestamp, PointStamp<u64>> {
1809    fn system_time(&mut self) -> &mut mz_repr::Timestamp {
1810        &mut self.outer
1811    }
1812    fn system_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1813        Product::new(delay, Default::default())
1814    }
1815    fn event_time(&self) -> mz_repr::Timestamp {
1816        self.outer
1817    }
1818    fn event_time_mut(&mut self) -> &mut mz_repr::Timestamp {
1819        &mut self.outer
1820    }
1821    fn event_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1822        Product::new(delay, Default::default())
1823    }
1824    fn step_back(&self) -> Self {
1825        // It is necessary to step back both coordinates of a product,
1826        // and when one is a `PointStamp` that also means all coordinates
1827        // of the pointstamp.
1828        let inner = self.inner.clone();
1829        let mut vec = inner.into_inner();
1830        for item in vec.iter_mut() {
1831            *item = item.saturating_sub(1);
1832        }
1833        Product::new(self.outer.saturating_sub(1), PointStamp::new(vec))
1834    }
1835}
1836
1837impl MaybeBucketByTime for Product<mz_repr::Timestamp, PointStamp<u64>> {
1838    fn maybe_apply_temporal_bucketing<'scope, D>(
1839        stream: Stream<'scope, Self, Column<(D, Self, Diff)>>,
1840        _as_of: Antichain<mz_repr::Timestamp>,
1841        _summary: mz_repr::Timestamp,
1842    ) -> Collection<'scope, Self, Column<(D, Self, Diff)>>
1843    where
1844        D: differential_dataflow::ExchangeData
1845            + crate::typedefs::MzData
1846            + differential_dataflow::Hashable
1847            + ColumnarData,
1848        for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1849        for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1850        for<'a> ::columnar::Ref<'a, Diff>: Ord,
1851        for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1852            ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1853    {
1854        // TODO: Implement bucketing on outer timestamp for iterative scopes.
1855        stream.as_collection()
1856    }
1857
1858    fn maybe_apply_temporal_bucketing_vec<'scope, D>(
1859        stream: StreamVec<'scope, Self, (D, Self, Diff)>,
1860        _as_of: Antichain<mz_repr::Timestamp>,
1861        _summary: mz_repr::Timestamp,
1862    ) -> VecCollection<'scope, Self, D, Diff>
1863    where
1864        D: differential_dataflow::ExchangeData
1865            + crate::typedefs::MzData
1866            + differential_dataflow::Hashable
1867            + ColumnarData,
1868        for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1869        for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1870        for<'a> ::columnar::Ref<'a, Diff>: Ord,
1871        for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1872            ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1873    {
1874        // TODO: Implement bucketing on outer timestamp for iterative scopes.
1875        stream.as_collection()
1876    }
1877}
1878
1879/// A signal that can be awaited by operators to suspend them prior to startup.
1880///
1881/// Creating a signal also yields a token, dropping of which causes the signal to fire.
1882///
1883/// `StartSignal` is designed to be usable by both async and sync Timely operators.
1884///
1885///  * Async operators can simply `await` it.
1886///  * Sync operators should register an [`ActivateOnDrop`] value via [`StartSignal::drop_on_fire`]
1887///    and then check `StartSignal::has_fired()` on each activation.
1888#[derive(Clone)]
1889pub(crate) struct StartSignal {
1890    /// A future that completes when the signal fires.
1891    ///
1892    /// The inner type is `Infallible` because no data is ever expected on this channel. Instead the
1893    /// signal is activated by dropping the corresponding `Sender`.
1894    fut: futures::future::Shared<oneshot::Receiver<Infallible>>,
1895    /// A weak reference to the token, to register drop-on-fire values.
1896    token_ref: Weak<RefCell<Box<dyn Any>>>,
1897}
1898
1899impl StartSignal {
1900    /// Create a new `StartSignal` and a corresponding token that activates the signal when
1901    /// dropped.
1902    pub fn new() -> (Self, Rc<dyn Any>) {
1903        let (tx, rx) = oneshot::channel::<Infallible>();
1904        let token: Rc<RefCell<Box<dyn Any>>> = Rc::new(RefCell::new(Box::new(tx)));
1905        let signal = Self {
1906            fut: rx.shared(),
1907            token_ref: Rc::downgrade(&token),
1908        };
1909        (signal, token)
1910    }
1911
1912    pub fn has_fired(&self) -> bool {
1913        self.token_ref.strong_count() == 0
1914    }
1915
1916    /// Returns a Send-safe future that completes when the signal fires.
1917    ///
1918    /// Unlike `StartSignal` itself, the returned future does not retain a reference to the token,
1919    /// so it cannot be used for `drop_on_fire` or `has_fired` checks.
1920    pub fn into_send_future(self) -> impl Future<Output = ()> + Send {
1921        use futures::FutureExt;
1922        self.fut.map(|_| ())
1923    }
1924
1925    pub fn drop_on_fire(&self, to_drop: Box<dyn Any>) {
1926        if let Some(token) = self.token_ref.upgrade() {
1927            let mut token = token.borrow_mut();
1928            let inner = std::mem::replace(&mut *token, Box::new(()));
1929            *token = Box::new((inner, to_drop));
1930        }
1931    }
1932}
1933
1934impl Future for StartSignal {
1935    type Output = ();
1936
1937    fn poll(mut self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Self::Output> {
1938        self.fut.poll_unpin(cx).map(|_| ())
1939    }
1940}
1941
1942/// Extension trait to attach a `StartSignal` to operator outputs.
1943pub(crate) trait WithStartSignal {
1944    /// Delays data and progress updates until the start signal has fired.
1945    ///
1946    /// Note that this operator needs to buffer all incoming data, so it has some memory footprint,
1947    /// depending on the amount and shape of its inputs.
1948    fn with_start_signal(self, signal: StartSignal) -> Self;
1949}
1950
1951impl<'scope, Tr> WithStartSignal for Arranged<'scope, Tr>
1952where
1953    Tr: TraceReader<Time: RenderTimestamp> + Clone,
1954{
1955    fn with_start_signal(self, signal: StartSignal) -> Self {
1956        Arranged {
1957            stream: self.stream.with_start_signal(signal),
1958            trace: self.trace,
1959        }
1960    }
1961}
1962
1963impl<'scope, T: Timestamp, D> WithStartSignal for Stream<'scope, T, D>
1964where
1965    D: timely::Container + Clone + 'static,
1966{
1967    fn with_start_signal(self, signal: StartSignal) -> Self {
1968        let activations = self.scope().activations();
1969        self.unary(Pipeline, "StartSignal", |_cap, info| {
1970            let token = Box::new(ActivateOnDrop::new((), info.address, activations));
1971            signal.drop_on_fire(token);
1972
1973            let mut stash = Vec::new();
1974
1975            move |input, output| {
1976                // Stash incoming updates as long as the start signal has not fired.
1977                if !signal.has_fired() {
1978                    input.for_each(|cap, data| stash.push((cap, std::mem::take(data))));
1979                    return;
1980                }
1981
1982                // Release any data we might still have stashed.
1983                for (cap, mut data) in std::mem::take(&mut stash) {
1984                    output.session(&cap).give_container(&mut data);
1985                }
1986
1987                // Pass through all remaining input data.
1988                input.for_each(|cap, data| {
1989                    output.session(&cap).give_container(data);
1990                });
1991            }
1992        })
1993    }
1994}
1995
1996/// Suppress progress messages for times before the given `as_of`.
1997///
1998/// This operator exists specifically to work around a memory spike we'd otherwise see when
1999/// hydrating arrangements (database-issues#6368). The memory spike happens because when the `arrange_core`
2000/// operator observes a frontier advancement without data it inserts an empty batch into the spine.
2001/// When it later inserts the snapshot batch into the spine, an empty batch is already there and
2002/// the spine initiates a merge of these batches, which requires allocating a new batch the size of
2003/// the snapshot batch.
2004///
2005/// The strategy to avoid the spike is to prevent the insertion of that initial empty batch by
2006/// ensuring that the first frontier advancement downstream `arrange_core` operators observe is
2007/// beyond the `as_of`, so the snapshot data has already been collected.
2008///
2009/// To ensure this, this operator needs to take two measures:
2010///  * Keep around a minimum capability until the input announces progress beyond the `as_of`.
2011///  * Reclock all updates emitted at times not beyond the `as_of` to the minimum time.
2012///
2013/// The second measure requires elaboration: If we wouldn't reclock snapshot updates, they might
2014/// still be upstream of `arrange_core` operators when those get to know about us dropping the
2015/// minimum capability. The in-flight snapshot updates would hold back the input frontiers of
2016/// `arrange_core` operators to the `as_of`, which would cause them to insert empty batches.
2017fn suppress_early_progress<'scope, T: Timestamp, D>(
2018    stream: Stream<'scope, T, D>,
2019    as_of: Antichain<T>,
2020) -> Stream<'scope, T, D>
2021where
2022    D: timely::Container + Clone,
2023{
2024    stream.unary_frontier(Pipeline, "SuppressEarlyProgress", |default_cap, _info| {
2025        let mut early_cap = Some(default_cap);
2026
2027        move |(input, frontier), output| {
2028            input.for_each_time(|data_cap, data| {
2029                if as_of.less_than(data_cap.time()) {
2030                    let mut session = output.session(&data_cap);
2031                    for data in data {
2032                        session.give_container(data);
2033                    }
2034                } else {
2035                    let cap = early_cap.as_ref().expect("early_cap can't be dropped yet");
2036                    let mut session = output.session(&cap);
2037                    for data in data {
2038                        session.give_container(data);
2039                    }
2040                }
2041            });
2042
2043            if !PartialOrder::less_equal(&frontier.frontier(), &as_of.borrow()) {
2044                early_cap.take();
2045            }
2046        }
2047    })
2048}
2049
2050/// Extension trait for [`Stream`] to selectively limit progress.
2051trait LimitProgress<T: Timestamp> {
2052    /// Limit the progress of the stream until its frontier reaches the given `upper` bound. Expects
2053    /// the implementation to observe times in data, and release capabilities based on the probe's
2054    /// frontier, after applying `slack` to round up timestamps.
2055    ///
2056    /// The implementation of this operator is subtle to avoid regressions in the rest of the
2057    /// system. Specifically joins hold back compaction on the other side of the join, so we need to
2058    /// make sure we release capabilities as soon as possible. This is why we only limit progress
2059    /// for times before the `upper`, which is the time until which the source can distinguish
2060    /// updates at the time of rendering. Once we make progress to the `upper`, we need to release
2061    /// our capability.
2062    ///
2063    /// This isn't perfect, and can result in regressions if on of the inputs lags behind. We could
2064    /// consider using the join of the uppers, i.e, use lower bound upper of all available inputs.
2065    ///
2066    /// Once the input frontier reaches `[]`, the implementation must release any capability to
2067    /// allow downstream operators to release resources.
2068    ///
2069    /// The implementation should limit the number of pending times to `limit` if it is `Some` to
2070    /// avoid unbounded memory usage.
2071    ///
2072    /// * `handle` is a probe installed on the dataflow's outputs as late as possible, but before
2073    ///   any timestamp rounding happens (c.f., `REFRESH EVERY` materialized views).
2074    /// * `slack_ms` is the number of milliseconds to round up timestamps to.
2075    /// * `name` is a human-readable name for the operator.
2076    /// * `limit` is the maximum number of pending times to keep around.
2077    /// * `upper` is the upper bound of the stream's frontier until which the implementation can
2078    ///   retain a capability.
2079    fn limit_progress(
2080        self,
2081        handle: MzProbeHandle<T>,
2082        slack_ms: u64,
2083        limit: Option<usize>,
2084        upper: Antichain<T>,
2085        name: String,
2086    ) -> Self;
2087}
2088
2089/// Reads the times of the records a container holds, for [`LimitProgress`].
2090trait RecordTimes {
2091    /// Call `f` once per record, with that record's time.
2092    fn for_each_time(&self, f: impl FnMut(mz_repr::Timestamp));
2093}
2094
2095impl<D, R> RecordTimes for Vec<(D, mz_repr::Timestamp, R)> {
2096    fn for_each_time(&self, mut f: impl FnMut(mz_repr::Timestamp)) {
2097        for (_, time, _) in self {
2098            f(*time);
2099        }
2100    }
2101}
2102
2103impl<D, R> RecordTimes for Column<(D, mz_repr::Timestamp, R)>
2104where
2105    D: ColumnarData,
2106    R: ColumnarData,
2107    (D, mz_repr::Timestamp, R): ColumnarData<
2108        Container = (
2109            D::Container,
2110            <mz_repr::Timestamp as ColumnarData>::Container,
2111            R::Container,
2112        ),
2113    >,
2114{
2115    fn for_each_time(&self, mut f: impl FnMut(mz_repr::Timestamp)) {
2116        for time in self.borrow().1.into_index_iter() {
2117            f(time);
2118        }
2119    }
2120}
2121
2122// TODO: We could make this generic over a `T` that can be converted to and from a u64 millisecond
2123// number.
2124impl<'scope, C> LimitProgress<mz_repr::Timestamp> for Stream<'scope, mz_repr::Timestamp, C>
2125where
2126    C: timely::Container + Clone + RecordTimes,
2127{
2128    fn limit_progress(
2129        self,
2130        handle: MzProbeHandle<mz_repr::Timestamp>,
2131        slack_ms: u64,
2132        limit: Option<usize>,
2133        upper: Antichain<mz_repr::Timestamp>,
2134        name: String,
2135    ) -> Self {
2136        let scope = self.scope();
2137        let stream =
2138            self.unary_frontier(Pipeline, &format!("LimitProgress({name})"), |_cap, info| {
2139                // Times that we've observed on our input.
2140                let mut pending_times: BTreeSet<mz_repr::Timestamp> = BTreeSet::new();
2141                // Capability for the lower bound of `pending_times`, if any.
2142                let mut retained_cap: Option<Capability<mz_repr::Timestamp>> = None;
2143
2144                let activator = scope.activator_for(info.address);
2145                handle.activate(activator.clone());
2146
2147                move |(input, frontier), output| {
2148                    input.for_each(|cap, data| {
2149                        data.for_each_time(|time| {
2150                            let Some(time) = u64::from(time).checked_add(slack_ms) else {
2151                                return;
2152                            };
2153                            // `slack_ms == 0` means no rounding; otherwise round up to the next
2154                            // multiple of `slack_ms`. Avoids a divide-by-zero panic when the
2155                            // operator is configured without slack.
2156                            let rounded_time = if slack_ms == 0 {
2157                                time
2158                            } else {
2159                                (time / slack_ms).saturating_add(1).saturating_mul(slack_ms)
2160                            };
2161                            if !upper.less_than(&rounded_time.into()) {
2162                                pending_times.insert(rounded_time.into());
2163                            }
2164                        });
2165                        output.session(&cap).give_container(data);
2166                        if retained_cap.as_ref().is_none_or(|c| {
2167                            !c.time().less_than(cap.time()) && !upper.less_than(cap.time())
2168                        }) {
2169                            retained_cap = Some(cap.retain(0));
2170                        }
2171                    });
2172
2173                    handle.with_frontier(|f| {
2174                        while pending_times
2175                            .first()
2176                            .map_or(false, |retained_time| !f.less_than(&retained_time))
2177                        {
2178                            let _ = pending_times.pop_first();
2179                        }
2180                    });
2181
2182                    while limit.map_or(false, |limit| pending_times.len() > limit) {
2183                        let _ = pending_times.pop_first();
2184                    }
2185
2186                    match (retained_cap.as_mut(), pending_times.first()) {
2187                        (Some(cap), Some(first)) => cap.downgrade(first),
2188                        (_, None) => retained_cap = None,
2189                        _ => {}
2190                    }
2191
2192                    if frontier.is_empty() {
2193                        retained_cap = None;
2194                        pending_times.clear();
2195                    }
2196
2197                    if !pending_times.is_empty() {
2198                        tracing::debug!(
2199                            name,
2200                            info.global_id,
2201                            pending_times = %PendingTimesDisplay(pending_times.iter().cloned()),
2202                            frontier = ?frontier.frontier().get(0),
2203                            probe = ?handle.with_frontier(|f| f.get(0).cloned()),
2204                            ?upper,
2205                            "pending times",
2206                        );
2207                    }
2208                }
2209            });
2210        stream
2211    }
2212}
2213
2214/// A formatter for an iterator of timestamps that displays the first element, and subsequently
2215/// the difference between timestamps.
2216struct PendingTimesDisplay<T>(T);
2217
2218impl<T> std::fmt::Display for PendingTimesDisplay<T>
2219where
2220    T: IntoIterator<Item = mz_repr::Timestamp> + Clone,
2221{
2222    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2223        let mut iter = self.0.clone().into_iter();
2224        write!(f, "[")?;
2225        if let Some(first) = iter.next() {
2226            write!(f, "{}", first)?;
2227            let mut last = u64::from(first);
2228            for time in iter {
2229                write!(f, ", +{}", u64::from(time) - last)?;
2230                last = u64::from(time);
2231            }
2232        }
2233        write!(f, "]")?;
2234        Ok(())
2235    }
2236}
2237
2238/// Helper to merge pairs of datum iterators into a row or split a datum iterator
2239/// into two rows, given the arity of the first component.
2240#[derive(Clone, Copy, Debug)]
2241struct Pairer {
2242    split_arity: usize,
2243}
2244
2245impl Pairer {
2246    /// Creates a pairer with knowledge of the arity of first component in the pair.
2247    fn new(split_arity: usize) -> Self {
2248        Self { split_arity }
2249    }
2250
2251    /// Merges a pair of datum iterators creating a `Row` instance.
2252    fn merge<'a, I1, I2>(&self, first: I1, second: I2) -> Row
2253    where
2254        I1: IntoIterator<Item = Datum<'a>>,
2255        I2: IntoIterator<Item = Datum<'a>>,
2256    {
2257        SharedRow::pack(first.into_iter().chain(second))
2258    }
2259
2260    /// Splits a datum iterator into a pair of `Row` instances.
2261    fn split<'a>(&self, datum_iter: impl IntoIterator<Item = Datum<'a>>) -> (Row, Row) {
2262        let mut datum_iter = datum_iter.into_iter();
2263        let mut row_builder = SharedRow::get();
2264        let first = row_builder.pack_using(datum_iter.by_ref().take(self.split_arity));
2265        let second = row_builder.pack_using(datum_iter);
2266        (first, second)
2267    }
2268}