1use std::borrow::Cow;
42#[cfg(test)]
43use std::cell::Cell;
44use std::cell::RefCell;
45use std::collections::VecDeque;
46use std::rc::Rc;
47use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
48
49use columnar::bytes::indexed;
50use columnar::{Borrow, BorrowedOf, Columnar, Container as _, Index, Len, Push as _};
51use differential_dataflow::difference::Semigroup;
52use differential_dataflow::lattice::Lattice;
53use differential_dataflow::trace::chunk::Chunk;
54use mz_ore::cast::CastFrom;
55use mz_ore::pool::{ChunkHandle, ChunkHints, ExtentCodec, IDENTITY_CODEC, Pool};
56use smallvec::SmallVec;
57use timely::Accountable;
58use timely::PartialOrder;
59use timely::container::{ContainerBuilder, PushInto};
60use timely::progress::Timestamp;
61use timely::progress::frontier::{Antichain, AntichainRef};
62
63use crate::columnar::batcher::{ColumnChunker, gallop};
64use crate::columnar::body::{ColumnBody, borrow_words};
65use crate::columnar::unload::UnloadChunk;
66use crate::columnar::{Column, at_serialized_capacity};
67
68pub mod metrics;
69
70static COMPUTE_SPILL_ENABLED: AtomicBool = AtomicBool::new(false);
72
73static STORAGE_SPILL_ENABLED: AtomicBool = AtomicBool::new(false);
75
76static SINK_SPILL_ENABLED: AtomicBool = AtomicBool::new(false);
79
80thread_local! {
81 static SPILL_OVERRIDE: RefCell<Option<Pool>> = const { RefCell::new(None) };
85
86 #[cfg(test)]
90 static COMPRESS_MIN_DEPTH_OVERRIDE: Cell<Option<u8>> = const { Cell::new(None) };
91
92 static READ_SCRATCH: RefCell<Vec<u64>> = const { RefCell::new(Vec::new()) };
94}
95
96pub fn set_compute_spill_enabled(enabled: bool) {
110 COMPUTE_SPILL_ENABLED.store(enabled, Ordering::Relaxed);
111}
112
113pub fn set_storage_spill_enabled(enabled: bool) {
117 STORAGE_SPILL_ENABLED.store(enabled, Ordering::Relaxed);
118}
119
120pub fn set_sink_spill_enabled(enabled: bool) {
128 SINK_SPILL_ENABLED.store(enabled, Ordering::Relaxed);
129}
130
131pub fn set_spill_override(pool: Option<Pool>) {
135 SPILL_OVERRIDE.with(|cell| *cell.borrow_mut() = pool);
136}
137
138static COMPRESS_MIN_DEPTH: AtomicU8 = AtomicU8::new(DEFAULT_COMPRESS_MIN_DEPTH);
141
142pub fn set_compress_min_depth(depth: u8) {
161 COMPRESS_MIN_DEPTH.store(depth, Ordering::Relaxed);
162}
163
164#[cfg(test)]
168pub fn set_compress_min_depth_override(depth: Option<u8>) {
169 COMPRESS_MIN_DEPTH_OVERRIDE.with(|cell| cell.set(depth));
170}
171
172fn compress_min_depth() -> u8 {
174 #[cfg(test)]
175 if let Some(depth) = COMPRESS_MIN_DEPTH_OVERRIDE.with(|cell| cell.get()) {
176 return depth;
177 }
178 COMPRESS_MIN_DEPTH.load(Ordering::Relaxed)
179}
180
181fn codec_for_depth(depth: u8) -> (&'static dyn ExtentCodec, bool) {
186 if depth < compress_min_depth() {
187 (&IDENTITY_CODEC, false)
188 } else {
189 (&LZ4_CODEC, true)
190 }
191}
192
193fn spill_target(enabled: bool, len_bytes: usize) -> Option<Pool> {
199 resolve_pool(enabled).filter(|_| len_bytes >= SPILL_MIN_BYTES)
200}
201
202fn spill_pool() -> Option<Pool> {
204 resolve_pool(chunk_spill_enabled())
205}
206
207fn chunk_spill_enabled() -> bool {
209 COMPUTE_SPILL_ENABLED.load(Ordering::Relaxed) || STORAGE_SPILL_ENABLED.load(Ordering::Relaxed)
210}
211
212fn resolve_pool(enabled: bool) -> Option<Pool> {
214 if let Some(pool) = SPILL_OVERRIDE.with(|cell| cell.borrow().clone()) {
215 return Some(pool);
216 }
217 if enabled {
218 crate::pool_config::active_pool()
219 } else {
220 None
221 }
222}
223
224const SCRATCH_RETAIN_WORDS: usize = 1 << 18;
228
229fn with_scratch<Out>(f: impl FnOnce(&mut Vec<u64>) -> Out) -> Out {
231 READ_SCRATCH.with(|cell| {
232 let mut scratch = cell.take();
233 scratch.clear();
234 let out = f(&mut scratch);
235 if scratch.capacity() > SCRATCH_RETAIN_WORDS {
236 scratch.clear();
237 scratch.shrink_to_fit();
238 }
239 cell.replace(scratch);
240 out
241 })
242}
243
244const COMMIT_BYTES: usize = 2 << 20;
247
248const SPILL_MIN_BYTES: usize = 64 << 10;
257
258const DEFAULT_COMPRESS_MIN_DEPTH: u8 = 1;
267
268fn cut_records(len: usize, bytes: usize, space: usize) -> usize {
274 let space = space.saturating_sub(COMMIT_BYTES / 20);
275 (len.saturating_mul(space) / bytes.max(1)).min(len)
276}
277
278fn at_commit_size<C: Columnar>(body: &ColumnBody<C>) -> bool {
280 body.length_in_bytes() >= COMMIT_BYTES - COMMIT_BYTES / 10
281}
282
283#[inline(always)]
287fn rr<'b, 'a: 'b, C: Columnar>(item: columnar::Ref<'a, C>) -> columnar::Ref<'b, C> {
288 columnar::ContainerOf::<C>::reborrow_ref(item)
289}
290
291pub struct SpilledBody<D: Columnar, T> {
297 records: usize,
299 fences: D::Container,
303 time_lower: Antichain<T>,
307 time_upper: SmallVec<[T; 1]>,
312 compressed: bool,
319 len_bytes: usize,
324 handle: ChunkHandle,
326}
327
328pub enum ColumnChunk<D: Columnar, T: Columnar, R: Columnar> {
345 Resident(Rc<ColumnBody<(D, T, R)>>, u8),
347 Spilled(Rc<SpilledBody<D, T>>, u8),
349}
350
351impl<D: Columnar, T: Columnar, R: Columnar> Clone for ColumnChunk<D, T, R> {
352 fn clone(&self) -> Self {
353 match self {
354 ColumnChunk::Resident(col, depth) => ColumnChunk::Resident(Rc::clone(col), *depth),
355 ColumnChunk::Spilled(body, depth) => ColumnChunk::Spilled(Rc::clone(body), *depth),
356 }
357 }
358}
359
360impl<D: Columnar, T: Columnar, R: Columnar> Default for ColumnChunk<D, T, R> {
361 fn default() -> Self {
362 ColumnChunk::Resident(Rc::new(ColumnBody::default()), 0)
363 }
364}
365
366impl<D: Columnar, T: Columnar, R: Columnar> Accountable for ColumnChunk<D, T, R> {
367 fn record_count(&self) -> i64 {
368 i64::try_from(self.records()).expect("record count fits i64")
369 }
370}
371
372impl<D: Columnar, T: Columnar, R: Columnar> ColumnChunk<D, T, R> {
373 pub fn from_body(body: ColumnBody<(D, T, R)>) -> Self {
376 mz_ore::soft_assert_no_log!(!body.is_empty(), "chunks must be non-empty");
377 ColumnChunk::Resident(Rc::new(body), 0)
378 }
379
380 fn push_bounded(body: ColumnBody<(D, T, R)>, depth: u8, out: &mut VecDeque<Self>) {
384 let len = body.len();
385 let bytes = body.length_in_bytes();
386 if len <= 1 || bytes <= COMMIT_BYTES {
387 if len > 0 {
388 out.push_back(Self::Resident(Rc::new(body), depth));
389 }
390 return;
391 }
392 Self::push_cuts(
393 &body,
394 cut_records(len, bytes, COMMIT_BYTES).max(1),
395 depth,
396 out,
397 );
398 }
399
400 fn push_cuts(
407 body: &ColumnBody<(D, T, R)>,
408 records: usize,
409 depth: u8,
410 out: &mut VecDeque<Self>,
411 ) {
412 let view = body.borrow();
413 Self::push_range(view, 0..view.len(), records, depth, out);
414 }
415
416 fn push_range(
421 view: BorrowedOf<'_, (D, T, R)>,
422 range: std::ops::Range<usize>,
423 records: usize,
424 depth: u8,
425 out: &mut VecDeque<Self>,
426 ) {
427 let mut start = range.start;
428 while start < range.end {
429 let end = range.end.min(start + records);
430 let mut part = <(D, T, R) as Columnar>::Container::default();
431 part.extend_from_self(view, start..end);
432 let part = ColumnBody::Typed(part);
433 if end - start > 1 && part.length_in_bytes() > COMMIT_BYTES {
434 drop(part);
435 Self::push_range(view, start..end, (end - start).div_ceil(2), depth, out);
436 } else {
437 out.push_back(Self::Resident(Rc::new(part), depth));
438 }
439 start = end;
440 }
441 }
442
443 pub fn into_body(self) -> ColumnBody<(D, T, R)> {
446 match self {
447 ColumnChunk::Resident(col, _) => {
448 Rc::try_unwrap(col).unwrap_or_else(|shared| shared.duplicate())
449 }
450 ColumnChunk::Spilled(body, _) => {
451 let mut words = Vec::new();
452 body.handle.read_into(&mut words);
453 ColumnBody::Words(words)
454 }
455 }
456 }
457
458 pub fn with_body<F, X>(&self, f: F) -> X
461 where
462 F: FnOnce(&ColumnBody<(D, T, R)>) -> X,
463 {
464 match self {
465 ColumnChunk::Resident(col, _) => f(col),
466 ColumnChunk::Spilled(body, _) => {
467 let mut words = Vec::new();
468 body.handle.read_into(&mut words);
469 f(&ColumnBody::Words(words))
470 }
471 }
472 }
473
474 pub fn is_spilled(&self) -> bool {
476 matches!(self, ColumnChunk::Spilled(_, _))
477 }
478
479 fn records(&self) -> usize {
481 match self {
482 ColumnChunk::Resident(col, _) => col.len(),
483 ColumnChunk::Spilled(body, _) => body.records,
484 }
485 }
486
487 fn depth(&self) -> u8 {
489 match self {
490 ColumnChunk::Resident(_, depth) | ColumnChunk::Spilled(_, depth) => *depth,
491 }
492 }
493
494 fn data_span(&self) -> (columnar::Ref<'_, D>, columnar::Ref<'_, D>) {
496 match self {
497 ColumnChunk::Resident(col, _) => {
498 let data = col.borrow().0;
499 (data.get(0), data.get(data.len() - 1))
500 }
501 ColumnChunk::Spilled(body, _) => {
502 let fences = body.fences.borrow();
503 (fences.get(0), fences.get(1))
504 }
505 }
506 }
507
508 fn commit(body: ColumnBody<(D, T, R)>, depth: u8) -> Self
512 where
513 T: Timestamp,
514 {
515 let len_bytes = body.length_in_bytes();
516 metrics::record(metrics::Stage::Commit, body.len(), len_bytes);
517 mz_ore::soft_assert_no_log!(!body.is_empty(), "chunks must be non-empty");
518 match spill_target(chunk_spill_enabled(), len_bytes) {
519 Some(pool) => Self::spill_body(body, &pool, depth),
520 None => ColumnChunk::Resident(Rc::new(body), depth),
521 }
522 }
523
524 fn spill_body(body: ColumnBody<(D, T, R)>, pool: &Pool, depth: u8) -> Self
532 where
533 T: Timestamp,
534 {
535 let (codec, compressed) = codec_for_depth(depth);
536 let len_bytes = body.length_in_bytes();
537 let (time_lower, time_upper) = Self::time_bounds(&body);
538 let view = body.borrow();
539 let records = view.len();
540 let mut fences = D::Container::default();
541 fences.push(view.0.get(0));
542 fences.push(view.0.get(records - 1));
543 let handle = spill_serialized(&body, pool, len_bytes, ChunkHints { depth }, codec);
544 ColumnChunk::Spilled(
545 Rc::new(SpilledBody {
546 records,
547 fences,
548 time_lower,
549 time_upper: time_upper.into(),
550 compressed,
551 len_bytes,
552 handle,
553 }),
554 depth,
555 )
556 }
557
558 fn survive_merge(self) -> Self
574 where
575 T: Timestamp,
576 {
577 let depth = self.depth().saturating_add(1);
578 match self {
579 ColumnChunk::Resident(col, _) => ColumnChunk::Resident(col, depth),
580 ColumnChunk::Spilled(body, was) => {
581 let migrate = !body.compressed && depth >= compress_min_depth();
582 if !migrate || Rc::strong_count(&body) > 1 {
583 return ColumnChunk::Spilled(body, depth);
584 }
585 match spill_pool() {
586 Some(pool) => {
587 let body = ColumnChunk::Spilled(body, was).into_body();
588 Self::spill_body(body, &pool, depth)
589 }
590 None => ColumnChunk::Spilled(body, depth),
591 }
592 }
593 }
594 }
595
596 fn chunk_time_bounds(&self) -> (Cow<'_, Antichain<T>>, Cow<'_, [T]>)
601 where
602 T: Timestamp,
603 {
604 match self {
605 ColumnChunk::Resident(col, _) => {
606 let (lower, upper) = Self::time_bounds(col);
607 (Cow::Owned(lower), Cow::Owned(upper))
608 }
609 ColumnChunk::Spilled(body, _) => (
610 Cow::Borrowed(&body.time_lower),
611 Cow::Borrowed(&body.time_upper[..]),
612 ),
613 }
614 }
615
616 fn time_bounds(body: &ColumnBody<(D, T, R)>) -> (Antichain<T>, Vec<T>)
621 where
622 T: Timestamp,
623 {
624 let (_, times, _) = body.borrow();
625 let mut lower = Antichain::new();
626 let mut upper: Vec<T> = Vec::new();
627 let mut time = T::minimum();
631 for i in 0..times.len() {
632 time.copy_from(rr::<T>(times.get(i)));
633 if !upper.iter().any(|u| PartialOrder::less_equal(&time, u)) {
634 upper.retain(|u| !PartialOrder::less_equal(u, &time));
635 upper.push(time.clone());
636 }
637 lower.insert_ref(&time);
638 }
639 (lower, upper)
640 }
641}
642
643#[derive(Debug)]
648pub struct Lz4Codec;
649
650pub static LZ4_CODEC: Lz4Codec = Lz4Codec;
653
654impl ExtentCodec for Lz4Codec {
655 fn encode(&self, body: &[u8], out: &mut Vec<u8>) {
656 let max_out = lz4_flex::block::get_maximum_output_size(body.len());
657 out.resize(4 + max_out, 0);
658 let len = u32::try_from(body.len()).expect("chunk bodies are bounded by the size classes");
659 out[..4].copy_from_slice(&len.to_le_bytes());
660 let compressed = lz4_flex::block::compress_into(body, &mut out[4..])
661 .expect("output sized to the maximum");
662 out.truncate(4 + compressed);
663 }
664
665 fn decode(&self, stored: &[u8], body: &mut [u8]) {
666 let prefix: [u8; 4] = stored[..4].try_into().expect("prefix length");
667 let len = usize::try_from(u32::from_le_bytes(prefix)).expect("length fits usize");
668 assert_eq!(
669 len,
670 body.len(),
671 "destination must match the encoded body length"
672 );
673 let written = lz4_flex::block::decompress_into(&stored[4..], body)
674 .expect("stored bytes hold a valid lz4 block");
675 assert_eq!(written, body.len(), "decoded length mismatch");
676 }
677}
678
679fn spill_serialized<C: Columnar>(
683 body: &ColumnBody<C>,
684 pool: &Pool,
685 len_bytes: usize,
686 hints: ChunkHints,
687 codec: &'static dyn ExtentCodec,
688) -> ChunkHandle {
689 mz_ore::soft_assert_eq_no_log!(len_bytes % 8, 0);
690 pool.insert_with(len_bytes / 8, hints, codec, |dst| {
691 let bytes: &mut [u8] = bytemuck::cast_slice_mut(dst);
692 let mut cursor = std::io::Cursor::new(bytes);
693 body.write_into(&mut cursor)
694 .expect("the slot is sized to the serialized body");
695 assert_eq!(
696 usize::try_from(cursor.position()).expect("usize position"),
697 len_bytes,
698 "serialized body must fill the chunk exactly",
699 );
700 })
701}
702
703pub fn try_spill_ref<C: Columnar>(body: &ColumnBody<C>, depth: u8) -> Option<ChunkHandle> {
716 let len_bytes = body.length_in_bytes();
717 let pool = spill_target(SINK_SPILL_ENABLED.load(Ordering::Relaxed), len_bytes)?;
718 let (codec, _) = codec_for_depth(depth);
719 Some(spill_serialized(
720 body,
721 &pool,
722 len_bytes,
723 ChunkHints { depth },
724 codec,
725 ))
726}
727
728impl<D, T, R> Chunk for ColumnChunk<D, T, R>
729where
730 D: Columnar,
731 for<'a> columnar::Ref<'a, D>: Copy + Ord,
732 T: Columnar + Default + Timestamp + Lattice + Ord,
733 for<'a> columnar::Ref<'a, T>: Copy + Ord,
734 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
735{
736 type Time = T;
737
738 const TARGET: usize = 65536;
744
745 fn len(&self) -> usize {
746 self.records()
747 }
748
749 fn merge(in1: &mut VecDeque<Self>, in2: &mut VecDeque<Self>, out: &mut VecDeque<Self>) {
756 let (a_first, a_last) = in1
761 .front()
762 .expect("caller guarantees non-empty input")
763 .data_span();
764 let (b_first, b_last) = in2
765 .front()
766 .expect("caller guarantees non-empty input")
767 .data_span();
768 let a_low = rr::<D>(a_last) < rr::<D>(b_first);
769 let b_low = rr::<D>(b_last) < rr::<D>(a_first);
770 if a_low {
771 let chunk = in1.pop_front().expect("front observed above");
772 out.push_back(chunk.survive_merge());
773 return;
774 }
775 if b_low {
776 let chunk = in2.pop_front().expect("front observed above");
777 out.push_back(chunk.survive_merge());
778 return;
779 }
780
781 let a = in1.pop_front().expect("caller guarantees non-empty input");
782 let b = in2.pop_front().expect("caller guarantees non-empty input");
783 let depths = [a.depth(), b.depth()];
786 let out_depth = depths[0].max(depths[1]).saturating_add(1);
787 let mut spill_a = match &a {
788 ColumnChunk::Spilled(body, _) => Some(Rc::clone(body)),
789 ColumnChunk::Resident(_, _) => None,
790 };
791 let mut spill_b = match &b {
792 ColumnChunk::Spilled(body, _) => Some(Rc::clone(body)),
793 ColumnChunk::Resident(_, _) => None,
794 };
795 let mut cols = [a.into_body(), b.into_body()];
796 let mut positions = [0usize, 0usize];
797 loop {
798 let mut result: ColumnBody<(D, T, R)> = ColumnBody::default();
799 let yielded = result.merge_from(&mut cols, &mut positions);
800 if !result.is_empty() {
801 out.push_back(ColumnChunk::Resident(Rc::new(result), out_depth));
802 }
803 if !yielded {
804 break;
805 }
806 }
807 metrics::record(metrics::Stage::Merge, positions[0] + positions[1], 0);
810 let [col_a, col_b] = &mut cols;
811 for (col, pos, depth, spilled, queue) in [
815 (col_a, positions[0], depths[0], &mut spill_a, in1),
816 (col_b, positions[1], depths[1], &mut spill_b, in2),
817 ] {
818 let len = col.len();
819 if pos == 0 && len > 0 {
820 let chunk = match spilled.take() {
823 Some(body) => ColumnChunk::Spilled(body, depth),
824 None => ColumnChunk::Resident(Rc::new(std::mem::take(col)), depth),
825 };
826 queue.push_front(chunk.survive_merge());
827 } else if pos < len {
828 let view = col.borrow();
829 let mut rest = <(D, T, R) as Columnar>::Container::default();
830 rest.extend_from_self(view, pos..len);
831 queue.push_front(ColumnChunk::Resident(
832 Rc::new(ColumnBody::Typed(rest)),
833 depth,
834 ));
835 }
836 }
837 }
838
839 fn extract(
843 input: &mut VecDeque<Self>,
844 frontier: AntichainRef<T>,
845 residual: &mut timely::progress::Antichain<T>,
846 keep: &mut VecDeque<Self>,
847 ship: &mut VecDeque<Self>,
848 ) {
849 let Some(chunk) = input.pop_front() else {
850 return;
851 };
852 let (time_lower, time_upper) = chunk.chunk_time_bounds();
858 if time_upper.iter().all(|t| !frontier.less_equal(t)) {
859 ship.push_back(chunk);
860 return;
861 }
862 if time_lower.elements().iter().all(|m| frontier.less_equal(m)) {
863 for m in time_lower.elements() {
866 residual.insert_ref(m);
867 }
868 keep.push_back(chunk);
869 return;
870 }
871 let depth = chunk.depth();
874 let mut col = chunk.into_body();
875 let len = col.len();
876 let mut pos = 0;
877 let mut keep_col: ColumnBody<(D, T, R)> = ColumnBody::default();
878 let mut ship_col: ColumnBody<(D, T, R)> = ColumnBody::default();
879 let cut = |col: &mut ColumnBody<(D, T, R)>, queue: &mut VecDeque<Self>, force: bool| {
885 if !col.is_empty() && (force || at_serialized_capacity(&col.borrow())) {
886 queue.push_back(ColumnChunk::Resident(Rc::new(std::mem::take(col)), depth));
887 }
888 };
889 while pos < len {
890 col.extract(&mut pos, frontier, residual, &mut keep_col, &mut ship_col);
891 if pos < len {
892 cut(&mut keep_col, keep, false);
893 cut(&mut ship_col, ship, false);
894 }
895 }
896 cut(&mut keep_col, keep, true);
897 cut(&mut ship_col, ship, true);
898 }
899
900 fn advance(
910 input: &mut VecDeque<Self>,
911 frontier: AntichainRef<T>,
912 done: bool,
913 out: &mut VecDeque<Self>,
914 ) {
915 let Some(front) = input.pop_front() else {
916 return;
917 };
918 let mut depth = front.depth();
921 let mut base = front.into_body();
925 {
926 let base_c = base.typed_mut();
927 for chunk in input.drain(..) {
928 depth = depth.max(chunk.depth());
929 let col = chunk.into_body();
930 let view = col.borrow();
931 base_c.extend_from_self(view, 0..view.len());
932 }
933 }
934 let view = base.borrow();
935 let total = view.len();
936 if total == 0 {
937 return;
938 }
939 let data = view.0;
940
941 if !done && data.get(0) == data.get(total - 1) {
944 input.push_front(ColumnChunk::Resident(Rc::new(base), depth));
945 return;
946 }
947
948 let end = if done {
951 total
952 } else {
953 let last = data.get(total - 1);
954 let mut end = total - 1;
955 while end > 0 && data.get(end - 1) == last {
956 end -= 1;
957 }
958 end
959 };
960 metrics::record(metrics::Stage::Advance, end, 0);
964
965 let mut result = <(D, T, R) as Columnar>::Container::default();
966 let mut scratch: Vec<(T, R)> = Vec::new();
968 let mut index = 0;
969 const CUT_CHECK_RECORDS: usize = 1024;
978 let mut records_since_check = 0usize;
979 while index < end {
986 let group_d = data.get(index);
987 scratch.clear();
988 while index < end && data.get(index) == group_d {
989 let (_, t, r) = view.get(index);
990 let mut owned_t = T::into_owned(t);
991 owned_t.advance_by(frontier);
992 scratch.push((owned_t, R::into_owned(r)));
993 index += 1;
994 }
995 scratch.sort_by(|a, b| a.0.cmp(&b.0));
996 let mut run = scratch.drain(..).peekable();
997 while let Some((t, mut r)) = run.next() {
998 while run.peek().is_some_and(|(t2, _)| *t2 == t) {
999 let (_, r2) = run.next().expect("peeked");
1000 r.plus_equals(&r2);
1001 }
1002 if !r.is_zero() {
1003 result.0.push(group_d);
1004 result.1.push(&t);
1005 result.2.push(&r);
1006 records_since_check += 1;
1007 if records_since_check >= CUT_CHECK_RECORDS {
1008 records_since_check = 0;
1009 if u64::cast_from(indexed::length_in_words(&result.borrow()))
1010 >= u64::cast_from(COMMIT_BYTES / 8)
1011 {
1012 out.push_back(ColumnChunk::Resident(
1013 Rc::new(ColumnBody::Typed(std::mem::take(&mut result))),
1014 depth,
1015 ));
1016 }
1017 }
1018 }
1019 }
1020 }
1021 if !result.is_empty() {
1022 out.push_back(ColumnChunk::Resident(
1023 Rc::new(ColumnBody::Typed(result)),
1024 depth,
1025 ));
1026 }
1027
1028 if end < total {
1030 let mut carry = <(D, T, R) as Columnar>::Container::default();
1031 carry.extend_from_self(view, end..total);
1032 input.push_front(ColumnChunk::Resident(
1033 Rc::new(ColumnBody::Typed(carry)),
1034 depth,
1035 ));
1036 }
1037 }
1038
1039 fn settle(input: &mut VecDeque<Self>, done: bool, out: &mut VecDeque<Self>) {
1045 Self::settle_graded(input, done, out, true)
1046 }
1047}
1048
1049impl<D, T, R> ColumnChunk<D, T, R>
1050where
1051 D: Columnar,
1052 T: Columnar + Timestamp,
1053 R: Columnar,
1054 for<'a> columnar::Ref<'a, D>: Ord,
1055 for<'a> columnar::Ref<'a, T>: Ord,
1056{
1057 fn settle_graded(
1064 input: &mut VecDeque<Self>,
1065 done: bool,
1066 out: &mut VecDeque<Self>,
1067 commit: bool,
1068 ) {
1069 let mut acc: Option<(ColumnBody<(D, T, R)>, u8)> = None;
1073 while let Some(chunk) = input.pop_front() {
1074 let (rc, depth) = match chunk {
1075 spilled @ ColumnChunk::Spilled(_, _) => {
1076 if let Some((body, depth)) = acc.take() {
1077 out.push_back(Self::grade(body, depth, commit));
1078 }
1079 out.push_back(spilled);
1080 continue;
1081 }
1082 ColumnChunk::Resident(rc, depth) => (rc, depth),
1083 };
1084 if rc.length_in_bytes() > COMMIT_BYTES && rc.len() > 1 {
1085 let records = cut_records(rc.len(), rc.length_in_bytes(), COMMIT_BYTES).max(1);
1088 let mut pieces = VecDeque::new();
1089 Self::push_cuts(&rc, records, depth, &mut pieces);
1090 for piece in pieces.into_iter().rev() {
1091 input.push_front(piece);
1092 }
1093 continue;
1094 }
1095 let full = at_commit_size(&rc);
1096 let fits = acc.as_ref().is_some_and(|(body, _)| {
1099 body.length_in_bytes().saturating_add(rc.length_in_bytes()) <= COMMIT_BYTES
1100 });
1101 if !full
1102 && fits
1103 && let Some((mut body, acc_depth)) = acc.take()
1104 {
1105 let view = rc.borrow();
1106 body.typed_mut().extend_from_self(view, 0..view.len());
1107 let acc_depth = acc_depth.max(depth);
1108 if at_commit_size(&body) {
1109 out.push_back(Self::grade(body, acc_depth, commit));
1110 } else {
1111 acc = Some((body, acc_depth));
1112 }
1113 continue;
1114 }
1115 if let Some((mut body, acc_depth)) = acc.take() {
1121 let len = rc.len();
1122 let prefix = if full {
1123 0
1124 } else {
1125 let space = COMMIT_BYTES.saturating_sub(body.length_in_bytes());
1126 cut_records(len, rc.length_in_bytes(), space)
1127 };
1128 if prefix == 0 || prefix >= len {
1129 out.push_back(Self::grade(body, acc_depth, commit));
1130 } else {
1131 let view = rc.borrow();
1132 body.typed_mut().extend_from_self(view, 0..prefix);
1133 let mut rest = <(D, T, R) as Columnar>::Container::default();
1134 rest.extend_from_self(view, prefix..len);
1135 input.push_front(ColumnChunk::Resident(
1136 Rc::new(ColumnBody::Typed(rest)),
1137 depth,
1138 ));
1139 let acc_depth = acc_depth.max(depth);
1140 if body.length_in_bytes() > COMMIT_BYTES {
1141 let mut pieces = VecDeque::new();
1143 Self::push_bounded(body, acc_depth, &mut pieces);
1144 for piece in pieces {
1145 let ColumnChunk::Resident(piece, piece_depth) = piece else {
1146 unreachable!("push_bounded emits resident pieces");
1147 };
1148 let piece =
1149 Rc::try_unwrap(piece).unwrap_or_else(|piece| piece.duplicate());
1150 out.push_back(Self::grade(piece, piece_depth, commit));
1151 }
1152 } else {
1153 out.push_back(Self::grade(body, acc_depth, commit));
1154 }
1155 continue;
1156 }
1157 }
1158 let body = Rc::try_unwrap(rc).unwrap_or_else(|rc| rc.duplicate());
1159 if full {
1160 out.push_back(Self::grade(body, depth, commit));
1161 } else {
1162 acc = Some((body, depth));
1163 }
1164 }
1165 if let Some((body, depth)) = acc {
1166 if done {
1167 out.push_back(Self::grade(body, depth, commit));
1168 } else {
1169 input.push_front(ColumnChunk::Resident(Rc::new(body), depth));
1170 }
1171 }
1172 }
1173
1174 fn grade(body: ColumnBody<(D, T, R)>, depth: u8, commit: bool) -> Self {
1177 if commit {
1178 Self::commit(body, depth)
1179 } else {
1180 ColumnChunk::Resident(Rc::new(body), depth)
1181 }
1182 }
1183}
1184
1185fn extract_view_into<'v, 'p, K, V, T, R>(
1190 view: BorrowedOf<'v, ((K, V), T, R)>,
1191 probes: BorrowedOf<'p, K>,
1192 probe_index: &mut usize,
1193 staging: &mut <((K, V), T, R) as Columnar>::Container,
1194) where
1195 K: Columnar,
1196 V: Columnar,
1197 T: Columnar,
1198 R: Columnar,
1199 for<'b> columnar::Ref<'b, K>: Copy + Ord,
1200{
1201 let keys = view.0.0;
1202 let len = keys.len();
1203 let last = keys.get(len - 1);
1204 let count = probes.len();
1205 let mut pos = 0;
1206 while *probe_index < count {
1207 let probe = probes.get(*probe_index);
1208 mz_ore::soft_assert_no_log!(
1209 *probe_index == 0 || rr::<K>(probes.get(*probe_index - 1)) < rr::<K>(probe),
1210 "probe keys must be sorted and deduplicated"
1211 );
1212 if rr::<K>(probe) > rr::<K>(last) {
1213 return;
1214 }
1215 gallop(len, &mut pos, |i| rr::<K>(keys.get(i)) < rr::<K>(probe));
1216 let start = pos;
1217 while pos < len && rr::<K>(keys.get(pos)) == rr::<K>(probe) {
1218 pos += 1;
1219 }
1220 staging.extend_from_self(view, start..pos);
1221 if rr::<K>(probe) == rr::<K>(last) {
1222 return;
1223 }
1224 *probe_index += 1;
1225 }
1226}
1227
1228impl<K, V, T, R> UnloadChunk for ColumnChunk<(K, V), T, R>
1229where
1230 K: Columnar,
1231 for<'a> columnar::Ref<'a, K>: Copy + Ord,
1232 V: Columnar,
1233 for<'a> columnar::Ref<'a, V>: Copy + Ord,
1234 T: Columnar + Default + Timestamp + Lattice + Ord,
1235 for<'a> columnar::Ref<'a, T>: Copy + Ord,
1236 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
1237{
1238 type Staging = <((K, V), T, R) as Columnar>::Container;
1241
1242 type Probes<'a> = BorrowedOf<'a, K>;
1245
1246 fn probe_count(probes: Self::Probes<'_>) -> usize {
1247 probes.len()
1248 }
1249
1250 fn locate(&self, probes: Self::Probes<'_>, probe_index: usize) -> std::cmp::Ordering {
1251 let probe = probes.get(probe_index);
1252 let (first, last) = self.data_span();
1255 let (first, last) = (first.0, last.0);
1256 if rr::<K>(probe) < rr::<K>(first) {
1257 std::cmp::Ordering::Less
1258 } else if rr::<K>(probe) > rr::<K>(last) {
1259 std::cmp::Ordering::Greater
1260 } else {
1261 std::cmp::Ordering::Equal
1262 }
1263 }
1264
1265 fn extract_into(
1266 &self,
1267 probes: Self::Probes<'_>,
1268 probe_index: &mut usize,
1269 staging: &mut Self::Staging,
1270 ) {
1271 match self {
1272 ColumnChunk::Resident(col, _) => {
1273 extract_view_into::<K, V, T, R>(col.borrow(), probes, probe_index, staging);
1274 }
1275 ColumnChunk::Spilled(body, _) => with_scratch(|scratch| {
1276 body.handle.read_into(scratch);
1282 let view = borrow_words::<((K, V), T, R)>(scratch);
1283 extract_view_into::<K, V, T, R>(view, probes, probe_index, staging);
1284 }),
1285 }
1286 }
1287
1288 fn fetch_into(&self, staging: &mut Self::Staging) {
1289 match self {
1290 ColumnChunk::Resident(col, _) => {
1291 let view = col.borrow();
1292 staging.extend_from_self(view, 0..view.len());
1293 }
1294 ColumnChunk::Spilled(body, _) => with_scratch(|scratch| {
1295 body.handle.read_into(scratch);
1296 let view = borrow_words::<((K, V), T, R)>(scratch);
1297 staging.extend_from_self(view, 0..view.len());
1298 }),
1299 }
1300 }
1301}
1302
1303pub struct UnchunkBuilder<Bu, D: Columnar, T: Columnar, R: Columnar> {
1313 inner: Bu,
1314 _marker: std::marker::PhantomData<(D, T, R)>,
1315}
1316
1317impl<Bu, D, T, R> differential_dataflow::trace::Builder for UnchunkBuilder<Bu, D, T, R>
1318where
1319 Bu: differential_dataflow::trace::Builder<Input = ColumnBody<(D, T, R)>> + ChainState,
1320 D: Columnar + 'static,
1321 T: Columnar + 'static,
1322 R: Columnar + 'static,
1323{
1324 type Input = ColumnChunk<D, T, R>;
1325 type Time = Bu::Time;
1326 type Output = Bu::Output;
1327
1328 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1329 Self {
1330 inner: Bu::with_capacity(keys, vals, upds),
1331 _marker: std::marker::PhantomData,
1332 }
1333 }
1334
1335 fn push(&mut self, chunk: &mut Self::Input) {
1336 let mut body = std::mem::take(chunk).into_body();
1337 self.inner.push(&mut body);
1338 }
1339
1340 fn done(
1341 self,
1342 description: differential_dataflow::trace::Description<Self::Time>,
1343 ) -> Self::Output {
1344 self.inner.done(description)
1345 }
1346
1347 fn seal(
1348 chain: &mut Vec<Self::Input>,
1349 description: differential_dataflow::trace::Description<Self::Time>,
1350 ) -> Self::Output {
1351 let mut state = Bu::State::default();
1352 for chunk in chain.iter() {
1353 if !chunk.is_spilled() || Bu::wants_bodies() {
1358 chunk.with_body(|body| Bu::observe(&mut state, body));
1359 } else {
1360 Bu::observe_records(&mut state, chunk.records());
1361 }
1362 }
1363 let mut builder = Self {
1364 inner: Bu::from_state(state),
1365 _marker: std::marker::PhantomData,
1366 };
1367 for chunk in chain.iter_mut() {
1370 builder.push(chunk);
1371 }
1372 chain.clear();
1373 builder.done(description)
1374 }
1375}
1376
1377pub trait ChainState: differential_dataflow::trace::Builder {
1394 type State: Default;
1396
1397 fn wants_bodies() -> bool;
1401
1402 fn observe(state: &mut Self::State, input: &Self::Input);
1404
1405 fn observe_records(state: &mut Self::State, records: usize);
1408
1409 fn from_state(state: Self::State) -> Self;
1411}
1412
1413pub struct ChunkChunker<D: Columnar, T: Columnar, R: Columnar> {
1416 inner: ColumnChunker<(D, T, R)>,
1417 ready: VecDeque<ColumnChunk<D, T, R>>,
1418 staged: ColumnChunk<D, T, R>,
1419}
1420
1421impl<D, T, R> Default for ChunkChunker<D, T, R>
1422where
1423 D: Columnar,
1424 T: Columnar,
1425 R: Columnar,
1426 ColumnChunker<(D, T, R)>: Default,
1427{
1428 fn default() -> Self {
1429 Self {
1430 inner: Default::default(),
1431 ready: VecDeque::new(),
1432 staged: Default::default(),
1433 }
1434 }
1435}
1436
1437impl<'a, D, T, R> PushInto<&'a mut Column<(D, T, R)>> for ChunkChunker<D, T, R>
1438where
1439 D: Columnar,
1440 T: Columnar,
1441 R: Columnar,
1442 ColumnChunker<(D, T, R)>: PushInto<&'a mut Column<(D, T, R)>>,
1443{
1444 fn push_into(&mut self, item: &'a mut Column<(D, T, R)>) {
1445 self.inner.push_into(item);
1446 }
1447}
1448
1449impl<D, T, R> ContainerBuilder for ChunkChunker<D, T, R>
1450where
1451 D: Columnar + 'static,
1452 T: Columnar + 'static,
1453 R: Columnar + 'static,
1454 ColumnChunker<(D, T, R)>: ContainerBuilder<Container = ColumnBody<(D, T, R)>>,
1455{
1456 type Container = ColumnChunk<D, T, R>;
1457
1458 fn extract(&mut self) -> Option<&mut Self::Container> {
1459 if self.ready.is_empty() {
1460 let body = self.inner.extract()?;
1461 ColumnChunk::push_bounded(std::mem::take(body), 0, &mut self.ready);
1462 }
1463 self.staged = self.ready.pop_front()?;
1464 Some(&mut self.staged)
1465 }
1466
1467 fn finish(&mut self) -> Option<&mut Self::Container> {
1468 if self.ready.is_empty() {
1469 let body = self.inner.finish()?;
1470 ColumnChunk::push_bounded(std::mem::take(body), 0, &mut self.ready);
1471 }
1472 self.staged = self.ready.pop_front()?;
1473 Some(&mut self.staged)
1474 }
1475}
1476
1477pub type AccountedChunkBatcher<D, T, R> =
1482 differential_dataflow::trace::implementations::merge_batcher::MergeBatcher<
1483 AccountedChunkMerger<D, T, R>,
1484 >;
1485
1486pub struct AccountedChunkMerger<D: Columnar, T: Columnar, R: Columnar> {
1499 inner: differential_dataflow::trace::chunk::ChunkMerger<ColumnChunk<D, T, R>>,
1500}
1501
1502impl<D: Columnar, T: Columnar, R: Columnar> Default for AccountedChunkMerger<D, T, R> {
1503 fn default() -> Self {
1504 Self {
1505 inner: Default::default(),
1506 }
1507 }
1508}
1509
1510impl<D, T, R> differential_dataflow::trace::implementations::merge_batcher::Merger
1511 for AccountedChunkMerger<D, T, R>
1512where
1513 D: Columnar + 'static,
1514 T: Columnar + Timestamp + 'static,
1515 R: Columnar + 'static,
1516 for<'a> columnar::Ref<'a, D>: Ord,
1517 for<'a> columnar::Ref<'a, T>: Ord,
1518 ColumnChunk<D, T, R>: Chunk,
1519 <ColumnChunk<D, T, R> as Chunk>::Time: Clone + PartialOrder + 'static,
1520{
1521 type Chunk = ColumnChunk<D, T, R>;
1522 type Time = <ColumnChunk<D, T, R> as Chunk>::Time;
1523
1524 fn merge(
1525 &mut self,
1526 list1: Vec<Self::Chunk>,
1527 list2: Vec<Self::Chunk>,
1528 output: &mut Vec<Self::Chunk>,
1529 stash: &mut Vec<Self::Chunk>,
1530 ) {
1531 self.inner.merge(list1, list2, output, stash)
1532 }
1533
1534 fn extract(
1535 &mut self,
1536 merged: Vec<Self::Chunk>,
1537 upper: AntichainRef<Self::Time>,
1538 frontier: &mut Antichain<Self::Time>,
1539 readied: &mut Vec<Self::Chunk>,
1540 kept: &mut Vec<Self::Chunk>,
1541 _stash: &mut Vec<Self::Chunk>,
1542 ) {
1543 let mut input: VecDeque<Self::Chunk> = merged.into();
1551 let (mut keep, mut shipped) = (VecDeque::new(), VecDeque::new());
1552 let (mut kept_q, mut shipped_q) = (VecDeque::new(), VecDeque::new());
1553 while !input.is_empty() {
1554 Chunk::extract(&mut input, upper, frontier, &mut keep, &mut shipped);
1555 Chunk::settle(&mut keep, false, &mut kept_q);
1556 ColumnChunk::settle_graded(&mut shipped, false, &mut shipped_q, false);
1557 }
1558 Chunk::settle(&mut keep, true, &mut kept_q);
1559 ColumnChunk::settle_graded(&mut shipped, true, &mut shipped_q, false);
1560 Extend::extend(kept, kept_q);
1561 Extend::extend(readied, shipped_q);
1562 }
1563
1564 fn len(chunk: &Self::Chunk) -> usize {
1565 Chunk::len(chunk)
1566 }
1567
1568 fn allocation(chunk: &Self::Chunk) -> (usize, usize, usize) {
1569 let bytes = match chunk {
1570 ColumnChunk::Resident(col, _) => col.length_in_bytes(),
1571 ColumnChunk::Spilled(body, _) => body.len_bytes,
1572 };
1573 (bytes, bytes, 1)
1574 }
1575}
1576
1577#[cfg(test)]
1578mod tests {
1579 use differential_dataflow::trace::chunk::{ChunkBatch, ChunkBatcher};
1589 use differential_dataflow::trace::{Batcher, Description};
1590 use mz_ore::pool::Pool;
1591 use proptest::prelude::*;
1592 use timely::container::PushInto;
1593 use timely::progress::Antichain;
1594
1595 use crate::columnar::unload::UnloadBatch;
1596
1597 use super::*;
1598
1599 type Tuple = ((u64, u64), u64, i64);
1600 type TestChunk = ColumnChunk<(u64, u64), u64, i64>;
1601
1602 #[mz_ore::test]
1607 fn lz4_codec_matches_the_previous_extent_framing() {
1608 let body: Vec<u8> = (0..100_000u32).flat_map(|i| i.to_le_bytes()).collect();
1609 let mut stored = Vec::new();
1610 LZ4_CODEC.encode(&body, &mut stored);
1611 assert_eq!(stored, lz4_flex::block::compress_prepend_size(&body));
1612 let mut round = vec![0u8; body.len()];
1613 LZ4_CODEC.decode(&stored, &mut round);
1614 assert_eq!(round, body);
1615 }
1616
1617 #[mz_ore::test]
1618 #[should_panic(expected = "destination must match")]
1619 fn lz4_codec_decode_length_mismatch_panics() {
1620 let mut stored = Vec::new();
1621 LZ4_CODEC.encode(&[7u8; 64], &mut stored);
1622 let mut short = vec![0u8; 32];
1623 LZ4_CODEC.decode(&stored, &mut short);
1624 }
1625
1626 fn consolidate(mut v: Vec<Tuple>) -> Vec<Tuple> {
1629 v.sort();
1630 let mut out: Vec<Tuple> = Vec::new();
1631 for (d, t, r) in v {
1632 if let Some(last) = out.last_mut() {
1633 if last.0 == d && last.1 == t {
1634 last.2 += r;
1635 continue;
1636 }
1637 }
1638 out.push((d, t, r));
1639 }
1640 out.retain(|x| x.2 != 0);
1641 out
1642 }
1643
1644 fn arb_consolidated() -> impl Strategy<Value = Vec<Tuple>> {
1645 prop::collection::vec(((0u64..5, 0u64..5), 0u64..4, -3i64..=3i64), 0..40)
1646 .prop_map(consolidate)
1647 }
1648
1649 fn build_column(v: &[Tuple]) -> ColumnBody<Tuple> {
1650 let mut col: ColumnBody<Tuple> = Default::default();
1651 for tup in v {
1652 col.push_into(*tup);
1653 }
1654 col
1655 }
1656
1657 fn collect_column(col: &ColumnBody<Tuple>) -> Vec<Tuple> {
1658 col.borrow()
1659 .into_index_iter()
1660 .map(|((k, v), t, r)| {
1661 (
1662 (u64::into_owned(k), u64::into_owned(v)),
1663 u64::into_owned(t),
1664 i64::into_owned(r),
1665 )
1666 })
1667 .collect()
1668 }
1669
1670 fn collect_chunks(chunks: impl IntoIterator<Item = TestChunk>) -> Vec<Tuple> {
1671 chunks
1672 .into_iter()
1673 .flat_map(|chunk| collect_column(&chunk.into_body()))
1674 .collect()
1675 }
1676
1677 fn collect_staging(staging: &<Tuple as Columnar>::Container) -> Vec<Tuple> {
1678 staging
1679 .borrow()
1680 .into_index_iter()
1681 .map(|((k, v), t, r)| {
1682 (
1683 (u64::into_owned(k), u64::into_owned(v)),
1684 u64::into_owned(t),
1685 i64::into_owned(r),
1686 )
1687 })
1688 .collect()
1689 }
1690
1691 fn chunked(data: &[Tuple], cuts: &[usize]) -> VecDeque<TestChunk> {
1693 let mut chunks = VecDeque::new();
1694 let mut start = 0;
1695 for cut in cuts {
1696 let end = (start + 1 + cut % 7).min(data.len());
1697 if end > start {
1698 chunks.push_back(ColumnChunk::from_body(build_column(&data[start..end])));
1699 start = end;
1700 }
1701 }
1702 if start < data.len() {
1703 chunks.push_back(ColumnChunk::from_body(build_column(&data[start..])));
1704 }
1705 chunks
1706 }
1707
1708 fn chunked_spilled(data: &[Tuple], cuts: &[usize], pool: &Pool) -> VecDeque<TestChunk> {
1711 chunked(data, cuts)
1712 .into_iter()
1713 .map(|chunk| force_spill(chunk, pool))
1714 .collect()
1715 }
1716
1717 fn body_compressed(chunk: &TestChunk) -> bool {
1721 match chunk {
1722 ColumnChunk::Spilled(body, _) => body.compressed,
1723 ColumnChunk::Resident(_, _) => panic!("chunk must be spilled"),
1724 }
1725 }
1726
1727 fn force_spill(chunk: TestChunk, pool: &Pool) -> TestChunk {
1730 let depth = chunk.depth();
1731 TestChunk::spill_body(chunk.into_body(), pool, depth)
1732 }
1733
1734 #[mz_ore::test]
1737 #[cfg_attr(miri, ignore)]
1738 fn allocation_reports_a_body_wherever_it_lives() {
1739 use differential_dataflow::trace::implementations::merge_batcher::Merger;
1740
1741 let data: Vec<Tuple> = (0..64u64).map(|i| ((i, i), 0, 1)).collect();
1742 let resident = TestChunk::from_body(build_column(&data));
1743 let (size, capacity, allocations) =
1744 AccountedChunkMerger::<(u64, u64), u64, i64>::allocation(&resident);
1745 assert!(size > 0, "a resident body reports its bytes");
1746 assert_eq!((capacity, allocations), (size, 1));
1747
1748 let spilled = force_spill(resident, &test_pool());
1749 assert_eq!(
1750 AccountedChunkMerger::<(u64, u64), u64, i64>::allocation(&spilled),
1751 (size, capacity, allocations),
1752 "spilling a body does not change what its owner holds"
1753 );
1754 }
1755
1756 #[mz_ore::test]
1765 #[cfg_attr(miri, ignore)]
1766 fn seal_readies_resident_chunks() {
1767 use differential_dataflow::trace::Batcher;
1768
1769 let data: Vec<Tuple> = (0..20_000u64).map(|i| ((i, i), 0, 1)).collect();
1773
1774 set_spill_override(Some(test_pool()));
1775 let mut batcher = AccountedChunkBatcher::<(u64, u64), u64, i64>::new(None, 0);
1776 batcher.push_into(TestChunk::from_body(build_column(&data)));
1777 let (chain, _description) = batcher.seal(Antichain::new());
1778 set_spill_override(None);
1779
1780 assert!(!chain.is_empty(), "the seal ships what was pushed");
1781 assert!(
1782 chain.iter().all(|chunk| !chunk.is_spilled()),
1783 "a chunk that arrived resident is readied resident"
1784 );
1785 }
1786
1787 fn test_pool() -> Pool {
1791 static POOL: std::sync::OnceLock<Pool> = std::sync::OnceLock::new();
1792 POOL.get_or_init(|| Pool::new().expect("pool creation"))
1793 .clone()
1794 }
1795
1796 proptest! {
1797 #[mz_ore::test]
1800 #[cfg_attr(miri, ignore)]
1801 fn batcher_round_trip(
1802 inputs in prop::collection::vec(arb_consolidated(), 1..6),
1803 cuts in prop::collection::vec(0usize..7, 0..8),
1804 ) {
1805 let mut batcher: ChunkBatcher<TestChunk> = Batcher::new(None, 0);
1806 let mut union = Vec::new();
1807 for input in &inputs {
1808 Extend::extend(&mut union, input.iter().copied());
1809 for chunk in chunked(input, &cuts) {
1810 batcher.push_into(chunk);
1811 }
1812 }
1813 let (sealed, _description) = batcher.seal(Antichain::new());
1815 prop_assert_eq!(collect_chunks(sealed), consolidate(union));
1816 }
1817
1818 #[mz_ore::test]
1821 #[cfg_attr(miri, ignore)]
1822 fn batcher_round_trip_spilled(
1823 inputs in prop::collection::vec(arb_consolidated(), 1..4),
1824 cuts in prop::collection::vec(0usize..7, 0..6),
1825 ) {
1826 let pool = test_pool();
1827 let mut batcher: ChunkBatcher<TestChunk> = Batcher::new(None, 0);
1828 let mut union = Vec::new();
1829 for input in &inputs {
1830 Extend::extend(&mut union, input.iter().copied());
1831 for chunk in chunked_spilled(input, &cuts, &pool) {
1832 batcher.push_into(chunk);
1833 }
1834 }
1835 let (sealed, _description) = batcher.seal(Antichain::new());
1836 prop_assert_eq!(collect_chunks(sealed), consolidate(union));
1837 }
1838
1839 #[mz_ore::test]
1842 #[cfg_attr(miri, ignore)]
1843 fn seal_partitions_by_time(
1844 input in arb_consolidated(),
1845 cuts in prop::collection::vec(0usize..7, 0..8),
1846 upper in 0u64..5,
1847 ) {
1848 let mut batcher: ChunkBatcher<TestChunk> = Batcher::new(None, 0);
1849 for chunk in chunked(&input, &cuts) {
1850 batcher.push_into(chunk);
1851 }
1852 let (shipped, _) = batcher.seal(Antichain::from_elem(upper));
1853 let expected_shipped: Vec<Tuple> =
1854 input.iter().copied().filter(|(_, t, _)| *t < upper).collect();
1855 prop_assert_eq!(collect_chunks(shipped), consolidate(expected_shipped));
1856
1857 let kept_min = input.iter().filter(|(_, t, _)| *t >= upper).map(|(_, t, _)| *t).min();
1858 let frontier = batcher.frontier().to_owned();
1859 prop_assert_eq!(frontier.elements().first().copied(), kept_min);
1860
1861 let (rest, _) = batcher.seal(Antichain::new());
1862 let expected_rest: Vec<Tuple> =
1863 input.iter().copied().filter(|(_, t, _)| *t >= upper).collect();
1864 prop_assert_eq!(collect_chunks(rest), consolidate(expected_rest));
1865 }
1866
1867 #[mz_ore::test]
1871 #[cfg_attr(miri, ignore)]
1872 fn seal_partitions_by_time_spilled(
1873 input in arb_consolidated(),
1874 cuts in prop::collection::vec(0usize..7, 0..8),
1875 upper in 0u64..5,
1876 ) {
1877 let pool = test_pool();
1878 let mut batcher: ChunkBatcher<TestChunk> = Batcher::new(None, 0);
1879 for chunk in chunked_spilled(&input, &cuts, &pool) {
1880 batcher.push_into(chunk);
1881 }
1882 let (shipped, _) = batcher.seal(Antichain::from_elem(upper));
1883 let expected_shipped: Vec<Tuple> =
1884 input.iter().copied().filter(|(_, t, _)| *t < upper).collect();
1885 prop_assert_eq!(collect_chunks(shipped), consolidate(expected_shipped));
1886
1887 let kept_min = input.iter().filter(|(_, t, _)| *t >= upper).map(|(_, t, _)| *t).min();
1888 let frontier = batcher.frontier().to_owned();
1889 prop_assert_eq!(frontier.elements().first().copied(), kept_min);
1890
1891 let (rest, _) = batcher.seal(Antichain::new());
1892 let expected_rest: Vec<Tuple> =
1893 input.iter().copied().filter(|(_, t, _)| *t >= upper).collect();
1894 prop_assert_eq!(collect_chunks(rest), consolidate(expected_rest));
1895 }
1896
1897 #[mz_ore::test]
1900 #[cfg_attr(miri, ignore)]
1901 fn advance_matches_reference(
1902 input in arb_consolidated(),
1903 cuts in prop::collection::vec(0usize..7, 0..8),
1904 frontier_elem in 0u64..5,
1905 ) {
1906 let frontier = Antichain::from_elem(frontier_elem);
1907 let mut chunks = chunked(&input, &cuts);
1908 let mut out = VecDeque::new();
1909 TestChunk::advance(&mut chunks, frontier.borrow(), false, &mut out);
1910 TestChunk::advance(&mut chunks, frontier.borrow(), true, &mut out);
1911 prop_assert!(chunks.is_empty());
1912
1913 let expected = consolidate(
1914 input
1915 .iter()
1916 .map(|&(d, mut t, r)| {
1917 t.advance_by(frontier.borrow());
1918 (d, t, r)
1919 })
1920 .collect(),
1921 );
1922 prop_assert_eq!(collect_chunks(out), expected);
1923 }
1924
1925 #[mz_ore::test]
1928 #[cfg_attr(miri, ignore)]
1929 fn settle_preserves_and_packs(
1930 input in arb_consolidated(),
1931 cuts in prop::collection::vec(0usize..7, 1..8),
1932 ) {
1933 let mut chunks = chunked(&input, &cuts);
1934 let mut out = VecDeque::new();
1935 TestChunk::settle(&mut chunks, true, &mut out);
1936 prop_assert!(chunks.is_empty());
1937 prop_assert!(out.len() <= 1);
1940 prop_assert_eq!(collect_chunks(out), input);
1941 }
1942
1943 #[mz_ore::test]
1947 #[cfg_attr(miri, ignore)]
1948 fn unload_extract_matches_filter(
1949 input in arb_consolidated(),
1950 cuts in prop::collection::vec(0usize..7, 0..8),
1951 probe_keys in prop::collection::btree_set(0u64..6, 0..6),
1952 spill in any::<bool>(),
1953 ) {
1954 prop_assume!(!input.is_empty());
1955 let pool = test_pool();
1956 let chunks: Vec<TestChunk> = if spill {
1957 chunked_spilled(&input, &cuts, &pool).into()
1958 } else {
1959 chunked(&input, &cuts).into()
1960 };
1961 let description = Description::new(
1962 Antichain::from_elem(0u64),
1963 Antichain::new(),
1964 Antichain::from_elem(0u64),
1965 );
1966 let batch = ChunkBatch::new(chunks, description);
1967
1968 let mut probe_col = <u64 as Columnar>::Container::default();
1969 for key in &probe_keys {
1970 probe_col.push(*key);
1971 }
1972 let mut staging = <Tuple as Columnar>::Container::default();
1973 batch.extract_into(probe_col.borrow(), &mut staging);
1974
1975 let expected: Vec<Tuple> = input
1976 .iter()
1977 .copied()
1978 .filter(|((k, _), _, _)| probe_keys.contains(k))
1979 .collect();
1980 prop_assert_eq!(collect_staging(&staging), expected);
1981
1982 let mut staging = <Tuple as Columnar>::Container::default();
1985 batch.fetch_into(&mut staging);
1986 prop_assert_eq!(collect_staging(&staging), input);
1987 }
1988 }
1989
1990 #[mz_ore::test]
1993 fn locate_spans_keys() {
1994 let chunk = ColumnChunk::from_body(build_column(&[
1995 ((2, 0), 0, 1),
1996 ((4, 0), 0, 1),
1997 ((6, 0), 0, 1),
1998 ]));
1999 let mut probe_col = <u64 as Columnar>::Container::default();
2000 for key in [0u64, 2, 3, 6, 9] {
2001 probe_col.push(key);
2002 }
2003 let probes = probe_col.borrow();
2004 use std::cmp::Ordering::*;
2005 let expected = [Less, Equal, Equal, Equal, Greater];
2006 for (index, expected) in expected.iter().enumerate() {
2007 assert_eq!(chunk.locate(probes, index), *expected, "probe {index}");
2008 }
2009 }
2010
2011 fn collect_bounded(chunks: impl IntoIterator<Item = TestChunk>, bound: usize) -> Vec<Tuple> {
2014 let mut collected = Vec::new();
2015 for chunk in chunks {
2016 let col = chunk.into_body();
2017 let bytes = col.length_in_bytes();
2018 assert!(bytes <= bound, "chunk of {bytes} bytes exceeds {bound}");
2019 Extend::extend(&mut collected, collect_column(&col));
2020 }
2021 collected
2022 }
2023
2024 type WideUpdate = ((u64, String), u64, i64);
2025 type WideChunk = ColumnChunk<(u64, String), u64, i64>;
2026
2027 fn wide_column(keys: impl Iterator<Item = u64>, bytes: usize) -> ColumnBody<WideUpdate> {
2028 let mut column = ColumnBody::default();
2029 for key in keys {
2030 column.push_into(&((key, "x".repeat(bytes)), 0, 1));
2031 }
2032 column
2033 }
2034
2035 fn assert_wide_byte_bound(chunks: VecDeque<WideChunk>, expected: usize, payload_bytes: usize) {
2036 let mut keys = Vec::new();
2037 for chunk in chunks {
2038 let column = chunk.into_body();
2039 assert!(
2040 column.length_in_bytes() <= COMMIT_BYTES || column.len() == 1,
2041 "{} bytes in a {}-record chunk",
2042 column.length_in_bytes(),
2043 column.len(),
2044 );
2045 let view = column.borrow();
2046 for index in 0..view.len() {
2047 let ((key, payload), time, diff) = view.get(index);
2048 assert_eq!(payload.len(), payload_bytes);
2049 assert!(payload.iter().all(|byte| *byte == b'x'));
2050 assert_eq!((*time, *diff), (0, 1));
2051 keys.push(*key);
2052 }
2053 }
2054 assert_eq!(keys, (0..u64::cast_from(expected)).collect::<Vec<_>>());
2055 }
2056
2057 #[mz_ore::test]
2058 fn chunker_enforces_byte_bound() {
2059 let mut chunker = ChunkChunker::default();
2060 let mut input: Column<WideUpdate> = wide_column((0..4000).rev(), 3000).into();
2061 chunker.push_into(&mut input);
2062 let mut chunks = VecDeque::new();
2063 if let Some(chunk) = chunker.extract() {
2064 chunks.push_back(std::mem::take(chunk));
2065 }
2066 while let Some(chunk) = chunker.finish() {
2067 chunks.push_back(std::mem::take(chunk));
2068 }
2069 assert_wide_byte_bound(chunks, 4000, 3000);
2070 }
2071
2072 #[mz_ore::test]
2073 fn merge_settle_enforces_byte_bound() {
2074 let mut left = VecDeque::from([WideChunk::from_body(wide_column(
2075 (0..2000).step_by(2),
2076 2100,
2077 ))]);
2078 let mut right = VecDeque::from([WideChunk::from_body(wide_column(
2079 (1..2000).step_by(2),
2080 2100,
2081 ))]);
2082 let mut merged = VecDeque::new();
2083 while !left.is_empty() && !right.is_empty() {
2084 WideChunk::merge(&mut left, &mut right, &mut merged);
2085 }
2086 merged.append(&mut left);
2087 merged.append(&mut right);
2088 let mut settled = VecDeque::new();
2089 WideChunk::settle(&mut merged, true, &mut settled);
2090 assert_wide_byte_bound(settled, 2000, 2100);
2091 }
2092
2093 #[mz_ore::test]
2094 fn settle_enforces_byte_bound_after_coalescing() {
2095 let mut input = VecDeque::from([
2096 WideChunk::from_body(wide_column(0..400, 3000)),
2097 WideChunk::from_body(wide_column(400..800, 3000)),
2098 ]);
2099 let mut settled = VecDeque::new();
2100 WideChunk::settle(&mut input, true, &mut settled);
2101 assert_wide_byte_bound(settled, 800, 3000);
2102 }
2103
2104 #[mz_ore::test]
2105 fn settle_byte_bound_allows_indivisible_update() {
2106 let mut input = VecDeque::from([WideChunk::from_body(wide_column(0..1, 2 * COMMIT_BYTES))]);
2107 let mut settled = VecDeque::new();
2108 WideChunk::settle(&mut input, true, &mut settled);
2109 assert_wide_byte_bound(settled, 1, 2 * COMMIT_BYTES);
2110 }
2111
2112 #[mz_ore::test]
2116 fn push_bounded_splits_uneven_widths() {
2117 let mut column: ColumnBody<WideUpdate> = ColumnBody::default();
2118 column.push_into(&((0, "x".repeat(COMMIT_BYTES - 4096)), 0, 1));
2119 for key in 1..20_000u64 {
2120 column.push_into(&((key, "x".to_string()), 0, 1));
2121 }
2122 assert!(column.length_in_bytes() > COMMIT_BYTES);
2123 let mut out = VecDeque::new();
2124 WideChunk::push_bounded(column, 3, &mut out);
2125 let mut keys = Vec::new();
2126 for chunk in out {
2127 let ColumnChunk::Resident(_, depth) = &chunk else {
2128 panic!("pieces are resident");
2129 };
2130 assert_eq!(*depth, 3, "pieces keep the input's depth");
2131 let column = chunk.into_body();
2132 assert!(
2133 column.length_in_bytes() <= COMMIT_BYTES || column.len() == 1,
2134 "{} bytes in a {}-record piece",
2135 column.length_in_bytes(),
2136 column.len(),
2137 );
2138 let view = column.borrow();
2139 for index in 0..view.len() {
2140 let ((key, _), _, _) = view.get(index);
2141 keys.push(*key);
2142 }
2143 }
2144 assert_eq!(keys, (0..20_000u64).collect::<Vec<_>>());
2145 }
2146
2147 #[mz_ore::test]
2151 fn merge_output_within_byte_bound_before_settle() {
2152 let interleaved = (
2153 wide_column((0..1600).step_by(2), 2100),
2154 wide_column((1..1600).step_by(2), 2100),
2155 );
2156 let runs = (
2157 wide_column((0..500).chain(1000..1100), 2100),
2158 wide_column(500..1000, 2100),
2159 );
2160 for (left, right) in [interleaved, runs] {
2161 let expected = left.len() + right.len();
2162 let mut left = VecDeque::from([WideChunk::from_body(left)]);
2163 let mut right = VecDeque::from([WideChunk::from_body(right)]);
2164 let mut merged = VecDeque::new();
2165 while !left.is_empty() && !right.is_empty() {
2166 WideChunk::merge(&mut left, &mut right, &mut merged);
2167 }
2168 assert!(merged.len() > 1, "merge never cut its output");
2169 merged.append(&mut left);
2170 merged.append(&mut right);
2171 assert_wide_byte_bound(merged, expected, 2100);
2172 }
2173 }
2174
2175 #[mz_ore::test]
2176 fn settle_split_preserves_depth() {
2177 let column = Rc::new(wide_column(0..2000, 2100));
2178 let shared = Rc::clone(&column);
2179 let mut input = VecDeque::from([WideChunk::Resident(column, 3)]);
2180 let mut settled = VecDeque::new();
2181 WideChunk::settle(&mut input, true, &mut settled);
2182 assert!(settled.len() > 1, "oversized chunk was not split");
2183 for chunk in &settled {
2184 assert_eq!(chunk.depth(), 3, "split piece changed generation");
2185 }
2186 assert_eq!(shared.len(), 2000, "shared body was mutated");
2187 assert_wide_byte_bound(settled, 2000, 2100);
2188 }
2189
2190 #[mz_ore::test]
2193 fn settle_split_then_coalesce_keeps_depth() {
2194 let mut input = VecDeque::from([
2195 WideChunk::Resident(Rc::new(wide_column(0..800, 3000)), 2),
2196 WideChunk::Resident(Rc::new(wide_column(800..810, 3000)), 0),
2197 ]);
2198 let mut settled = VecDeque::new();
2199 WideChunk::settle(&mut input, true, &mut settled);
2200 assert!(settled.len() > 1, "oversized chunk was not split");
2201 for chunk in &settled {
2202 assert_eq!(chunk.depth(), 2, "coalescing lost the deeper generation");
2203 }
2204 assert_wide_byte_bound(settled, 810, 3000);
2205 }
2206
2207 #[mz_ore::test]
2211 fn settle_fills_accumulated_column_from_next_chunk() {
2212 let per_chunk = 370;
2213 let mut input: VecDeque<WideChunk> = (0..8u64)
2214 .map(|i| {
2215 let start = i * per_chunk;
2216 WideChunk::from_body(wide_column(start..start + per_chunk, 3000))
2217 })
2218 .collect();
2219 let chunk_bytes = input[0].clone().into_body().length_in_bytes();
2220 assert!(2 * chunk_bytes > COMMIT_BYTES && chunk_bytes < COMMIT_BYTES);
2221 let mut settled = VecDeque::new();
2222 WideChunk::settle(&mut input, true, &mut settled);
2223 let sizes: Vec<usize> = settled
2224 .iter()
2225 .map(|chunk| chunk.clone().into_body().length_in_bytes())
2226 .collect();
2227 let (last, rest) = sizes.split_last().expect("settled output");
2228 assert!(*last <= COMMIT_BYTES);
2229 for size in rest {
2230 assert!(
2231 *size >= COMMIT_BYTES - COMMIT_BYTES / 10 && *size <= COMMIT_BYTES,
2232 "committed {size} bytes, sizes {sizes:?}",
2233 );
2234 }
2235 assert_wide_byte_bound(settled, 8 * 370, 3000);
2236 }
2237
2238 #[mz_ore::test]
2241 #[cfg_attr(miri, ignore)]
2242 fn advance_cuts_large_output() {
2243 let records: Vec<Tuple> = (0..300_000u64).map(|k| ((k, 0), 0, 1)).collect();
2244 let mut input = VecDeque::from([ColumnChunk::from_body(build_column(&records))]);
2245 let frontier = Antichain::from_elem(0u64);
2246 let mut out = VecDeque::new();
2247 TestChunk::advance(&mut input, frontier.borrow(), true, &mut out);
2248 assert!(input.is_empty());
2249 assert!(
2250 out.len() >= 2,
2251 "expected a cut output, got {} chunk(s)",
2252 out.len()
2253 );
2254 assert_eq!(collect_bounded(out, 2 * COMMIT_BYTES), records);
2255 }
2256
2257 #[mz_ore::test]
2260 #[cfg_attr(miri, ignore)]
2261 fn advance_withholds_giant_group() {
2262 let records: Vec<Tuple> = (0..100u64).map(|t| ((7, 7), t, 1)).collect();
2263 let mut input: VecDeque<TestChunk> = VecDeque::new();
2264 for piece in records.chunks(30) {
2265 input.push_back(ColumnChunk::from_body(build_column(piece)));
2266 }
2267 let frontier = Antichain::from_elem(50u64);
2268 let mut out = VecDeque::new();
2269 TestChunk::advance(&mut input, frontier.borrow(), false, &mut out);
2270 assert!(out.is_empty(), "nothing may ship from a single open group");
2271 assert_eq!(input.len(), 1, "the whole input becomes one carry chunk");
2272 TestChunk::advance(&mut input, frontier.borrow(), true, &mut out);
2274 assert!(input.is_empty());
2275 let advanced = records.iter().map(|&(d, t, r)| (d, t.max(50), r)).collect();
2276 assert_eq!(collect_chunks(out), consolidate(advanced));
2277 }
2278
2279 #[mz_ore::test]
2284 #[cfg_attr(miri, ignore)]
2285 fn extract_passes_frontier_disjoint_chunks_through() {
2286 set_spill_override(Some(test_pool()));
2287 let low: Vec<Tuple> = (0..20_000u64).map(|i| ((i, 0), i % 4, 1)).collect();
2288 let high: Vec<Tuple> = (0..20_000u64).map(|i| ((i, 0), 6 + i % 4, 1)).collect();
2289 let spilled_chunk = |data: &[Tuple]| {
2290 let chunk = TestChunk::commit(build_column(&consolidate(data.to_vec())), 1);
2291 assert!(chunk.is_spilled());
2292 chunk
2293 };
2294
2295 let mut input = VecDeque::from([spilled_chunk(&low), spilled_chunk(&high)]);
2300 let frontier = Antichain::from_elem(5u64);
2301 let mut residual = Antichain::new();
2302 let (mut keep, mut ship) = (VecDeque::new(), VecDeque::new());
2303 while !input.is_empty() {
2304 TestChunk::extract(
2305 &mut input,
2306 frontier.borrow(),
2307 &mut residual,
2308 &mut keep,
2309 &mut ship,
2310 );
2311 }
2312 assert_eq!(ship.len(), 1);
2313 assert!(ship[0].is_spilled(), "shipped whole: body untouched");
2314 assert_eq!(keep.len(), 1);
2315 assert!(keep[0].is_spilled(), "kept whole: body untouched");
2316 assert_eq!(residual, Antichain::from_elem(6));
2317 let shipped = ship.pop_front().unwrap().into_body();
2318 assert_eq!(collect_column(&shipped), consolidate(low));
2319 let kept = keep.pop_front().unwrap().into_body();
2320 assert_eq!(collect_column(&kept), consolidate(high));
2321 set_spill_override(None);
2322 }
2323
2324 #[mz_ore::test]
2327 #[cfg_attr(miri, ignore)]
2328 fn extract_cuts_large_output() {
2329 let records: Vec<Tuple> = (0..300_000u64).map(|k| ((k, 0), k % 2, 1)).collect();
2330 let mut input = VecDeque::from([ColumnChunk::from_body(build_column(&records))]);
2331 let frontier = Antichain::from_elem(1u64);
2332 let mut residual = Antichain::new();
2333 let (mut keep, mut ship) = (VecDeque::new(), VecDeque::new());
2334 while !input.is_empty() {
2335 TestChunk::extract(
2336 &mut input,
2337 frontier.borrow(),
2338 &mut residual,
2339 &mut keep,
2340 &mut ship,
2341 );
2342 }
2343 assert!(
2344 keep.len() >= 2,
2345 "expected a cut keep side, got {} chunk(s)",
2346 keep.len()
2347 );
2348 assert!(
2349 ship.len() >= 2,
2350 "expected a cut ship side, got {} chunk(s)",
2351 ship.len()
2352 );
2353 let kept: Vec<Tuple> = records.iter().copied().filter(|r| r.1 >= 1).collect();
2354 let shipped: Vec<Tuple> = records.iter().copied().filter(|r| r.1 < 1).collect();
2355 assert_eq!(collect_bounded(keep, 2 * COMMIT_BYTES), kept);
2356 assert_eq!(collect_bounded(ship, 2 * COMMIT_BYTES), shipped);
2357 assert_eq!(residual, Antichain::from_elem(1));
2358 }
2359
2360 #[mz_ore::test]
2363 fn locate_uses_resident_bounds() {
2364 let pool = test_pool();
2365 let data: Vec<Tuple> = vec![((2, 0), 0, 1), ((4, 0), 0, 1)];
2366 let chunk = force_spill(ColumnChunk::from_body(build_column(&data)), &pool);
2367
2368 let mut probe_col = <u64 as Columnar>::Container::default();
2369 for key in [1u64, 3, 5] {
2370 probe_col.push(key);
2371 }
2372 let probes = probe_col.borrow();
2373 assert_eq!(chunk.locate(probes, 0), std::cmp::Ordering::Less);
2374 assert_eq!(chunk.locate(probes, 1), std::cmp::Ordering::Equal);
2375 assert_eq!(chunk.locate(probes, 2), std::cmp::Ordering::Greater);
2376 }
2377
2378 #[mz_ore::test]
2381 #[cfg_attr(miri, ignore)] fn spill_round_trip() {
2383 set_spill_override(Some(test_pool()));
2384
2385 let data: Vec<Tuple> = (0..40_000u64)
2386 .map(|i| ((i / 4, i % 4), i % 8, 1i64))
2387 .collect();
2388 let data = consolidate(data);
2389
2390 let column = build_column(&data);
2391 let committed = TestChunk::commit(column, 0);
2392 assert!(committed.is_spilled(), "large body must spill");
2393 assert_eq!(committed.len(), data.len());
2394 assert_eq!(collect_column(&committed.clone().into_body()), data);
2395
2396 let mut batcher: ChunkBatcher<TestChunk> = Batcher::new(None, 0);
2397 for piece in data.chunks(10_000) {
2398 batcher.push_into(ColumnChunk::from_body(build_column(piece)));
2399 }
2400 let (sealed, _) = batcher.seal(Antichain::new());
2401 assert!(
2402 sealed.iter().any(ColumnChunk::is_spilled),
2403 "sealed output should contain spilled chunks",
2404 );
2405 assert_eq!(collect_chunks(sealed), data);
2406
2407 set_spill_override(None);
2408 }
2409
2410 #[mz_ore::test]
2413 #[cfg_attr(miri, ignore)] fn merge_spilled_chains() {
2415 set_spill_override(Some(test_pool()));
2416
2417 let a: Vec<Tuple> = (0..20_000u64).map(|i| ((i, 0), 0, 1i64)).collect();
2418 let b: Vec<Tuple> = (0..20_000u64).map(|i| ((i, 0), 0, 2i64)).collect();
2419
2420 let mut in1 = VecDeque::from([TestChunk::commit(build_column(&a), 0)]);
2421 let mut in2 = VecDeque::from([TestChunk::commit(build_column(&b), 0)]);
2422 assert!(in1[0].is_spilled() && in2[0].is_spilled());
2423
2424 let mut out = VecDeque::new();
2425 while !in1.is_empty() && !in2.is_empty() {
2426 TestChunk::merge(&mut in1, &mut in2, &mut out);
2427 }
2428 for tail in in1.drain(..).chain(in2.drain(..)) {
2429 out.push_back(tail);
2430 }
2431
2432 let expected: Vec<Tuple> = (0..20_000u64).map(|i| ((i, 0), 0, 3i64)).collect();
2433 assert_eq!(collect_chunks(out), expected);
2434
2435 set_spill_override(None);
2436 }
2437
2438 #[mz_ore::test]
2441 fn merge_untouched_survivor_stays_spilled() {
2442 let pool = test_pool();
2443 let low: Vec<Tuple> = (0..100u64).map(|i| ((i, 0), 0, 1i64)).collect();
2444 let high: Vec<Tuple> = (1000..1100u64).map(|i| ((i, 0), 0, 1i64)).collect();
2445
2446 let mut in1 = VecDeque::from([force_spill(
2447 ColumnChunk::from_body(build_column(&low)),
2448 &pool,
2449 )]);
2450 let mut in2 = VecDeque::from([force_spill(
2451 ColumnChunk::from_body(build_column(&high)),
2452 &pool,
2453 )]);
2454 let mut out = VecDeque::new();
2455 TestChunk::merge(&mut in1, &mut in2, &mut out);
2456
2457 assert!(in1.is_empty());
2460 assert_eq!(in2.len(), 1);
2461 assert!(in2[0].is_spilled(), "untouched survivor must stay spilled");
2462 let mut all = collect_chunks(out);
2463 Extend::extend(&mut all, collect_chunks(in2.drain(..)));
2464 let mut expected = low;
2465 Extend::extend(&mut expected, high);
2466 assert_eq!(all, expected);
2467 }
2468
2469 #[mz_ore::test]
2473 fn merge_derives_generational_depth() {
2474 let low: Vec<Tuple> = (0..100u64).map(|i| ((i, 0), 0, 1i64)).collect();
2475 let high: Vec<Tuple> = (50..150u64).map(|i| ((i, 0), 0, 1i64)).collect();
2476 let mut in1 = VecDeque::from([ColumnChunk::from_body(build_column(&low))]);
2477 let mut in2 = VecDeque::from([ColumnChunk::from_body(build_column(&high))]);
2478 assert_eq!(in1[0].depth(), 0, "fresh chunks start at depth 0");
2479 let mut out = VecDeque::new();
2480 TestChunk::merge(&mut in1, &mut in2, &mut out);
2481 assert!(!out.is_empty());
2482 for chunk in &out {
2483 assert_eq!(chunk.depth(), 1, "merge output is one past its inputs");
2484 }
2485 assert!(in1.is_empty());
2488 assert_eq!(in2.len(), 1);
2489 assert_eq!(in2[0].depth(), 0, "rewritten survivor keeps its depth");
2490
2491 let mut in1 = VecDeque::from([ColumnChunk::Resident(Rc::new(build_column(&low)), 3)]);
2494 let far: Vec<Tuple> = (1000..1100u64).map(|i| ((i, 0), 0, 1i64)).collect();
2495 let mut in2 = VecDeque::from([ColumnChunk::from_body(build_column(&far))]);
2496 let mut out = VecDeque::new();
2497 TestChunk::merge(&mut in1, &mut in2, &mut out);
2498 assert_eq!(out.len(), 1);
2499 assert_eq!(out[0].depth(), 4, "pass-through ages a generation");
2500 assert_eq!(collect_chunks(out), low);
2501 }
2502
2503 #[mz_ore::test]
2506 fn advance_preserves_depth() {
2507 let data: Vec<Tuple> = (0..100u64).map(|i| ((i, 0), 1, 1i64)).collect();
2508 let mut input = VecDeque::from([
2509 ColumnChunk::Resident(Rc::new(build_column(&data[..50])), 2),
2510 ColumnChunk::Resident(Rc::new(build_column(&data[50..])), 1),
2511 ]);
2512 let frontier = Antichain::from_elem(5u64);
2513 let mut out = VecDeque::new();
2514 TestChunk::advance(&mut input, frontier.borrow(), false, &mut out);
2515 for chunk in out.iter().chain(input.iter()) {
2516 assert_eq!(chunk.depth(), 2);
2517 }
2518 TestChunk::advance(&mut input, frontier.borrow(), true, &mut out);
2519 assert!(input.is_empty());
2520 assert!(!out.is_empty());
2521 for chunk in &out {
2522 assert_eq!(chunk.depth(), 2);
2523 }
2524 }
2525
2526 #[mz_ore::test]
2530 #[cfg_attr(miri, ignore)] fn settle_commits_at_accumulated_depth() {
2532 set_spill_override(Some(test_pool()));
2533 let big: Vec<Tuple> = (0..60_000u64).map(|i| ((i, 0), 0, 1i64)).collect();
2534 let mut input = VecDeque::from([
2535 ColumnChunk::Resident(Rc::new(build_column(&big)), 1),
2536 ColumnChunk::Resident(Rc::new(build_column(&[((0, 0), 0, 1)])), 0),
2537 ColumnChunk::Resident(Rc::new(build_column(&[((1, 0), 0, 1)])), 2),
2538 ]);
2539 let mut out = VecDeque::new();
2540 TestChunk::settle(&mut input, true, &mut out);
2541 assert!(input.is_empty());
2542 assert_eq!(out.len(), 2);
2543 assert!(out[0].is_spilled(), "large commit must spill");
2544 assert_eq!(out[0].depth(), 1, "sole commit keeps its depth");
2545 assert!(!out[1].is_spilled(), "small commit stays resident");
2546 assert_eq!(out[1].depth(), 2, "coalesced commit takes the max depth");
2547 set_spill_override(None);
2548 }
2549
2550 #[mz_ore::test]
2554 #[cfg_attr(miri, ignore)] fn settle_carry_commits_at_target() {
2556 let chunk_rows = u64::cast_from(800_000usize / 32);
2558 let mut input: VecDeque<TestChunk> = (0..4u64)
2559 .map(|c| {
2560 let data: Vec<Tuple> = (0..chunk_rows)
2561 .map(|i| ((c * chunk_rows + i, 0), 0, 1i64))
2562 .collect();
2563 ColumnChunk::from_body(build_column(&data))
2564 })
2565 .collect();
2566 let mut out = VecDeque::new();
2567 TestChunk::settle(&mut input, true, &mut out);
2568 assert!(out.len() < 4, "nothing coalesced");
2571 for chunk in &out {
2572 let col = chunk.clone().into_body();
2573 assert!(
2574 col.length_in_bytes() <= COMMIT_BYTES,
2575 "settled chunk of {} bytes exceeds the commit target",
2576 col.length_in_bytes(),
2577 );
2578 }
2579 assert_eq!(
2580 collect_chunks(out).len(),
2581 usize::try_from(4 * chunk_rows).unwrap(),
2582 );
2583 }
2584
2585 #[mz_ore::test]
2586 fn small_chunks_stay_resident() {
2587 set_spill_override(Some(test_pool()));
2588 let committed = TestChunk::commit(build_column(&[((1, 1), 0, 1)]), 0);
2589 assert!(!committed.is_spilled());
2590 set_spill_override(None);
2591 }
2592
2593 fn column_at_spill_floor() -> (ColumnBody<Tuple>, u64) {
2596 let mut col: ColumnBody<Tuple> = ColumnBody::default();
2597 let mut n = 0u64;
2598 while col.length_in_bytes() < SPILL_MIN_BYTES {
2599 col.push_into(((n, n), 0, 1));
2600 n += 1;
2601 }
2602 (col, n)
2603 }
2604
2605 #[mz_ore::test]
2608 fn spill_floor_boundary() {
2609 set_spill_override(Some(test_pool()));
2610 let (col, n) = column_at_spill_floor();
2611 let mut under: ColumnBody<Tuple> = ColumnBody::default();
2612 for m in 0..n - 1 {
2613 under.push_into(((m, m), 0, 1));
2614 }
2615 assert!(under.length_in_bytes() < SPILL_MIN_BYTES);
2616 assert!(!TestChunk::commit(under, 0).is_spilled());
2617 assert!(TestChunk::commit(col, 0).is_spilled());
2618 set_spill_override(None);
2619 }
2620
2621 #[mz_ore::test]
2625 fn spill_codec_depth_floor() {
2626 set_spill_override(Some(test_pool()));
2627 set_compress_min_depth_override(Some(2));
2628 let codec_name = |depth: u8| {
2632 let (codec, compressed) = codec_for_depth(depth);
2633 let name = format!("{:?}", codec);
2634 assert_eq!(compressed, name == "Lz4Codec", "flag tracks the codec");
2635 name
2636 };
2637 assert_eq!(codec_name(0), "IdentityCodec");
2638 assert_eq!(codec_name(1), "IdentityCodec");
2639 assert_eq!(codec_name(2), "Lz4Codec");
2640 assert_eq!(codec_name(u8::MAX), "Lz4Codec");
2641
2642 let data: Vec<Tuple> = (0..20_000u64).map(|i| ((i, 0), 0, 1i64)).collect();
2643 let data = consolidate(data);
2644 let column = build_column(&data);
2645 for depth in [0u8, 1, 2, 3] {
2646 let chunk = TestChunk::commit(column.clone(), depth);
2647 assert!(chunk.is_spilled(), "depth {depth} must spill");
2648 assert_eq!(collect_column(&chunk.into_body()), data);
2649 }
2650 set_spill_override(None);
2651 set_compress_min_depth_override(None);
2652
2653 set_compress_min_depth_override(Some(DEFAULT_COMPRESS_MIN_DEPTH));
2655 assert_eq!(codec_name(0), "IdentityCodec");
2656 assert_eq!(codec_name(1), "Lz4Codec");
2657 set_compress_min_depth_override(None);
2658 }
2659
2660 #[mz_ore::test]
2666 #[cfg_attr(miri, ignore)] fn merge_survivor_crosses_compression_floor() {
2668 set_spill_override(Some(test_pool()));
2669 set_compress_min_depth_override(Some(1));
2670
2671 let low = consolidate((0..20_000u64).map(|i| ((i, 0), 0, 1i64)).collect());
2672 let far = consolidate((100_000..120_000u64).map(|i| ((i, 0), 0, 1i64)).collect());
2673 let fresh_far = || VecDeque::from([TestChunk::commit(build_column(&far), 0)]);
2674
2675 let mut in1 = VecDeque::from([TestChunk::commit(build_column(&low), 0)]);
2678 let mut in2 = fresh_far();
2679 assert!(in1[0].is_spilled() && in2[0].is_spilled());
2680 assert!(
2681 !body_compressed(&in1[0]),
2682 "a fresh body below the floor is identity coded"
2683 );
2684
2685 let mut out = VecDeque::new();
2686 TestChunk::merge(&mut in1, &mut in2, &mut out);
2687 assert_eq!(out.len(), 1);
2688 let survived = out.pop_front().expect("the lower front passes through");
2689 assert_eq!(survived.depth(), 1, "survival ages across the floor");
2690 assert!(
2691 survived.is_spilled(),
2692 "the crossing re-spills, it does not evict"
2693 );
2694 assert!(
2695 body_compressed(&survived),
2696 "the survivor is re-spilled under the compressing codec"
2697 );
2698
2699 let mut in1 = VecDeque::from([survived]);
2702 let mut in2 = fresh_far();
2703 let mut out = VecDeque::new();
2704 TestChunk::merge(&mut in1, &mut in2, &mut out);
2705 assert_eq!(out.len(), 1);
2706 assert_eq!(out[0].depth(), 2, "an aged survivor keeps aging");
2707 assert!(out[0].is_spilled());
2708 assert_eq!(
2709 collect_chunks(out),
2710 low,
2711 "the body reads back intact across both survivals"
2712 );
2713
2714 set_spill_override(None);
2715 set_compress_min_depth_override(None);
2716 }
2717
2718 #[mz_ore::test]
2725 #[cfg_attr(miri, ignore)] fn merge_survivor_ages_while_shared() {
2727 set_spill_override(Some(test_pool()));
2728 set_compress_min_depth_override(Some(1));
2729
2730 let low = consolidate((0..20_000u64).map(|i| ((i, 0), 0, 1i64)).collect());
2731 let far = consolidate((100_000..120_000u64).map(|i| ((i, 0), 0, 1i64)).collect());
2732
2733 let source = TestChunk::commit(build_column(&low), 0);
2736 let ColumnChunk::Spilled(source_body, 0) = &source else {
2737 panic!("a fresh commit above the spill floor is spilled at depth 0");
2738 };
2739 let source_body = Rc::clone(source_body);
2740
2741 let mut in1 = VecDeque::from([source.clone()]);
2742 let mut in2 = VecDeque::from([TestChunk::commit(build_column(&far), 0)]);
2743 let mut out = VecDeque::new();
2744 TestChunk::merge(&mut in1, &mut in2, &mut out);
2745
2746 assert_eq!(out.len(), 1);
2747 assert_eq!(out[0].depth(), 1, "a shared body ages all the same");
2748 let ColumnChunk::Spilled(survived_body, _) = &out[0] else {
2749 panic!("the survivor stays spilled");
2750 };
2751 assert!(
2752 Rc::ptr_eq(&source_body, survived_body),
2753 "a shared body is aged in place, not re-spilled"
2754 );
2755 assert_eq!(source.depth(), 0, "the other holder is left as it was");
2756
2757 let mut in1 = VecDeque::from([out.pop_front().expect("survivor observed above")]);
2760 let mut in2 = VecDeque::from([TestChunk::commit(build_column(&far), 0)]);
2761 let mut out = VecDeque::new();
2762 TestChunk::merge(&mut in1, &mut in2, &mut out);
2763 assert_eq!(out.len(), 1);
2764 assert_eq!(out[0].depth(), 2, "aging past the floor is not pinned");
2765 assert_eq!(collect_chunks(out), low);
2766
2767 set_spill_override(None);
2768 set_compress_min_depth_override(None);
2769 }
2770
2771 #[mz_ore::test]
2777 #[cfg_attr(miri, ignore)] fn survive_merge_retries_missed_migrations() {
2779 let low = consolidate((0..20_000u64).map(|i| ((i, 0), 0, 1i64)).collect());
2780 let far = consolidate((100_000..120_000u64).map(|i| ((i, 0), 0, 1i64)).collect());
2781
2782 let survive = |chunk: TestChunk| {
2785 let mut in1 = VecDeque::from([chunk]);
2786 let mut in2 = VecDeque::from([TestChunk::commit(build_column(&far), 0)]);
2787 let mut out = VecDeque::new();
2788 TestChunk::merge(&mut in1, &mut in2, &mut out);
2789 out.pop_front().expect("the lower front passes through")
2790 };
2791
2792 set_spill_override(Some(test_pool()));
2795 set_compress_min_depth_override(Some(1));
2796 let chunk = TestChunk::commit(build_column(&low), 0);
2797 assert!(!body_compressed(&chunk));
2798 set_spill_override(None);
2799 let chunk = survive(chunk);
2800 assert_eq!(chunk.depth(), 1, "aging does not need a pool");
2801 assert!(!body_compressed(&chunk), "no pool, no migration");
2802 set_spill_override(Some(test_pool()));
2803 let chunk = survive(chunk);
2804 assert!(
2805 body_compressed(&chunk),
2806 "the migration retries once a pool is back"
2807 );
2808
2809 let chunk = TestChunk::commit(build_column(&low), 0);
2812 let held = chunk.clone();
2813 let chunk = survive(chunk);
2814 assert!(!body_compressed(&chunk), "shared, so not migrated");
2815 drop(held);
2816 let chunk = survive(chunk);
2817 assert!(
2818 body_compressed(&chunk),
2819 "the migration retries once the body is unshared"
2820 );
2821
2822 set_compress_min_depth_override(Some(8));
2827 let chunk = TestChunk::commit(build_column(&low), 3);
2828 assert!(!body_compressed(&chunk));
2829 set_compress_min_depth_override(Some(1));
2830 let chunk = survive(chunk);
2831 assert_eq!(chunk.depth(), 4);
2832 assert!(
2833 body_compressed(&chunk),
2834 "lowering the floor migrates bodies already past it"
2835 );
2836
2837 set_spill_override(None);
2838 set_compress_min_depth_override(None);
2839 }
2840
2841 #[mz_ore::test]
2847 #[cfg_attr(miri, ignore)]
2848 fn spill_gates_compose() {
2849 let installed =
2850 crate::pool_config::apply_pool_config(crate::pool_config::PoolPagerConfig {
2851 budget_bytes: 32 << 20,
2852 spill_threads: 1,
2853 eager_backing: false,
2854 rss_target_bytes: 16 << 20,
2855 });
2856 assert!(installed, "pool reservation failed");
2857 let (col, _) = column_at_spill_floor();
2859 let commit = |col: &ColumnBody<Tuple>| TestChunk::commit(col.clone(), 0).is_spilled();
2860 let spill_ref = |col: &ColumnBody<Tuple>| try_spill_ref(col, 0).is_some();
2861
2862 assert!(!commit(&col), "both gates off");
2863 assert!(!spill_ref(&col), "all gates off");
2864 set_storage_spill_enabled(true);
2865 assert!(commit(&col), "the storage gate alone spills");
2866 assert!(
2867 !spill_ref(&col),
2868 "the chunk gates must not spill sink bodies"
2869 );
2870 set_compute_spill_enabled(false);
2871 assert!(
2872 commit(&col),
2873 "the compute setter must not clobber the storage gate"
2874 );
2875 set_compute_spill_enabled(true);
2876 set_storage_spill_enabled(false);
2877 assert!(commit(&col), "the compute gate alone spills");
2878 set_compute_spill_enabled(false);
2879 assert!(!commit(&col), "both gates off again");
2880 set_sink_spill_enabled(true);
2881 assert!(spill_ref(&col), "the sink gate alone spills sink bodies");
2882 assert!(!commit(&col), "the sink gate must not spill chunks");
2883 set_sink_spill_enabled(false);
2884 set_compress_min_depth_override(None);
2885 }
2886
2887 #[mz_ore::test]
2890 fn spill_align_round_trip() {
2891 let pool = test_pool();
2892 let data: Vec<Tuple> = (0..64u64).map(|k| ((k, k), 0, 1)).collect();
2893 let spilled = force_spill(ColumnChunk::from_body(build_column(&data)), &pool);
2894 let column = spilled.into_body();
2895 let ColumnBody::Words(words) = &column else {
2896 panic!("a spilled body reads back as ColumnBody::Words");
2897 };
2898 let words = words.clone();
2899 let respilled = force_spill(ColumnChunk::from_body(column), &pool);
2900 let reread = respilled.into_body();
2901 let ColumnBody::Words(words2) = &reread else {
2902 panic!("a spilled body reads back as ColumnBody::Words");
2903 };
2904 assert_eq!(&words, words2, "byte-identical round trip");
2905 assert_eq!(collect_column(&reread), data);
2906 }
2907
2908 #[mz_ore::test]
2909 fn merge_depth_saturates() {
2910 let a = ColumnChunk::Resident(
2911 Rc::new(build_column(&[((1, 0), 0, 1), ((3, 0), 0, 1)])),
2912 u8::MAX,
2913 );
2914 let b = ColumnChunk::Resident(
2915 Rc::new(build_column(&[((2, 0), 0, 1), ((4, 0), 0, 1)])),
2916 u8::MAX,
2917 );
2918 let mut in1 = VecDeque::from([a]);
2919 let mut in2 = VecDeque::from([b]);
2920 let mut out = VecDeque::new();
2921 TestChunk::merge(&mut in1, &mut in2, &mut out);
2922 for chunk in out.iter().chain(in1.iter()).chain(in2.iter()) {
2923 assert_eq!(chunk.depth(), u8::MAX, "depth saturates");
2924 }
2925 }
2926
2927 #[mz_ore::test]
2928 fn into_body_copies_shared_resident() {
2929 let data: Vec<Tuple> = vec![((1, 1), 0, 1), ((2, 2), 0, 1)];
2930 let a = ColumnChunk::from_body(build_column(&data));
2931 let b = a.clone();
2932 assert_eq!(collect_column(&a.into_body()), data);
2933 assert_eq!(collect_column(&b.into_body()), data);
2934 }
2935}