Skip to main content

mz_timely_util/columnar/
batcher.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License in the LICENSE file at the
6// root of this repository, or online at
7//
8//     http://www.apache.org/licenses/LICENSE-2.0
9//
10// Unless required by applicable law or agreed to in writing, software
11// distributed under the License is distributed on an "AS IS" BASIS,
12// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13// See the License for the specific language governing permissions and
14// limitations under the License.
15
16//! Types for consolidating, merging, and extracting columnar update collections.
17
18use std::collections::VecDeque;
19use std::marker::PhantomData;
20
21use crate::columnation::ColumnationStack;
22use columnar::Container as _;
23use columnar::Push as _;
24use columnar::bytes::indexed;
25use columnar::{BorrowedOf, Clear, Columnar, Index, Len};
26use columnation::Columnation;
27use differential_dataflow::difference::Semigroup;
28use differential_dataflow::trace::implementations::merge_batcher::Merger;
29use timely::Container;
30use timely::PartialOrder;
31use timely::container::{ContainerBuilder, PushInto};
32use timely::progress::frontier::{Antichain, AntichainRef};
33
34use crate::columnar::Column;
35use crate::columnar::body::ColumnBody;
36
37/// A chunker to transform input data into sorted columns.
38#[derive(Default)]
39pub struct Chunker<C> {
40    /// Buffer into which we'll consolidate.
41    ///
42    /// Also the buffer where we'll stage responses to `extract` and `finish`.
43    /// When these calls return, the buffer is available for reuse.
44    target: C,
45    /// Consolidated buffers ready to go.
46    ready: VecDeque<C>,
47}
48
49impl<C: Container + Clone + 'static> ContainerBuilder for Chunker<C> {
50    type Container = C;
51
52    fn extract(&mut self) -> Option<&mut Self::Container> {
53        if let Some(ready) = self.ready.pop_front() {
54            self.target = ready;
55            Some(&mut self.target)
56        } else {
57            None
58        }
59    }
60
61    fn finish(&mut self) -> Option<&mut Self::Container> {
62        self.extract()
63    }
64}
65
66impl<'a, D, T, R> PushInto<&'a mut Column<(D, T, R)>> for Chunker<ColumnationStack<(D, T, R)>>
67where
68    D: Columnar + Columnation,
69    for<'b> columnar::Ref<'b, D>: Ord + Copy,
70    T: Columnar + Columnation,
71    for<'b> columnar::Ref<'b, T>: Ord + Copy,
72    R: Columnar + Columnation + Semigroup + for<'b> Semigroup<columnar::Ref<'b, R>>,
73    for<'b> columnar::Ref<'b, R>: Ord,
74{
75    fn push_into(&mut self, container: &'a mut Column<(D, T, R)>) {
76        // Sort input data
77        // TODO: consider `Vec<usize>` that we retain, containing indexes.
78        let borrowed = container.borrow();
79        let mut permutation = Vec::with_capacity(borrowed.len());
80        Extend::extend(&mut permutation, borrowed.into_index_iter());
81        permutation.sort();
82
83        self.target.clear();
84        // Iterate over the data, accumulating diffs for like keys.
85        let mut iter = permutation.drain(..);
86        if let Some((data, time, diff)) = iter.next() {
87            let mut owned_data = D::into_owned(data);
88            let mut owned_time = T::into_owned(time);
89
90            let mut prev_data = data;
91            let mut prev_time = time;
92            let mut prev_diff = <R as Columnar>::into_owned(diff);
93
94            for (data, time, diff) in iter {
95                if (&prev_data, &prev_time) == (&data, &time) {
96                    prev_diff.plus_equals(&diff);
97                } else {
98                    if !prev_diff.is_zero() {
99                        D::copy_from(&mut owned_data, prev_data);
100                        T::copy_from(&mut owned_time, prev_time);
101                        let tuple = (owned_data, owned_time, prev_diff);
102                        self.target.push_into(&tuple);
103                        (owned_data, owned_time, prev_diff) = tuple;
104                    }
105                    prev_data = data;
106                    prev_time = time;
107                    R::copy_from(&mut prev_diff, diff);
108                }
109            }
110
111            if !prev_diff.is_zero() {
112                D::copy_from(&mut owned_data, prev_data);
113                T::copy_from(&mut owned_time, prev_time);
114                let tuple = (owned_data, owned_time, prev_diff);
115                self.target.push_into(&tuple);
116            }
117        }
118
119        if !self.target.is_empty() {
120            self.ready.push_back(std::mem::take(&mut self.target));
121        }
122    }
123}
124
125/// A chunker that consolidates `Column<(D, T, R)>` updates into sorted
126/// [`ColumnBody`] chunks, without round-tripping through columnation.
127///
128/// Drop-in counterpart to [`Chunker`] for the merge-batcher path: same control
129/// flow (sort borrowed refs, fold equal `(data, time)` runs, drop zero diffs).
130/// This is where records leave the edge container: the input is whatever the
131/// edge delivered, the output is a body the batcher chains.
132pub struct ColumnChunker<U: Columnar> {
133    /// Body we consolidate into and present to extract/finish callers.
134    target: ColumnBody<U>,
135    /// Sorted, consolidated chunks pending extraction.
136    ready: VecDeque<ColumnBody<U>>,
137}
138
139// Manual impl rather than `#[derive(Default)]`: the derive would synthesize
140// `impl<U: Columnar + Default>`, but `ColumnBody<U>: Default` only requires
141// `U: Columnar`, and adding a spurious `U: Default` bound would propagate
142// through every `ContainerBuilder for ColumnChunker<U>` impl.
143impl<U: Columnar> Default for ColumnChunker<U> {
144    fn default() -> Self {
145        Self {
146            target: ColumnBody::default(),
147            ready: VecDeque::new(),
148        }
149    }
150}
151
152impl<U: Columnar> ContainerBuilder for ColumnChunker<U>
153where
154    U::Container: Clone + 'static,
155{
156    type Container = ColumnBody<U>;
157
158    fn extract(&mut self) -> Option<&mut Self::Container> {
159        if let Some(ready) = self.ready.pop_front() {
160            self.target = ready;
161            Some(&mut self.target)
162        } else {
163            None
164        }
165    }
166
167    fn finish(&mut self) -> Option<&mut Self::Container> {
168        self.extract()
169    }
170}
171
172impl<'a, D, T, R> PushInto<&'a mut Column<(D, T, R)>> for ColumnChunker<(D, T, R)>
173where
174    D: Columnar,
175    for<'b> columnar::Ref<'b, D>: Copy + Ord,
176    T: Columnar,
177    for<'b> columnar::Ref<'b, T>: Copy + Ord,
178    R: Columnar + Default + Semigroup + for<'b> Semigroup<columnar::Ref<'b, R>>,
179    for<'b> columnar::Ref<'b, R>: Ord,
180    for<'b> <D as Columnar>::Container: columnar::Push<columnar::Ref<'b, D>>,
181    for<'b> <T as Columnar>::Container: columnar::Push<columnar::Ref<'b, T>>,
182    for<'b> <R as Columnar>::Container: columnar::Push<&'b R>,
183{
184    fn push_into(&mut self, container: &'a mut Column<(D, T, R)>) {
185        // Clearing keeps the allocations of a typed target, which is the
186        // steady state: possibly a chunk just handed back via `extract`.
187        self.target.clear();
188
189        // Sort input by columnar ref order.
190        let borrowed = container.borrow();
191        let mut permutation = Vec::with_capacity(borrowed.len());
192        Extend::extend(&mut permutation, borrowed.into_index_iter());
193        permutation.sort();
194
195        // Sweep sorted refs, accumulating diffs over equal `(data, time)`
196        // pairs and pushing non-zero results to the target's leaves. Refs
197        // from the input borrow are valid through the sweep, so D and T
198        // are pushed directly via each leaf's `Push<Ref<_>>` impl. Only R
199        // needs an owned scratch since it carries the consolidated sum.
200        {
201            let (target_d, target_t, target_r) = self.target.typed_mut();
202
203            let mut iter = permutation.drain(..);
204            if let Some((data, time, diff)) = iter.next() {
205                let mut prev_data = data;
206                let mut prev_time = time;
207                let mut prev_diff = <R as Columnar>::into_owned(diff);
208
209                for (data, time, diff) in iter {
210                    if (&prev_data, &prev_time) == (&data, &time) {
211                        prev_diff.plus_equals(&diff);
212                    } else {
213                        if !prev_diff.is_zero() {
214                            target_d.push(prev_data);
215                            target_t.push(prev_time);
216                            target_r.push(&prev_diff);
217                        }
218                        prev_data = data;
219                        prev_time = time;
220                        R::copy_from(&mut prev_diff, diff);
221                    }
222                }
223
224                if !prev_diff.is_zero() {
225                    target_d.push(prev_data);
226                    target_t.push(prev_time);
227                    target_r.push(&prev_diff);
228                }
229            }
230        }
231
232        if !self.target.is_empty() {
233            self.ready.push_back(std::mem::take(&mut self.target));
234        }
235    }
236}
237
238/// Advance `*lower` past every position in `[*lower, upper)` where `cmp`
239/// returns true.
240///
241/// On return, `*lower` is the first index `>= initial *lower` where `cmp`
242/// returns false, or `upper` if `cmp` holds through the end.
243///
244/// Takes the predicate as `FnMut(usize) -> bool` rather than a value-bearing
245/// closure so callers can index whichever subset of the input columns they
246/// actually need to compare — for the merger's `(d, t)`-keyed sort, this lets
247/// each probe touch only the D and T leaf views, skipping the diff column.
248///
249/// Compared to a linear scan, this is `O(log K)` for a run of length `K`
250/// satisfying `cmp` — useful when one side of a sorted merge has long runs
251/// dominated by the other side.
252pub(crate) fn gallop(upper: usize, lower: &mut usize, mut cmp: impl FnMut(usize) -> bool) {
253    if *lower < upper && cmp(*lower) {
254        let mut step = 1;
255        while *lower + step < upper && cmp(*lower + step) {
256            *lower += step;
257            step <<= 1;
258        }
259
260        step >>= 1;
261        while step > 0 {
262            if *lower + step < upper && cmp(*lower + step) {
263                *lower += step;
264            }
265            step >>= 1;
266        }
267
268        // `*lower` is the last index where `cmp` holds; step to the first
269        // where it does not.
270        *lower += 1;
271    }
272}
273
274/// Counterpart to `ColInternalMerger` (which merges `ColumnationStack` chunks).
275/// Drives the merge batcher with [`ColumnBody`] chunks, no columnation
276/// detour, by way of the inherent `merge_from` / `extract` methods on
277/// `ColumnBody<(D, T, R)>` below.
278pub struct ColumnMerger<D, T, R> {
279    _marker: PhantomData<(D, T, R)>,
280}
281
282impl<D, T, R> Default for ColumnMerger<D, T, R> {
283    fn default() -> Self {
284        Self {
285            _marker: PhantomData,
286        }
287    }
288}
289
290/// A chunk the merge and extract bodies below write to: the merger's
291/// [`ColumnBody`], and the column pager's [`Column`] until that pager is
292/// retired.
293pub trait MergeChunk<C: Columnar>: Default {
294    /// The typed containers, for writing to.
295    fn typed(&mut self) -> &mut C::Container;
296    /// The chunk as a columnar view.
297    fn view(&self) -> BorrowedOf<'_, C>;
298    /// True when the chunk holds no records.
299    fn is_empty(&self) -> bool;
300}
301
302impl<C: Columnar> MergeChunk<C> for ColumnBody<C> {
303    fn typed(&mut self) -> &mut C::Container {
304        self.typed_mut()
305    }
306    fn view(&self) -> BorrowedOf<'_, C> {
307        self.borrow()
308    }
309    fn is_empty(&self) -> bool {
310        self.is_empty()
311    }
312}
313
314impl<C: Columnar> MergeChunk<C> for Column<C> {
315    fn typed(&mut self) -> &mut C::Container {
316        // The pager merges into recycled typed chunks, so this is a move. A
317        // serialized chunk is copied, as `ColumnBody::typed_mut` does.
318        if !matches!(self, Column::Typed(_)) {
319            let view = self.borrow();
320            let mut fresh = C::Container::default();
321            fresh.extend_from_self(view, 0..view.len());
322            *self = Column::Typed(fresh);
323        }
324        let Column::Typed(typed) = self else {
325            unreachable!("a serialized chunk was materialized above");
326        };
327        typed
328    }
329    fn view(&self) -> BorrowedOf<'_, C> {
330        self.borrow()
331    }
332    fn is_empty(&self) -> bool {
333        self.is_empty()
334    }
335}
336
337/// Per-chunk merge and extract for sorted [`ColumnBody`] chunks.
338///
339/// These are the building blocks that [`Merger for ColumnMerger`] orchestrates
340/// over chains of chunks. They're inherent methods rather than a trait impl
341/// so the merger can call them without going through any wrapper indirection.
342/// The output side of each is written to, so a serialized output body is
343/// materialized first (see [`ColumnBody::typed_mut`]).
344impl<D, T, R> ColumnBody<(D, T, R)>
345where
346    D: Columnar,
347    for<'a> columnar::Ref<'a, D>: Copy + Ord,
348    T: Columnar + Default + Clone + PartialOrder,
349    for<'a> columnar::Ref<'a, T>: Copy + Ord,
350    R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
351{
352    /// Merge items from sorted inputs into `self`, advancing positions. See
353    /// [`merge_from`].
354    #[must_use]
355    pub fn merge_from(&mut self, others: &mut [Self], positions: &mut [usize]) -> bool {
356        merge_from(self, others, positions)
357    }
358
359    /// Partition records starting at `*position` into `keep` and `ship`. See
360    /// [`extract`].
361    pub fn extract(
362        &mut self,
363        position: &mut usize,
364        upper: AntichainRef<T>,
365        frontier: &mut Antichain<T>,
366        keep: &mut Self,
367        ship: &mut Self,
368    ) {
369        extract(self, position, upper, frontier, keep, ship)
370    }
371}
372
373/// The same merge and extract for the column pager's [`Column`] chunks.
374impl<D, T, R> Column<(D, T, R)>
375where
376    D: Columnar,
377    for<'a> columnar::Ref<'a, D>: Copy + Ord,
378    T: Columnar + Default + Clone + PartialOrder,
379    for<'a> columnar::Ref<'a, T>: Copy + Ord,
380    R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
381{
382    /// Merge items from sorted inputs into `self`, advancing positions. See
383    /// [`merge_from`].
384    #[must_use]
385    pub fn merge_from(&mut self, others: &mut [Self], positions: &mut [usize]) -> bool {
386        merge_from(self, others, positions)
387    }
388
389    /// Partition records starting at `*position` into `keep` and `ship`. See
390    /// [`extract`].
391    pub fn extract(
392        &mut self,
393        position: &mut usize,
394        upper: AntichainRef<T>,
395        frontier: &mut Antichain<T>,
396        keep: &mut Self,
397        ship: &mut Self,
398    ) {
399        extract(self, position, upper, frontier, keep, ship)
400    }
401}
402
403/// Merge items from sorted inputs into `target`, advancing positions.
404///
405/// Mirrors the dispatch shape used by the merge-batcher framework:
406/// - **0**: no-op
407/// - **1**: bulk copy (or swap, if `target` is empty and `*pos == 0`)
408/// - **2**: merge two sorted streams, with diff consolidation on equal
409///   `(data, time)` keys and gallop bulk-copy of long single-side runs.
410///
411/// Returns `true` if the merge stopped because the amortized ship-threshold
412/// check inside the inner loop fired (the caller should ship `target` before
413/// the next call). Returns `false` if the merge stopped because at least
414/// one input was exhausted at its position (the caller should refill that
415/// side; `target` may still be at capacity from accumulation across short
416/// calls and the caller should also check `at_capacity` in that case).
417///
418/// The 0- and 1-input dispatches always return `false`: 0 does no work,
419/// 1 is a bulk copy or swap that runs to completion.
420#[must_use]
421pub fn merge_from<D, T, R, Ch>(target: &mut Ch, others: &mut [Ch], positions: &mut [usize]) -> bool
422where
423    D: Columnar,
424    for<'a> columnar::Ref<'a, D>: Copy + Ord,
425    T: Columnar + Default + Clone + PartialOrder,
426    for<'a> columnar::Ref<'a, T>: Copy + Ord,
427    R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
428    Ch: MergeChunk<(D, T, R)>,
429{
430    match others.len() {
431        0 => false,
432        1 => {
433            let other = &mut others[0];
434            let pos = &mut positions[0];
435            if target.is_empty() && *pos == 0 {
436                std::mem::swap(target, other);
437                return false;
438            }
439            let self_c = target.typed();
440            let src_c = other.view();
441            self_c.extend_from_self(src_c, *pos..other.view().len());
442            *pos = other.view().len();
443            false
444        }
445        2 => {
446            let (left, right) = others.split_at(1);
447            let (left_pos, right_pos) = positions.split_at_mut(1);
448            let left_borrow = left[0].view();
449            let right_borrow = right[0].view();
450
451            let self_c = target.typed();
452
453            // Split the input borrows into per-leaf views.
454            //
455            // A columnar tuple `Borrow::Ref` is recursive: indexing the
456            // tuple borrow walks every leaf (and reconstructs the nested
457            // ref tuple) regardless of which leaves the caller actually
458            // reads. Indexing each leaf view directly cuts probe-path
459            // work to the columns we consult — for the merge step
460            // that's the `(D, T)` key, with the diff column read only
461            // when we push.
462            let l_d = left_borrow.0;
463            let l_t = left_borrow.1;
464            let l_r = left_borrow.2;
465            let r_d = right_borrow.0;
466            let r_t = right_borrow.1;
467            let r_r = right_borrow.2;
468            let upper_l = l_d.len();
469            let upper_r = r_d.len();
470
471            // Mirror the split on the output container. Tuple
472            // containers split into per-leaf containers, which lets us
473            // address each leaf independently — both gallop bulk-copies
474            // and single-record pushes resolve to a primitive operation
475            // per leaf. The leaves stay length-synchronized as long as
476            // every record path pushes exactly one element to each.
477            let (sd, st, sr) = self_c;
478
479            // Pre-size each output leaf for the worst-case merge (no
480            // consolidation). `reserve_for` walks each input's
481            // `as_bytes`, which is accurate for variable-length leaves,
482            // where reserving a record count wouldn't size the byte
483            // buffer correctly. The reservation covers both whole
484            // inputs, however little of them this call consumes, so
485            // reserve only when the inputs fit within two ship-sized
486            // chunks, which bounds it to 4 MiB.
487            let input_bytes = 8
488                * (indexed::length_in_words(&left_borrow)
489                    + indexed::length_in_words(&right_borrow));
490            if input_bytes <= 2 * crate::columnar::SHIP_WORDS * 8 {
491                let inputs = [left_borrow, right_borrow];
492                sd.reserve_for(inputs.iter().map(|b| b.0));
493                st.reserve_for(inputs.iter().map(|b| b.1));
494                sr.reserve_for(inputs.iter().map(|b| b.2));
495            }
496
497            let mut stash = R::default();
498
499            // The size check walks every leaf slice, which is not cheap on
500            // variable-length leaves, so it runs only after spending the
501            // ship threshold's headroom at the inputs' average width.
502            // That bounds uniform rows to just under the ship size. For
503            // mixed widths it is a cadence hint: a stretch of wide rows
504            // can overshoot. `ColumnChunk::settle` enforces the bound for
505            // chunks that reach the pool. The column pager's merge-batcher
506            // output has no size classes and ships as cut.
507            let average_bytes = input_bytes / (upper_l + upper_r).max(1);
508            let check_records =
509                (crate::columnar::SHIP_WORDS * 8 / 10 / average_bytes.max(1)).clamp(1, 1024);
510            let mut next_check = sd.len() + check_records;
511            let mut yielded = false;
512
513            while left_pos[0] < upper_l && right_pos[0] < upper_r {
514                let d1 = l_d.get(left_pos[0]);
515                let t1 = l_t.get(left_pos[0]);
516                let d2 = r_d.get(right_pos[0]);
517                let t2 = r_t.get(right_pos[0]);
518                match (d1, t1).cmp(&(d2, t2)) {
519                    std::cmp::Ordering::Less => {
520                        // Common case (interleaved data): single-record
521                        // advance. Skip the gallop call entirely — its
522                        // setup plus the first cmp probe is more
523                        // expensive than just pushing this record and
524                        // re-entering the outer loop. Galloping is only
525                        // worthwhile when there's an actual run, which
526                        // we detect with the peek check below.
527                        sd.push(d1);
528                        st.push(t1);
529                        sr.push(l_r.get(left_pos[0]));
530                        left_pos[0] += 1;
531                        // Long-run case: peek at the next record; if
532                        // it's still strictly less than `(d2, t2)`,
533                        // we have a run worth galloping (and bulk-
534                        // copying).
535                        if left_pos[0] < upper_l
536                            && (l_d.get(left_pos[0]), l_t.get(left_pos[0])) < (d2, t2)
537                        {
538                            let start = left_pos[0];
539                            gallop(
540                                upper_l.min(left_pos[0] + next_check.saturating_sub(sd.len())),
541                                &mut left_pos[0],
542                                |i| (l_d.get(i), l_t.get(i)) < (d2, t2),
543                            );
544                            // Per-leaf bulk copy of the run: each call
545                            // resolves to an `extend_from_slice` on its
546                            // leaf (recursively for nested leaves).
547                            sd.extend_from_self(l_d, start..left_pos[0]);
548                            st.extend_from_self(l_t, start..left_pos[0]);
549                            sr.extend_from_self(l_r, start..left_pos[0]);
550                        }
551                    }
552                    std::cmp::Ordering::Greater => {
553                        // Symmetric on the right side.
554                        sd.push(d2);
555                        st.push(t2);
556                        sr.push(r_r.get(right_pos[0]));
557                        right_pos[0] += 1;
558                        if right_pos[0] < upper_r
559                            && (r_d.get(right_pos[0]), r_t.get(right_pos[0])) < (d1, t1)
560                        {
561                            let start = right_pos[0];
562                            gallop(
563                                upper_r.min(right_pos[0] + next_check.saturating_sub(sd.len())),
564                                &mut right_pos[0],
565                                |i| (r_d.get(i), r_t.get(i)) < (d1, t1),
566                            );
567                            sd.extend_from_self(r_d, start..right_pos[0]);
568                            st.extend_from_self(r_t, start..right_pos[0]);
569                            sr.extend_from_self(r_r, start..right_pos[0]);
570                        }
571                    }
572                    std::cmp::Ordering::Equal => {
573                        let r1 = l_r.get(left_pos[0]);
574                        let r2 = r_r.get(right_pos[0]);
575                        R::copy_from(&mut stash, r1);
576                        stash.plus_equals(&r2);
577                        if !stash.is_zero() {
578                            sd.push(d1);
579                            st.push(t1);
580                            sr.push(&stash);
581                        }
582                        left_pos[0] += 1;
583                        right_pos[0] += 1;
584                    }
585                }
586
587                if sd.len() >= next_check {
588                    use columnar::Borrow as _;
589                    if crate::columnar::at_serialized_capacity(&(
590                        sd.borrow(),
591                        st.borrow(),
592                        sr.borrow(),
593                    )) {
594                        yielded = true;
595                        break;
596                    }
597                    next_check = sd.len() + check_records;
598                }
599            }
600            yielded
601        }
602        // `Merger::merge` only ever calls `merge_from` with 0/1/2-input
603        // slices (k-way merge isn't part of the merge-batcher contract).
604        n => unreachable!("merge_from called with {n} inputs; expected 0, 1, or 2"),
605    }
606}
607
608/// Partition records of `source` starting at `*position` into `keep` (times
609/// beyond `upper`, retained for the next round) and `ship` (times not beyond
610/// `upper`, sealed into the output batch). Updates `frontier` with the
611/// times of kept records.
612///
613/// The caller invokes `extract` repeatedly until `*position >= source.len()`,
614/// swapping out a full output buffer between calls. This shape exists
615/// because the framework only checks `at_capacity()` between calls, so
616/// without an inner-loop yield a single call could quietly produce
617/// oversized output chunks.
618pub fn extract<D, T, R, Ch>(
619    source: &Ch,
620    position: &mut usize,
621    upper: AntichainRef<T>,
622    frontier: &mut Antichain<T>,
623    keep: &mut Ch,
624    ship: &mut Ch,
625) where
626    D: Columnar,
627    for<'a> columnar::Ref<'a, D>: Copy + Ord,
628    T: Columnar + Default + Clone + PartialOrder,
629    for<'a> columnar::Ref<'a, T>: Copy + Ord,
630    R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
631    Ch: MergeChunk<(D, T, R)>,
632{
633    let keep_c = keep.typed();
634    let ship_c = ship.typed();
635
636    let self_view = source.view();
637    let len = self_view.len();
638
639    use columnar::Borrow as _;
640    let mut owned_t = T::default();
641    while *position < len
642        && !crate::columnar::at_serialized_capacity(&keep_c.borrow())
643        && !crate::columnar::at_serialized_capacity(&ship_c.borrow())
644    {
645        let (_, time, _) = self_view.get(*position);
646        T::copy_from(&mut owned_t, time);
647        if upper.less_equal(&owned_t) {
648            // `insert_with` only clones when the time isn't already
649            // present in the antichain.
650            frontier.insert_with(&owned_t, |t| t.clone());
651            keep_c.extend_from_self(self_view, *position..*position + 1);
652        } else {
653            ship_c.extend_from_self(self_view, *position..*position + 1);
654        }
655        *position += 1;
656    }
657}
658
659/// `Merger` impl driving [`MergeBatcher`] over [`ColumnBody`] chunks,
660/// built on the inherent `ColumnBody::merge_from` and `ColumnBody::extract`
661/// methods.
662/// Exhausted input chunks are recycled through `stash`, and remaining full
663/// chunks on a drained side move to the output directly, with no per-element
664/// copy.
665///
666/// [`MergeBatcher`]: differential_dataflow::trace::implementations::merge_batcher::MergeBatcher
667impl<D, T, R> Merger for ColumnMerger<D, T, R>
668where
669    D: Columnar,
670    for<'a> columnar::Ref<'a, D>: Copy + Ord,
671    T: Columnar + Default + Clone + Ord + PartialOrder,
672    for<'a> columnar::Ref<'a, T>: Copy + Ord,
673    R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
674{
675    type Time = T;
676    type Chunk = ColumnBody<(D, T, R)>;
677
678    fn merge(
679        &mut self,
680        list1: Vec<Self::Chunk>,
681        list2: Vec<Self::Chunk>,
682        output: &mut Vec<Self::Chunk>,
683        stash: &mut Vec<Self::Chunk>,
684    ) {
685        let mut list1 = list1.into_iter();
686        let mut list2 = list2.into_iter();
687
688        let mut heads = [
689            list1.next().unwrap_or_default(),
690            list2.next().unwrap_or_default(),
691        ];
692        let mut positions = [0usize, 0usize];
693
694        let mut result = empty_chunk(stash);
695
696        // Main merge loop: both sides have data.
697        loop {
698            let upper_l = heads[0].borrow().len();
699            let upper_r = heads[1].borrow().len();
700            if positions[0] >= upper_l || positions[1] >= upper_r {
701                break;
702            }
703
704            // Whole-chunk passthrough fast path: when one head's tail (from
705            // its current position) sorts entirely before the other head's
706            // current record, the head moves to `output` wholesale, with two
707            // probe records in place of per-record compares and per-leaf
708            // byte copies. Restricted to `positions[i] == 0` so the head can
709            // be handed off intact; a partial tail is what gallop already
710            // handles inside the merge loop.
711            let lhs_passthrough = positions[0] == 0 && upper_l > 0 && {
712                let lhs = heads[0].borrow();
713                let rhs = heads[1].borrow();
714                let last_l = (lhs.0.get(upper_l - 1), lhs.1.get(upper_l - 1));
715                let cur_r = (rhs.0.get(positions[1]), rhs.1.get(positions[1]));
716                last_l < cur_r
717            };
718            if lhs_passthrough {
719                if !result.is_empty() {
720                    output.push(std::mem::take(&mut result));
721                    result = empty_chunk(stash);
722                }
723                let head = std::mem::replace(&mut heads[0], list1.next().unwrap_or_default());
724                output.push(head);
725                positions[0] = 0;
726                continue;
727            }
728
729            let rhs_passthrough = positions[1] == 0 && upper_r > 0 && {
730                let lhs = heads[0].borrow();
731                let rhs = heads[1].borrow();
732                let last_r = (rhs.0.get(upper_r - 1), rhs.1.get(upper_r - 1));
733                let cur_l = (lhs.0.get(positions[0]), lhs.1.get(positions[0]));
734                last_r < cur_l
735            };
736            if rhs_passthrough {
737                if !result.is_empty() {
738                    output.push(std::mem::take(&mut result));
739                    result = empty_chunk(stash);
740                }
741                let head = std::mem::replace(&mut heads[1], list2.next().unwrap_or_default());
742                output.push(head);
743                positions[1] = 0;
744                continue;
745            }
746
747            // Per-record merge. `merge_from` returns `true` when its inner
748            // amortized ship-threshold check fires — short-circuit the
749            // outer `at_capacity` walk in that case.
750            let yielded = result.merge_from(&mut heads, &mut positions);
751
752            if positions[0] >= heads[0].borrow().len() {
753                let old = std::mem::replace(&mut heads[0], list1.next().unwrap_or_default());
754                recycle_chunk(old, stash);
755                positions[0] = 0;
756            }
757            if positions[1] >= heads[1].borrow().len() {
758                let old = std::mem::replace(&mut heads[1], list2.next().unwrap_or_default());
759                recycle_chunk(old, stash);
760                positions[1] = 0;
761            }
762            if yielded || result.at_capacity() {
763                output.push(std::mem::take(&mut result));
764                result = empty_chunk(stash);
765            }
766        }
767
768        // Drain remaining from each side: copy partial head, then append
769        // full chunks directly to output (no per-element copy).
770        drain_side(
771            &mut heads[0],
772            &mut positions[0],
773            &mut list1,
774            &mut result,
775            output,
776            stash,
777        );
778        drain_side(
779            &mut heads[1],
780            &mut positions[1],
781            &mut list2,
782            &mut result,
783            output,
784            stash,
785        );
786        if !result.is_empty() {
787            output.push(result);
788        }
789    }
790
791    fn extract(
792        &mut self,
793        merged: Vec<Self::Chunk>,
794        upper: AntichainRef<Self::Time>,
795        frontier: &mut Antichain<Self::Time>,
796        ship: &mut Vec<Self::Chunk>,
797        kept: &mut Vec<Self::Chunk>,
798        stash: &mut Vec<Self::Chunk>,
799    ) {
800        let mut keep = empty_chunk(stash);
801        let mut ready = empty_chunk(stash);
802
803        for mut buffer in merged {
804            let mut position = 0;
805            let len = buffer.borrow().len();
806            while position < len {
807                buffer.extract(&mut position, upper, frontier, &mut keep, &mut ready);
808                if keep.at_capacity() {
809                    kept.push(std::mem::take(&mut keep));
810                    keep = empty_chunk(stash);
811                }
812                if ready.at_capacity() {
813                    ship.push(std::mem::take(&mut ready));
814                    ready = empty_chunk(stash);
815                }
816            }
817            recycle_chunk(buffer, stash);
818        }
819        if !keep.is_empty() {
820            kept.push(keep);
821        }
822        if !ready.is_empty() {
823            ship.push(ready);
824        }
825    }
826
827    fn len(chunk: &Self::Chunk) -> usize {
828        chunk.len()
829    }
830
831    fn allocation(chunk: &Self::Chunk) -> (usize, usize, usize) {
832        // Serialized footprint stands in for both `size` and `capacity`: the
833        // chunk owns one logical allocation worth of leaf storage, and we
834        // ship/recycle the whole thing rather than tracking per-leaf
835        // capacities. Treating `size == capacity` matches how the framework
836        // accounts already-shipped chunks (no slack to absorb).
837        let bytes = chunk.length_in_bytes();
838        (bytes, bytes, 1)
839    }
840}
841
842/// Pop a chunk from `stash` or allocate a fresh one. Stashed chunks are
843/// already cleared via `recycle_chunk`, so they're ready for push.
844#[inline]
845pub(crate) fn empty_chunk<C: Columnar>(stash: &mut Vec<ColumnBody<C>>) -> ColumnBody<C> {
846    stash.pop().unwrap_or_default()
847}
848
849/// Reset `chunk` to an empty `Typed` and push it to `stash` for reuse.
850///
851/// Chunks recycled here come from the merger and chunker, both of which
852/// produce `Typed`; only the typed allocations are worth caching for reuse.
853/// A serialized chunk has no typed-side allocation to preserve, so we simply
854/// drop it: `empty_chunk` will produce a fresh default just as cheaply, and
855/// pushing it onto `stash` would only displace useful recycled allocations.
856#[inline]
857pub(crate) fn recycle_chunk<C: Columnar>(mut chunk: ColumnBody<C>, stash: &mut Vec<ColumnBody<C>>) {
858    if let ColumnBody::Typed(c) = &mut chunk {
859        c.clear();
860        stash.push(chunk);
861    }
862}
863
864/// Drain remaining items from one side into `result` / `output`.
865///
866/// Copies the partially-consumed head into `result` via `merge_from`'s 1-input
867/// path, then appends remaining full chunks directly to `output` without
868/// per-element copy.
869fn drain_side<D, T, R>(
870    head: &mut ColumnBody<(D, T, R)>,
871    pos: &mut usize,
872    list: &mut std::vec::IntoIter<ColumnBody<(D, T, R)>>,
873    result: &mut ColumnBody<(D, T, R)>,
874    output: &mut Vec<ColumnBody<(D, T, R)>>,
875    stash: &mut Vec<ColumnBody<(D, T, R)>>,
876) where
877    D: Columnar,
878    for<'a> columnar::Ref<'a, D>: Copy + Ord,
879    T: Columnar + Default + Clone + PartialOrder,
880    for<'a> columnar::Ref<'a, T>: Copy + Ord,
881    R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
882{
883    if *pos < head.borrow().len() {
884        // 1-input dispatch — bulk copy that runs to completion; the yield
885        // signal is unused.
886        let _ = result.merge_from(std::slice::from_mut(head), std::slice::from_mut(pos));
887    }
888    if !result.is_empty() {
889        output.push(std::mem::take(result));
890        *result = empty_chunk(stash);
891    }
892    Extend::extend(output, list);
893}
894
895#[cfg(test)]
896mod tests {
897    use timely::dataflow::channels::ContainerBytes;
898
899    use super::*;
900
901    /// Re-encode `column` as the `Align` variant, to drive the chunker over serialized
902    /// input.
903    ///
904    /// `into_bytes` writes whole `u64` words, read back in native byte order.
905    fn serialize<C: Columnar>(column: &Column<C>) -> Column<C> {
906        let mut bytes: Vec<u8> = Vec::new();
907        column.into_bytes(&mut bytes);
908        assert_eq!(bytes.len() % 8, 0);
909        let words = bytes
910            .chunks_exact(8)
911            .map(|w| u64::from_ne_bytes(w.try_into().expect("chunk is 8 bytes")))
912            .collect();
913        Column::Align(words)
914    }
915
916    /// Drive a single `push_into` call with `inputs` and collect the
917    /// consolidated output (if any) as owned tuples.
918    fn run_chunker<D, T, R>(inputs: &[(D, T, R)]) -> Vec<(D, T, R)>
919    where
920        D: Columnar + Clone,
921        for<'a> columnar::Ref<'a, D>: Copy + Ord,
922        T: Columnar + Clone,
923        for<'a> columnar::Ref<'a, T>: Copy + Ord,
924        R: Columnar + Clone + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
925        for<'a> columnar::Ref<'a, R>: Ord,
926        <(D, T, R) as Columnar>::Container: Clone,
927        for<'a> <(D, T, R) as Columnar>::Container: columnar::Push<&'a (D, T, R)>,
928        <(D, T, R) as Columnar>::Container: columnar::Push<(D, T, R)>,
929        for<'a> <D as Columnar>::Container: columnar::Push<columnar::Ref<'a, D>>,
930        for<'a> <T as Columnar>::Container: columnar::Push<columnar::Ref<'a, T>>,
931        for<'a> <R as Columnar>::Container: columnar::Push<&'a R>,
932    {
933        let mut input: Column<(D, T, R)> = Default::default();
934        for tuple in inputs.iter().cloned() {
935            input.push_into(tuple);
936        }
937
938        let mut chunker: ColumnChunker<(D, T, R)> = Default::default();
939        chunker.push_into(&mut input);
940
941        let mut out = Vec::new();
942        while let Some(chunk) = chunker.extract() {
943            for (d, t, r) in chunk.borrow().into_index_iter() {
944                out.push((D::into_owned(d), T::into_owned(t), R::into_owned(r)));
945            }
946        }
947        out
948    }
949
950    #[mz_ore::test]
951    fn empty_input_yields_no_chunk() {
952        let mut chunker: ColumnChunker<(u64, u64, i64)> = Default::default();
953        let mut input: Column<(u64, u64, i64)> = Default::default();
954        chunker.push_into(&mut input);
955        assert!(chunker.extract().is_none());
956        assert!(chunker.finish().is_none());
957    }
958
959    #[mz_ore::test]
960    fn unsorted_input_is_sorted() {
961        let out = run_chunker(&[(3u64, 0u64, 1i64), (1u64, 0u64, 1i64), (2u64, 0u64, 1i64)]);
962        assert_eq!(out, vec![(1, 0, 1), (2, 0, 1), (3, 0, 1)]);
963    }
964
965    #[mz_ore::test]
966    fn serialized_input_is_consolidated() {
967        // Callers hand the chunker whatever the upstream edge delivered, which is
968        // serialized once the data crossed an exchange.
969        let mut input: Column<(u64, u64, i64)> = Default::default();
970        for tuple in [
971            (2u64, 0u64, 1i64),
972            (1, 0, 1),
973            (2, 0, 2),
974            (3, 0, 1),
975            (3, 0, -1),
976        ] {
977            input.push_into(tuple);
978        }
979        let mut input = serialize(&input);
980
981        let mut chunker: ColumnChunker<(u64, u64, i64)> = Default::default();
982        chunker.push_into(&mut input);
983
984        let mut out = Vec::new();
985        while let Some(chunk) = chunker.extract() {
986            for (d, t, r) in chunk.borrow().into_index_iter() {
987                out.push((*d, *t, *r));
988            }
989        }
990        assert_eq!(out, vec![(1, 0, 1), (2, 0, 3)]);
991    }
992
993    #[mz_ore::test]
994    fn duplicate_keys_consolidate() {
995        let out = run_chunker(&[(1u64, 0u64, 1i64), (1u64, 0u64, 2i64), (1u64, 0u64, -1i64)]);
996        assert_eq!(out, vec![(1, 0, 2)]);
997    }
998
999    #[mz_ore::test]
1000    fn diffs_summing_to_zero_are_dropped() {
1001        let out = run_chunker(&[(1u64, 0u64, 1i64), (1u64, 0u64, -1i64)]);
1002        assert!(out.is_empty());
1003    }
1004
1005    #[mz_ore::test]
1006    fn mixed_consolidation() {
1007        // (1, 0): 1 + 2 + (-3) = 0  -> dropped
1008        // (2, 0): 1            = 1  -> kept
1009        // (1, 1): 5            = 5  -> kept (different time from the (1, 0) group)
1010        let out = run_chunker(&[
1011            (1u64, 0u64, 1i64),
1012            (2u64, 0u64, 1i64),
1013            (1u64, 0u64, 2i64),
1014            (1u64, 1u64, 5i64),
1015            (1u64, 0u64, -3i64),
1016        ]);
1017        assert_eq!(out, vec![(1, 1, 5), (2, 0, 1)]);
1018    }
1019
1020    #[mz_ore::test]
1021    fn key_val_tuple_data() {
1022        // Exercise the actual val-batcher shape: `D = (K, V)`.
1023        let out = run_chunker(&[
1024            ((1u64, 10u64), 0u64, 1i64),
1025            ((1u64, 10u64), 0u64, 1i64),
1026            ((1u64, 11u64), 0u64, 1i64),
1027            ((2u64, 10u64), 0u64, 1i64),
1028        ]);
1029        assert_eq!(
1030            out,
1031            vec![((1, 10), 0, 2), ((1, 11), 0, 1), ((2, 10), 0, 1),]
1032        );
1033    }
1034
1035    #[mz_ore::test]
1036    fn buffer_reuse_across_calls() {
1037        // Two sequential push_into calls; second runs after extract returned
1038        // the first chunk, exercising the in-place clear path.
1039        let mut input1: Column<(u64, u64, i64)> = Default::default();
1040        input1.push_into((1u64, 0u64, 1i64));
1041        input1.push_into((2u64, 0u64, 1i64));
1042
1043        let mut input2: Column<(u64, u64, i64)> = Default::default();
1044        input2.push_into((3u64, 0u64, 1i64));
1045        input2.push_into((1u64, 0u64, 1i64));
1046
1047        let mut chunker: ColumnChunker<(u64, u64, i64)> = Default::default();
1048        chunker.push_into(&mut input1);
1049
1050        // Hand back the first chunk via extract, simulating the merge batcher
1051        // taking ownership of the &mut and then returning.
1052        {
1053            let _ = chunker.extract().expect("first chunk");
1054        }
1055
1056        chunker.push_into(&mut input2);
1057
1058        let chunk = chunker.extract().expect("second chunk");
1059        let collected: Vec<_> = chunk
1060            .borrow()
1061            .into_index_iter()
1062            .map(|(d, t, r)| (u64::into_owned(d), u64::into_owned(t), i64::into_owned(r)))
1063            .collect();
1064        assert_eq!(collected, vec![(1, 0, 1), (3, 0, 1)]);
1065    }
1066
1067    /// Build a `ColumnBody<((u64, u64), u64, i64)>` from a slice of tuples.
1068    fn col(rows: &[((u64, u64), u64, i64)]) -> ColumnBody<((u64, u64), u64, i64)> {
1069        let mut c: ColumnBody<((u64, u64), u64, i64)> = Default::default();
1070        for &t in rows {
1071            c.push_into(t);
1072        }
1073        c
1074    }
1075
1076    fn collect_chunks(
1077        chunks: &[ColumnBody<((u64, u64), u64, i64)>],
1078    ) -> Vec<((u64, u64), u64, i64)> {
1079        chunks
1080            .iter()
1081            .flat_map(|c| {
1082                c.borrow().into_index_iter().map(|((k, v), t, r)| {
1083                    (
1084                        (u64::into_owned(k), u64::into_owned(v)),
1085                        u64::into_owned(t),
1086                        i64::into_owned(r),
1087                    )
1088                })
1089            })
1090            .collect()
1091    }
1092
1093    /// Disjoint-range chains exercise the whole-chunk passthrough fast path:
1094    /// every chunk in chain1 is sortable-before every chunk in chain2, so
1095    /// each outer-loop iteration should hand a chunk straight to `output`
1096    /// without recursing through the per-record merge.
1097    #[mz_ore::test]
1098    fn merger_disjoint_chains_passthrough() {
1099        let chain1 = vec![
1100            col(&[((0, 0), 0, 1), ((1, 0), 0, 1)]),
1101            col(&[((2, 0), 0, 1), ((3, 0), 0, 1)]),
1102        ];
1103        let chain2 = vec![
1104            col(&[((10, 0), 0, 1), ((11, 0), 0, 1)]),
1105            col(&[((12, 0), 0, 1), ((13, 0), 0, 1)]),
1106        ];
1107
1108        let mut merger: ColumnMerger<(u64, u64), u64, i64> = Default::default();
1109        let mut output = Vec::new();
1110        let mut stash = Vec::new();
1111        Merger::merge(&mut merger, chain1, chain2, &mut output, &mut stash);
1112
1113        let collected = collect_chunks(&output);
1114        let expected: Vec<_> = (0..4u64)
1115            .map(|d| ((d, 0u64), 0u64, 1i64))
1116            .chain((10..14u64).map(|d| ((d, 0u64), 0u64, 1i64)))
1117            .collect();
1118        assert_eq!(collected, expected);
1119    }
1120
1121    /// Interleaved chains never satisfy the passthrough condition; each
1122    /// outer iteration falls through to `merge_from`. Same correctness
1123    /// expectation, exercises the non-passthrough path under
1124    /// `Merger::merge`.
1125    #[mz_ore::test]
1126    fn merger_interleaved_chains() {
1127        // Even keys on one chain, odd on the other; chunks alternate so the
1128        // per-record path is the only viable route.
1129        let chain1 = vec![
1130            col(&[((0, 0), 0, 1), ((2, 0), 0, 1)]),
1131            col(&[((4, 0), 0, 1), ((6, 0), 0, 1)]),
1132        ];
1133        let chain2 = vec![
1134            col(&[((1, 0), 0, 1), ((3, 0), 0, 1)]),
1135            col(&[((5, 0), 0, 1), ((7, 0), 0, 1)]),
1136        ];
1137
1138        let mut merger: ColumnMerger<(u64, u64), u64, i64> = Default::default();
1139        let mut output = Vec::new();
1140        let mut stash = Vec::new();
1141        Merger::merge(&mut merger, chain1, chain2, &mut output, &mut stash);
1142
1143        let collected = collect_chunks(&output);
1144        let expected: Vec<_> = (0..8u64).map(|d| ((d, 0u64), 0u64, 1i64)).collect();
1145        assert_eq!(collected, expected);
1146    }
1147
1148    /// Passthrough must consolidate adjacent equal keys at chunk
1149    /// boundaries — i.e., must NOT fire when `chain1`'s last record's
1150    /// `(d, t)` equals `chain2`'s first.
1151    #[mz_ore::test]
1152    fn merger_passthrough_respects_equal_boundary() {
1153        // chain1's last == chain2's first key: equal-key consolidation
1154        // must kick in (sum of diffs would be 2). If passthrough fired
1155        // erroneously, both records would land in different output chunks
1156        // unconsolidated.
1157        let chain1 = vec![col(&[((0, 0), 0, 1), ((5, 0), 0, 1)])];
1158        let chain2 = vec![col(&[((5, 0), 0, 1), ((10, 0), 0, 1)])];
1159
1160        let mut merger: ColumnMerger<(u64, u64), u64, i64> = Default::default();
1161        let mut output = Vec::new();
1162        let mut stash = Vec::new();
1163        Merger::merge(&mut merger, chain1, chain2, &mut output, &mut stash);
1164
1165        let collected = collect_chunks(&output);
1166        assert_eq!(
1167            collected,
1168            vec![((0, 0), 0, 1), ((5, 0), 0, 2), ((10, 0), 0, 1)]
1169        );
1170    }
1171}
1172
1173#[cfg(test)]
1174mod proptests {
1175    //! Property tests for `ColumnBody::merge_from` and `ColumnBody::extract`.
1176    //!
1177    //! Strategy: generate sorted+consolidated inputs (the merger's input
1178    //! contract), drive `merge_from` / `extract` the same way the framework
1179    //! would, and compare against a brute-force reference impl.
1180    //!
1181    //! Test types are `D = (u64, u64)`, `T = u64`, `R = i64` drawn from small
1182    //! ranges so that equal-key collisions are common and the consolidation
1183    //! path actually runs.
1184    use super::*;
1185    use mz_ore::cast::CastFrom;
1186    use proptest::prelude::*;
1187    use timely::progress::frontier::Antichain;
1188
1189    type Tuple = ((u64, u64), u64, i64);
1190
1191    /// Reference consolidation: sort by `(data, time)`, sum diffs over equal
1192    /// pairs, drop zeros.
1193    fn consolidate(mut v: Vec<Tuple>) -> Vec<Tuple> {
1194        v.sort();
1195        let mut out: Vec<Tuple> = Vec::new();
1196        for (d, t, r) in v {
1197            if let Some(last) = out.last_mut() {
1198                if last.0 == d && last.1 == t {
1199                    last.2 += r;
1200                    continue;
1201                }
1202            }
1203            out.push((d, t, r));
1204        }
1205        out.retain(|x| x.2 != 0);
1206        out
1207    }
1208
1209    /// Strategy for sorted+consolidated input lists. Ranges are small to
1210    /// encourage equal-key collisions.
1211    fn arb_consolidated() -> impl Strategy<Value = Vec<Tuple>> {
1212        prop::collection::vec(((0u64..5, 0u64..5), 0u64..3, -3i64..=3i64), 0..30)
1213            .prop_map(consolidate)
1214    }
1215
1216    fn build_column(v: &[Tuple]) -> ColumnBody<Tuple> {
1217        let mut col: ColumnBody<Tuple> = Default::default();
1218        for tup in v {
1219            col.push_into(*tup);
1220        }
1221        col
1222    }
1223
1224    /// The body in `ColumnBody::Words` form when `words` is set, as a store hands it back.
1225    fn in_form(col: ColumnBody<Tuple>, words: bool) -> ColumnBody<Tuple> {
1226        if !words {
1227            return col;
1228        }
1229        let mut bytes = Vec::new();
1230        col.write_into(&mut bytes).expect("vec writes");
1231        ColumnBody::Words(bytemuck::allocation::pod_collect_to_vec(&bytes))
1232    }
1233
1234    fn collect_column(col: &ColumnBody<Tuple>) -> Vec<Tuple> {
1235        col.borrow()
1236            .into_index_iter()
1237            .map(|((k, v), t, r)| {
1238                (
1239                    (u64::into_owned(k), u64::into_owned(v)),
1240                    u64::into_owned(t),
1241                    i64::into_owned(r),
1242                )
1243            })
1244            .collect()
1245    }
1246
1247    /// Drive a 2-way merge the same way `Merger::merge` would: a 2-input
1248    /// call until one side exhausts, then a 1-input drain for whichever
1249    /// side still has data.
1250    fn drive_merge(left: ColumnBody<Tuple>, right: ColumnBody<Tuple>) -> ColumnBody<Tuple> {
1251        let mut self_col: ColumnBody<Tuple> = Default::default();
1252        let mut others = [left, right];
1253        let mut positions = [0usize, 0];
1254        let _ = self_col.merge_from(&mut others, &mut positions);
1255
1256        let [left_done, right_done] = others;
1257        let [left_pos, right_pos] = positions;
1258
1259        if left_pos < left_done.borrow().len() {
1260            let mut tail = [left_done];
1261            let mut p = [left_pos];
1262            let _ = self_col.merge_from(&mut tail, &mut p);
1263        } else if right_pos < right_done.borrow().len() {
1264            let mut tail = [right_done];
1265            let mut p = [right_pos];
1266            let _ = self_col.merge_from(&mut tail, &mut p);
1267        }
1268
1269        self_col
1270    }
1271
1272    proptest! {
1273        /// `merge_from` with two sorted+consolidated inputs equals the
1274        /// reference consolidate(union).
1275        #[mz_ore::test]
1276        #[cfg_attr(miri, ignore)]
1277        fn merge_from_equals_consolidated_union(
1278            a in arb_consolidated(),
1279            b in arb_consolidated(),
1280            a_words in any::<bool>(),
1281            b_words in any::<bool>(),
1282        ) {
1283            let merged = drive_merge(
1284                in_form(build_column(&a), a_words),
1285                in_form(build_column(&b), b_words),
1286            );
1287
1288            let mut union = a.clone();
1289            Extend::extend(&mut union, b.iter().copied());
1290            let expected = consolidate(union);
1291
1292            prop_assert_eq!(collect_column(&merged), expected);
1293        }
1294
1295        /// `merge_from` 1-input bulk-copy from a non-zero position equals
1296        /// `other[*pos..]`.
1297        #[mz_ore::test]
1298        #[cfg_attr(miri, ignore)]
1299        fn merge_from_one_input_drains_tail(
1300            data in arb_consolidated(),
1301            pos_frac in 0u32..=100,
1302            self_words in any::<bool>(),
1303            other_words in any::<bool>(),
1304        ) {
1305            // Cap at len so we always have a valid position.
1306            let len = data.len();
1307            let start_pos = if len == 0 { 0 } else {
1308                (usize::cast_from(pos_frac) * len) / 101
1309            };
1310
1311            // Self starts non-empty so we exercise the bulk-copy path, not the
1312            // empty-self swap shortcut.
1313            let mut self_col: ColumnBody<Tuple> = Default::default();
1314            let sentinel: Tuple = ((u64::MAX, u64::MAX), 0, 1);
1315            self_col.push_into(sentinel);
1316            // A serialized target is materialized before the copy appends to it.
1317            let mut self_col = in_form(self_col, self_words);
1318
1319            let mut others = [in_form(build_column(&data), other_words)];
1320            let mut positions = [start_pos];
1321            let _ = self_col.merge_from(&mut others, &mut positions);
1322
1323            let mut expected = vec![sentinel];
1324            Extend::extend(&mut expected, data[start_pos..].iter().copied());
1325
1326            prop_assert_eq!(collect_column(&self_col), expected);
1327            prop_assert_eq!(positions[0], len);
1328        }
1329
1330        /// `merge_from` 1-input swap shortcut: empty self + pos=0 should
1331        /// produce a column equal to the input.
1332        #[mz_ore::test]
1333        #[cfg_attr(miri, ignore)]
1334        fn merge_from_empty_self_swap(data in arb_consolidated(), words in any::<bool>()) {
1335            let mut self_col: ColumnBody<Tuple> = Default::default();
1336            let mut others = [in_form(build_column(&data), words)];
1337            let mut positions = [0usize];
1338            let _ = self_col.merge_from(&mut others, &mut positions);
1339
1340            prop_assert_eq!(collect_column(&self_col), data);
1341        }
1342
1343        /// `extract` partitions correctly:
1344        ///   - keep ∪ ship multiset-equals self
1345        ///   - upper.less_equal(t) for every kept time
1346        ///   - !upper.less_equal(t) for every shipped time
1347        ///   - frontier covers every kept time
1348        #[mz_ore::test]
1349        #[cfg_attr(miri, ignore)]
1350        fn extract_partitions_by_frontier(
1351            data in arb_consolidated(),
1352            upper_time in 0u64..=4,
1353            words in any::<bool>(),
1354        ) {
1355            let mut self_col = in_form(build_column(&data), words);
1356            let upper = Antichain::from_elem(upper_time);
1357            let mut frontier: Antichain<u64> = Antichain::new();
1358            let mut keep: ColumnBody<Tuple> = Default::default();
1359            let mut ship: ColumnBody<Tuple> = Default::default();
1360            let mut position = 0;
1361
1362            self_col.extract(
1363                &mut position,
1364                upper.borrow(),
1365                &mut frontier,
1366                &mut keep,
1367                &mut ship,
1368            );
1369
1370            // Single call drains the input (we removed the at_capacity yield).
1371            prop_assert_eq!(position, data.len());
1372
1373            let kept = collect_column(&keep);
1374            let shipped = collect_column(&ship);
1375
1376            // Partition predicate: kept times >= upper, shipped times < upper.
1377            for (_, t, _) in &kept {
1378                prop_assert!(
1379                    upper.borrow().less_equal(t),
1380                    "kept time {} should satisfy upper.less_equal", t,
1381                );
1382            }
1383            for (_, t, _) in &shipped {
1384                prop_assert!(
1385                    !upper.borrow().less_equal(t),
1386                    "shipped time {} should NOT satisfy upper.less_equal", t,
1387                );
1388            }
1389
1390            // Union (multiset) equals input.
1391            let mut union = kept.clone();
1392            Extend::extend(&mut union, shipped.iter().copied());
1393            union.sort();
1394            let mut expected_sorted = data.clone();
1395            expected_sorted.sort();
1396            prop_assert_eq!(union, expected_sorted);
1397
1398            // Frontier dominates every kept time.
1399            for (_, t, _) in &kept {
1400                prop_assert!(
1401                    frontier.less_equal(t),
1402                    "frontier should dominate kept time {}", t,
1403                );
1404            }
1405        }
1406
1407        /// Empty input → no work, frontier untouched, position = 0.
1408        #[mz_ore::test]
1409        #[cfg_attr(miri, ignore)]
1410        fn extract_empty_input(upper_time in 0u64..=4) {
1411            let mut self_col: ColumnBody<Tuple> = Default::default();
1412            let upper = Antichain::from_elem(upper_time);
1413            let mut frontier: Antichain<u64> = Antichain::new();
1414            let mut keep: ColumnBody<Tuple> = Default::default();
1415            let mut ship: ColumnBody<Tuple> = Default::default();
1416            let mut position = 0;
1417
1418            self_col.extract(
1419                &mut position,
1420                upper.borrow(),
1421                &mut frontier,
1422                &mut keep,
1423                &mut ship,
1424            );
1425
1426            prop_assert_eq!(position, 0);
1427            prop_assert!(collect_column(&keep).is_empty());
1428            prop_assert!(collect_column(&ship).is_empty());
1429            prop_assert!(frontier.elements().is_empty());
1430        }
1431    }
1432}