1use std::collections::VecDeque;
19use std::marker::PhantomData;
20
21use crate::columnation::ColumnationStack;
22use columnar::Container as _;
23use columnar::Push as _;
24use columnar::{Clear, Columnar, Index, Len};
25use columnation::Columnation;
26use differential_dataflow::difference::Semigroup;
27use differential_dataflow::trace::implementations::merge_batcher::Merger;
28use timely::Accountable;
29use timely::Container;
30use timely::PartialOrder;
31use timely::container::{ContainerBuilder, PushInto, SizableContainer};
32use timely::progress::frontier::{Antichain, AntichainRef};
33
34use crate::columnar::Column;
35
36#[derive(Default)]
38pub struct Chunker<C> {
39 target: C,
44 ready: VecDeque<C>,
46}
47
48impl<C: Container + Clone + 'static> ContainerBuilder for Chunker<C> {
49 type Container = C;
50
51 fn extract(&mut self) -> Option<&mut Self::Container> {
52 if let Some(ready) = self.ready.pop_front() {
53 self.target = ready;
54 Some(&mut self.target)
55 } else {
56 None
57 }
58 }
59
60 fn finish(&mut self) -> Option<&mut Self::Container> {
61 self.extract()
62 }
63}
64
65impl<'a, D, T, R> PushInto<&'a mut Column<(D, T, R)>> for Chunker<ColumnationStack<(D, T, R)>>
66where
67 D: Columnar + Columnation,
68 for<'b> columnar::Ref<'b, D>: Ord + Copy,
69 T: Columnar + Columnation,
70 for<'b> columnar::Ref<'b, T>: Ord + Copy,
71 R: Columnar + Columnation + Semigroup + for<'b> Semigroup<columnar::Ref<'b, R>>,
72 for<'b> columnar::Ref<'b, R>: Ord,
73{
74 fn push_into(&mut self, container: &'a mut Column<(D, T, R)>) {
75 let borrowed = container.borrow();
78 let mut permutation = Vec::with_capacity(borrowed.len());
79 Extend::extend(&mut permutation, borrowed.into_index_iter());
80 permutation.sort();
81
82 self.target.clear();
83 let mut iter = permutation.drain(..);
85 if let Some((data, time, diff)) = iter.next() {
86 let mut owned_data = D::into_owned(data);
87 let mut owned_time = T::into_owned(time);
88
89 let mut prev_data = data;
90 let mut prev_time = time;
91 let mut prev_diff = <R as Columnar>::into_owned(diff);
92
93 for (data, time, diff) in iter {
94 if (&prev_data, &prev_time) == (&data, &time) {
95 prev_diff.plus_equals(&diff);
96 } else {
97 if !prev_diff.is_zero() {
98 D::copy_from(&mut owned_data, prev_data);
99 T::copy_from(&mut owned_time, prev_time);
100 let tuple = (owned_data, owned_time, prev_diff);
101 self.target.push_into(&tuple);
102 (owned_data, owned_time, prev_diff) = tuple;
103 }
104 prev_data = data;
105 prev_time = time;
106 R::copy_from(&mut prev_diff, diff);
107 }
108 }
109
110 if !prev_diff.is_zero() {
111 D::copy_from(&mut owned_data, prev_data);
112 T::copy_from(&mut owned_time, prev_time);
113 let tuple = (owned_data, owned_time, prev_diff);
114 self.target.push_into(&tuple);
115 }
116 }
117
118 if !self.target.is_empty() {
119 self.ready.push_back(std::mem::take(&mut self.target));
120 }
121 }
122}
123
124pub struct ColumnChunker<U: Columnar> {
131 target: Column<U>,
134 ready: VecDeque<Column<U>>,
136}
137
138impl<U: Columnar> Default for ColumnChunker<U> {
143 fn default() -> Self {
144 Self {
145 target: Column::default(),
146 ready: VecDeque::new(),
147 }
148 }
149}
150
151impl<U: Columnar> ContainerBuilder for ColumnChunker<U>
152where
153 U::Container: Clone + 'static,
154{
155 type Container = Column<U>;
156
157 fn extract(&mut self) -> Option<&mut Self::Container> {
158 if let Some(ready) = self.ready.pop_front() {
159 self.target = ready;
160 Some(&mut self.target)
161 } else {
162 None
163 }
164 }
165
166 fn finish(&mut self) -> Option<&mut Self::Container> {
167 self.extract()
168 }
169}
170
171impl<'a, D, T, R> PushInto<&'a mut Column<(D, T, R)>> for ColumnChunker<(D, T, R)>
172where
173 D: Columnar,
174 for<'b> columnar::Ref<'b, D>: Copy + Ord,
175 T: Columnar,
176 for<'b> columnar::Ref<'b, T>: Copy + Ord,
177 R: Columnar + Default + Semigroup + for<'b> Semigroup<columnar::Ref<'b, R>>,
178 for<'b> columnar::Ref<'b, R>: Ord,
179 for<'b> <D as Columnar>::Container: columnar::Push<columnar::Ref<'b, D>>,
180 for<'b> <T as Columnar>::Container: columnar::Push<columnar::Ref<'b, T>>,
181 for<'b> <R as Columnar>::Container: columnar::Push<&'b R>,
182{
183 fn push_into(&mut self, container: &'a mut Column<(D, T, R)>) {
184 match &mut self.target {
189 Column::Typed(c) => c.clear(),
190 Column::Bytes(_) | Column::Align(_) => {
191 self.target = Column::Typed(Default::default());
192 }
193 }
194
195 let borrowed = container.borrow();
197 let mut permutation = Vec::with_capacity(borrowed.len());
198 Extend::extend(&mut permutation, borrowed.into_index_iter());
199 permutation.sort();
200
201 {
207 let Column::Typed(target_c) = &mut self.target else {
208 unreachable!("target reset to Typed above");
209 };
210 let (target_d, target_t, target_r) = target_c;
211
212 let mut iter = permutation.drain(..);
213 if let Some((data, time, diff)) = iter.next() {
214 let mut prev_data = data;
215 let mut prev_time = time;
216 let mut prev_diff = <R as Columnar>::into_owned(diff);
217
218 for (data, time, diff) in iter {
219 if (&prev_data, &prev_time) == (&data, &time) {
220 prev_diff.plus_equals(&diff);
221 } else {
222 if !prev_diff.is_zero() {
223 target_d.push(prev_data);
224 target_t.push(prev_time);
225 target_r.push(&prev_diff);
226 }
227 prev_data = data;
228 prev_time = time;
229 R::copy_from(&mut prev_diff, diff);
230 }
231 }
232
233 if !prev_diff.is_zero() {
234 target_d.push(prev_data);
235 target_t.push(prev_time);
236 target_r.push(&prev_diff);
237 }
238 }
239 }
240
241 if !self.target.is_empty() {
242 let chunk = std::mem::replace(&mut self.target, Column::Typed(Default::default()));
243 self.ready.push_back(chunk);
244 }
245 }
246}
247
248pub(crate) fn gallop(upper: usize, lower: &mut usize, mut cmp: impl FnMut(usize) -> bool) {
263 if *lower < upper && cmp(*lower) {
264 let mut step = 1;
265 while *lower + step < upper && cmp(*lower + step) {
266 *lower += step;
267 step <<= 1;
268 }
269
270 step >>= 1;
271 while step > 0 {
272 if *lower + step < upper && cmp(*lower + step) {
273 *lower += step;
274 }
275 step >>= 1;
276 }
277
278 *lower += 1;
281 }
282}
283
284pub struct ColumnMerger<D, T, R> {
289 _marker: PhantomData<(D, T, R)>,
290}
291
292impl<D, T, R> Default for ColumnMerger<D, T, R> {
293 fn default() -> Self {
294 Self {
295 _marker: PhantomData,
296 }
297 }
298}
299
300impl<D, T, R> Column<(D, T, R)>
306where
307 D: Columnar,
308 for<'a> columnar::Ref<'a, D>: Copy + Ord,
309 T: Columnar + Default + Clone + PartialOrder,
310 for<'a> columnar::Ref<'a, T>: Copy + Ord,
311 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
312{
313 #[must_use]
331 pub fn merge_from(&mut self, others: &mut [Self], positions: &mut [usize]) -> bool {
332 match others.len() {
333 0 => false,
334 1 => {
335 let other = &mut others[0];
336 let pos = &mut positions[0];
337 if self.is_empty() && *pos == 0 {
338 std::mem::swap(self, other);
339 return false;
340 }
341 let Column::Typed(self_c) = self else {
342 unreachable!("merger chunks are always Column::Typed");
343 };
344 let src_c = other.borrow();
345 self_c.extend_from_self(src_c, *pos..other.borrow().len());
346 *pos = other.borrow().len();
347 false
348 }
349 2 => {
350 let (left, right) = others.split_at(1);
351 let (left_pos, right_pos) = positions.split_at_mut(1);
352 let left_borrow = left[0].borrow();
353 let right_borrow = right[0].borrow();
354
355 let Column::Typed(self_c) = self else {
356 unreachable!("merger chunks are always Column::Typed");
357 };
358
359 let l_d = left_borrow.0;
369 let l_t = left_borrow.1;
370 let l_r = left_borrow.2;
371 let r_d = right_borrow.0;
372 let r_t = right_borrow.1;
373 let r_r = right_borrow.2;
374 let upper_l = l_d.len();
375 let upper_r = r_d.len();
376
377 let (sd, st, sr) = self_c;
384
385 const RESERVE_RECORD_THRESHOLD: usize = 1_000_000;
399 if upper_l + upper_r <= RESERVE_RECORD_THRESHOLD {
400 use columnar::Container as _;
401 let inputs = [left_borrow, right_borrow];
402 sd.reserve_for(inputs.iter().map(|b| b.0));
403 st.reserve_for(inputs.iter().map(|b| b.1));
404 sr.reserve_for(inputs.iter().map(|b| b.2));
405 }
406
407 let mut stash = R::default();
408
409 let at_ship_threshold =
426 |sd: &D::Container, st: &T::Container, sr: &R::Container| {
427 use columnar::Borrow as _;
428 crate::columnar::at_serialized_capacity(&(
429 sd.borrow(),
430 st.borrow(),
431 sr.borrow(),
432 ))
433 };
434 const THRESHOLD_PERIOD_MASK: u32 = 1023;
435 let mut iter: u32 = 0;
436 let mut yielded = false;
437
438 while left_pos[0] < upper_l && right_pos[0] < upper_r {
439 let d1 = l_d.get(left_pos[0]);
440 let t1 = l_t.get(left_pos[0]);
441 let d2 = r_d.get(right_pos[0]);
442 let t2 = r_t.get(right_pos[0]);
443 match (d1, t1).cmp(&(d2, t2)) {
444 std::cmp::Ordering::Less => {
445 sd.push(d1);
453 st.push(t1);
454 sr.push(l_r.get(left_pos[0]));
455 left_pos[0] += 1;
456 if left_pos[0] < upper_l
461 && (l_d.get(left_pos[0]), l_t.get(left_pos[0])) < (d2, t2)
462 {
463 let start = left_pos[0];
464 gallop(upper_l, &mut left_pos[0], |i| {
465 (l_d.get(i), l_t.get(i)) < (d2, t2)
466 });
467 sd.extend_from_self(l_d, start..left_pos[0]);
471 st.extend_from_self(l_t, start..left_pos[0]);
472 sr.extend_from_self(l_r, start..left_pos[0]);
473 }
474 }
475 std::cmp::Ordering::Greater => {
476 sd.push(d2);
478 st.push(t2);
479 sr.push(r_r.get(right_pos[0]));
480 right_pos[0] += 1;
481 if right_pos[0] < upper_r
482 && (r_d.get(right_pos[0]), r_t.get(right_pos[0])) < (d1, t1)
483 {
484 let start = right_pos[0];
485 gallop(upper_r, &mut right_pos[0], |i| {
486 (r_d.get(i), r_t.get(i)) < (d1, t1)
487 });
488 sd.extend_from_self(r_d, start..right_pos[0]);
489 st.extend_from_self(r_t, start..right_pos[0]);
490 sr.extend_from_self(r_r, start..right_pos[0]);
491 }
492 }
493 std::cmp::Ordering::Equal => {
494 let r1 = l_r.get(left_pos[0]);
495 let r2 = r_r.get(right_pos[0]);
496 R::copy_from(&mut stash, r1);
497 stash.plus_equals(&r2);
498 if !stash.is_zero() {
499 sd.push(d1);
500 st.push(t1);
501 sr.push(&stash);
502 }
503 left_pos[0] += 1;
504 right_pos[0] += 1;
505 }
506 }
507
508 iter = iter.wrapping_add(1);
511 if iter & THRESHOLD_PERIOD_MASK == 0 && at_ship_threshold(sd, st, sr) {
512 yielded = true;
513 break;
514 }
515 }
516 yielded
517 }
518 n => unreachable!("merge_from called with {n} inputs; expected 0, 1, or 2"),
521 }
522 }
523
524 pub fn extract(
535 &mut self,
536 position: &mut usize,
537 upper: AntichainRef<T>,
538 frontier: &mut Antichain<T>,
539 keep: &mut Self,
540 ship: &mut Self,
541 ) {
542 let Column::Typed(keep_c) = keep else {
543 unreachable!("merger chunks are always Column::Typed");
544 };
545 let Column::Typed(ship_c) = ship else {
546 unreachable!("merger chunks are always Column::Typed");
547 };
548
549 let self_view = self.borrow();
550 let len = self_view.len();
551
552 use columnar::Borrow as _;
553 let mut owned_t = T::default();
554 while *position < len
555 && !crate::columnar::at_serialized_capacity(&keep_c.borrow())
556 && !crate::columnar::at_serialized_capacity(&ship_c.borrow())
557 {
558 let (_, time, _) = self_view.get(*position);
559 T::copy_from(&mut owned_t, time);
560 if upper.less_equal(&owned_t) {
561 frontier.insert_with(&owned_t, |t| t.clone());
564 keep_c.extend_from_self(self_view, *position..*position + 1);
565 } else {
566 ship_c.extend_from_self(self_view, *position..*position + 1);
567 }
568 *position += 1;
569 }
570 }
571}
572
573impl<D, T, R> Merger for ColumnMerger<D, T, R>
581where
582 D: Columnar,
583 for<'a> columnar::Ref<'a, D>: Copy + Ord,
584 T: Columnar + Default + Clone + Ord + PartialOrder,
585 for<'a> columnar::Ref<'a, T>: Copy + Ord,
586 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
587{
588 type Time = T;
589 type Chunk = Column<(D, T, R)>;
590
591 fn merge(
592 &mut self,
593 list1: Vec<Self::Chunk>,
594 list2: Vec<Self::Chunk>,
595 output: &mut Vec<Self::Chunk>,
596 stash: &mut Vec<Self::Chunk>,
597 ) {
598 let mut list1 = list1.into_iter();
599 let mut list2 = list2.into_iter();
600
601 let mut heads = [
602 list1.next().unwrap_or_default(),
603 list2.next().unwrap_or_default(),
604 ];
605 let mut positions = [0usize, 0usize];
606
607 let mut result = empty_chunk(stash);
608
609 loop {
611 let upper_l = heads[0].borrow().len();
612 let upper_r = heads[1].borrow().len();
613 if positions[0] >= upper_l || positions[1] >= upper_r {
614 break;
615 }
616
617 let lhs_passthrough = positions[0] == 0 && upper_l > 0 && {
625 let lhs = heads[0].borrow();
626 let rhs = heads[1].borrow();
627 let last_l = (lhs.0.get(upper_l - 1), lhs.1.get(upper_l - 1));
628 let cur_r = (rhs.0.get(positions[1]), rhs.1.get(positions[1]));
629 last_l < cur_r
630 };
631 if lhs_passthrough {
632 if !result.is_empty() {
633 output.push(std::mem::take(&mut result));
634 result = empty_chunk(stash);
635 }
636 let head = std::mem::replace(&mut heads[0], list1.next().unwrap_or_default());
637 output.push(head);
638 positions[0] = 0;
639 continue;
640 }
641
642 let rhs_passthrough = positions[1] == 0 && upper_r > 0 && {
643 let lhs = heads[0].borrow();
644 let rhs = heads[1].borrow();
645 let last_r = (rhs.0.get(upper_r - 1), rhs.1.get(upper_r - 1));
646 let cur_l = (lhs.0.get(positions[0]), lhs.1.get(positions[0]));
647 last_r < cur_l
648 };
649 if rhs_passthrough {
650 if !result.is_empty() {
651 output.push(std::mem::take(&mut result));
652 result = empty_chunk(stash);
653 }
654 let head = std::mem::replace(&mut heads[1], list2.next().unwrap_or_default());
655 output.push(head);
656 positions[1] = 0;
657 continue;
658 }
659
660 let yielded = result.merge_from(&mut heads, &mut positions);
664
665 if positions[0] >= heads[0].borrow().len() {
666 let old = std::mem::replace(&mut heads[0], list1.next().unwrap_or_default());
667 recycle_chunk(old, stash);
668 positions[0] = 0;
669 }
670 if positions[1] >= heads[1].borrow().len() {
671 let old = std::mem::replace(&mut heads[1], list2.next().unwrap_or_default());
672 recycle_chunk(old, stash);
673 positions[1] = 0;
674 }
675 if yielded || result.at_capacity() {
676 output.push(std::mem::take(&mut result));
677 result = empty_chunk(stash);
678 }
679 }
680
681 drain_side(
684 &mut heads[0],
685 &mut positions[0],
686 &mut list1,
687 &mut result,
688 output,
689 stash,
690 );
691 drain_side(
692 &mut heads[1],
693 &mut positions[1],
694 &mut list2,
695 &mut result,
696 output,
697 stash,
698 );
699 if !result.is_empty() {
700 output.push(result);
701 }
702 }
703
704 fn extract(
705 &mut self,
706 merged: Vec<Self::Chunk>,
707 upper: AntichainRef<Self::Time>,
708 frontier: &mut Antichain<Self::Time>,
709 ship: &mut Vec<Self::Chunk>,
710 kept: &mut Vec<Self::Chunk>,
711 stash: &mut Vec<Self::Chunk>,
712 ) {
713 let mut keep = empty_chunk(stash);
714 let mut ready = empty_chunk(stash);
715
716 for mut buffer in merged {
717 let mut position = 0;
718 let len = buffer.borrow().len();
719 while position < len {
720 buffer.extract(&mut position, upper, frontier, &mut keep, &mut ready);
721 if keep.at_capacity() {
722 kept.push(std::mem::take(&mut keep));
723 keep = empty_chunk(stash);
724 }
725 if ready.at_capacity() {
726 ship.push(std::mem::take(&mut ready));
727 ready = empty_chunk(stash);
728 }
729 }
730 recycle_chunk(buffer, stash);
731 }
732 if !keep.is_empty() {
733 kept.push(keep);
734 }
735 if !ready.is_empty() {
736 ship.push(ready);
737 }
738 }
739
740 fn len(chunk: &Self::Chunk) -> usize {
741 usize::try_from(chunk.record_count()).expect("record_count is non-negative")
742 }
743
744 fn allocation(chunk: &Self::Chunk) -> (usize, usize, usize) {
745 use timely::dataflow::channels::ContainerBytes;
746 let bytes = chunk.length_in_bytes();
752 (bytes, bytes, 1)
753 }
754}
755
756#[inline]
759pub(crate) fn empty_chunk<C: Columnar>(stash: &mut Vec<Column<C>>) -> Column<C> {
760 stash.pop().unwrap_or_default()
761}
762
763#[inline]
772pub(crate) fn recycle_chunk<C: Columnar>(mut chunk: Column<C>, stash: &mut Vec<Column<C>>) {
773 if let Column::Typed(c) = &mut chunk {
774 c.clear();
775 stash.push(chunk);
776 }
777}
778
779fn drain_side<D, T, R>(
785 head: &mut Column<(D, T, R)>,
786 pos: &mut usize,
787 list: &mut std::vec::IntoIter<Column<(D, T, R)>>,
788 result: &mut Column<(D, T, R)>,
789 output: &mut Vec<Column<(D, T, R)>>,
790 stash: &mut Vec<Column<(D, T, R)>>,
791) where
792 D: Columnar,
793 for<'a> columnar::Ref<'a, D>: Copy + Ord,
794 T: Columnar + Default + Clone + PartialOrder,
795 for<'a> columnar::Ref<'a, T>: Copy + Ord,
796 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
797{
798 if *pos < head.borrow().len() {
799 let _ = result.merge_from(std::slice::from_mut(head), std::slice::from_mut(pos));
802 }
803 if !result.is_empty() {
804 output.push(std::mem::take(result));
805 *result = empty_chunk(stash);
806 }
807 Extend::extend(output, list);
808}
809
810#[cfg(test)]
811mod tests {
812 use timely::dataflow::channels::ContainerBytes;
813
814 use super::*;
815
816 fn serialize<C: Columnar>(column: &Column<C>) -> Column<C> {
821 let mut bytes: Vec<u8> = Vec::new();
822 column.into_bytes(&mut bytes);
823 assert_eq!(bytes.len() % 8, 0);
824 let words = bytes
825 .chunks_exact(8)
826 .map(|w| u64::from_ne_bytes(w.try_into().expect("chunk is 8 bytes")))
827 .collect();
828 Column::Align(words)
829 }
830
831 fn run_chunker<D, T, R>(inputs: &[(D, T, R)]) -> Vec<(D, T, R)>
834 where
835 D: Columnar + Clone,
836 for<'a> columnar::Ref<'a, D>: Copy + Ord,
837 T: Columnar + Clone,
838 for<'a> columnar::Ref<'a, T>: Copy + Ord,
839 R: Columnar + Clone + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
840 for<'a> columnar::Ref<'a, R>: Ord,
841 <(D, T, R) as Columnar>::Container: Clone,
842 for<'a> <(D, T, R) as Columnar>::Container: columnar::Push<&'a (D, T, R)>,
843 <(D, T, R) as Columnar>::Container: columnar::Push<(D, T, R)>,
844 for<'a> <D as Columnar>::Container: columnar::Push<columnar::Ref<'a, D>>,
845 for<'a> <T as Columnar>::Container: columnar::Push<columnar::Ref<'a, T>>,
846 for<'a> <R as Columnar>::Container: columnar::Push<&'a R>,
847 {
848 let mut input: Column<(D, T, R)> = Default::default();
849 for tuple in inputs.iter().cloned() {
850 input.push_into(tuple);
851 }
852
853 let mut chunker: ColumnChunker<(D, T, R)> = Default::default();
854 chunker.push_into(&mut input);
855
856 let mut out = Vec::new();
857 while let Some(chunk) = chunker.extract() {
858 for (d, t, r) in chunk.borrow().into_index_iter() {
859 out.push((D::into_owned(d), T::into_owned(t), R::into_owned(r)));
860 }
861 }
862 out
863 }
864
865 #[mz_ore::test]
866 fn empty_input_yields_no_chunk() {
867 let mut chunker: ColumnChunker<(u64, u64, i64)> = Default::default();
868 let mut input: Column<(u64, u64, i64)> = Default::default();
869 chunker.push_into(&mut input);
870 assert!(chunker.extract().is_none());
871 assert!(chunker.finish().is_none());
872 }
873
874 #[mz_ore::test]
875 fn unsorted_input_is_sorted() {
876 let out = run_chunker(&[(3u64, 0u64, 1i64), (1u64, 0u64, 1i64), (2u64, 0u64, 1i64)]);
877 assert_eq!(out, vec![(1, 0, 1), (2, 0, 1), (3, 0, 1)]);
878 }
879
880 #[mz_ore::test]
881 fn serialized_input_is_consolidated() {
882 let mut input: Column<(u64, u64, i64)> = Default::default();
885 for tuple in [
886 (2u64, 0u64, 1i64),
887 (1, 0, 1),
888 (2, 0, 2),
889 (3, 0, 1),
890 (3, 0, -1),
891 ] {
892 input.push_into(tuple);
893 }
894 let mut input = serialize(&input);
895
896 let mut chunker: ColumnChunker<(u64, u64, i64)> = Default::default();
897 chunker.push_into(&mut input);
898
899 let mut out = Vec::new();
900 while let Some(chunk) = chunker.extract() {
901 for (d, t, r) in chunk.borrow().into_index_iter() {
902 out.push((*d, *t, *r));
903 }
904 }
905 assert_eq!(out, vec![(1, 0, 1), (2, 0, 3)]);
906 }
907
908 #[mz_ore::test]
909 fn duplicate_keys_consolidate() {
910 let out = run_chunker(&[(1u64, 0u64, 1i64), (1u64, 0u64, 2i64), (1u64, 0u64, -1i64)]);
911 assert_eq!(out, vec![(1, 0, 2)]);
912 }
913
914 #[mz_ore::test]
915 fn diffs_summing_to_zero_are_dropped() {
916 let out = run_chunker(&[(1u64, 0u64, 1i64), (1u64, 0u64, -1i64)]);
917 assert!(out.is_empty());
918 }
919
920 #[mz_ore::test]
921 fn mixed_consolidation() {
922 let out = run_chunker(&[
926 (1u64, 0u64, 1i64),
927 (2u64, 0u64, 1i64),
928 (1u64, 0u64, 2i64),
929 (1u64, 1u64, 5i64),
930 (1u64, 0u64, -3i64),
931 ]);
932 assert_eq!(out, vec![(1, 1, 5), (2, 0, 1)]);
933 }
934
935 #[mz_ore::test]
936 fn key_val_tuple_data() {
937 let out = run_chunker(&[
939 ((1u64, 10u64), 0u64, 1i64),
940 ((1u64, 10u64), 0u64, 1i64),
941 ((1u64, 11u64), 0u64, 1i64),
942 ((2u64, 10u64), 0u64, 1i64),
943 ]);
944 assert_eq!(
945 out,
946 vec![((1, 10), 0, 2), ((1, 11), 0, 1), ((2, 10), 0, 1),]
947 );
948 }
949
950 #[mz_ore::test]
951 fn buffer_reuse_across_calls() {
952 let mut input1: Column<(u64, u64, i64)> = Default::default();
955 input1.push_into((1u64, 0u64, 1i64));
956 input1.push_into((2u64, 0u64, 1i64));
957
958 let mut input2: Column<(u64, u64, i64)> = Default::default();
959 input2.push_into((3u64, 0u64, 1i64));
960 input2.push_into((1u64, 0u64, 1i64));
961
962 let mut chunker: ColumnChunker<(u64, u64, i64)> = Default::default();
963 chunker.push_into(&mut input1);
964
965 {
968 let _ = chunker.extract().expect("first chunk");
969 }
970
971 chunker.push_into(&mut input2);
972
973 let chunk = chunker.extract().expect("second chunk");
974 let collected: Vec<_> = chunk
975 .borrow()
976 .into_index_iter()
977 .map(|(d, t, r)| (u64::into_owned(d), u64::into_owned(t), i64::into_owned(r)))
978 .collect();
979 assert_eq!(collected, vec![(1, 0, 1), (3, 0, 1)]);
980 }
981
982 fn col(rows: &[((u64, u64), u64, i64)]) -> Column<((u64, u64), u64, i64)> {
984 let mut c: Column<((u64, u64), u64, i64)> = Default::default();
985 for &t in rows {
986 c.push_into(t);
987 }
988 c
989 }
990
991 fn collect_chunks(chunks: &[Column<((u64, u64), u64, i64)>]) -> Vec<((u64, u64), u64, i64)> {
992 chunks
993 .iter()
994 .flat_map(|c| {
995 c.borrow().into_index_iter().map(|((k, v), t, r)| {
996 (
997 (u64::into_owned(k), u64::into_owned(v)),
998 u64::into_owned(t),
999 i64::into_owned(r),
1000 )
1001 })
1002 })
1003 .collect()
1004 }
1005
1006 #[mz_ore::test]
1011 fn merger_disjoint_chains_passthrough() {
1012 let chain1 = vec![
1013 col(&[((0, 0), 0, 1), ((1, 0), 0, 1)]),
1014 col(&[((2, 0), 0, 1), ((3, 0), 0, 1)]),
1015 ];
1016 let chain2 = vec![
1017 col(&[((10, 0), 0, 1), ((11, 0), 0, 1)]),
1018 col(&[((12, 0), 0, 1), ((13, 0), 0, 1)]),
1019 ];
1020
1021 let mut merger: ColumnMerger<(u64, u64), u64, i64> = Default::default();
1022 let mut output = Vec::new();
1023 let mut stash = Vec::new();
1024 Merger::merge(&mut merger, chain1, chain2, &mut output, &mut stash);
1025
1026 let collected = collect_chunks(&output);
1027 let expected: Vec<_> = (0..4u64)
1028 .map(|d| ((d, 0u64), 0u64, 1i64))
1029 .chain((10..14u64).map(|d| ((d, 0u64), 0u64, 1i64)))
1030 .collect();
1031 assert_eq!(collected, expected);
1032 }
1033
1034 #[mz_ore::test]
1039 fn merger_interleaved_chains() {
1040 let chain1 = vec![
1043 col(&[((0, 0), 0, 1), ((2, 0), 0, 1)]),
1044 col(&[((4, 0), 0, 1), ((6, 0), 0, 1)]),
1045 ];
1046 let chain2 = vec![
1047 col(&[((1, 0), 0, 1), ((3, 0), 0, 1)]),
1048 col(&[((5, 0), 0, 1), ((7, 0), 0, 1)]),
1049 ];
1050
1051 let mut merger: ColumnMerger<(u64, u64), u64, i64> = Default::default();
1052 let mut output = Vec::new();
1053 let mut stash = Vec::new();
1054 Merger::merge(&mut merger, chain1, chain2, &mut output, &mut stash);
1055
1056 let collected = collect_chunks(&output);
1057 let expected: Vec<_> = (0..8u64).map(|d| ((d, 0u64), 0u64, 1i64)).collect();
1058 assert_eq!(collected, expected);
1059 }
1060
1061 #[mz_ore::test]
1065 fn merger_passthrough_respects_equal_boundary() {
1066 let chain1 = vec![col(&[((0, 0), 0, 1), ((5, 0), 0, 1)])];
1071 let chain2 = vec![col(&[((5, 0), 0, 1), ((10, 0), 0, 1)])];
1072
1073 let mut merger: ColumnMerger<(u64, u64), u64, i64> = Default::default();
1074 let mut output = Vec::new();
1075 let mut stash = Vec::new();
1076 Merger::merge(&mut merger, chain1, chain2, &mut output, &mut stash);
1077
1078 let collected = collect_chunks(&output);
1079 assert_eq!(
1080 collected,
1081 vec![((0, 0), 0, 1), ((5, 0), 0, 2), ((10, 0), 0, 1)]
1082 );
1083 }
1084}
1085
1086#[cfg(test)]
1087mod proptests {
1088 use super::*;
1098 use mz_ore::cast::CastFrom;
1099 use proptest::prelude::*;
1100 use timely::progress::frontier::Antichain;
1101
1102 type Tuple = ((u64, u64), u64, i64);
1103
1104 fn consolidate(mut v: Vec<Tuple>) -> Vec<Tuple> {
1107 v.sort();
1108 let mut out: Vec<Tuple> = Vec::new();
1109 for (d, t, r) in v {
1110 if let Some(last) = out.last_mut() {
1111 if last.0 == d && last.1 == t {
1112 last.2 += r;
1113 continue;
1114 }
1115 }
1116 out.push((d, t, r));
1117 }
1118 out.retain(|x| x.2 != 0);
1119 out
1120 }
1121
1122 fn arb_consolidated() -> impl Strategy<Value = Vec<Tuple>> {
1125 prop::collection::vec(((0u64..5, 0u64..5), 0u64..3, -3i64..=3i64), 0..30)
1126 .prop_map(consolidate)
1127 }
1128
1129 fn build_column(v: &[Tuple]) -> Column<Tuple> {
1130 let mut col: Column<Tuple> = Default::default();
1131 for tup in v {
1132 col.push_into(*tup);
1133 }
1134 col
1135 }
1136
1137 fn collect_column(col: &Column<Tuple>) -> Vec<Tuple> {
1138 col.borrow()
1139 .into_index_iter()
1140 .map(|((k, v), t, r)| {
1141 (
1142 (u64::into_owned(k), u64::into_owned(v)),
1143 u64::into_owned(t),
1144 i64::into_owned(r),
1145 )
1146 })
1147 .collect()
1148 }
1149
1150 fn drive_merge(left: Column<Tuple>, right: Column<Tuple>) -> Column<Tuple> {
1154 let mut self_col: Column<Tuple> = Default::default();
1155 let mut others = [left, right];
1156 let mut positions = [0usize, 0];
1157 let _ = self_col.merge_from(&mut others, &mut positions);
1158
1159 let [left_done, right_done] = others;
1160 let [left_pos, right_pos] = positions;
1161
1162 if left_pos < left_done.borrow().len() {
1163 let mut tail = [left_done];
1164 let mut p = [left_pos];
1165 let _ = self_col.merge_from(&mut tail, &mut p);
1166 } else if right_pos < right_done.borrow().len() {
1167 let mut tail = [right_done];
1168 let mut p = [right_pos];
1169 let _ = self_col.merge_from(&mut tail, &mut p);
1170 }
1171
1172 self_col
1173 }
1174
1175 proptest! {
1176 #[mz_ore::test]
1179 #[cfg_attr(miri, ignore)]
1180 fn merge_from_equals_consolidated_union(
1181 a in arb_consolidated(),
1182 b in arb_consolidated(),
1183 ) {
1184 let merged = drive_merge(build_column(&a), build_column(&b));
1185
1186 let mut union = a.clone();
1187 Extend::extend(&mut union, b.iter().copied());
1188 let expected = consolidate(union);
1189
1190 prop_assert_eq!(collect_column(&merged), expected);
1191 }
1192
1193 #[mz_ore::test]
1196 #[cfg_attr(miri, ignore)]
1197 fn merge_from_one_input_drains_tail(
1198 data in arb_consolidated(),
1199 pos_frac in 0u32..=100,
1200 ) {
1201 let len = data.len();
1203 let start_pos = if len == 0 { 0 } else {
1204 (usize::cast_from(pos_frac) * len) / 101
1205 };
1206
1207 let mut self_col: Column<Tuple> = Default::default();
1210 let sentinel: Tuple = ((u64::MAX, u64::MAX), 0, 1);
1211 self_col.push_into(sentinel);
1212
1213 let mut others = [build_column(&data)];
1214 let mut positions = [start_pos];
1215 let _ = self_col.merge_from(&mut others, &mut positions);
1216
1217 let mut expected = vec![sentinel];
1218 Extend::extend(&mut expected, data[start_pos..].iter().copied());
1219
1220 prop_assert_eq!(collect_column(&self_col), expected);
1221 prop_assert_eq!(positions[0], len);
1222 }
1223
1224 #[mz_ore::test]
1227 #[cfg_attr(miri, ignore)]
1228 fn merge_from_empty_self_swap(data in arb_consolidated()) {
1229 let mut self_col: Column<Tuple> = Default::default();
1230 let mut others = [build_column(&data)];
1231 let mut positions = [0usize];
1232 let _ = self_col.merge_from(&mut others, &mut positions);
1233
1234 prop_assert_eq!(collect_column(&self_col), data);
1235 }
1236
1237 #[mz_ore::test]
1243 #[cfg_attr(miri, ignore)]
1244 fn extract_partitions_by_frontier(
1245 data in arb_consolidated(),
1246 upper_time in 0u64..=4,
1247 ) {
1248 let mut self_col = build_column(&data);
1249 let upper = Antichain::from_elem(upper_time);
1250 let mut frontier: Antichain<u64> = Antichain::new();
1251 let mut keep: Column<Tuple> = Default::default();
1252 let mut ship: Column<Tuple> = Default::default();
1253 let mut position = 0;
1254
1255 self_col.extract(
1256 &mut position,
1257 upper.borrow(),
1258 &mut frontier,
1259 &mut keep,
1260 &mut ship,
1261 );
1262
1263 prop_assert_eq!(position, data.len());
1265
1266 let kept = collect_column(&keep);
1267 let shipped = collect_column(&ship);
1268
1269 for (_, t, _) in &kept {
1271 prop_assert!(
1272 upper.borrow().less_equal(t),
1273 "kept time {} should satisfy upper.less_equal", t,
1274 );
1275 }
1276 for (_, t, _) in &shipped {
1277 prop_assert!(
1278 !upper.borrow().less_equal(t),
1279 "shipped time {} should NOT satisfy upper.less_equal", t,
1280 );
1281 }
1282
1283 let mut union = kept.clone();
1285 Extend::extend(&mut union, shipped.iter().copied());
1286 union.sort();
1287 let mut expected_sorted = data.clone();
1288 expected_sorted.sort();
1289 prop_assert_eq!(union, expected_sorted);
1290
1291 for (_, t, _) in &kept {
1293 prop_assert!(
1294 frontier.less_equal(t),
1295 "frontier should dominate kept time {}", t,
1296 );
1297 }
1298 }
1299
1300 #[mz_ore::test]
1302 #[cfg_attr(miri, ignore)]
1303 fn extract_empty_input(upper_time in 0u64..=4) {
1304 let mut self_col: Column<Tuple> = Default::default();
1305 let upper = Antichain::from_elem(upper_time);
1306 let mut frontier: Antichain<u64> = Antichain::new();
1307 let mut keep: Column<Tuple> = Default::default();
1308 let mut ship: Column<Tuple> = Default::default();
1309 let mut position = 0;
1310
1311 self_col.extract(
1312 &mut position,
1313 upper.borrow(),
1314 &mut frontier,
1315 &mut keep,
1316 &mut ship,
1317 );
1318
1319 prop_assert_eq!(position, 0);
1320 prop_assert!(collect_column(&keep).is_empty());
1321 prop_assert!(collect_column(&ship).is_empty());
1322 prop_assert!(frontier.elements().is_empty());
1323 }
1324 }
1325}