Skip to main content

mz_compute/render/
context.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//! Management of dataflow-local state, like arrangements, while building a
11//! dataflow.
12
13use std::collections::BTreeMap;
14use std::rc::Rc;
15
16use columnar::{Columnar, Index};
17use differential_dataflow::consolidation::ConsolidatingContainerBuilder;
18use differential_dataflow::operators::arrange::Arranged;
19use differential_dataflow::trace::cursor::{BatchCursor, BatchKey, BatchVal};
20use differential_dataflow::trace::implementations::BatchContainer;
21use differential_dataflow::trace::{Cursor, Navigable, TraceReader};
22use differential_dataflow::{AsCollection, VecCollection};
23use mz_compute_types::dataflows::DataflowDescription;
24use mz_compute_types::dyncfgs::{ENABLE_COMPUTE_TEMPORAL_BUCKETING, TEMPORAL_BUCKETING_SUMMARY};
25use mz_compute_types::plan::scalar::{LirScalarExpr, mfp_mir_to_lir_plan, mfp_plan_lir_to_mir};
26use mz_compute_types::plan::{ArrangementStrategy, AvailableCollections};
27use mz_dyncfg::ConfigSet;
28use mz_expr::{Eval, Id, MfpPlan};
29use mz_ore::soft_assert_or_log;
30use mz_repr::fixed_length::ExtendDatums;
31use mz_repr::{DatumVec, DatumVecBorrow, Diff, GlobalId, Row, RowArena, SharedRow, StableRow};
32use mz_storage_types::controller::CollectionMetadata;
33use mz_timely_util::columnar::Column;
34use mz_timely_util::columnar::batcher;
35use mz_timely_util::columnar::builder::ColumnBuilder;
36use mz_timely_util::columnar::chunk::{AccountedChunkBatcher, ChunkChunker, UnchunkBuilder};
37use mz_timely_util::columnar::consolidate::ConsolidatingColumnBuilder;
38use mz_timely_util::columnar::{Col2ValBatcher, Col2ValColBatcher, columnar_exchange};
39use mz_timely_util::columnation::ColumnationChunker;
40use timely::ContainerBuilder;
41use timely::container::NoopBuilder;
42use timely::dataflow::channels::pact::{ExchangeCore, Pipeline};
43use timely::dataflow::operators::Capability;
44use timely::dataflow::operators::generic::builder_rc::OperatorBuilder;
45use timely::dataflow::operators::generic::{OutputBuilder, OutputBuilderSession};
46use timely::dataflow::{Scope, Stream};
47use timely::progress::operate::FrontierInterest;
48use timely::progress::{Antichain, Timestamp};
49
50use crate::compute_state::ComputeState;
51use crate::extensions::arrange::{ArrangementBatcher, KeyCollection, MzArrange, MzArrangeCore};
52use crate::extensions::reduce::MzReduce;
53use crate::render::columnar::{ColCollection, flat_map_datums};
54use crate::render::errors::{DataflowErrorSer, ErrorLogger};
55use crate::render::{LinearJoinSpec, MaybeBucketByTime, RenderTimestamp};
56use crate::typedefs::{
57    ErrAgent, ErrBatcher, ErrBuilder, ErrEnter, ErrSpine, RowRowAgent, RowRowEnter, RowRowSpine,
58};
59use mz_row_spine::{RowRowBuilder, RowRowColPagedBuilder};
60
61/// Dataflow-local collections and arrangements.
62///
63/// A context means to wrap available data assets and present them in an easy-to-use manner.
64/// These assets include dataflow-local collections and arrangements, as well as imported
65/// arrangements from outside the dataflow.
66///
67/// Context has a timestamp type `T`, which is the timestamp used by the scope in question.
68pub struct Context<'scope, T: RenderTimestamp> {
69    /// The scope within which all managed collections exist.
70    ///
71    /// It is an error to add any collections not contained in this scope.
72    pub(crate) scope: Scope<'scope, T>,
73    /// The debug name of the dataflow associated with this context.
74    pub debug_name: String,
75    /// The Timely ID of the dataflow associated with this context.
76    pub dataflow_id: usize,
77    /// The collection IDs of exports of the dataflow associated with this context.
78    pub export_ids: Vec<GlobalId>,
79    /// Frontier before which updates should not be emitted.
80    ///
81    /// We *must* apply it to sinks, to ensure correct outputs.
82    /// We *should* apply it to sources and imported traces, because it improves performance.
83    pub as_of_frontier: Antichain<mz_repr::Timestamp>,
84    /// Frontier after which updates should not be emitted.
85    /// Used to limit the amount of work done when appropriate.
86    pub until: Antichain<mz_repr::Timestamp>,
87    /// Bindings of identifiers to collections.
88    pub bindings: BTreeMap<Id, CollectionBundle<'scope, T>>,
89    /// The logger, from Timely's logging framework, if logs are enabled.
90    pub(super) compute_logger: Option<crate::logging::compute::Logger>,
91    /// Specification for rendering linear joins.
92    pub(super) linear_join_spec: LinearJoinSpec,
93    /// The expiration time for dataflows in this context. The output's frontier should never advance
94    /// past this frontier, except the empty frontier.
95    pub dataflow_expiration: Antichain<mz_repr::Timestamp>,
96    /// The config set for this context.
97    pub config_set: Rc<ConfigSet>,
98}
99
100impl<'scope, T: RenderTimestamp> Context<'scope, T> {
101    /// Creates a new empty Context.
102    pub fn for_dataflow_in<Plan>(
103        dataflow: &DataflowDescription<Plan, CollectionMetadata>,
104        scope: Scope<'scope, T>,
105        compute_state: &ComputeState,
106        until: Antichain<mz_repr::Timestamp>,
107        dataflow_expiration: Antichain<mz_repr::Timestamp>,
108    ) -> Self {
109        use mz_ore::collections::CollectionExt as IteratorExt;
110        let dataflow_id = *scope.addr().into_first();
111        let as_of_frontier = dataflow
112            .as_of
113            .clone()
114            .unwrap_or_else(|| Antichain::from_elem(Timestamp::minimum()));
115
116        let export_ids = dataflow.export_ids().collect();
117
118        // Skip compute event logging for transient dataflows. We do this to avoid overhead for
119        // slow-path peeks, but it also affects subscribes. For now that seems fine, but we may
120        // want to reconsider in the future.
121        let compute_logger = if dataflow.is_transient() {
122            None
123        } else {
124            compute_state.compute_logger.clone()
125        };
126
127        Self {
128            scope,
129            debug_name: dataflow.debug_name.clone(),
130            dataflow_id,
131            export_ids,
132            as_of_frontier,
133            until,
134            bindings: BTreeMap::new(),
135            compute_logger,
136            linear_join_spec: compute_state.linear_join_spec,
137            dataflow_expiration,
138            config_set: Rc::clone(&compute_state.worker_config),
139        }
140    }
141}
142
143impl<'scope, T: RenderTimestamp> Context<'scope, T> {
144    /// Insert a collection bundle by an identifier.
145    ///
146    /// This is expected to be used to install external collections (sources, indexes, other views),
147    /// as well as for `Let` bindings of local collections.
148    pub fn insert_id(
149        &mut self,
150        id: Id,
151        collection: CollectionBundle<'scope, T>,
152    ) -> Option<CollectionBundle<'scope, T>> {
153        self.bindings.insert(id, collection)
154    }
155    /// Remove a collection bundle by an identifier.
156    ///
157    /// The primary use of this method is uninstalling `Let` bindings.
158    pub fn remove_id(&mut self, id: Id) -> Option<CollectionBundle<'scope, T>> {
159        self.bindings.remove(&id)
160    }
161    /// Melds a collection bundle to whatever exists.
162    pub fn update_id(&mut self, id: Id, collection: CollectionBundle<'scope, T>) {
163        if !self.bindings.contains_key(&id) {
164            self.bindings.insert(id, collection);
165        } else {
166            let binding = self
167                .bindings
168                .get_mut(&id)
169                .expect("Binding verified to exist");
170            if collection.collection.is_some() {
171                binding.collection = collection.collection;
172            }
173            for (key, flavor) in collection.arranged.into_iter() {
174                binding.arranged.insert(key, flavor);
175            }
176        }
177    }
178    /// Look up a collection bundle by an identifier.
179    pub fn lookup_id(&self, id: Id) -> Option<CollectionBundle<'scope, T>> {
180        self.bindings.get(&id).cloned()
181    }
182
183    pub(super) fn error_logger(&self) -> ErrorLogger {
184        ErrorLogger::new(self.debug_name.clone())
185    }
186}
187
188impl<'scope, T: RenderTimestamp> Context<'scope, T> {
189    /// Brings the underlying arrangements and collections into a region.
190    pub fn enter_region<'a>(
191        &self,
192        region: Scope<'a, T>,
193        bindings: Option<&std::collections::BTreeSet<Id>>,
194    ) -> Context<'a, T> {
195        let bindings = self
196            .bindings
197            .iter()
198            .filter(|(key, _)| bindings.as_ref().map(|b| b.contains(key)).unwrap_or(true))
199            .map(|(key, bundle)| (*key, bundle.enter_region(region)))
200            .collect();
201
202        Context {
203            scope: region,
204            debug_name: self.debug_name.clone(),
205            dataflow_id: self.dataflow_id.clone(),
206            export_ids: self.export_ids.clone(),
207            as_of_frontier: self.as_of_frontier.clone(),
208            until: self.until.clone(),
209            compute_logger: self.compute_logger.clone(),
210            linear_join_spec: self.linear_join_spec.clone(),
211            bindings,
212            dataflow_expiration: self.dataflow_expiration.clone(),
213            config_set: Rc::clone(&self.config_set),
214        }
215    }
216}
217
218/// Describes flavor of arrangement: local or imported trace.
219#[derive(Clone)]
220pub enum ArrangementFlavor<'scope, T: RenderTimestamp> {
221    /// A dataflow-local arrangement.
222    Local(
223        Arranged<'scope, RowRowAgent<T, Diff>>,
224        Arranged<'scope, ErrAgent<T, Diff>>,
225    ),
226    /// An imported trace from outside the dataflow.
227    ///
228    /// The `GlobalId` identifier exists so that exports of this same trace
229    /// can refer back to and depend on the original instance.
230    Trace(
231        GlobalId,
232        Arranged<'scope, RowRowEnter<mz_repr::Timestamp, Diff, T>>,
233        Arranged<'scope, ErrEnter<mz_repr::Timestamp, T>>,
234    ),
235}
236
237impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
238    /// Constructs and applies logic to elements of `self` and returns the results.
239    ///
240    /// The `logic` callback receives a borrow of the decoded datum vector, a timestamp, a
241    /// diff, and two output sessions: one for `ok` updates of type `(D, T, Diff)` and one for
242    /// MFP-style `DataflowErrorSer` updates. It must return the number of records *produced*
243    /// (written to either session), not the number of input tuples consumed.
244    ///
245    /// # Fuel
246    ///
247    /// The operator accumulates the returned counts as fuel and yields when the total reaches
248    /// an internal refuel threshold. The metric is output-produced (not input-consumed) on
249    /// purpose: it regulates two asymmetric pressures.
250    ///
251    /// * **Drain inputs.** The operator holds a clone of each pending `Batch` until its work
252    ///   item pops; we want to release that memory back to the upstream arrangement as soon
253    ///   as possible. A `filter(false)` MFP returns 0 for every tuple, so fuel never trips
254    ///   and the cursor runs to end-of-batch in one activation.
255    /// * **Throttle outputs.** A `map("1KB-string")` MFP produces large records per input;
256    ///   stopping when emit count hits the threshold caps how much data a single activation
257    ///   dumps on the next operator.
258    ///
259    /// The refuel constant is a pragmatic compromise: large enough to be a non-event in
260    /// steady-state, small enough that one activation can't flood downstream. There is no
261    /// universal value across MFP shapes.
262    ///
263    /// If `key` is set, this is a promise that `logic` will produce no results on
264    /// records for which the key does not evaluate to the value. This is used to
265    /// leap directly to exactly those records.
266    ///
267    /// The `max_demand` parameter limits the number of columns decoded from the
268    /// input. Only the first `max_demand` columns are decoded. Pass `usize::MAX` to
269    /// decode all columns.
270    pub fn flat_map<DCB, L>(
271        &self,
272        key: Option<&Row>,
273        max_demand: usize,
274        logic: L,
275    ) -> (
276        Stream<'scope, T, DCB::Container>,
277        VecCollection<'scope, T, DataflowErrorSer, Diff>,
278    )
279    where
280        DCB: ContainerBuilder,
281        L: for<'a, 'b> FnMut(
282                &'a mut DatumVecBorrow<'b>,
283                T,
284                Diff,
285                &mut Session<T, DCB>,
286                &mut Session<T, ECB<T>>,
287            ) -> usize
288            + 'static,
289    {
290        // `logic` is passed straight through to `flat_map_core_fallible`, which owns the per-row
291        // decode (and the activation-scoped arena it decodes into).
292        match &self {
293            ArrangementFlavor::Local(oks, errs) => {
294                let (oks, mfp_errs) = CollectionBundle::<T>::flat_map_core_fallible::<_, DCB, _>(
295                    oks.clone(),
296                    key,
297                    max_demand,
298                    logic,
299                    REFUEL,
300                );
301                let errs = errs.clone().as_collection(|k, &()| k.clone());
302                let errs = errs.concat(mfp_errs.as_collection());
303                (oks, errs)
304            }
305            ArrangementFlavor::Trace(_, oks, errs) => {
306                let (oks, mfp_errs) = CollectionBundle::<T>::flat_map_core_fallible::<_, DCB, _>(
307                    oks.clone(),
308                    key,
309                    max_demand,
310                    logic,
311                    REFUEL,
312                );
313                let errs = errs.clone().as_collection(|k, &()| k.clone());
314                let errs = errs.concat(mfp_errs.as_collection());
315                (oks, errs)
316            }
317        }
318    }
319
320    /// Ok-only variant of [`Self::flat_map`]. The `logic` callback receives a single output
321    /// session, cannot produce errors, and returns the number of records produced (see
322    /// [`Self::flat_map`] for fuel semantics). The returned err collection comes solely from
323    /// the arrangement; no extra operator is built to carry an empty MFP-error stream.
324    pub fn flat_map_ok<DCB, L>(
325        &self,
326        key: Option<&Row>,
327        max_demand: usize,
328        logic: L,
329    ) -> (
330        Stream<'scope, T, DCB::Container>,
331        VecCollection<'scope, T, DataflowErrorSer, Diff>,
332    )
333    where
334        // No push bound here: it lives at `logic`'s `give` call site, so a caller can push
335        // borrowed records into a columnar builder that has no owned-tuple `Push`.
336        DCB: ContainerBuilder,
337        L: for<'a, 'b> FnMut(&'a mut DatumVecBorrow<'b>, T, Diff, &mut Session<T, DCB>) -> usize
338            + 'static,
339    {
340        match &self {
341            ArrangementFlavor::Local(oks, errs) => {
342                let oks = CollectionBundle::<T>::flat_map_core_ok::<_, DCB, _>(
343                    oks.clone(),
344                    key,
345                    max_demand,
346                    logic,
347                    REFUEL,
348                );
349                let errs = errs.clone().as_collection(|k, &()| k.clone());
350                (oks, errs)
351            }
352            ArrangementFlavor::Trace(_, oks, errs) => {
353                let oks = CollectionBundle::<T>::flat_map_core_ok::<_, DCB, _>(
354                    oks.clone(),
355                    key,
356                    max_demand,
357                    logic,
358                    REFUEL,
359                );
360                let errs = errs.clone().as_collection(|k, &()| k.clone());
361                (oks, errs)
362            }
363        }
364    }
365}
366impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
367    /// The scope containing the collection bundle.
368    pub fn scope(&self) -> Scope<'scope, T> {
369        match self {
370            ArrangementFlavor::Local(oks, _errs) => oks.stream.scope(),
371            ArrangementFlavor::Trace(_gid, oks, _errs) => oks.stream.scope(),
372        }
373    }
374
375    /// Brings the arrangement flavor into a region.
376    pub fn enter_region<'a>(&self, region: Scope<'a, T>) -> ArrangementFlavor<'a, T> {
377        match self {
378            ArrangementFlavor::Local(oks, errs) => ArrangementFlavor::Local(
379                oks.clone().enter_region(region),
380                errs.clone().enter_region(region),
381            ),
382            ArrangementFlavor::Trace(gid, oks, errs) => ArrangementFlavor::Trace(
383                *gid,
384                oks.clone().enter_region(region),
385                errs.clone().enter_region(region),
386            ),
387        }
388    }
389}
390impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
391    /// Extracts the arrangement flavor from a region.
392    pub fn leave_region<'outer>(&self, outer: Scope<'outer, T>) -> ArrangementFlavor<'outer, T> {
393        match self {
394            ArrangementFlavor::Local(oks, errs) => ArrangementFlavor::Local(
395                oks.clone().leave_region(outer),
396                errs.clone().leave_region(outer),
397            ),
398            ArrangementFlavor::Trace(gid, oks, errs) => ArrangementFlavor::Trace(
399                *gid,
400                oks.clone().leave_region(outer),
401                errs.clone().leave_region(outer),
402            ),
403        }
404    }
405}
406/// Rewrites an arranged error collection to hold each of its errors once.
407///
408/// Sound because query semantics depend only on whether an error is present, and correct under
409/// retraction only because it reads the accumulated collection: no pointwise function of the input
410/// diffs (a saturating add, a sign) can collapse multiplicity and still cancel when the errors
411/// retract.
412///
413/// NOTE: One consumer does read the multiplicity. Error-count introspection reports it as the
414/// number of failing rows, so `log_dataflow_errors` must see the collection before this collapses
415/// it, or a dataflow reports one error however many rows failed. Collapse after the logging, never
416/// before.
417pub(crate) fn distinct_arranged_errs<'a, T: RenderTimestamp>(
418    errs: Arranged<'a, ErrAgent<T, Diff>>,
419    name: &str,
420) -> Arranged<'a, ErrAgent<T, Diff>> {
421    errs.mz_reduce_abelian::<_, ErrBuilder<_, _>, ErrSpine<_, _>, _>(
422        name,
423        |_err, _input, output| output.push(((), Diff::ONE)),
424    )
425}
426
427/// Rewrites an error collection to hold each of its errors once.
428///
429/// Costs an arrangement more than [`distinct_arranged_errs`], which reuses the arrangement its
430/// input already has.
431pub(crate) fn distinct_errs_collection<'a, T: RenderTimestamp>(
432    errs: VecCollection<'a, T, DataflowErrorSer, Diff>,
433) -> VecCollection<'a, T, DataflowErrorSer, Diff> {
434    let errs: KeyCollection<_, _, _> = errs.into();
435    let errs = errs
436        .mz_arrange::<ColumnationChunker<_>, ErrBatcher<_, _>, ErrBuilder<_, _>, ErrSpine<_, _>>(
437            "Arrange errors",
438        );
439    distinct_arranged_errs(errs, "Distinct errors").as_collection(|err, _| err.clone())
440}
441
442/// A bundle of the various ways a collection can be represented.
443///
444/// This type maintains the invariant that it does contain at least one valid
445/// source of data, either a collection or at least one arrangement.
446#[derive(Clone)]
447pub struct CollectionBundle<'scope, T: RenderTimestamp> {
448    pub collection: Option<(
449        ColCollection<'scope, T>,
450        VecCollection<'scope, T, DataflowErrorSer, Diff>,
451    )>,
452    pub arranged: BTreeMap<Vec<LirScalarExpr>, ArrangementFlavor<'scope, T>>,
453}
454
455impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
456    /// Construct a new collection bundle from a [`ColCollection`] and an error stream.
457    pub fn from_edge(
458        oks: ColCollection<'scope, T>,
459        errs: VecCollection<'scope, T, DataflowErrorSer, Diff>,
460    ) -> Self {
461        Self {
462            collection: Some((oks, errs)),
463            arranged: BTreeMap::default(),
464        }
465    }
466
467    /// Inserts arrangements by the expressions on which they are keyed.
468    pub fn from_expressions(
469        exprs: Vec<LirScalarExpr>,
470        arrangements: ArrangementFlavor<'scope, T>,
471    ) -> Self {
472        let mut arranged = BTreeMap::new();
473        arranged.insert(exprs, arrangements);
474        Self {
475            collection: None,
476            arranged,
477        }
478    }
479
480    /// Inserts arrangements by the columns on which they are keyed.
481    pub fn from_columns<I: IntoIterator<Item = usize>>(
482        columns: I,
483        arrangements: ArrangementFlavor<'scope, T>,
484    ) -> Self {
485        let mut keys = Vec::new();
486        for column in columns {
487            keys.push(LirScalarExpr::column(column));
488        }
489        Self::from_expressions(keys, arrangements)
490    }
491
492    /// The scope containing the collection bundle.
493    pub fn scope(&self) -> Scope<'scope, T> {
494        if let Some((oks, _errs)) = &self.collection {
495            oks.scope()
496        } else {
497            self.arranged
498                .values()
499                .next()
500                .expect("Must contain a valid collection")
501                .scope()
502        }
503    }
504
505    /// Collapses the multiplicity of every error this bundle carries to one.
506    ///
507    /// Belongs at the definition of a binding that more than one `Get` reads. Each reader
508    /// propagates the binding's errors independently, so a binding read `f` times contributes its
509    /// errors `f` times to the dataflow's error output, and those factors apply again at every
510    /// further level of sharing: a chain of diamonds multiplies rather than adds, and reaches
511    /// `Diff` overflow at a depth plans really do have. Collapsing at each definition holds the
512    /// dataflow's error multiplicity to the fan-out of a single level.
513    ///
514    /// Sound because error semantics depend only on whether an error is present, and correct under
515    /// retraction only because it reads the accumulated collection: no pointwise function of the
516    /// input diffs (a saturating add, a sign) can collapse multiplicity and still cancel when the
517    /// errors retract.
518    ///
519    /// Collapses every form the bundle offers, not just one. Each form carries its own error
520    /// stream, and those streams differ in content as well as identity: an arrangement's errors
521    /// include the key-formation errors that the raw collection's do not. Which form a consumer
522    /// reads is the consumer's choice, and a delta join reads both within one operator, so a
523    /// binding's definition cannot know which form to collapse.
524    ///
525    /// NOTE: Leaves imported arrangements (`ArrangementFlavor::Trace`) alone, whose error traces
526    /// this dataflow cannot rewrite in place. Their errors arrive bounded by the exporting
527    /// dataflow's last level of sharing rather than collapsed to one, since nothing collapses at an
528    /// export. A global read more than once within one dataflow is not collapsed either, because
529    /// only local bindings reach this.
530    pub fn distinct_errs(mut self) -> Self {
531        if let Some((oks, errs)) = self.collection.take() {
532            self.collection = Some((oks, distinct_errs_collection(errs)));
533        }
534        for (key, flavor) in std::mem::take(&mut self.arranged) {
535            let flavor = match flavor {
536                ArrangementFlavor::Local(oks, errs) => {
537                    // Names the key, not the binding: an operator name carrying a `LocalId` would
538                    // churn the introspection goldens every time the optimizer renumbers locals.
539                    let name = format!("Distinct errors[{key:?}]");
540                    ArrangementFlavor::Local(oks, distinct_arranged_errs(errs, &name))
541                }
542                flavor @ ArrangementFlavor::Trace(..) => flavor,
543            };
544            self.arranged.insert(key, flavor);
545        }
546        self
547    }
548
549    /// Brings the collection bundle into a region.
550    pub fn enter_region<'inner>(&self, region: Scope<'inner, T>) -> CollectionBundle<'inner, T> {
551        CollectionBundle {
552            collection: self.collection.as_ref().map(|(oks, errs)| {
553                (
554                    oks.clone().enter_region(region),
555                    errs.clone().enter_region(region),
556                )
557            }),
558            arranged: self
559                .arranged
560                .iter()
561                .map(|(key, bundle)| (key.clone(), bundle.enter_region(region)))
562                .collect(),
563        }
564    }
565}
566
567impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
568    /// Extracts the collection bundle from a region.
569    pub fn leave_region<'outer>(&self, outer: Scope<'outer, T>) -> CollectionBundle<'outer, T> {
570        CollectionBundle {
571            collection: self.collection.as_ref().map(|(oks, errs)| {
572                (
573                    oks.clone().leave_region(outer),
574                    errs.clone().leave_region(outer),
575                )
576            }),
577            arranged: self
578                .arranged
579                .iter()
580                .map(|(key, bundle)| (key.clone(), bundle.leave_region(outer)))
581                .collect(),
582        }
583    }
584}
585
586impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
587    /// Asserts that the arrangement for a specific key
588    /// (or the raw collection for no key) exists,
589    /// and returns the corresponding collection.
590    ///
591    /// This returns the collection as-is, without
592    /// doing any unthinning transformation.
593    /// Therefore, it should be used when the appropriate transformation
594    /// was planned as part of a following MFP.
595    ///
596    /// If `key` is specified, the function converts the arrangement to a collection using a
597    /// fueled `flat_map` operator.
598    ///
599    /// The keyed path materializes the arrangement as the columnar edge. The unkeyed path
600    /// returns the unarranged `.collection` edge with its variant intact.
601    pub fn as_specific_collection(
602        &self,
603        key: Option<&[LirScalarExpr]>,
604    ) -> (
605        ColCollection<'scope, T>,
606        VecCollection<'scope, T, DataflowErrorSer, Diff>,
607    ) {
608        // Any operator that uses this method was told to use a particular
609        // collection during LIR planning, where we should have made
610        // sure that that collection exists.
611        //
612        // If it doesn't, we panic.
613        match key {
614            None => self
615                .collection
616                .clone()
617                .expect("The unarranged collection doesn't exist."),
618            Some(key) => {
619                let arranged = self.arranged.get(key).unwrap_or_else(|| {
620                    panic!("The collection arranged by {:?} doesn't exist.", key)
621                });
622                // Output is 1:1 from the already-consolidated cursor, so a
623                // non-consolidating `ColumnBuilder` suffices. `max_demand` is
624                // `usize::MAX` because the materialized collection carries every column.
625                let (ok, err) =
626                    arranged.flat_map_ok::<ColumnBuilder<(Row, T, Diff)>, _>(None, usize::MAX, {
627                        // `give` copies the bytes into the column, so one buffer
628                        // serves every record.
629                        let mut row_buf = Row::default();
630                        move |borrow, t, r, ok_session| {
631                            row_buf.packer().extend(borrow.iter());
632                            ok_session.give((&row_buf, &t, &r));
633                            1
634                        }
635                    });
636                (ok.as_collection(), err)
637            }
638        }
639    }
640
641    /// Constructs and applies logic to elements of a collection and returns the results.
642    ///
643    /// The function applies `logic` on elements. The logic conceptually receives
644    /// `(&Row, &Row)` pairs in the form of a datum vec in the expected order.
645    ///
646    /// If `key_val` is set, this is a promise that `logic` will produce no results on
647    /// records for which the key does not evaluate to the value. This is used when we
648    /// have an arrangement by that key to leap directly to exactly those records.
649    /// It is important that `logic` still guard against data that does not satisfy
650    /// this constraint, as this method does not statically know that it will have
651    /// that arrangement.
652    ///
653    /// The `max_demand` parameter limits the number of columns decoded from the
654    /// input. Only the first `max_demand` columns are decoded. Pass `usize::MAX` to
655    /// decode all columns.
656    pub fn flat_map<DCB, L>(
657        &self,
658        key_val: Option<(Vec<LirScalarExpr>, Option<Row>)>,
659        max_demand: usize,
660        logic: L,
661    ) -> (
662        Stream<'scope, T, DCB::Container>,
663        VecCollection<'scope, T, DataflowErrorSer, Diff>,
664    )
665    where
666        DCB: ContainerBuilder,
667        L: for<'a> FnMut(
668                &'a mut DatumVecBorrow<'_>,
669                T,
670                Diff,
671                &mut Session<T, DCB>,
672                &mut Session<T, ECB<T>>,
673            ) -> usize
674            + 'static,
675    {
676        // If `key_val` is set, we should have to use the corresponding arrangement.
677        // If there isn't one, that implies an error in the contract between
678        // key-production and available arrangements.
679        if let Some((key, val)) = key_val {
680            self.arrangement(&key)
681                .expect("Should have ensured during planning that this arrangement exists.")
682                .flat_map::<DCB, _>(val.as_ref(), max_demand, logic)
683        } else {
684            let (oks, errs) = self
685                .collection
686                .clone()
687                .expect("Invariant violated: CollectionBundle contains no collection.");
688            let (ok_stream, err_stream) =
689                flat_map_datums::<_, DCB, _>(oks, "CollectionFlatMap", max_demand, logic);
690            let errs = errs.concat(err_stream.as_collection());
691            (ok_stream, errs)
692        }
693    }
694
695    /// Factored out common logic for using literal keys in general traces.
696    ///
697    /// This logic is sufficiently interesting that we want to write it only
698    /// once, and thereby avoid any skew in the two uses of the logic.
699    ///
700    /// The function presents the contents of the trace as `(key, value, time, delta)` tuples,
701    /// where key and value are potentially specialized, but convertible into rows. The `logic`
702    /// callback writes ok results into the first session and errors into the second, returning
703    /// the number of records produced. See [`ArrangementFlavor::flat_map`] for the fuel
704    /// rationale.
705    fn flat_map_core_fallible<Tr, DCB, L>(
706        trace: Arranged<'scope, Tr>,
707        key: Option<&<<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned>,
708        max_demand: usize,
709        mut logic: L,
710        refuel: usize,
711    ) -> (
712        Stream<'scope, T, DCB::Container>,
713        Stream<'scope, T, Vec<(DataflowErrorSer, T, Diff)>>,
714    )
715    where
716        Tr: TraceReader<Batch: Navigable, Time = T> + Clone + 'static,
717        for<'a> BatchCursor<Tr>:
718            Cursor<Key<'a>: ExtendDatums, Val<'a>: ExtendDatums, Time = T, Diff = mz_repr::Diff>,
719        <<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned: PartialEq,
720        DCB: ContainerBuilder,
721        // `logic` receives the key and value already decoded into a `DatumVecBorrow`. The decode
722        // (and its arena/`DatumVec`) lives in the per-activation closure below, so it is scoped to
723        // a single scheduling invocation rather than to the operator.
724        L: for<'a, 'b> FnMut(
725                &'a mut DatumVecBorrow<'b>,
726                T,
727                mz_repr::Diff,
728                &mut Session<T, DCB>,
729                &mut Session<T, ECB<T>>,
730            ) -> usize
731            + 'static,
732    {
733        let scope = trace.stream.scope();
734
735        let mut key_con = <BatchCursor<Tr> as Cursor>::KeyContainer::with_capacity(1);
736        if let Some(key) = &key {
737            key_con.push_own(key);
738        }
739        let mode = if key.is_some() { "index" } else { "scan" };
740        let name = format!("ArrangementFlatMap({})", mode);
741
742        let mut builder = OperatorBuilder::new(name, scope.clone());
743        let (ok_output, ok_stream) = builder.new_output();
744        let mut ok_output = OutputBuilder::<_, DCB>::from(ok_output);
745        let (err_output, err_stream) = builder.new_output();
746        let mut err_output = OutputBuilder::<_, ECB<T>>::from(err_output);
747        let mut input = builder.new_input(trace.stream.clone(), Pipeline);
748        let operator_info = builder.operator_info();
749
750        builder.build(move |_capabilities| {
751            // Acquire an activator to reschedule the operator when it has unfinished work.
752            let activator = scope.activator_for(operator_info.address);
753            // Maintain a list of work to do, cursor to navigate and process.
754            let mut todo = std::collections::VecDeque::new();
755            move |_frontiers| {
756                let key = key_con.get(0);
757                let mut ok_output = ok_output.activate();
758                let mut err_output = err_output.activate();
759
760                // First, dequeue all batches.
761                input.for_each(|time, data| {
762                    // Retain a capability for each output, as the work may complete across
763                    // multiple activations.
764                    let ok_cap = time.retain(0);
765                    let err_cap = time.retain(1);
766                    for batch in data.iter() {
767                        todo.push_back(PendingWork::new(
768                            ok_cap.clone(),
769                            err_cap.clone(),
770                            batch.cursor(),
771                            batch.clone(),
772                        ));
773                    }
774                });
775
776                // Decode the key/value of each record into datums for `logic`. The arena and datum
777                // buffer are created here, so they are scoped to this activation (dropped when it
778                // returns) rather than retained for the operator's lifetime; both are reused across
779                // the records processed within the activation.
780                let mut temp_storage = RowArena::new();
781                let mut datums = DatumVec::new();
782                let mut decode_logic =
783                    |k: BatchKey<'_, Tr>,
784                     v: BatchVal<'_, Tr>,
785                     t: T,
786                     d: mz_repr::Diff,
787                     ok_session: &mut Session<T, DCB>,
788                     err_session: &mut Session<T, ECB<T>>| {
789                        temp_storage.clear();
790                        let mut datums_borrow = datums.borrow();
791                        k.extend_datums(&temp_storage, &mut datums_borrow, Some(max_demand));
792                        let remaining = max_demand.saturating_sub(datums_borrow.len());
793                        v.extend_datums(&temp_storage, &mut datums_borrow, Some(remaining));
794                        logic(&mut datums_borrow, t, d, ok_session, err_session)
795                    };
796
797                // Second, make progress on `todo`.
798                let mut fuel = refuel;
799                while !todo.is_empty() && fuel > 0 {
800                    todo.front_mut().unwrap().do_work(
801                        key.as_ref(),
802                        &mut decode_logic,
803                        &mut fuel,
804                        &mut ok_output,
805                        &mut err_output,
806                    );
807                    if fuel > 0 {
808                        todo.pop_front();
809                    }
810                }
811                // If we have not finished all work, re-activate the operator.
812                if !todo.is_empty() {
813                    activator.activate();
814                }
815            }
816        });
817
818        (ok_stream, err_stream)
819    }
820
821    /// Ok-only variant of [`Self::flat_map_core_fallible`]. The `logic` callback writes results
822    /// into a single output session and returns the number of records produced (see the
823    /// fallible variant for fuel semantics). Use this when the caller statically knows it
824    /// will never produce `DataflowErrorSer` records, to avoid building a second output port
825    /// and the empty err stream that would follow it.
826    fn flat_map_core_ok<Tr, DCB, L>(
827        trace: Arranged<'scope, Tr>,
828        key: Option<&<<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned>,
829        max_demand: usize,
830        mut logic: L,
831        refuel: usize,
832    ) -> Stream<'scope, T, DCB::Container>
833    where
834        Tr: TraceReader<Batch: Navigable, Time = T> + Clone + 'static,
835        for<'a> BatchCursor<Tr>:
836            Cursor<Key<'a>: ExtendDatums, Val<'a>: ExtendDatums, Time = T, Diff = mz_repr::Diff>,
837        <<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned: PartialEq,
838        DCB: ContainerBuilder,
839        L: for<'a, 'b> FnMut(
840                &'a mut DatumVecBorrow<'b>,
841                T,
842                mz_repr::Diff,
843                &mut Session<T, DCB>,
844            ) -> usize
845            + 'static,
846    {
847        let scope = trace.stream.scope();
848
849        let mut key_con = <BatchCursor<Tr> as Cursor>::KeyContainer::with_capacity(1);
850        if let Some(key) = &key {
851            key_con.push_own(key);
852        }
853        let mode = if key.is_some() { "index" } else { "scan" };
854        let name = format!("ArrangementFlatMapOk({})", mode);
855
856        let mut builder = OperatorBuilder::new(name, scope.clone());
857        let (ok_output, ok_stream) = builder.new_output();
858        let mut ok_output = OutputBuilder::<_, DCB>::from(ok_output);
859        let mut input = builder.new_input(trace.stream.clone(), Pipeline);
860        let operator_info = builder.operator_info();
861
862        builder.build(move |_capabilities| {
863            let activator = scope.activator_for(operator_info.address);
864            let mut todo = std::collections::VecDeque::new();
865            move |_frontiers| {
866                let key = key_con.get(0);
867                let mut ok_output = ok_output.activate();
868
869                input.for_each(|time, data| {
870                    let cap = time.retain(0);
871                    for batch in data.iter() {
872                        todo.push_back(PendingWorkOk::new(
873                            cap.clone(),
874                            batch.cursor(),
875                            batch.clone(),
876                        ));
877                    }
878                });
879
880                // Activation-scoped decode storage; see `flat_map_core_fallible`.
881                let mut temp_storage = RowArena::new();
882                let mut datums = DatumVec::new();
883                let mut decode_logic =
884                    |k: BatchKey<'_, Tr>,
885                     v: BatchVal<'_, Tr>,
886                     t: T,
887                     d: mz_repr::Diff,
888                     ok_session: &mut Session<T, DCB>| {
889                        temp_storage.clear();
890                        let mut datums_borrow = datums.borrow();
891                        k.extend_datums(&temp_storage, &mut datums_borrow, Some(max_demand));
892                        let remaining = max_demand.saturating_sub(datums_borrow.len());
893                        v.extend_datums(&temp_storage, &mut datums_borrow, Some(remaining));
894                        logic(&mut datums_borrow, t, d, ok_session)
895                    };
896
897                let mut fuel = refuel;
898                while !todo.is_empty() && fuel > 0 {
899                    todo.front_mut().unwrap().do_work(
900                        key.as_ref(),
901                        &mut decode_logic,
902                        &mut fuel,
903                        &mut ok_output,
904                    );
905                    if fuel > 0 {
906                        todo.pop_front();
907                    }
908                }
909                if !todo.is_empty() {
910                    activator.activate();
911                }
912            }
913        });
914
915        ok_stream
916    }
917
918    /// Look up an arrangement by the expressions that form the key.
919    ///
920    /// The result may be `None` if no such arrangement exists, or it may be one of many
921    /// "arrangement flavors" that represent the types of arranged data we might have.
922    pub fn arrangement(&self, key: &[LirScalarExpr]) -> Option<ArrangementFlavor<'scope, T>> {
923        self.arranged.get(key).map(|x| x.clone())
924    }
925}
926
927impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
928    /// Presents `self` as a stream of updates, having been subjected to `mfp`.
929    ///
930    /// This operator is able to apply the logic of `mfp` early, which can substantially
931    /// reduce the amount of data produced when `mfp` is non-trivial.
932    ///
933    /// The `key_val` argument, when present, indicates that a specific arrangement should
934    /// be used, and if, in addition, the `val` component is present,
935    /// that we can seek to the supplied row.
936    pub fn as_collection_core(
937        &self,
938        mfp_plan: MfpPlan<LirScalarExpr>,
939        key_val: Option<(Vec<LirScalarExpr>, Option<StableRow>)>,
940        until: Antichain<mz_repr::Timestamp>,
941    ) -> (
942        ColCollection<'scope, T>,
943        VecCollection<'scope, T, DataflowErrorSer, Diff>,
944    ) {
945        // Unwrap the stable-serialization row wrapper, seeking works on
946        // plain rows.
947        let key_val = key_val.map(|(key, val)| (key, val.map(|val| val.0)));
948        // If the MFP is trivial, we can just return the collection.
949        // In the case that we weren't going to apply the `key_val` optimization,
950        // this path results in a slightly smaller and faster
951        // dataflow graph, and is intended to fix
952        // https://github.com/MaterializeInc/database-issues/issues/3111
953        let has_key_val = if let Some((_key, Some(_val))) = &key_val {
954            true
955        } else {
956            false
957        };
958
959        if mfp_plan.is_identity() && !has_key_val {
960            let key = key_val.map(|(k, _v)| k);
961            return match key {
962                // Unarranged identity hands the edge straight through, so a columnar
963                // producer stays columnar.
964                None => self
965                    .collection
966                    .clone()
967                    .expect("The unarranged collection doesn't exist."),
968                Some(key) => self.as_specific_collection(Some(&key)),
969            };
970        }
971
972        // Apply demand-based column pruning. We round-trip through MIR
973        // so temporal bounds are folded back as mz_now() predicates —
974        // this way demand() sees all column references (including those
975        // in temporal bounds), and permute_fn applies uniformly.
976        let (mfp_plan, max_demand) = {
977            let mut mir_mfp = mfp_plan_lir_to_mir(mfp_plan).into_map_filter_project();
978            let max_demand = mir_mfp.demand().last().map(|x| *x + 1).unwrap_or(0);
979            mir_mfp.permute_fn(|c| c, max_demand);
980            mir_mfp.optimize();
981            let plan = mfp_mir_to_lir_plan(mir_mfp);
982            (plan, max_demand)
983        };
984
985        let mut datum_vec = DatumVec::new();
986        // Wrap in an `Rc` so that lifetimes work out.
987        let until = std::rc::Rc::new(until);
988
989        // `ConsolidatingColumnBuilder` folds within-batch duplicates. It consolidates in
990        // place, so records are given owned, which costs nothing here because
991        // `mfp_plan.evaluate` already produces a fresh `Row` per result.
992        let (stream, errors) = self.flat_map::<ConsolidatingColumnBuilder<Row, T, Diff>, _>(
993            key_val,
994            max_demand,
995            move |row_datums, time, diff, ok_session, err_session| {
996                let mut row_builder = SharedRow::get();
997                let until = std::rc::Rc::clone(&until);
998                let temp_storage = RowArena::new();
999                let row_iter = row_datums.iter();
1000                let mut datums_local = datum_vec.borrow();
1001                datums_local.extend(row_iter);
1002                let event_time = time.event_time();
1003                let mut work: usize = 0;
1004                for result in mfp_plan.evaluate(
1005                    &mut datums_local,
1006                    &temp_storage,
1007                    event_time,
1008                    diff.clone(),
1009                    move |time| !until.less_equal(time),
1010                    &mut row_builder,
1011                ) {
1012                    work += 1;
1013                    match result {
1014                        Ok((row, event_time, diff)) => {
1015                            // Copy the whole time, and re-populate event time.
1016                            let mut time: T = time.clone();
1017                            *time.event_time_mut() = event_time;
1018                            ok_session.give((row, time, diff));
1019                        }
1020                        Err((e, event_time, diff)) => {
1021                            // Copy the whole time, and re-populate event time.
1022                            let mut time: T = time.clone();
1023                            *time.event_time_mut() = event_time;
1024                            err_session.give((e, time, diff));
1025                        }
1026                    }
1027                }
1028                work
1029            },
1030        );
1031
1032        (stream.as_collection(), errors)
1033    }
1034    pub fn ensure_collections(
1035        mut self,
1036        collections: AvailableCollections,
1037        input_key: Option<Vec<LirScalarExpr>>,
1038        input_mfp: MfpPlan<LirScalarExpr>,
1039        as_of: Antichain<mz_repr::Timestamp>,
1040        until: Antichain<mz_repr::Timestamp>,
1041        config_set: &ConfigSet,
1042        strategy: ArrangementStrategy,
1043    ) -> Self
1044    where
1045        T: MaybeBucketByTime,
1046    {
1047        if collections == Default::default() {
1048            return self;
1049        }
1050        // Cache collection to avoid reforming it each time.
1051        //
1052        // TODO(mcsherry): In theory this could be faster run out of another arrangement,
1053        // as the `map_fallible` that follows could be run against an arrangement itself.
1054        //
1055        // Note(btv): If we ever do that, we would then only need to make the raw collection here
1056        // if `collections.raw` is true.
1057
1058        for (key, _, _) in collections.arranged.iter() {
1059            soft_assert_or_log!(
1060                !self.arranged.contains_key(key),
1061                "LIR ArrangeBy tried to create an existing arrangement"
1062            );
1063        }
1064
1065        // Track whether we already applied temporal bucketing in this call, to
1066        // avoid bucketing the same updates twice.
1067        let mut bucketed = false;
1068
1069        // True iff at least one new arrangement will actually be built below. Bucketing only
1070        // pays off when something downstream merges/compacts the future-stamped updates; on a
1071        // pure raw collection (no new arrangement) the work is wasted.
1072        let will_create_arrangement = collections
1073            .arranged
1074            .iter()
1075            .any(|(key, _, _)| !self.arranged.contains_key(key));
1076
1077        // We need the collection if either (1) it is explicitly demanded, or (2) we are going to render any arrangement
1078        let form_raw_collection = collections.raw || will_create_arrangement;
1079        if form_raw_collection && self.collection.is_none() {
1080            let (oks, errs) =
1081                self.as_collection_core(input_mfp, input_key.map(|k| (k, None)), until);
1082            // Apply temporal bucketing when the lowering selected `TemporalBucketing` and
1083            // we will build at least one arrangement. This path fires when the collection
1084            // must be formed from scratch (e.g., from an arrangement via as_collection_core).
1085            let effective_strategy = if will_create_arrangement {
1086                strategy
1087            } else {
1088                ArrangementStrategy::Direct
1089            };
1090            let oks = if matches!(effective_strategy, ArrangementStrategy::TemporalBucketing)
1091                && ENABLE_COMPUTE_TEMPORAL_BUCKETING.get(config_set)
1092            {
1093                let summary: mz_repr::Timestamp = TEMPORAL_BUCKETING_SUMMARY
1094                    .get(config_set)
1095                    .try_into()
1096                    .expect("must fit");
1097                bucketed = true;
1098                // Temporal bucketing is columnar throughout, so no round trip here.
1099                T::maybe_apply_temporal_bucketing(oks.inner, as_of.clone(), summary)
1100            } else {
1101                oks
1102            };
1103            self.collection = Some((oks, errs));
1104        }
1105        for (key, _, thinning) in collections.arranged {
1106            if !self.arranged.contains_key(&key) {
1107                // TODO: Consider allowing more expressive names.
1108                let name = format!("ArrangeBy[{:?}]", key);
1109
1110                let (oks, errs) = self
1111                    .collection
1112                    .take()
1113                    .expect("Collection constructed above");
1114                // Apply temporal bucketing if the collection already existed on
1115                // the bundle (e.g., from an upstream temporal Mfp or Get) and we
1116                // haven't bucketed yet. This is the common path for temporal-MFP
1117                // → ArrangeBy flows.
1118                let effective_strategy = if bucketed {
1119                    ArrangementStrategy::Direct
1120                } else {
1121                    strategy
1122                };
1123                let oks = if matches!(effective_strategy, ArrangementStrategy::TemporalBucketing)
1124                    && ENABLE_COMPUTE_TEMPORAL_BUCKETING.get(config_set)
1125                {
1126                    let summary: mz_repr::Timestamp = TEMPORAL_BUCKETING_SUMMARY
1127                        .get(config_set)
1128                        .try_into()
1129                        .expect("must fit");
1130                    bucketed = true;
1131                    // Temporal bucketing is columnar throughout, so no round trip here.
1132                    T::maybe_apply_temporal_bucketing(oks.inner, as_of.clone(), summary)
1133                } else {
1134                    oks
1135                };
1136                let batcher = ArrangementBatcher::from_config(config_set);
1137                let (oks, errs_keyed, passthrough) =
1138                    Self::arrange_collection(&name, oks, key.clone(), thinning.clone(), batcher);
1139                let errs_concat: KeyCollection<_, _, _> = errs.clone().concat(errs_keyed).into();
1140                self.collection = Some((passthrough, errs));
1141                let errs =
1142                    errs_concat.mz_arrange::<
1143                        ColumnationChunker<_>,
1144                        ErrBatcher<_, _>,
1145                        ErrBuilder<_, _>,
1146                        ErrSpine<_, _>,
1147                    >(
1148                        &format!("{}-errors", name),
1149                    );
1150                self.arranged
1151                    .insert(key, ArrangementFlavor::Local(oks, errs));
1152            }
1153        }
1154        self
1155    }
1156
1157    /// Builds an arrangement from a collection, using the specified key and value thinning.
1158    ///
1159    /// The arrangement's key is based on the `key` expressions, and the value the input with
1160    /// the `thinning` applied to it. It selects which of the input columns are included in the
1161    /// value of the arrangement. The thinning is in support of permuting arrangements such that
1162    /// columns in the key are not included in the value.
1163    ///
1164    /// In addition to the ok and err streams, we produce a passthrough stream that forwards
1165    /// the input as-is, which allows downstream consumers to reuse the collection without
1166    /// teeing the stream.
1167    fn arrange_collection(
1168        name: &String,
1169        oks: ColCollection<'scope, T>,
1170        key: Vec<LirScalarExpr>,
1171        thinning: Vec<usize>,
1172        batcher: ArrangementBatcher,
1173    ) -> (
1174        Arranged<'scope, RowRowAgent<T, Diff>>,
1175        VecCollection<'scope, T, DataflowErrorSer, Diff>,
1176        ColCollection<'scope, T>,
1177    ) {
1178        // Spelled out rather than `map_fallible`, whose closure cannot return the references
1179        // a columnar stream is pushed from.
1180        //
1181        // The arena is per-activation and cleared per row, so the capacity it retains never
1182        // outlives a scheduling invocation.
1183        let (ok_stream, err_stream, passthrough) = {
1184            let mut builder =
1185                OperatorBuilder::new("FormArrangementKey".to_string(), oks.inner.scope());
1186            let (ok_output, ok_stream) = builder.new_output();
1187            let mut ok_output =
1188                OutputBuilder::<_, ColumnBuilder<((Row, Row), T, Diff)>>::from(ok_output);
1189            let (err_output, err_stream) = builder.new_output();
1190            let mut err_output = OutputBuilder::from(err_output);
1191            let (passthrough_output, passthrough_stream) = builder.new_output();
1192            let mut passthrough_output =
1193                OutputBuilder::<_, NoopBuilder<Column<(Row, T, Diff)>>>::from(passthrough_output);
1194            let mut input = builder.new_input(oks.inner, Pipeline);
1195            builder.set_notify_for(0, FrontierInterest::Never);
1196            builder.build(move |_capabilities| {
1197                let mut key_buf = Row::default();
1198                let mut val_buf = Row::default();
1199                let mut datums = DatumVec::new();
1200                move |_frontiers| {
1201                    // Scoped to the activation so the arena's retained capacity does not
1202                    // outlive a single scheduling invocation; cleared per row to reuse it
1203                    // within the batch.
1204                    let mut temp_storage = RowArena::new();
1205                    let mut ok_output = ok_output.activate();
1206                    let mut err_output = err_output.activate();
1207                    let mut passthrough_output = passthrough_output.activate();
1208                    input.for_each(|time, data| {
1209                        let mut ok_session = ok_output.session_with_builder(&time);
1210                        let mut err_session = err_output.session(&time);
1211                        // Rows are read from the borrowed column, never materialized as
1212                        // owned `Row`s. Times and diffs are owned only on the error path.
1213                        for (row, t, d) in data.borrow().into_index_iter() {
1214                            temp_storage.clear();
1215                            let datums = datums.borrow_with(row);
1216                            let key_iter = key.iter().map(|k| k.eval(&datums, &temp_storage));
1217                            match key_buf.packer().try_extend(key_iter) {
1218                                Ok(()) => {
1219                                    let val_datum_iter = thinning.iter().map(|c| datums[*c]);
1220                                    val_buf.packer().extend(val_datum_iter);
1221                                    ok_session.give(((&*key_buf, &*val_buf), t, d));
1222                                }
1223                                Err(e) => {
1224                                    err_session.give((
1225                                        e.into(),
1226                                        Columnar::into_owned(t),
1227                                        Columnar::into_owned(d),
1228                                    ));
1229                                }
1230                            }
1231                        }
1232                        passthrough_output
1233                            .session_with_builder(&time)
1234                            .give_container(data);
1235                    });
1236                }
1237            });
1238            (ok_stream, err_stream, passthrough_stream.as_collection())
1239        };
1240
1241        let exchange =
1242            ExchangeCore::<ColumnBuilder<_>, _>::new_core(columnar_exchange::<Row, Row, T, Diff>);
1243        let oks = match batcher {
1244            ArrangementBatcher::Chunked => ok_stream.mz_arrange_core::<
1245                _,
1246                ChunkChunker<(Row, Row), T, Diff>,
1247                AccountedChunkBatcher<(Row, Row), T, Diff>,
1248                UnchunkBuilder<RowRowColPagedBuilder<T, Diff>, (Row, Row), T, Diff>,
1249                RowRowSpine<_, _>,
1250            >(exchange, name),
1251            ArrangementBatcher::Columnar => ok_stream.mz_arrange_core::<
1252                _,
1253                batcher::ColumnChunker<_>,
1254                Col2ValColBatcher<_, _, _, _>,
1255                RowRowColPagedBuilder<_, _>,
1256                RowRowSpine<_, _>,
1257            >(exchange, name),
1258            ArrangementBatcher::Columnation => ok_stream.mz_arrange_core::<
1259                _,
1260                batcher::Chunker<_>,
1261                Col2ValBatcher<_, _, _, _>,
1262                RowRowBuilder<_, _>,
1263                RowRowSpine<_, _>,
1264            >(exchange, name),
1265        };
1266        (oks, err_stream.as_collection(), passthrough)
1267    }
1268}
1269
1270/// Type alias for a timely output `Session` whose capability is a `Capability<T>`. The container
1271/// builder `CB` is left to the caller; sessions can therefore drive consolidating, capacity, or
1272/// (in the future) columnar output builders without changing call sites.
1273pub(crate) type Session<'a, 'b, T, CB> =
1274    timely::dataflow::operators::generic::Session<'a, 'b, T, CB, Capability<T>>;
1275
1276/// Container builder used for the err output of every flat_map variant. Pre-refactor the
1277/// merged Ok/Err stream flowed through a [`ConsolidatingContainerBuilder`] before the
1278/// `map_fallible` demux split it; we preserve that consolidation here so errors with the
1279/// same `(error, time)` cancel within a batch rather than propagating to downstream.
1280pub(crate) type ECB<T> = ConsolidatingContainerBuilder<Vec<(DataflowErrorSer, T, Diff)>>;
1281
1282/// Number of output records the arrangement flat_map operators may produce before yielding.
1283/// See [`ArrangementFlavor::flat_map`] for the fuel rationale; the constant is a pragmatic
1284/// compromise and not tuned empirically.
1285const REFUEL: usize = 1_000_000;
1286
1287struct PendingWork<C>
1288where
1289    C: Cursor,
1290{
1291    /// Capability for the `ok` output (output port 0).
1292    ok_capability: Capability<C::Time>,
1293    /// Capability for the `err` output (output port 1).
1294    err_capability: Capability<C::Time>,
1295    cursor: C,
1296    batch: C::Storage,
1297}
1298
1299impl<C> PendingWork<C>
1300where
1301    C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1302{
1303    /// Create a new bundle of pending work, from a pair of capabilities (one per output),
1304    /// a cursor, and backing storage.
1305    fn new(
1306        ok_capability: Capability<C::Time>,
1307        err_capability: Capability<C::Time>,
1308        cursor: C,
1309        batch: C::Storage,
1310    ) -> Self {
1311        Self {
1312            ok_capability,
1313            err_capability,
1314            cursor,
1315            batch,
1316        }
1317    }
1318    /// Perform roughly `fuel` work through the cursor, applying `logic` and sending results to
1319    /// the two output sessions.
1320    fn do_work<DCB, L>(
1321        &mut self,
1322        key: Option<&C::Key<'_>>,
1323        logic: &mut L,
1324        fuel: &mut usize,
1325        ok_output: &mut OutputBuilderSession<'_, C::Time, DCB>,
1326        err_output: &mut OutputBuilderSession<'_, C::Time, ECB<C::Time>>,
1327    ) where
1328        DCB: ContainerBuilder,
1329        L: FnMut(
1330            C::Key<'_>,
1331            C::Val<'_>,
1332            C::Time,
1333            C::Diff,
1334            &mut Session<C::Time, DCB>,
1335            &mut Session<C::Time, ECB<C::Time>>,
1336        ) -> usize,
1337    {
1338        let mut ok_session = ok_output.session_with_builder(&self.ok_capability);
1339        let mut err_session = err_output.session_with_builder(&self.err_capability);
1340        walk_cursor(&mut self.cursor, &self.batch, key, fuel, |k, v, t, d| {
1341            logic(k, v, t, d, &mut ok_session, &mut err_session)
1342        });
1343    }
1344}
1345
1346/// Pending work for the Ok-only variant of `flat_map_core_fallible`. Holds a single capability since
1347/// the operator has only one output port.
1348struct PendingWorkOk<C>
1349where
1350    C: Cursor,
1351{
1352    capability: Capability<C::Time>,
1353    cursor: C,
1354    batch: C::Storage,
1355}
1356
1357impl<C> PendingWorkOk<C>
1358where
1359    C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1360{
1361    fn new(capability: Capability<C::Time>, cursor: C, batch: C::Storage) -> Self {
1362        Self {
1363            capability,
1364            cursor,
1365            batch,
1366        }
1367    }
1368
1369    /// Perform roughly `fuel` work through the cursor, applying `logic` and sending results to
1370    /// the single output session.
1371    fn do_work<DCB, L>(
1372        &mut self,
1373        key: Option<&C::Key<'_>>,
1374        logic: &mut L,
1375        fuel: &mut usize,
1376        ok_output: &mut OutputBuilderSession<'_, C::Time, DCB>,
1377    ) where
1378        DCB: ContainerBuilder,
1379        L: FnMut(C::Key<'_>, C::Val<'_>, C::Time, C::Diff, &mut Session<C::Time, DCB>) -> usize,
1380    {
1381        let mut ok_session = ok_output.session_with_builder(&self.capability);
1382        walk_cursor(&mut self.cursor, &self.batch, key, fuel, |k, v, t, d| {
1383            logic(k, v, t, d, &mut ok_session)
1384        });
1385    }
1386}
1387
1388/// Walk a cursor, calling `emit` for each consolidated `(key, val, time, diff)` tuple. If
1389/// `key` is set, the cursor is seeked to it and only values for that key are produced.
1390///
1391/// `emit` returns the number of records it produced for the given input tuple. The cursor
1392/// stops as soon as the accumulated emit count reaches `*fuel`, leaving the cursor in place
1393/// so work can resume on a later call. Within a batch, both the inner val loop and the
1394/// outer key loop are bounded only by emit count, so selective filters (`emit` returns 0)
1395/// run to batch completion in a single activation — see [`ArrangementFlavor::flat_map`]
1396/// for why fuel counts output rather than input.
1397fn walk_cursor<C, F>(
1398    cursor: &mut C,
1399    batch: &C::Storage,
1400    key: Option<&C::Key<'_>>,
1401    fuel: &mut usize,
1402    mut emit: F,
1403) where
1404    C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1405    F: FnMut(C::Key<'_>, C::Val<'_>, C::Time, C::Diff) -> usize,
1406{
1407    use differential_dataflow::consolidation::consolidate;
1408
1409    let mut work: usize = 0;
1410    let mut buffer = Vec::new();
1411    if let Some(key) = key {
1412        let key = C::KeyContainer::reborrow(*key);
1413        if cursor.get_key(batch).map(|k| k == key) != Some(true) {
1414            cursor.seek_key(batch, key);
1415        }
1416        if cursor.get_key(batch).map(|k| k == key) == Some(true) {
1417            let key = cursor.key(batch);
1418            while let Some(val) = cursor.get_val(batch) {
1419                cursor.map_times(batch, |time, diff| {
1420                    buffer.push((C::owned_time(time), C::owned_diff(diff)));
1421                });
1422                consolidate(&mut buffer);
1423                for (time, diff) in buffer.drain(..) {
1424                    work += emit(key, val, time, diff);
1425                }
1426                cursor.step_val(batch);
1427                if work >= *fuel {
1428                    *fuel = 0;
1429                    return;
1430                }
1431            }
1432        }
1433    } else {
1434        while let Some(key) = cursor.get_key(batch) {
1435            while let Some(val) = cursor.get_val(batch) {
1436                cursor.map_times(batch, |time, diff| {
1437                    buffer.push((C::owned_time(time), C::owned_diff(diff)));
1438                });
1439                consolidate(&mut buffer);
1440                for (time, diff) in buffer.drain(..) {
1441                    work += emit(key, val, time, diff);
1442                }
1443                cursor.step_val(batch);
1444                if work >= *fuel {
1445                    *fuel = 0;
1446                    return;
1447                }
1448            }
1449            cursor.step_key(batch);
1450        }
1451    }
1452    *fuel -= work;
1453}
1454
1455#[cfg(test)]
1456mod tests {
1457    use differential_dataflow::input::Input;
1458    use mz_expr::{EvalError, MapFilterProject};
1459    use mz_repr::{Datum, ReprScalarType, Timestamp};
1460    use timely::dataflow::operators::Capture;
1461    use timely::dataflow::operators::capture::{Event, Extract};
1462
1463    use super::*;
1464    use crate::render::columnar::{columnar_to_vec, vec_to_columnar};
1465
1466    type OkUpdate = ((Row, Row), Timestamp, Diff);
1467    type ErrUpdate = (DataflowErrorSer, Timestamp, Diff);
1468    type Captured<D> = std::sync::mpsc::Receiver<Event<Timestamp, Vec<D>>>;
1469
1470    fn extract_ok(captured: Captured<OkUpdate>) -> Vec<OkUpdate> {
1471        let mut updates: Vec<_> = captured
1472            .extract()
1473            .into_iter()
1474            .flat_map(|(_, data)| data)
1475            .collect();
1476        updates.sort();
1477        updates
1478    }
1479
1480    // `DataflowErrorSer` is not `Ord`, so order by the error's debug string.
1481    fn extract_err(captured: Captured<ErrUpdate>) -> Vec<(String, Timestamp, Diff)> {
1482        let mut updates: Vec<_> = captured
1483            .extract()
1484            .into_iter()
1485            .flat_map(|(_, data)| data)
1486            .map(|(e, t, d)| (format!("{e:?}"), t, d))
1487            .collect();
1488        updates.sort();
1489        updates
1490    }
1491
1492    /// Arranges `rows`, returning the sorted ok and err output.
1493    fn arrange_columnar(
1494        rows: Vec<(Row, u64)>,
1495        key: Vec<LirScalarExpr>,
1496    ) -> (Vec<OkUpdate>, Vec<(String, Timestamp, Diff)>) {
1497        let thinning = vec![0, 1];
1498        let (ok, err) = timely::execute_directly(move |worker| {
1499            worker.dataflow::<Timestamp, _, _>(|scope| {
1500                let (mut input, collection) = scope.new_collection();
1501                let (arranged, errs, _passthrough) =
1502                    CollectionBundle::<Timestamp>::arrange_collection(
1503                        &"col".to_string(),
1504                        vec_to_columnar(collection),
1505                        key,
1506                        thinning,
1507                        ArrangementBatcher::Columnation,
1508                    );
1509                let ok = arranged
1510                    .as_collection(|k, v| (k.to_row(), v.to_row()))
1511                    .inner
1512                    .capture();
1513                let err = errs.inner.capture();
1514
1515                let max_time = rows.iter().map(|(_, t)| *t).max().unwrap_or(0);
1516                for (row, time) in rows {
1517                    input.update_at(row, Timestamp::from(time), Diff::ONE);
1518                }
1519                input.advance_to(Timestamp::from(max_time + 1));
1520                input.flush();
1521                (ok, err)
1522            })
1523        });
1524        (extract_ok(ok), extract_err(err))
1525    }
1526
1527    // Uniform two columns so `Column(0)` and full-row thinning are in bounds, over three
1528    // distinct times so per-record time handling is exercised.
1529    fn test_rows() -> Vec<(Row, u64)> {
1530        vec![
1531            (Row::pack_slice(&[Datum::Int32(1), Datum::String("a")]), 0),
1532            (Row::pack_slice(&[Datum::Int32(2), Datum::String("b")]), 1),
1533            (Row::pack_slice(&[Datum::Int32(1), Datum::String("a")]), 1),
1534            (Row::pack_slice(&[Datum::Int32(3), Datum::Null]), 2),
1535        ]
1536    }
1537
1538    /// Agreeing contents do not rule out a silent `columnar_to_vec` on the ok path. That the
1539    /// operator never decodes holds by inspection, not by this test.
1540    #[mz_ore::test]
1541    fn arrange_collection_keys_correctly() {
1542        let rows = test_rows();
1543        let mut expected: Vec<OkUpdate> = rows
1544            .iter()
1545            .map(|(row, t)| {
1546                let key = Row::pack_slice(&[row.iter().next().unwrap()]);
1547                ((key, row.clone()), Timestamp::from(*t), Diff::ONE)
1548            })
1549            .collect();
1550        expected.sort();
1551
1552        let (ok, err) = arrange_columnar(rows, vec![LirScalarExpr::column(0)]);
1553        assert_eq!(ok, expected);
1554        assert!(err.is_empty());
1555    }
1556
1557    /// A key expression that always errors drives every record onto the error path.
1558    #[mz_ore::test]
1559    fn arrange_collection_error_path() {
1560        let key = vec![LirScalarExpr::literal(
1561            Err(EvalError::DivisionByZero),
1562            ReprScalarType::Int32,
1563        )];
1564        let (ok, err) = arrange_columnar(test_rows(), key);
1565
1566        assert!(ok.is_empty());
1567        assert!(!err.is_empty());
1568    }
1569
1570    /// The passthrough forwards every input record, including those whose key
1571    /// evaluation errored, so a consumer of the bundle's collection sees the
1572    /// unarranged input rather than the ok side of the arrangement.
1573    #[mz_ore::test]
1574    fn arrange_collection_passthrough_forwards_input() {
1575        let rows = test_rows();
1576        let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1577            .iter()
1578            .map(|(row, t)| (row.clone(), Timestamp::from(*t), Diff::ONE))
1579            .collect();
1580        expected.sort();
1581
1582        let key = vec![LirScalarExpr::literal(
1583            Err(EvalError::DivisionByZero),
1584            ReprScalarType::Int32,
1585        )];
1586        let captured = timely::execute_directly(move |worker| {
1587            worker.dataflow::<Timestamp, _, _>(|scope| {
1588                let (mut input, collection) = scope.new_collection();
1589                let (_arranged, _errs, passthrough) =
1590                    CollectionBundle::<Timestamp>::arrange_collection(
1591                        &"col".to_string(),
1592                        vec_to_columnar(collection),
1593                        key,
1594                        vec![0, 1],
1595                        ArrangementBatcher::Columnation,
1596                    );
1597                let captured = columnar_to_vec(passthrough).inner.capture();
1598
1599                let max_time = rows.iter().map(|(_, t)| *t).max().unwrap_or(0);
1600                for (row, time) in rows {
1601                    input.update_at(row, Timestamp::from(time), Diff::ONE);
1602                }
1603                input.advance_to(Timestamp::from(max_time + 1));
1604                input.flush();
1605                captured
1606            })
1607        });
1608
1609        assert_eq!(extract_row_updates(captured), expected);
1610    }
1611
1612    fn extract_row_updates(
1613        captured: Captured<(Row, Timestamp, Diff)>,
1614    ) -> Vec<(Row, Timestamp, Diff)> {
1615        let mut updates: Vec<_> = captured
1616            .extract()
1617            .into_iter()
1618            .flat_map(|(_, data)| data)
1619            .collect();
1620        updates.sort();
1621        updates
1622    }
1623
1624    /// Arrange correctness itself is covered by `arrange_collection_arms_agree`; this only
1625    /// asserts the columnar variant survives the hand-off.
1626    #[mz_ore::test]
1627    fn get_arrange_by_produces_projected_rows() {
1628        let rows = vec![
1629            (Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]), 0u64),
1630            (Row::pack_slice(&[Datum::Int64(2), Datum::Int64(20)]), 1),
1631            (Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]), 1),
1632        ];
1633        // A projection is non-identity, so `as_collection_core` takes the producer path
1634        // rather than the identity passthrough.
1635        let mfp = MapFilterProject::<LirScalarExpr>::new(2)
1636            .project(vec![0])
1637            .into_plan()
1638            .expect("mfp");
1639        let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1640            .iter()
1641            .map(|(r, t)| {
1642                let col0 = r.iter().next().unwrap();
1643                (Row::pack_slice(&[col0]), Timestamp::from(*t), Diff::ONE)
1644            })
1645            .collect();
1646        expected.sort();
1647
1648        let produced = timely::execute_directly(move |worker| {
1649            worker.dataflow::<Timestamp, _, _>(|scope| {
1650                let (mut input, collection) = scope.new_collection();
1651                let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1652                let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1653                let (edge, _errs) = bundle.as_collection_core(mfp, None, Antichain::new());
1654                let produced = columnar_to_vec(edge.clone()).inner.capture();
1655                let (_arranged, _arrange_errs, _passthrough) =
1656                    CollectionBundle::<Timestamp>::arrange_collection(
1657                        &"arrange".to_string(),
1658                        edge,
1659                        vec![LirScalarExpr::column(0)],
1660                        vec![0],
1661                        ArrangementBatcher::Columnation,
1662                    );
1663
1664                let max_time = rows.iter().map(|(_, t)| *t).max().unwrap();
1665                for (row, time) in rows {
1666                    input.update_at(row, Timestamp::from(time), Diff::ONE);
1667                }
1668                input.advance_to(Timestamp::from(max_time + 1));
1669                input.flush();
1670                produced
1671            })
1672        });
1673
1674        assert_eq!(extract_row_updates(produced), expected);
1675    }
1676
1677    /// The identity fast-path returns the unarranged input edge unchanged, with no
1678    /// repack, so its contents pass straight through.
1679    #[mz_ore::test]
1680    fn as_collection_core_identity_passes_edge_through() {
1681        let expected = vec![(
1682            Row::pack_slice(&[Datum::Int64(1)]),
1683            Timestamp::from(0u64),
1684            Diff::ONE,
1685        )];
1686        let captured = timely::execute_directly(move |worker| {
1687            worker.dataflow::<Timestamp, _, _>(|scope| {
1688                let (mut input, collection) = scope.new_collection::<Row, Diff>();
1689                let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1690                let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1691                let identity = MapFilterProject::<LirScalarExpr>::new(1)
1692                    .into_plan()
1693                    .expect("identity mfp");
1694                let (out, _errs) = bundle.as_collection_core(identity, None, Antichain::new());
1695                let captured = columnar_to_vec(out).inner.capture();
1696                input.update_at(
1697                    Row::pack_slice(&[Datum::Int64(1)]),
1698                    Timestamp::from(0u64),
1699                    Diff::ONE,
1700                );
1701                input.advance_to(Timestamp::from(1u64));
1702                input.flush();
1703                captured
1704            })
1705        });
1706        assert_eq!(extract_row_updates(captured), expected);
1707    }
1708
1709    /// Input rows that project to the same output row at the same time collapse to one
1710    /// record with summed diff. A plain `ColumnBuilder` would emit both.
1711    #[mz_ore::test]
1712    fn as_collection_core_consolidates_within_batch() {
1713        // `[1, 10]` and `[1, 20]` both project (dropping column 1) to `[1]` at
1714        // t=0, so their `+1` diffs must fold to `+2`. `[2, 30]` projects to a
1715        // distinct record.
1716        let rows = vec![
1717            Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]),
1718            Row::pack_slice(&[Datum::Int64(1), Datum::Int64(20)]),
1719            Row::pack_slice(&[Datum::Int64(2), Datum::Int64(30)]),
1720        ];
1721        let mfp = MapFilterProject::<LirScalarExpr>::new(2)
1722            .project(vec![0])
1723            .into_plan()
1724            .expect("mfp");
1725        let expected = vec![
1726            (
1727                Row::pack_slice(&[Datum::Int64(1)]),
1728                Timestamp::from(0u64),
1729                Diff::from(2),
1730            ),
1731            (
1732                Row::pack_slice(&[Datum::Int64(2)]),
1733                Timestamp::from(0u64),
1734                Diff::ONE,
1735            ),
1736        ];
1737
1738        let captured = timely::execute_directly(move |worker| {
1739            worker.dataflow::<Timestamp, _, _>(|scope| {
1740                let (mut input, collection) = scope.new_collection();
1741                let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1742                let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1743                let (edge, _errs) = bundle.as_collection_core(mfp, None, Antichain::new());
1744                let captured = columnar_to_vec(edge).inner.capture();
1745                // Feed all rows at the same time in one batch so the fold is
1746                // within-batch, not a downstream re-consolidation.
1747                for row in rows {
1748                    input.update_at(row, Timestamp::from(0u64), Diff::ONE);
1749                }
1750                input.advance_to(Timestamp::from(1u64));
1751                input.flush();
1752                captured
1753            })
1754        });
1755
1756        assert_eq!(extract_row_updates(captured), expected);
1757    }
1758
1759    /// Keying by column 0 and thinning the value to column 1 reconstructs the original
1760    /// two-column row. The decode below belongs to the capture harness.
1761    #[mz_ore::test]
1762    fn as_specific_collection_materializes_columnar() {
1763        let rows = test_rows();
1764        let key = vec![LirScalarExpr::column(0)];
1765        let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1766            .iter()
1767            .map(|(r, t)| (r.clone(), Timestamp::from(*t), Diff::ONE))
1768            .collect();
1769        expected.sort();
1770
1771        let captured = timely::execute_directly(move |worker| {
1772            worker.dataflow::<Timestamp, _, _>(|scope| {
1773                let (mut input, collection) = scope.new_collection();
1774                let (arranged, arr_errs, _passthrough) =
1775                    CollectionBundle::<Timestamp>::arrange_collection(
1776                        &"agg".to_string(),
1777                        vec_to_columnar(collection),
1778                        key.clone(),
1779                        vec![1],
1780                        ArrangementBatcher::Columnation,
1781                    );
1782                let err_arranged = {
1783                    let kc: KeyCollection<_, _, _> = arr_errs.into();
1784                    kc.mz_arrange::<
1785                        ColumnationChunker<_>,
1786                        ErrBatcher<_, _>,
1787                        ErrBuilder<_, _>,
1788                        ErrSpine<_, _>,
1789                    >("agg-errs")
1790                };
1791                // An arrangement-only bundle, as Reduce/Threshold/TopK produce.
1792                let bundle = CollectionBundle::from_columns(
1793                    0..1,
1794                    ArrangementFlavor::Local(arranged, err_arranged),
1795                );
1796                let (edge, _errs) = bundle.as_specific_collection(Some(&key));
1797                let captured = columnar_to_vec(edge).inner.capture();
1798
1799                let max_time = rows.iter().map(|(_, t)| *t).max().unwrap();
1800                for (row, time) in rows {
1801                    input.update_at(row, Timestamp::from(time), Diff::ONE);
1802                }
1803                input.advance_to(Timestamp::from(max_time + 1));
1804                input.flush();
1805                captured
1806            })
1807        });
1808
1809        assert_eq!(extract_row_updates(captured), expected);
1810    }
1811}