Skip to main content

mz_timely_util/columnar/
chunk.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! [`ColumnChunk`]: differential's [`Chunk`] over [`ColumnBody`] updates.
11//!
12//! A chunk is a sorted, consolidated run of `(D, T, R)` updates in the flat
13//! columnar layout, in one of two homes:
14//!
15//! * **Resident**: an `Rc`-shared [`ColumnBody`] on the heap. Fresh input,
16//!   merge output, and small tails live here.
17//! * **Spilled**: the serialized body in the process [`Pool`], with the record
18//!   count and the first and last data items resident. The pool owns residency
19//!   from there, with slots under a memory budget and compression and device
20//!   pageout under pressure, and a body that dies before pressure reaches it
21//!   is freed without I/O.
22//!
23//! Reads of a spilled body are copy-out and scoped to the call that needs
24//! them, the contract that lets the pool evict with no reader accounting
25//! (see [`mz_ore::pool`]).
26//!
27//! Spilling happens in [`Chunk::settle`], the trait's designated commit point:
28//! chunks moved to settled output are handed to the pool when spilling is
29//! enabled (see [`set_compute_spill_enabled`] and [`set_storage_spill_enabled`]
30//! for how the per-commit destination resolves). Grading is by serialized
31//! bytes, the ship size [`Column`] already targets, rather than by the
32//! record-count `TARGET`, since record count does not bound bytes for
33//! variable-width data.
34//!
35//! Chunks whose data is a `(key, val)` pair additionally implement
36//! [`UnloadChunk`], the bulk-read capability: sorted probe keys in, matching
37//! updates appended to caller-owned staging, with `locate` answered from the
38//! resident fence metadata so a probe set faults only the chunk bodies it
39//! actually touches.
40
41use 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
70/// Compute's leg of the process spill gate. See [`set_compute_spill_enabled`].
71static COMPUTE_SPILL_ENABLED: AtomicBool = AtomicBool::new(false);
72
73/// Storage's leg of the process spill gate. See [`set_storage_spill_enabled`].
74static STORAGE_SPILL_ENABLED: AtomicBool = AtomicBool::new(false);
75
76/// The gate for bodies spilled through [`try_spill_ref`]. See
77/// [`set_sink_spill_enabled`].
78static SINK_SPILL_ENABLED: AtomicBool = AtomicBool::new(false);
79
80thread_local! {
81    /// A thread-scoped pool override, taking precedence over the global
82    /// enable flag and pool. Lets tests and benches spill through a private
83    /// pool without touching process-global state.
84    static SPILL_OVERRIDE: RefCell<Option<Pool>> = const { RefCell::new(None) };
85
86    /// A thread-scoped depth-floor override, taking precedence over the
87    /// global value. Lets tests pin the floor without racing concurrently
88    /// running tests on the process-global state.
89    #[cfg(test)]
90    static COMPRESS_MIN_DEPTH_OVERRIDE: Cell<Option<u8>> = const { Cell::new(None) };
91
92    /// Reusable staging for call-scoped reads of spilled bodies.
93    static READ_SCRATCH: RefCell<Vec<u64>> = const { RefCell::new(Vec::new()) };
94}
95
96/// Enable or disable chunk spilling on behalf of compute's arrangement
97/// batchers.
98///
99/// Chunks carry no subsystem identity, so the spill decision is process-wide:
100/// committed chunks spill while *either* the compute or the storage gate is
101/// set. Each subsystem's config application writes only its own gate, so the
102/// two dyncfg flags compose as an OR instead of clobbering each other.
103///
104/// Takes effect at the next `settle`. Already-spilled chunks are unaffected
105/// either way. The pool is resolved per commit through
106/// [`crate::pool_config::active_pool`], so chunks spill only once
107/// `apply_pool_config` has installed and budgeted the pool. With no pool
108/// installed chunks stay resident regardless of the gates.
109pub fn set_compute_spill_enabled(enabled: bool) {
110    COMPUTE_SPILL_ENABLED.store(enabled, Ordering::Relaxed);
111}
112
113/// Enable or disable chunk spilling on behalf of storage's upsert dataflows.
114///
115/// See [`set_compute_spill_enabled`] for the shared-gate semantics.
116pub fn set_storage_spill_enabled(enabled: bool) {
117    STORAGE_SPILL_ENABLED.store(enabled, Ordering::Relaxed);
118}
119
120/// Enable or disable spilling of bodies offered through [`try_spill_ref`],
121/// which is how compute's MV sink correction buffer spills.
122///
123/// Independent of the compute and storage legs: those gate [`ColumnChunk`]s
124/// only, and this gates [`try_spill_ref`] only, so enabling one subsystem's
125/// spilling does not spill another's state. Like them, it takes effect only
126/// once `apply_pool_config` has installed the pool.
127pub fn set_sink_spill_enabled(enabled: bool) {
128    SINK_SPILL_ENABLED.store(enabled, Ordering::Relaxed);
129}
130
131/// Set or unset the pool through which this thread's chunk spills are
132/// routed, taking precedence over the gates and the process pool. `None`
133/// restores the global resolution.
134pub fn set_spill_override(pool: Option<Pool>) {
135    SPILL_OVERRIDE.with(|cell| *cell.borrow_mut() = pool);
136}
137
138/// The youngest generational depth whose spilled bodies are compressed. See
139/// [`set_compress_min_depth`].
140static COMPRESS_MIN_DEPTH: AtomicU8 = AtomicU8::new(DEFAULT_COMPRESS_MIN_DEPTH);
141
142/// Set the youngest generational depth whose spilled bodies are compressed.
143///
144/// A chunk at depth `d` is rewritten (merged, extracted, advanced) with
145/// frequency proportional to `2^-d` under geometric merging, so compressing
146/// a shallow chunk buys a short stay in the pool at the cost of a guaranteed
147/// near-term codec round-trip: the body is encoded only to be read back and
148/// decoded by the next rewrite. Generations below the floor spill under the
149/// identity codec instead: still budgeted and swap-backed like every extent,
150/// but encode and decode are copies. The floor never exempts a body from the
151/// pool, so it cannot grow unbudgeted resident state.
152///
153/// The floor cannot strand a long-lived body uncompressed: a chunk that a
154/// merge carries forward untouched also ages a generation, and its body is
155/// re-spilled under the compressing codec when it crosses the floor (see
156/// `survive_merge`).
157///
158/// `0` compresses every spilled body. Consulted at every commit, so changes
159/// apply to running dataflows.
160pub fn set_compress_min_depth(depth: u8) {
161    COMPRESS_MIN_DEPTH.store(depth, Ordering::Relaxed);
162}
163
164/// Set or unset a thread-scoped depth-floor override, taking precedence over
165/// [`set_compress_min_depth`]. Tests run concurrently and must not race on
166/// the process-global floor.
167#[cfg(test)]
168pub fn set_compress_min_depth_override(depth: Option<u8>) {
169    COMPRESS_MIN_DEPTH_OVERRIDE.with(|cell| cell.set(depth));
170}
171
172/// The depth floor in effect for this thread's commits.
173fn 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
181/// The codec a body at `depth` stores under, identity below the compression
182/// floor and lz4 at and past it, paired with whether that codec compresses.
183/// One read of the floor, so the pair cannot disagree with itself when the
184/// floor moves under a concurrent commit.
185fn 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
193/// The pool a body of `len_bytes` spills into while `enabled`, or `None` when
194/// it stays resident.
195///
196/// The one place the spill decision lives: the thread override, the gate, the
197/// installed pool, and the size floor.
198fn spill_target(enabled: bool, len_bytes: usize) -> Option<Pool> {
199    resolve_pool(enabled).filter(|_| len_bytes >= SPILL_MIN_BYTES)
200}
201
202/// The pool committed chunks spill to, if any.
203fn spill_pool() -> Option<Pool> {
204    resolve_pool(chunk_spill_enabled())
205}
206
207/// Whether [`ColumnChunk`]s spill: the OR of the compute and storage legs.
208fn chunk_spill_enabled() -> bool {
209    COMPUTE_SPILL_ENABLED.load(Ordering::Relaxed) || STORAGE_SPILL_ENABLED.load(Ordering::Relaxed)
210}
211
212/// The thread's override pool, else the installed pool while `enabled`.
213fn 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
224/// Scratch capacity retained across reads, in words. A read larger than this
225/// releases the buffer afterward, so a thread's scratch does not ratchet to
226/// the largest body it ever carried (heap no pool gauge can see).
227const SCRATCH_RETAIN_WORDS: usize = 1 << 18;
228
229/// Run `f` with this thread's read scratch, cleared of any previous use.
230fn 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
244/// The serialized-byte size committed chunks aim for, matching the ship size
245/// of the columnar merge machinery.
246const COMMIT_BYTES: usize = 2 << 20;
247
248/// Bodies smaller than this stay resident: the pool's smallest size class is
249/// 64 KiB, so spilling below it trades no meaningful memory for slot waste.
250///
251/// Sub-floor bodies are invisible to the pool's budget, which is safe only
252/// while they are rare. `settle` coalesces toward `COMMIT_BYTES` before
253/// committing, so in the harness only a final `done` tail commits below the
254/// floor. A caller that commits many small chunks directly accumulates
255/// unbudgeted heap, and no accounting here would catch it.
256const SPILL_MIN_BYTES: usize = 64 << 10;
257
258/// The default compression depth floor: fresh (depth 0) bodies spill
259/// uncompressed.
260///
261/// A fresh chunk is consumed by its first merge with certainty, so
262/// compressing it can never save pool bytes for longer than one merge
263/// cadence and always costs a full encode plus decode. Depth 1 and beyond
264/// have survived a merge and wait geometrically longer for the next, so
265/// their compression amortizes.
266const DEFAULT_COMPRESS_MIN_DEPTH: u8 = 1;
267
268/// Records of a body with `len` records in `bytes` that fit in `space`
269/// bytes at the body's average width, keeping 5% of `COMMIT_BYTES` as a
270/// margin for per-piece framing. Uniform rows cut this way fill their size
271/// class and leave one short remainder that `settle` can coalesce. Uneven
272/// widths can still overflow, so callers check the cut piece.
273fn 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
278/// Whether a body is big enough to commit on its own.
279fn at_commit_size<C: Columnar>(body: &ColumnBody<C>) -> bool {
280    body.length_in_bytes() >= COMMIT_BYTES - COMMIT_BYTES / 10
281}
282
283/// Narrow a columnar ref to a shorter lifetime, so refs from different
284/// borrows, such as a probe column and a chunk's own columns, can be compared
285/// (the refs are lifetime-invariant).
286#[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
291/// A spilled chunk body: the serialized column in the pool, plus the resident
292/// metadata every [`Chunk`] must answer without fetching. That metadata is
293/// the record count, the first and last data items (the fence entries
294/// [`UnloadChunk::locate`] consults), and the time bounds `extract` consults
295/// to pass frontier-disjoint chunks through without loading them.
296pub struct SpilledBody<D: Columnar, T> {
297    /// Number of updates in the body.
298    records: usize,
299    /// The first and last data items, as a two-element container. One
300    /// container rather than two singletons, so the leaf allocations are not
301    /// duplicated per fence.
302    fences: D::Container,
303    /// The minimal times in the body: a lower bound antichain every
304    /// contained time is greater-or-equal to. Folded into `extract`'s
305    /// residual frontier when the chunk is kept whole.
306    time_lower: Antichain<T>,
307    /// The maximal times in the body. Some contained time is
308    /// greater-or-equal to a frontier exactly when some maximal time is,
309    /// which is `extract`'s ship-whole test. A single element for totally
310    /// ordered times, hence the inline capacity.
311    time_upper: SmallVec<[T; 1]>,
312    /// Whether the body was inserted under the compressing codec. The pool
313    /// stores the codec itself and reads decode through it, so this is the
314    /// only handle chunk code has on what a body is stored as, and it is
315    /// what `survive_merge` consults to decide a body wants migrating.
316    /// Deriving that from depth instead would tie it to a single transition
317    /// and miss every path that skips it.
318    compressed: bool,
319    /// The body's serialized size, as it was before the pool's codec saw
320    /// it. Retained because the pool reports no per-chunk figure, and a
321    /// chunk that reported nothing would drop its operator's share of the
322    /// batcher's memory out of the introspection tables.
323    len_bytes: usize,
324    /// The pool chunk holding the serialized column.
325    handle: ChunkHandle,
326}
327
328/// A sorted, consolidated run of `(D, T, R)` updates, resident or spilled.
329///
330/// Every chunk carries a generational depth counting the merge cadences it
331/// has lived through: fresh chunks are depth 0, a merge output is one
332/// generation past its deepest input (saturating at `u8::MAX`, where
333/// remerged long-lived chunks stay), a chunk a merge carries forward
334/// untouched also gains a generation (see `survive_merge`), and rewrites
335/// within a generation (extract, advance, settle coalescing) preserve
336/// depth.
337///
338/// Depth belongs to the chunk, not to the body: a body outlives the chunks
339/// that share it, and aging must not depend on whether a caller happens to
340/// hold the only reference. At spill time the depth becomes the pool's
341/// [`ChunkHints`], so repeatedly merged (older, colder) data lands in deeper
342/// eviction bands. Hints are fixed at insert, so a chunk aged without a
343/// re-spill keeps the band it spilled into.
344pub enum ColumnChunk<D: Columnar, T: Columnar, R: Columnar> {
345    /// Body on the heap, shared via `Rc`, with its generational depth.
346    Resident(Rc<ColumnBody<(D, T, R)>>, u8),
347    /// Body in the pool, with its generational depth. See [`SpilledBody`].
348    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    /// Wrap a sorted, consolidated, non-empty body as a resident chunk of
374    /// the youngest generation.
375    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    /// Append pieces of at most `COMMIT_BYTES`, preserving order and
381    /// generational depth. A single update may exceed the bound because it
382    /// cannot be split.
383    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    /// Append `body` in consecutive pieces of `records` records. A piece
401    /// that still exceeds `COMMIT_BYTES`, from uneven widths, is halved until
402    /// it fits or holds one record. Halving bounds the recursion depth by
403    /// log2 of `records`, where cutting again at the piece's own average
404    /// width can shrink a piece holding one wide record by only a few
405    /// percent per level.
406    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    /// [`Self::push_cuts`] over `range` of `view`. An oversized piece is
417    /// dropped before descending, and the halves are cut again from `view`,
418    /// so transient memory is one piece rather than one copy of a wide
419    /// update per level.
420    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    /// The body, owned. A spilled body is copied out of the pool within this
444    /// call. A shared resident body is copied in its current form.
445    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    /// The body for the duration of `f`: a resident body is borrowed, a
459    /// spilled body is loaded from the pool for the call and dropped after it.
460    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    /// True when the body lives in the pool.
475    pub fn is_spilled(&self) -> bool {
476        matches!(self, ColumnChunk::Spilled(_, _))
477    }
478
479    /// The number of updates, from resident state only.
480    fn records(&self) -> usize {
481        match self {
482            ColumnChunk::Resident(col, _) => col.len(),
483            ColumnChunk::Spilled(body, _) => body.records,
484        }
485    }
486
487    /// The generational depth, from resident state only.
488    fn depth(&self) -> u8 {
489        match self {
490            ColumnChunk::Resident(_, depth) | ColumnChunk::Spilled(_, depth) => *depth,
491        }
492    }
493
494    /// The first and last data items, from resident state only.
495    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    /// Commit a non-empty body at the given generational depth: spill it to
509    /// the pool when spilling is on and the body is worth a slot, else keep it
510    /// resident.
511    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    /// Spill a non-empty body into `pool` unconditionally, capturing the
525    /// resident fence metadata.
526    ///
527    /// Generations below the compression depth floor store under the
528    /// identity codec: rewritten too soon for compression to amortize, they
529    /// stay budgeted and swap-backed while encode and decode reduce to
530    /// copies.
531    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    /// Age a chunk that a merge carried forward untouched by one generation.
559    /// Depth counts merge cadences lived through, not rewrites, so a
560    /// pass-through survivor ages like a merged chunk; the bump rides on the
561    /// chunk, so it is free whether or not the body is shared.
562    ///
563    /// An identity-coded body at or past the floor wants migrating, so a
564    /// sole owner re-spills it compressed; without the re-spill,
565    /// key-disjoint input would keep its whole spilled backlog
566    /// identity-coded for as long as it lived. The test is the body's stored
567    /// codec against the floor, not a depth transition, so a migration that
568    /// cannot happen now (a shared body, no pool installed, a floor lowered
569    /// long after the spill) is retried at the next survival rather than
570    /// consumed. A shared body is skipped because re-spilling this reference
571    /// cannot change what the other holder stores, and the compaction merger
572    /// that shares bodies rewrites its clones immediately.
573    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    /// The chunk's time bounds: borrowed from the stored metadata for
597    /// spilled bodies, computed by a time-column scan for resident ones. The
598    /// scan costs less than the copy it lets `extract` avoid when the chunk
599    /// passes through whole.
600    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    /// The time bounds of a non-empty body: the antichain of minimal times
617    /// (every contained time is greater-or-equal to some element) and the
618    /// set of maximal times (some contained time is greater-or-equal to a
619    /// frontier exactly when some maximal one is).
620    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        // One owned time reused across the scan, so times with owned
628        // allocations do not allocate per element; the bound sets clone only
629        // the elements they retain.
630        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/// The chunk-side [`ExtentCodec`]: a little-endian `u32` body-length prefix
644/// followed by one lz4 block, the framing
645/// `lz4_flex::block::compress_prepend_size` produces. Every chunk consumer
646/// passes [`LZ4_CODEC`] at insert; the pool itself has no codec opinion.
647#[derive(Debug)]
648pub struct Lz4Codec;
649
650/// The [`Lz4Codec`] instance chunk consumers pass to
651/// [`Pool::insert_with`].
652pub 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
679/// Serialize a body into a pool slot, writing its serialized form through a cursor over the
680/// slot memory. A serialized body's encoding is one copy. Sizing is exact, so a short or
681/// overlong write is a contract violation and panics.
682fn 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
703/// Spill a serialized copy of `body` into the process pool, leaving `body` untouched, or
704/// `None` when the body stays resident.
705///
706/// A body stays resident when the gate [`set_sink_spill_enabled`] sets is off, when no pool is
707/// installed, or when it is smaller than `SPILL_MIN_BYTES`. `depth` is the body's generational
708/// depth: it selects the stored codec and the pool's eviction band.
709///
710/// For consumers that keep their own chunk representation and so cannot use [`ColumnChunk`];
711/// the two share the spill decision through `spill_target`. The encode writes straight into
712/// the pool slot, so a spilled body costs no intermediate allocation and the caller can
713/// [`ColumnBody::clear`] and keep its allocation. A caller that gets `None` back still owns
714/// the body and must keep it resident itself.
715pub 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    /// A nominal record count for the harness's fuel and ladder accounting,
739    /// not a bound. Actual chunk sizing is by serialized bytes: `merge` and
740    /// `extract` cut output at the [`Column`] ship threshold, and `settle`
741    /// grades by `at_commit_size`, so a chunk of narrow records can hold more
742    /// records than this and nothing here consults it.
743    const TARGET: usize = 65536;
744
745    fn len(&self) -> usize {
746        self.records()
747    }
748
749    /// [`ColumnBody::merge_from`] does the work: gallop bulk-copies for disjoint
750    /// runs, semigroup consolidation on equal `(data, time)`, output cut at
751    /// the ship threshold.
752    ///
753    /// Fronts whose data ranges are disjoint never load at all: the resident
754    /// fence entries decide, and the lower front moves to the output whole.
755    fn merge(in1: &mut VecDeque<Self>, in2: &mut VecDeque<Self>, out: &mut VecDeque<Self>) {
756        // Disjoint fast path: when one front lies strictly below the other's
757        // first data item (equal boundary data could still interleave on
758        // time), the merged prefix through the shared horizon is exactly that
759        // front, unchanged.
760        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        // Merged output is one generation past its deepest input. A survivor
784        // (untouched or rewritten from its remainder) keeps its own depth.
785        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        // The disjoint fast paths above move fronts without reading a row, so
808        // only this path counts as merge work.
809        metrics::record(metrics::Stage::Merge, positions[0] + positions[1], 0);
810        let [col_a, col_b] = &mut cols;
811        // Per input side: the loaded column and the merge's consumed position
812        // within it, the side's pre-merge depth, its original spilled body
813        // when it had one, and the deque a survivor returns to.
814        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                // Untouched survivor: restore it as it was (the loaded copy
821                // is dropped), aged one generation by its survival.
822                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    /// Partition one front chunk by `frontier`, folding kept times into
840    /// `residual`. One chunk per call, so the harness settles both sides
841    /// between chunks. Output is cut at the ship threshold.
842    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        // Whole-chunk pass-through from the resident time bounds: a chunk
853        // the frontier is entirely past ships unchanged, one entirely at or
854        // past the frontier keeps unchanged. Spilled bodies pass through
855        // without a load, a re-commit, or any codec work; only chunks the
856        // frontier actually splits are loaded below.
857        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            // The residual must lower-bound every kept time, which is the
864            // chunk's lower bound antichain by construction.
865            for m in time_lower.elements() {
866                residual.insert_ref(m);
867            }
868            keep.push_back(chunk);
869            return;
870        }
871        // Partitioning rewrites within a generation, so both sides keep the
872        // input chunk's depth.
873        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        // TODO: rewrite the underlying `ColumnBody::extract` as two passes, the
880        // time column first to find run boundaries, then bulk per-range
881        // copies of the remaining leaves.
882        // Move a side's accumulation to its queue, at the ship threshold
883        // mid-loop, or any non-empty remainder at the end.
884        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    /// Advance times by `frontier` and consolidate, withholding the trailing
901    /// `D` group as the carry unless `done` (its updates may continue in input
902    /// this call has not seen).
903    ///
904    /// The input concatenates into the carry's container, so a group that
905    /// grows across many calls is appended to, not rebuilt. Each record is
906    /// copied once on arrival, keeping the run linear. Advancing is
907    /// lattice-monotone but not order-monotone, so each group's advanced
908    /// times are re-sorted before adjacent equal times fold.
909    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        // Advancing rewrites within a generation, so output and carry keep
919        // the deepest input depth. Only merges increment.
920        let mut depth = front.depth();
921        // Concatenate the input into one body, reusing the front chunk's
922        // storage when it is exclusively owned (the usual case: it is last
923        // call's carry).
924        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        // Giant-group early-out: if the whole input is one `D` group, nothing
942        // is provably complete. Unless `done`, push it all back as the carry.
943        if !done && data.get(0) == data.get(total - 1) {
944            input.push_front(ColumnChunk::Resident(Rc::new(base), depth));
945            return;
946        }
947
948        // The processing bound: everything, or everything before the trailing
949        // `D` group when it must be withheld.
950        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        // Count only the rows this call advances. The withheld carry is
961        // counted by the call that finally processes it, so a group spanning
962        // many calls counts once.
963        metrics::record(metrics::Stage::Advance, end, 0);
964
965        let mut result = <(D, T, R) as Columnar>::Container::default();
966        // Per-group scratch: advanced owned times with owned diffs.
967        let mut scratch: Vec<(T, R)> = Vec::new();
968        let mut index = 0;
969        // Cut output at the commit size, checked amortized by emitted records
970        // (the size test walks the container's leaves, so probing it per
971        // record would be quadratic). Records, not groups: a single group may
972        // carry arbitrarily many advanced times, and a cut is legal anywhere
973        // in the sorted sequence, so bounding by records keeps the largest
974        // possible output chunk within one check period of the target. It
975        // must not outgrow the pool's largest size class, past which a body
976        // degrades to a permanently resident heap chunk.
977        const CUT_CHECK_RECORDS: usize = 1024;
978        let mut records_since_check = 0usize;
979        // TODO: the output leaves are addressed independently, so a group
980        // that folds nothing (no time collisions, no zeroed diffs) could bulk
981        // `extend_from_self` the D leaf over the whole group range and push
982        // only the advanced times and diffs per record, and a singleton group
983        // (the common case for mostly-unique D) could skip the scratch and
984        // sort round trip entirely.
985        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        // Rebuild the withheld trailing group as the carry.
1029        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    /// Grade by serialized bytes and commit: spilled chunks pass through
1040    /// untouched, resident chunks at the commit size commit as they are, and
1041    /// smaller neighbors coalesce until the accumulation reaches it.
1042    /// Committing is the spill hook (see `ColumnChunk::commit`). Unless
1043    /// `done`, a tail below the commit size returns to the front of `input`.
1044    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    /// [`Chunk::settle`], with `commit` choosing what a graded chunk becomes.
1058    ///
1059    /// A caller that reads its output back immediately passes `false`: the
1060    /// chunk is graded resident, so the body is not serialized into a pool
1061    /// slot only to be copied straight back out, and the pool's budget is
1062    /// not charged for a body whose lifetime ends at the next read.
1063    fn settle_graded(
1064        input: &mut VecDeque<Self>,
1065        done: bool,
1066        out: &mut VecDeque<Self>,
1067        commit: bool,
1068    ) {
1069        // Resident chunks below the commit size accumulate into `acc` until
1070        // it reaches that size. Its depth is the deepest among the chunks it
1071        // holds, since coalescing rewrites within a generation.
1072        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                // Cut by borrow, so a shared body is copied once, into its
1086                // pieces.
1087                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            // A chunk is copied into `acc` from its borrow, so a shared body
1097            // is never unwrapped.
1098            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            // The chunk does not fit in `acc`. If the chunk is below the
1116            // commit size, as many of its records as fit move into `acc` and
1117            // the rest return to the front of `input`. Without that move, a
1118            // stream of chunks just over half the bound would commit each one
1119            // at about half the bound. Either way `acc` then goes to `out`.
1120            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                        // Uneven record widths pushed `acc` past the bound.
1142                        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    /// A graded body as a chunk: committed, which may spill it, or left
1175    /// resident.
1176    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
1185/// Append every update in `view` whose key matches a probe at or after
1186/// `*probe_index` into `staging`, per the [`UnloadChunk`] consume-index
1187/// protocol: probes strictly below the view's last key are consumed, a probe
1188/// equal to it is extracted but left for the next chunk.
1189fn 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    /// The flat columnar accumulation. Appends are bulk column-range copies,
1239    /// and a group straddling chunks stitches by plain concatenation.
1240    type Staging = <((K, V), T, R) as Columnar>::Container;
1241
1242    /// A borrowed key column, e.g. of a `Column<K>` the consumer assembled
1243    /// from its sorted, deduplicated probe keys.
1244    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        // A data ref is a `(key ref, val ref)` tuple, so the key fences are a
1253        // projection of the data fences.
1254        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                // NOTE: deliberately the non-admitting read. One probe set
1277                // touching a chunk is weak evidence it will be touched again,
1278                // and probing a spilled trace must not accrete it back into
1279                // residency. The cost is a full decode per probe set against
1280                // an evicted chunk.
1281                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
1303/// A batch builder over [`ColumnChunk`] input that delegates to a builder
1304/// over [`ColumnBody`] input, loading each chunk's body as it is pushed.
1305///
1306/// This is the adapter that lets a [`ChunkBatcher`] feed the body-input batch
1307/// builders (and through them the existing spine layouts):
1308/// the batcher's chains carry pool-spillable chunks, and bodies are read back
1309/// copy-out only at the seal, one chunk at a time.
1310///
1311/// [`ChunkBatcher`]: differential_dataflow::trace::chunk::ChunkBatcher
1312pub 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            // A resident body folds for the price of a borrow, so it goes to
1354            // `observe` whatever the builder asked for: the state gets exact
1355            // figures for free. Folding a spilled body costs a second pool
1356            // read on top of the push below, so it waits on `wants_bodies`.
1357            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        // One chunk at a time through `push`, so peak transient memory is a
1368        // single loaded body rather than the whole chain at once.
1369        for chunk in chain.iter_mut() {
1370            builder.push(chunk);
1371        }
1372        chain.clear();
1373        builder.done(description)
1374    }
1375}
1376
1377/// A builder whose batches carry state derived from a whole chain, computed
1378/// before the chain's first push.
1379///
1380/// [`Builder::seal`] receives the chain at once, which a chain of pool-backed
1381/// chunks cannot supply without holding every body resident at the same time.
1382/// An implementor splits the derivation instead: a caller folds the chain into
1383/// [`State`](Self::State) one entry at a time, then builds from it.
1384///
1385/// A caller folds each chain entry exactly once, through
1386/// [`observe`](Self::observe) when it holds the entry's contents and through
1387/// [`observe_records`](Self::observe_records) when it does not, so an
1388/// implementor may accumulate the same figure in both without double
1389/// counting. Which one an entry takes is the caller's choice, so an
1390/// implementor reads no meaning into it beyond what each carries.
1391///
1392/// [`Builder::seal`]: differential_dataflow::trace::Builder::seal
1393pub trait ChainState: differential_dataflow::trace::Builder {
1394    /// State accumulated across a chain.
1395    type State: Default;
1396
1397    /// Whether the state is worth reading a chain entry's contents for when
1398    /// they are not already at hand. A caller that holds them passes them on
1399    /// regardless.
1400    fn wants_bodies() -> bool;
1401
1402    /// Fold one chain entry's contents into `state`.
1403    fn observe(state: &mut Self::State, input: &Self::Input);
1404
1405    /// Fold one chain entry's update count into `state`, for a caller that
1406    /// reads no bodies.
1407    fn observe_records(state: &mut Self::State, records: usize);
1408
1409    /// Build with `state` installed.
1410    fn from_state(state: Self::State) -> Self;
1411}
1412
1413/// A chunker for `arrange_core` over [`ColumnChunk`]s: sorts and consolidates
1414/// raw input columns through a [`ColumnChunker`] and wraps its output chunks.
1415pub 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
1477/// The [`ChunkBatcher`] of a chunk chain, reporting resident bytes to the
1478/// batcher size logger.
1479///
1480/// [`ChunkBatcher`]: differential_dataflow::trace::chunk::ChunkBatcher
1481pub type AccountedChunkBatcher<D, T, R> =
1482    differential_dataflow::trace::implementations::merge_batcher::MergeBatcher<
1483        AccountedChunkMerger<D, T, R>,
1484    >;
1485
1486/// The chunk merger of [`AccountedChunkBatcher`]: differential's merger with
1487/// a resident readied side, plus the [`Merger::allocation`] figures the
1488/// `mz_arrangement_batcher_*_raw` introspection tables are assembled from.
1489///
1490/// A chunk reports the bytes of the body it holds whether that body sits on
1491/// the heap or in the pool. A spilled body is not free: the pool keeps it in
1492/// a slot, or compressed in an extent, or on the device, and only the last
1493/// of those is off the process's books. Which tier holds it is the pool's
1494/// own business and its metrics report that split, so attributing the body
1495/// to the operator that owns it is what keeps these tables whole.
1496///
1497/// [`Merger::allocation`]: differential_dataflow::trace::implementations::merge_batcher::Merger::allocation
1498pub 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        // Differential's merger settles both sides with committing settles.
1544        // The readied side is handed to a builder that loads every chunk
1545        // back, so committing it spills bodies the next call copies out
1546        // again, and charges the pool's budget for them in between. Drive
1547        // the loop here instead and grade that side resident. The kept side
1548        // stays committed: it is the side that lives on, and spilling it is
1549        // the point.
1550        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    //! Property tests for the [`Chunk`] and [`UnloadChunk`] contracts on
1580    //! [`ColumnChunk`].
1581    //!
1582    //! Strategy: generate sorted+consolidated inputs (the chunk invariant),
1583    //! drive the trait methods the way the differential harness does, and
1584    //! compare against brute-force references on owned tuples. Test types are
1585    //! `D = (u64, u64)`, `T = u64`, `R = i64` from small ranges so equal-key
1586    //! collisions are common and consolidation actually runs.
1587
1588    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    /// The delegated codec's stored form is byte-identical to the extent
1603    /// store's previous hard-coded framing: a little-endian `u32`
1604    /// body-length prefix followed by one lz4 block, which is exactly what
1605    /// `compress_prepend_size` produces.
1606    #[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    /// Reference consolidation: sort by `(data, time)`, sum diffs over equal
1627    /// pairs, drop zeros.
1628    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    /// Cut consolidated data into non-empty chunks at the given points.
1692    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    /// The chunked cut, with every chunk force-spilled through a private pool
1709    /// (bounds captured, bodies in the pool) regardless of size thresholds.
1710    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    /// Whether a spilled chunk's body is stored compressed. Deliberately
1718    /// does not retain the body: an `Rc` held across a `survive_merge` would
1719    /// itself make the body shared and suppress the migration under test.
1720    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    /// Spill one chunk through `pool`, bypassing the size threshold and
1728    /// keeping the chunk's depth.
1729    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    /// The batcher size logger reports a body's bytes wherever the body
1735    /// lives, so an operator's share does not vanish when a chunk spills.
1736    #[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    /// A chunk that reaches the seal resident is handed back resident, even
1757    /// with a pool installed. The readied side goes straight to a builder
1758    /// that loads every chunk, so committing it there would serialize a body
1759    /// the next call copies out again.
1760    ///
1761    /// A body the batcher spilled earlier, while merging chains it was
1762    /// holding, passes through as it is: that spill is the one worth paying
1763    /// for, and undoing it would be another copy.
1764    #[mz_ore::test]
1765    #[cfg_attr(miri, ignore)]
1766    fn seal_readies_resident_chunks() {
1767        use differential_dataflow::trace::Batcher;
1768
1769        // One push, so nothing merges before the seal and the chunk reaches
1770        // `extract` resident. Comfortably past the 64 KiB spill floor at 32
1771        // bytes per record, so a committing settle would spill it.
1772        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    /// A single pool shared by every test in the module. A pool reserves a
1788    /// large slab of address space, so one per test (let alone per proptest
1789    /// case) exhausts the VM map under parallel test threads.
1790    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        /// A full batcher round trip: push chunked inputs, seal everything,
1798        /// and compare with the reference consolidation of the union.
1799        #[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            // An empty upper ships everything.
1814            let (sealed, _description) = batcher.seal(Antichain::new());
1815            prop_assert_eq!(collect_chunks(sealed), consolidate(union));
1816        }
1817
1818        /// The same round trip over force-spilled inputs: merge and extract
1819        /// read bodies back from the pool call-scoped.
1820        #[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        /// Sealing at an intermediate upper partitions by time and reports
1840        /// the kept lower envelope as the frontier.
1841        #[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        /// The intermediate-upper partition of `seal_partitions_by_time`, over
1868        /// force-spilled inputs: bodies read back from the pool and split by
1869        /// time in one seal.
1870        #[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        /// `advance` equals per-record time advancement plus reference
1898        /// consolidation, including across a `done = false` carry.
1899        #[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        /// `settle` preserves contents and order, moves everything on `done`,
1926        /// and coalesces small neighbors.
1927        #[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            // Test chunks are far below the byte threshold, so maximal
1938            // packing coalesces everything into a single chunk.
1939            prop_assert!(out.len() <= 1);
1940            prop_assert_eq!(collect_chunks(out), input);
1941        }
1942
1943        /// `ChunkBatch::extract_into` over sorted, deduplicated probe keys
1944        /// equals the reference filter, resident and spilled alike, straddled
1945        /// keys included.
1946        #[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            // The scan path reproduces the batch exactly, resident and
1983            // spilled alike.
1984            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    /// `locate` answers the three-way span comparison for every probe
1991    /// placement: below, within, and past the chunk's keys.
1992    #[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    /// Collect chunk contents while asserting each chunk's serialized size
2012    /// stays within `bound` bytes.
2013    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    /// One record just under the bound among many narrow ones, so a cut at
2113    /// the average width leaves the piece holding the wide record over the
2114    /// bound.
2115    #[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    /// Each side is under the bound, so any oversized output would have to
2148    /// come from `merge_from` itself. One input interleaves record by record,
2149    /// the other has runs long enough to gallop.
2150    #[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    /// The last piece of a split coalesces with a shallower neighbor and
2191    /// commits at the deeper depth.
2192    #[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    /// Chunks of just over half the bound: unless settle moves part of each
2208    /// chunk into the column it is accumulating, every chunk commits at about
2209    /// half a slot.
2210    #[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    /// Advancing a large input cuts the output into several chunks near the
2239    /// ship threshold, and their concatenation is the reference result.
2240    #[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    /// An input that is entirely one `D` group is withheld whole as the
2258    /// carry unless `done`: none of it is provably complete.
2259    #[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        // Sealing the carry advances and consolidates it.
2273        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    /// Chunks the extract frontier does not split pass through whole from
2280    /// their resident time bounds: spilled bodies land on their side still
2281    /// spilled, with no load or re-commit, and a kept chunk's minimal times
2282    /// feed the residual frontier.
2283    #[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        // A frontier between the two chunks' time ranges: the low chunk
2296        // ships whole and the high chunk keeps whole, both still spilled
2297        // (no load, no re-commit), and the residual is the kept chunk's
2298        // minimal time.
2299        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    /// Extracting a large chunk at an intermediate frontier cuts both sides
2325    /// into several chunks and partitions exactly by time.
2326    #[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    /// `locate` answers from resident metadata on spilled chunks and follows
2361    /// the probe-relative-to-span convention.
2362    #[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    /// A body large enough to spill round-trips through the pool with resident
2379    /// metadata intact, and the batcher produces spilled sealed output.
2380    #[mz_ore::test]
2381    #[cfg_attr(miri, ignore)] // too slow
2382    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    /// Merging spilled chains loads bodies call-scoped and consolidates
2411    /// correctly, and an untouched survivor keeps its spilled body.
2412    #[mz_ore::test]
2413    #[cfg_attr(miri, ignore)] // too slow
2414    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    /// A merge whose fronts have disjoint key ranges pushes the untouched
2439    /// survivor back in its original (spilled) form rather than rewriting it.
2440    #[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        // `low` is fully consumed. `high` was never touched and must come
2458        // back spilled.
2459        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    /// Merge output is one generation past its deepest input, a survivor
2470    /// rewritten from its remainder keeps its own depth, and a chunk passed
2471    /// through the disjoint fast path ages by its survival.
2472    #[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        // The merge runs through the shared horizon, so `high` survives with
2486        // its unmerged remainder at its original depth.
2487        assert!(in1.is_empty());
2488        assert_eq!(in2.len(), 1);
2489        assert_eq!(in2[0].depth(), 0, "rewritten survivor keeps its depth");
2490
2491        // A disjoint merge moves the lower front to the output with its data
2492        // unchanged, one generation older for having outlived the merge.
2493        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    /// Advance output and carry keep the deepest input depth, since
2504    /// compaction rewrites within a generation.
2505    #[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    /// Settle commits at the deepest depth among coalesced chunks, and a
2527    /// commit large enough to spill carries the depth into its spilled
2528    /// metadata (and thus into the pool hints).
2529    #[mz_ore::test]
2530    #[cfg_attr(miri, ignore)] // too slow
2531    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    /// The settle carry commits at a monotone size threshold rather than the
2551    /// periodic ship window, so mid-window chunk sizes cannot make it grow
2552    /// past the target unbounded.
2553    #[mz_ore::test]
2554    #[cfg_attr(miri, ignore)] // too slow
2555    fn settle_carry_commits_at_target() {
2556        // Two inputs fit in one slot. A third must start another chunk.
2557        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        // Catches the fixture drifting above `at_commit_size`, where settle
2569        // commits each chunk as-is and the size cap below holds vacuously.
2570        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    /// The smallest column whose serialized size reaches `SPILL_MIN_BYTES`.
2594    /// One record less sits under the spill floor.
2595    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    /// Bodies straddling the spill floor: one record under stays resident,
2606    /// at the floor spills.
2607    #[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    /// The compression depth floor picks the codec, not whether a body
2622    /// spills: shallow generations store at identity, the floor and deeper
2623    /// at lz4, and every depth spills and round-trips.
2624    #[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        // Codec identity via Debug: ZST statics and dyn vtables make
2629        // pointer comparison unreliable. The flag must agree with the codec,
2630        // since it is what decides whether a body wants migrating.
2631        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        // The default floor stores only fresh (depth 0) bodies at identity.
2654        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    /// A chunk a merge carries forward untouched ages a generation, and a
2661    /// spilled body crossing the compression floor by doing so is re-spilled
2662    /// under the compressing codec. Key-disjoint input takes that path on
2663    /// every merge, so without the crossing its backlog would stay
2664    /// identity-coded for as long as it lived.
2665    #[mz_ore::test]
2666    #[cfg_attr(miri, ignore)] // too slow
2667    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        // Fresh spilled chunks sit below the floor, so both store identity
2676        // coded, and their data ranges are disjoint.
2677        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        // Past the floor the next survival is a metadata-only bump: the body
2700        // is already compressed and stays where it is.
2701        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    /// Aging does not depend on holding the only reference to a body. The
2719    /// trace's compaction merger feeds `merge` clones of a source batch's
2720    /// chunks and keeps the batch alive throughout, so a shared body must
2721    /// still age. It must not be re-spilled: the other holder goes on
2722    /// storing the original whatever this reference does, and the merger
2723    /// rewrites its clone immediately.
2724    #[mz_ore::test]
2725    #[cfg_attr(miri, ignore)] // too slow
2726    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        // The source batch's chunk, held for the whole merge as the spine
2734        // holds it.
2735        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        // Past the floor, where no re-spill is in question, a shared body
2758        // goes on aging rather than pinning at the crossing depth.
2759        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    /// A migration that cannot happen when a body first qualifies is retried
2772    /// at the next survival, never consumed. Each case leaves an
2773    /// identity-coded body at or past the floor, which would be stranded
2774    /// uncompressed for the rest of its life if the test were a depth
2775    /// transition rather than the body's stored codec.
2776    #[mz_ore::test]
2777    #[cfg_attr(miri, ignore)] // too slow
2778    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        // Age a chunk one generation through a disjoint merge, which passes
2783        // the lower front through `survive_merge`.
2784        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        // No pool installed when the body qualifies: spilling can be toggled
2793        // off at runtime while existing handles stay valid.
2794        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        // Shared when the body qualifies: the compaction merger holds the
2810        // source batch while merging clones of its chunks.
2811        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        // The floor lowered long after the body spilled, which is what an
2823        // operator reaches for under pool pressure. Nothing here is a
2824        // transition: the body is already several generations past the new
2825        // floor when it moves.
2826        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    /// The compute and storage spill gates compose as an OR: either gate
2842    /// routes commits to the installed pool, and each setter writes only its
2843    /// own gate. The sink gate is independent of both, in both directions.
2844    ///
2845    /// One test, because the gates are process-global.
2846    #[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        // A body at the spill floor, so the gates alone decide.
2858        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    /// Re-spilling an already-serialized body exercises the `ColumnBody::Words`
2888    /// branch of `spill_serialized` and round-trips byte-identically.
2889    #[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}