1use 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#[derive(Default)]
39pub struct Chunker<C> {
40 target: C,
45 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 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 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
125pub struct ColumnChunker<U: Columnar> {
133 target: ColumnBody<U>,
135 ready: VecDeque<ColumnBody<U>>,
137}
138
139impl<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 self.target.clear();
188
189 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 {
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
238pub(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 += 1;
271 }
272}
273
274pub 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
290pub trait MergeChunk<C: Columnar>: Default {
294 fn typed(&mut self) -> &mut C::Container;
296 fn view(&self) -> BorrowedOf<'_, C>;
298 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 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
337impl<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 #[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 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
373impl<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 #[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 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#[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 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 let (sd, st, sr) = self_c;
478
479 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 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 sd.push(d1);
528 st.push(t1);
529 sr.push(l_r.get(left_pos[0]));
530 left_pos[0] += 1;
531 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 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 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 n => unreachable!("merge_from called with {n} inputs; expected 0, 1, or 2"),
605 }
606}
607
608pub 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 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
659impl<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 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 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 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_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 let bytes = chunk.length_in_bytes();
838 (bytes, bytes, 1)
839 }
840}
841
842#[inline]
845pub(crate) fn empty_chunk<C: Columnar>(stash: &mut Vec<ColumnBody<C>>) -> ColumnBody<C> {
846 stash.pop().unwrap_or_default()
847}
848
849#[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
864fn 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 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 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 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 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 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 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 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 {
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 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 #[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 #[mz_ore::test]
1126 fn merger_interleaved_chains() {
1127 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 #[mz_ore::test]
1152 fn merger_passthrough_respects_equal_boundary() {
1153 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 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 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 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 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 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 #[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 #[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 let len = data.len();
1307 let start_pos = if len == 0 { 0 } else {
1308 (usize::cast_from(pos_frac) * len) / 101
1309 };
1310
1311 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 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 #[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 #[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 prop_assert_eq!(position, data.len());
1372
1373 let kept = collect_column(&keep);
1374 let shipped = collect_column(&ship);
1375
1376 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 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 for (_, t, _) in &kept {
1400 prop_assert!(
1401 frontier.less_equal(t),
1402 "frontier should dominate kept time {}", t,
1403 );
1404 }
1405 }
1406
1407 #[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}