Skip to main content

mz_compute/sink/
correction_v2.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//! An implementation of the `Correction` data structure used by the MV sink's `write_batches`
11//! operator to stash updates before they are written.
12//!
13//! The `Correction` data structure provides methods to:
14//!  * insert new updates
15//!  * advance the compaction frontier (called `since`)
16//!  * obtain an iterator over consolidated updates before some `upper`
17//!  * force consolidation of updates before some `upper`
18//!
19//! The goal is to provide good performance for each of these operations, even in the presence of
20//! future updates. MVs downstream of temporal filters might have to deal with large amounts of
21//! retractions for future times and we want those to be handled efficiently as well.
22//!
23//! Note that `Correction` does not provide a method to directly remove updates. Instead updates
24//! are removed by inserting their retractions so that they consolidate away to nothing.
25//!
26//! ## Storage of Updates
27//!
28//! Stored updates are of the form `(data, time, diff)`, where `time` and `diff` are fixed to
29//! [`mz_repr::Timestamp`] and [`mz_repr::Diff`], respectively.
30//!
31//! [`CorrectionV2`] holds onto a list of `Chain`s containing `Chunk`s of stashed updates. Each
32//! `Chunk` is a columnation region containing a fixed maximum number of updates. All updates in
33//! a chunk, and all updates in a chain, are ordered by (time, data) and consolidated.
34//!
35//! Chains live in three places:
36//!
37//!  * A [`BucketChain`] partitions times at or beyond the `boundary` (the largest read `upper`
38//!    seen so far) into buckets of exponentially growing time ranges, each holding a list of
39//!    chains. Reads only touch the buckets below their `upper`, so the bulk of the buffered
40//!    updates — in particular far-future retractions produced by temporal filters — is left
41//!    alone.
42//!  * `pending_low` holds chains at times below the `boundary`, mostly insertions arriving
43//!    through the persist feedback.
44//!  * `emitted` is a single chain holding the updates returned by the last read. Updates must
45//!    stay in the buffer until their feedback retractions arrive, and keeping them separate from
46//!    the bucket chain means reads never have to re-merge future updates.
47//!
48//! ```text
49//!       chain[0]   |   chain[1]   |   chain[2]
50//!                  |              |
51//!     chunk[0]     | chunk[0]     | chunk[0]
52//!       (a, 1, +1) |   (a, 1, +1) |   (d, 3, +1)
53//!       (b, 1, +1) |   (b, 2, -1) |   (d, 4, -1)
54//!     chunk[1]     | chunk[1]     |
55//!       (c, 1, +1) |   (c, 2, -2) |
56//!       (a, 2, -1) |   (c, 4, -1) |
57//!     chunk[2]     |              |
58//!       (b, 2, +1) |              |
59//!       (c, 2, +1) |              |
60//!     chunk[3]     |              |
61//!       (b, 3, -1) |              |
62//!       (c, 3, +1) |              |
63//! ```
64//!
65//! The "chain invariant" states that each chain in a bucket has at least `chain_proportionality` times as
66//! many updates as the next one. This means that chain sizes will often be powers of
67//! `chain_proportionality`, but they don't have to be. For example, for a proportionality of 2,
68//! the chain sizes `[11, 5, 2, 1]` would satisfy the chain invariant.
69//!
70//! Note that the invariant is maintained on update counts, not chunk counts. Chunks are
71//! byte-bounded (see `ChunkBuilder`), so chunk count is not proportional to update count and
72//! would be a poor proxy: any chain below the chunk byte boundary is a single chunk regardless
73//! of how many updates it holds, which would let the geometric invariant collapse and break the
74//! O(log N) amortization of inserts.
75//!
76//! Choosing the `chain_proportionality` value allows tuning the trade-off between memory and CPU
77//! resources required to maintain corrections. A higher proportionality forces more frequent chain
78//! merges, and therefore consolidation, reducing memory usage but increasing CPU usage.
79//!
80//! ## Inserting Updates
81//!
82//! A batch of updates is routed by time: updates below the `boundary` become a `pending_low`
83//! chain, the rest is appended as new chains to their respective buckets. Appending to a bucket
84//! merges chains until the chain invariant is restored.
85//!
86//! Inserting an update into the correction buffer can be expensive: It involves allocating a new
87//! chunk, copying the update in, and then likely merging with an existing chain to restore the
88//! chain invariant. If updates trickle in in small batches, this can cause a considerable
89//! overhead. To amortize this overhead, new updates aren't immediately inserted into the sorted
90//! chains but instead stored in a `Stage` buffer. Once the staged updates reach the configured
91//! byte size, they are sorted and routed.
92//!
93//! The insert operation has an amortized complexity of O(log N), with N being the current number
94//! of updates stored.
95//!
96//! ## Retrieving Consolidated Updates
97//!
98//! Retrieving consolidated updates before a given `upper` works by peeling all buckets below the
99//! `upper` off the bucket chain, splitting their chains, the pending low chains, and the previous
100//! `emitted` chain at the `upper`, merging the parts below the `upper` into the new `emitted`
101//! chain, and returning an iterator over that chain.
102//!
103//! Because each chain contains updates ordered by time first, splitting a chain at the `upper`
104//! reuses whole chunks and copies at most one chunk straddling the split point. Updates at times
105//! at or beyond the `upper` are never touched, no matter how many the buffer holds. The
106//! complexity of a read is O(U log K), with U being the number of updates before `upper` and K
107//! the number of chains containing them.
108//!
109//! ## Merging Chains
110//!
111//! Merging multiple chains into a single chain is done using a k-way merge. As the input chains
112//! are sorted by (time, data) and consolidated, the same properties hold for the output chain. The
113//! complexity of a merge of K chains containing N updates is O(N log K).
114//!
115//! There is a twist though: Merging also has to respect the `since` frontier, which determines how
116//! far the times of updates should be advanced. Advancing times in a sorted chain of updates
117//! can make them become unsorted, so we cannot just merge the chains from top to bottom.
118//!
119//! For example, consider these two chains, assuming `since = [2]`:
120//!   chain 1: [(c, 1, +1), (b, 2, -1), (a, 3, -1)]
121//!   chain 2: [(b, 1, +1), (a, 2, +1), (c, 2, -1)]
122//! After time advancement, the chains look like this:
123//!   chain 1: [(c, 2, +1), (b, 2, -1), (a, 3, -1)]
124//!   chain 2: [(b, 2, +1), (a, 2, +1), (c, 2, -1)]
125//! Merging them naively yields [(b, 2, +1), (a, 2, +1), (b, 2, -1), (a, 3, -1)], a chain that's
126//! neither sorted nor consolidated.
127//!
128//! Times below the `since` can only exist in chains read by `consolidate_before`, and only if
129//! the `since` advanced past buffered times since the previous read. For few distinct stale
130//! times — the steady state, where the previously emitted chain was written just before the
131//! since advanced past it — we merge sub-chains, one for each distinct time that's before or at
132//! the `since`. Each of these sub-chains retains the (time, data) ordering after the time
133//! advancement to `since`, so merging those yields the expected result.
134//!
135//! For the above example, the chains we would merge are:
136//!   chain 1.a: [(c, 2, +1)]
137//!   chain 1.b: [(b, 2, -1), (a, 3, -1)]
138//!   chain 2.a: [(b, 2, +1)],
139//!   chain 2.b: [(a, 2, +1), (c, 2, -1)]
140//!
141//! For many distinct stale times — e.g. a since jump across many buffered timestamps when a sink
142//! restarts with an old as-of — the number of sub-chains grows with the number of distinct times,
143//! so we instead materialize the affected updates, advance their times, and sort and consolidate
144//! them in one O(U log U) pass.
145
146use std::cmp::{Ordering, Reverse};
147use std::collections::{BinaryHeap, VecDeque};
148use std::fmt;
149use std::rc::Rc;
150use std::sync::atomic::{self, AtomicUsize};
151use std::sync::{Arc, Mutex, OnceLock};
152
153use columnar::bytes::indexed;
154use columnar::{Columnar, Index, Len, Ref};
155use itertools::Itertools;
156use mz_ore::cast::CastLossy;
157use mz_ore::pool::ChunkHandle;
158use mz_ore::soft_assert_or_log;
159use mz_persist_client::metrics::{SinkMetrics, SinkWorkerMetrics, UpdateDelta};
160use mz_repr::{Diff, Row, Timestamp};
161use mz_timely_util::columnar::body::ColumnBody;
162use mz_timely_util::columnar::chunk;
163use mz_timely_util::temporal::{Bucket, BucketChain};
164use timely::PartialOrder;
165use timely::container::PushInto;
166use timely::progress::Antichain;
167
168use crate::sink::correction::{ChannelLogging, SizeMetrics};
169
170/// Convenient alias for use in data trait bounds.
171///
172/// `D` is constrained to be `Columnar`, so that updates can be stored in a single columnar
173/// region per chunk, and the variable-length payload (e.g. `Row` bytes) lives in the same
174/// allocation as the rest of the chunk. [`DataContainer`] carries the bounds on that container.
175pub trait Data:
176    differential_dataflow::Data + Columnar<Container: DataContainer> + DataBytes + Send + Sync
177{
178}
179impl<D> Data for D where
180    D: differential_dataflow::Data + Columnar<Container: DataContainer> + DataBytes + Send + Sync
181{
182}
183
184/// The bytes a datum occupies, so the staging area can be sized in bytes.
185///
186/// Counts `size_of::<Self>()` plus what the datum owns on the heap. An estimate suffices: the
187/// figure paces when staged updates ship and feeds the buffer's size metrics, nothing else.
188pub trait DataBytes {
189    /// The bytes this datum occupies, heap included.
190    fn data_bytes(&self) -> usize;
191}
192
193impl DataBytes for Row {
194    fn data_bytes(&self) -> usize {
195        self.byte_len()
196    }
197}
198
199/// The bounds [`Data`] places on its columnar container.
200///
201/// The `Ref`-level `Eq + Ord` bounds let the merge/heap code compare updates directly through
202/// the columnar borrow, avoiding `into_owned` clones on the hot path. The `Borrowed`-level
203/// `Send` bound lets a hoisted `Chunk::view` travel with the iterators that
204/// [`CorrectionV2::updates_before`] hands across the persist writer's `await`.
205pub trait DataContainer:
206    Send + Sync + Clone + for<'a> columnar::Borrow<Ref<'a>: Eq + Ord, Borrowed<'a>: Send>
207{
208}
209impl<C> DataContainer for C where
210    C: Send + Sync + Clone + for<'a> columnar::Borrow<Ref<'a>: Eq + Ord, Borrowed<'a>: Send>
211{
212}
213
214/// A borrowed view over a [`Chunk`]'s column.
215///
216/// Obtained from [`Chunk::view`] and indexed with `get`.
217type ChunkView<'a, D> =
218    <<(D, Timestamp, Diff) as Columnar>::Container as columnar::Borrow>::Borrowed<'a>;
219
220/// A data structure used to store corrections in the MV sink implementation.
221///
222/// In contrast to `CorrectionV1`, this implementation stores updates in columnation regions,
223/// allowing their memory to be transparently spilled to disk.
224#[derive(Debug)]
225pub struct CorrectionV2<D: Data> {
226    /// Bucketed storage for updates at times at or beyond `boundary`.
227    ///
228    /// Buckets cover exponentially growing time ranges, so reads only touch the buckets below
229    /// their `upper`, and far-future updates (e.g. retractions produced by temporal filters) are
230    /// rarely touched.
231    chain: BucketChain<ChainBucket<D>>,
232    /// Chains at times below `boundary` that were not yet emitted.
233    ///
234    /// Filled by inserts at times below the boundary (mostly persist feedback) and by the
235    /// remainders of `emitted` when a read uses a smaller `upper` than the previous one. Merged
236    /// into `emitted` by the next read.
237    pending_low: Vec<Chain<D>>,
238    /// Updates that were emitted by `updates_before` but not yet cancelled by persist feedback.
239    ///
240    /// Sorted and consolidated, with all times advanced to the `since`.
241    emitted: Chain<D>,
242    /// A staging area for updates, to speed up small inserts.
243    stage: Stage<D>,
244    /// The lower bound of times stored in `chain`. Only ever advances.
245    ///
246    /// Times below the boundary have been peeled off the bucket chain and can only be stored in
247    /// `pending_low` or `emitted`.
248    boundary: Antichain<Timestamp>,
249    /// The frontier by which all contained times are advanced.
250    since: Antichain<Timestamp>,
251
252    /// Total count of updates last reported to metrics.
253    ///
254    /// Tracked to compute deltas in `update_metrics`.
255    prev_update_count: usize,
256    /// Total size last reported to metrics.
257    ///
258    /// Tracked to compute deltas in `update_metrics`.
259    prev_size: SizeMetrics,
260    /// Global persist sink metrics.
261    metrics: SinkMetrics,
262    /// Per-worker persist sink metrics.
263    worker_metrics: SinkWorkerMetrics,
264    /// Running totals and introspection logging.
265    accounting: Accounting,
266}
267
268/// Fuel for restoring the bucket chain invariant after peeling.
269///
270/// Bounds the restoration work per buffer operation. The bucket chain remains functional when
271/// restoration is incomplete -- peeling and finding work on ill-formed chains, at the cost of
272/// more in-line splitting -- so leftover restoration is simply picked up by the next operation.
273///
274/// `restore` spends one unit of fuel per bucket split, and a single `peel` leaves at most
275/// `BucketTimestamp::DOMAIN` (64) buckets to re-split, so this budget completes restoration in one
276/// call for any realistic buffer. It is deliberately generous: the "incomplete restoration is
277/// picked up next op" path is a correctness safety net for pathological bucket counts, not a hot
278/// path we expect to exercise. Lower it if restoration ever needs to interleave with other work.
279const RESTORE_FUEL: i64 = 1_000_000;
280
281impl<D: Data> CorrectionV2<D> {
282    /// Construct a new [`CorrectionV2`] instance.
283    pub fn new(
284        metrics: SinkMetrics,
285        worker_metrics: SinkWorkerMetrics,
286        logging: Option<ChannelLogging>,
287        chain_proportionality: f64,
288        chunk_size: usize,
289    ) -> Self {
290        let accounting = Accounting::new(logging);
291
292        Self {
293            chain: BucketChain::new(ChainBucket::new(chain_proportionality, accounting.clone())),
294            pending_low: Vec::new(),
295            emitted: Chain::new(),
296            stage: Stage::new(accounting.clone(), chunk_size),
297            boundary: Antichain::from_elem(Timestamp::MIN),
298            since: Antichain::from_elem(Timestamp::MIN),
299            prev_update_count: 0,
300            prev_size: Default::default(),
301            metrics,
302            worker_metrics,
303            accounting,
304        }
305    }
306
307    /// Insert a batch of updates.
308    pub fn insert(&mut self, updates: &mut Vec<(D, Timestamp, Diff)>) {
309        let Some(since_ts) = self.since.as_option() else {
310            // If the since is the empty frontier, discard all updates.
311            updates.clear();
312            return;
313        };
314
315        for (_, time, _) in &mut *updates {
316            *time = std::cmp::max(*time, *since_ts);
317        }
318
319        self.insert_inner(updates);
320    }
321
322    /// Insert a batch of updates, after negating their diffs.
323    pub fn insert_negated(&mut self, updates: &mut Vec<(D, Timestamp, Diff)>) {
324        let Some(since_ts) = self.since.as_option() else {
325            // If the since is the empty frontier, discard all updates.
326            updates.clear();
327            return;
328        };
329
330        for (_, time, diff) in &mut *updates {
331            *time = std::cmp::max(*time, *since_ts);
332            *diff = -*diff;
333        }
334
335        self.insert_inner(updates);
336    }
337
338    /// Insert a batch of updates into the stage, flushing it when full.
339    ///
340    /// All times are expected to be >= the `since`.
341    fn insert_inner(&mut self, updates: &mut Vec<(D, Timestamp, Diff)>) {
342        debug_assert!(updates.iter().all(|(_, t, _)| self.since.less_equal(t)));
343
344        if let Some(mut ready) = self.stage.insert(updates) {
345            self.route(&mut ready);
346        }
347
348        self.update_metrics();
349    }
350
351    /// Route a batch of sorted, consolidated updates to `pending_low` or their chain buckets.
352    fn route(&mut self, updates: &mut Vec<(D, Timestamp, Diff)>) {
353        // Updates at times below the boundary become a pending low chain.
354        let idx = updates.partition_point(|(_, t, _)| !self.boundary.less_equal(t));
355        if idx > 0 {
356            let mut builder = ChainBuilder::default();
357            builder.extend(updates.drain(..idx));
358            let chain = builder.finish();
359            if !chain.is_empty() {
360                self.account_chain_created(&chain);
361                self.pending_low.push(chain);
362            }
363        }
364
365        // Updates at times at or beyond the boundary go into their chain buckets. Walk ranges of
366        // times that fall into the same bucket, to push batches of updates at once.
367        let mut drain = updates.drain(..).peekable();
368        while let Some(update) = drain.next() {
369            let time = update.1;
370            let range = self
371                .chain
372                .range_of(&time)
373                .expect("bucket chain covers all times at or beyond the boundary");
374            let mut builder = ChainBuilder::default();
375            builder.extend(std::iter::once(update));
376            while let Some(update) = drain.next_if(|(_, t, _)| range.contains(t)) {
377                builder.extend(std::iter::once(update));
378            }
379            let bucket = self
380                .chain
381                .find_mut(&range.start)
382                .expect("bucket chain covers all times at or beyond the boundary");
383            bucket.push_chain(builder.finish());
384        }
385    }
386
387    /// Return consolidated updates before the given `upper`.
388    pub fn updates_before<'a>(
389        &'a mut self,
390        upper: &Antichain<Timestamp>,
391    ) -> impl Iterator<Item = (D, Timestamp, Diff)> + Send + 'a {
392        self.consolidate_before(upper);
393        self.consolidated_updates_before(upper)
394    }
395
396    /// Return the updates before the given `upper`, as consolidated by a preceding
397    /// [`CorrectionV2::consolidate_before`] call.
398    ///
399    /// The caller must have invoked `consolidate_before` with the same `upper` and must not have
400    /// mutated the buffer since. Otherwise the returned updates are neither consolidated nor
401    /// necessarily complete.
402    pub fn consolidated_updates_before<'a>(
403        &'a self,
404        upper: &Antichain<Timestamp>,
405    ) -> impl Iterator<Item = (D, Timestamp, Diff)> + Send + use<'a, D> {
406        // All contained times are advanced to at least the `since`, so a read at an `upper` that
407        // is not beyond the `since` is always empty. This mirrors the short-circuit in
408        // `consolidate_before`, which leaves `emitted` untouched in that case.
409        if !PartialOrder::less_than(&self.since, upper) {
410            return None.into_iter().flatten();
411        }
412
413        // After `consolidate_before`, `emitted` holds exactly the updates before `upper`: every
414        // path that populates it splits at `upper` (pushing the remainder to `pending_low`), and
415        // the guard above guarantees `upper > since`, so advancing stale times to the `since`
416        // cannot lift them to or beyond `upper`. We can therefore yield all of `emitted`. Guard
417        // the invariant: a violation would write updates beyond the batch upper to persist.
418        soft_assert_or_log!(
419            self.emitted
420                .chunks
421                .last()
422                .is_none_or(|c| !upper.less_equal(&c.last_time())),
423            "emitted contains times at or beyond the upper",
424        );
425        Some(self.emitted.iter()).into_iter().flatten()
426    }
427
428    /// Consolidate all updates before the given `upper` into the `emitted` chain.
429    ///
430    /// Once this method returns, `emitted` contains all updates at times before `upper`,
431    /// consolidated.
432    ///
433    /// Does nothing if `upper` is not beyond the `since`: all contained times are advanced to at
434    /// least the `since`, so such a read is empty anyway, and skipping avoids an eager peel,
435    /// merge, and `boundary` advancement. Normal reads and `consolidate_at_since` always pass an
436    /// `upper` beyond the `since`.
437    pub fn consolidate_before(&mut self, upper: &Antichain<Timestamp>) {
438        if !PartialOrder::less_than(&self.since, upper) {
439            return;
440        }
441
442        if let Some(mut ready) = self.stage.flush() {
443            self.route(&mut ready);
444        }
445
446        let Some(&since_ts) = self.since.as_option() else {
447            // If the since is the empty frontier, discard all updates.
448            let peeled = self.chain.peel(Antichain::new().borrow());
449            for bucket in peeled {
450                for chain in bucket.into_chains() {
451                    self.account_chain_dropped(&chain);
452                }
453            }
454            for chain in std::mem::take(&mut self.pending_low) {
455                self.account_chain_dropped(&chain);
456            }
457            let emitted = std::mem::replace(&mut self.emitted, Chain::new());
458            if !emitted.is_empty() {
459                self.account_chain_dropped(&emitted);
460            }
461            self.update_metrics();
462            return;
463        };
464
465        // Peel the buckets below the upper off the bucket chain. Bucket splits during the peel
466        // only touch chunks around the upper; chunks wholly on either side are reused.
467        let peeled = self.chain.peel(upper.borrow());
468        if PartialOrder::less_than(&self.boundary, upper) {
469            self.boundary = upper.clone();
470        }
471
472        // Collect candidate chains: peeled bucket contents, pending low chains, and the previous
473        // emitted chain. All contain only times below the boundary.
474        let emitted = std::mem::replace(&mut self.emitted, Chain::new());
475        let mut candidates: Vec<Chain<D>> = Vec::new();
476        for bucket in peeled {
477            candidates.extend(bucket.into_chains());
478        }
479        candidates.append(&mut self.pending_low);
480        if !emitted.is_empty() {
481            candidates.push(emitted);
482        }
483
484        if candidates.is_empty() {
485            self.restore_chain();
486            self.update_metrics();
487            return;
488        }
489
490        candidates
491            .iter()
492            .for_each(|c| self.account_chain_dropped(c));
493
494        // Split the candidates at the upper. Parts at or beyond the upper (possible when `upper`
495        // regresses below a previous one) stay pending.
496        let mut lowers = Vec::new();
497        for chain in candidates {
498            match upper.as_option() {
499                Some(&upper_ts) => {
500                    let (lower, remainder) = chain.split_at_time(upper_ts);
501                    if !lower.is_empty() {
502                        lowers.push(lower);
503                    }
504                    if !remainder.is_empty() {
505                        self.account_chain_created(&remainder);
506                        self.pending_low.push(remainder);
507                    }
508                }
509                // The empty upper is greater than all times.
510                None => lowers.push(chain),
511            }
512        }
513
514        // Merge the lower parts into the new emitted chain, advancing times below the since.
515        // Advancing times in a (time, data)-sorted chain can break its sort order, so chains
516        // containing stale times cannot be merged as they are. Stale times are expected in steady
517        // state: the previous emitted chain was written before the since advanced past it.
518        //
519        // Count the distinct stale times, up to a small cap. For few distinct stale times -- the
520        // steady state -- split cursors into runs that remain sorted under advancement and merge
521        // those. For many distinct stale times -- e.g. a since jump across many buffered
522        // timestamps when a sink restarts with an old as-of -- the number of runs and the cost of
523        // cloning cursor state per run grow with the number of distinct times, so materialize,
524        // advance, and consolidate in one O(U log U) pass instead.
525        const MAX_STALE_RUNS: usize = 32;
526        let mut stale_times = 0;
527        for chain in &lowers {
528            stale_times += chain.distinct_times_before(since_ts, MAX_STALE_RUNS - stale_times);
529            if stale_times >= MAX_STALE_RUNS {
530                break;
531            }
532        }
533
534        // The merged chain becomes `emitted`, which the caller reads next and the next
535        // consolidation merges with the feedback retractions, so the chunks the merge writes are
536        // the youngest generation whatever the depth of their inputs. A lone input chain is
537        // reused as is and keeps its chunks' depths, since re-spilling them would cost a copy.
538        let merged = if stale_times == 0 {
539            let cursors: Vec<_> = lowers.into_iter().filter_map(Chain::into_cursor).collect();
540            merge_cursors(cursors, 0)
541        } else if stale_times < MAX_STALE_RUNS {
542            let mut runs = Vec::new();
543            for chain in lowers {
544                if let Some(cursor) = chain.into_cursor() {
545                    runs.append(&mut cursor.advance_by(since_ts));
546                }
547            }
548            merge_cursors(runs, 0)
549        } else {
550            let mut updates: Vec<_> = lowers.iter().flat_map(|c| c.iter()).collect();
551            for (_, time, _) in &mut updates {
552                *time = std::cmp::max(*time, since_ts);
553            }
554            consolidate(&mut updates);
555            let mut builder = ChainBuilder::default();
556            builder.extend(updates);
557            let chain = builder.finish();
558
559            // Advancement can move updates to or beyond the upper; such updates stay pending.
560            match upper.as_option() {
561                Some(&upper_ts) => {
562                    let (lower, remainder) = chain.split_at_time(upper_ts);
563                    if !remainder.is_empty() {
564                        self.account_chain_created(&remainder);
565                        self.pending_low.push(remainder);
566                    }
567                    lower
568                }
569                None => chain,
570            }
571        };
572
573        if !merged.is_empty() {
574            self.account_chain_created(&merged);
575        }
576        self.emitted = merged;
577
578        self.restore_chain();
579        self.update_metrics();
580    }
581
582    /// Perform a bounded amount of work towards restoring the bucket chain invariant.
583    ///
584    /// Restoration is allowed to remain incomplete: the bucket chain supports peeling and finding
585    /// on ill-formed chains, so any leftover work is picked up by subsequent operations. The fuel
586    /// bound keeps individual buffer operations from stalling the operator that owns the buffer.
587    fn restore_chain(&mut self) {
588        let mut fuel = RESTORE_FUEL;
589        self.chain.restore(&mut fuel);
590    }
591
592    /// Advance the since frontier.
593    ///
594    /// Time advancement of updates in the bucket chain is lazy: it happens when the updates are
595    /// consolidated by a read.
596    ///
597    /// # Panics
598    ///
599    /// Panics if the given `since` is less than the current since frontier.
600    pub fn advance_since(&mut self, since: Antichain<Timestamp>) {
601        assert!(PartialOrder::less_equal(&self.since, &since));
602        self.stage.advance_times(&since);
603        self.since = since;
604    }
605
606    /// Consolidate all updates at the current `since`.
607    pub fn consolidate_at_since(&mut self) {
608        let upper_ts = self.since.as_option().and_then(|t| t.try_step_forward());
609        if let Some(upper_ts) = upper_ts {
610            let upper = Antichain::from_elem(upper_ts);
611            self.consolidate_before(&upper);
612        }
613    }
614
615    fn account_chain_created(&self, chain: &Chain<D>) {
616        self.accounting.chain_created(chain);
617    }
618
619    fn account_chain_dropped(&self, chain: &Chain<D>) {
620        self.accounting.chain_dropped(chain);
621    }
622
623    /// Update persist sink metrics.
624    ///
625    /// Reads the running totals maintained by [`Accounting`], so its cost is independent of how
626    /// much the buffer holds. Nothing here walks chains or pages a chunk in.
627    fn update_metrics(&mut self) {
628        let (new_length, new_size) = self.accounting.totals();
629        self.update_metrics_inner(new_size, new_length);
630    }
631
632    /// Update persist sink metrics to the given new size and length.
633    fn update_metrics_inner(&mut self, new_size: SizeMetrics, new_length: usize) {
634        let old_size = self.prev_size;
635        let old_length = self.prev_update_count;
636        let len_delta = UpdateDelta::new(new_length, old_length);
637        let cap_delta = UpdateDelta::new(new_size.capacity, old_size.capacity);
638        self.metrics
639            .report_correction_update_deltas(len_delta, cap_delta);
640        self.worker_metrics
641            .report_correction_update_totals(new_length, new_size.capacity);
642
643        self.accounting.report_size_metrics(new_size, old_size);
644
645        self.prev_size = new_size;
646        self.prev_update_count = new_length;
647    }
648}
649
650/// Merge the given cursors into one chain, writing new chunks at `depth`.
651fn merge_cursors<D: Data>(cursors: Vec<Cursor<D>>, depth: u8) -> Chain<D> {
652    match cursors.len() {
653        0 => Chain::new(),
654        1 => {
655            let [cur] = cursors.try_into().unwrap();
656            cur.into_chain(depth)
657        }
658        2 => {
659            let [a, b] = cursors.try_into().unwrap();
660            merge_2(a, b, depth)
661        }
662        _ => merge_many(cursors, depth),
663    }
664}
665
666/// Merge the given two cursors using a 2-way merge.
667///
668/// This function is a specialization of `merge_many` that avoids the overhead of a binary heap.
669fn merge_2<D: Data>(cursor1: Cursor<D>, cursor2: Cursor<D>, depth: u8) -> Chain<D> {
670    let mut rest1 = Some(cursor1);
671    let mut rest2 = Some(cursor2);
672    let mut merged = ChainBuilder::at_depth(depth);
673
674    // One borrow per chunk pair, not per update: `Chunk::view` re-decodes the column header on
675    // every call. The inner loop runs until either cursor crosses into its next chunk, at which
676    // point the outer loop re-borrows both.
677    while rest1.is_some() && rest2.is_some() {
678        let chunk1 = rest1.as_ref().expect("checked above").chunk_handle();
679        let chunk2 = rest2.as_ref().expect("checked above").chunk_handle();
680        let view1 = chunk1.view();
681        let view2 = chunk2.view();
682
683        loop {
684            let (Some(c1), Some(c2)) = (rest1.as_ref(), rest2.as_ref()) else {
685                break;
686            };
687            if !c1.reads_from(&chunk1) || !c2.reads_from(&chunk2) {
688                break;
689            }
690
691            let (d1, t1, r1) = c1.get_with(&view1);
692            let (d2, t2, r2) = c2.get_with(&view2);
693
694            match refs_cmp::<D>((t1, d1), (t2, d2)) {
695                Ordering::Less => {
696                    merged.push_ref((d1, t1, r1));
697                    rest1 = rest1.take().expect("checked above").step();
698                }
699                Ordering::Greater => {
700                    merged.push_ref((d2, t2, r2));
701                    rest2 = rest2.take().expect("checked above").step();
702                }
703                Ordering::Equal => {
704                    let r = r1 + r2;
705                    if r != Diff::ZERO {
706                        merged.push_ref((d1, t1, r));
707                    }
708                    rest1 = rest1.take().expect("checked above").step();
709                    rest2 = rest2.take().expect("checked above").step();
710                }
711            }
712        }
713    }
714
715    match (rest1, rest2) {
716        (Some(c), None) | (None, Some(c)) => merged.push_cursor(c),
717        (Some(_), Some(_)) => unreachable!("loop runs while both cursors are live"),
718        (None, None) => (),
719    }
720
721    merged.finish()
722}
723
724/// Merge the given cursors using a k-way merge with a binary heap.
725fn merge_many<D: Data>(cursors: Vec<Cursor<D>>, depth: u8) -> Chain<D> {
726    let mut cursors: Vec<Option<Cursor<D>>> = cursors.into_iter().map(Some).collect();
727    let mut merged = ChainBuilder::at_depth(depth);
728
729    // One borrow per chunk, not per update, as in `merge_2`. Each round borrows the current chunk
730    // of every live cursor, and merges until a cursor crosses into its next chunk, whose updates
731    // the round's heap cannot hold. The next round re-borrows.
732    while cursors.iter().any(Option::is_some) {
733        let chunks: Vec<Option<Rc<Chunk<D>>>> = cursors
734            .iter()
735            .map(|c| c.as_ref().map(Cursor::chunk_handle))
736            .collect();
737        let views: Vec<Option<ChunkView<'_, D>>> = chunks
738            .iter()
739            .map(|c| c.as_deref().map(Chunk::view))
740            .collect();
741
742        // Keyed by `(time, data)`, the merge order, with the cursor index as a tie breaker.
743        let mut heap = BinaryHeap::new();
744        for (i, (cursor, view)) in cursors.iter().zip_eq(&views).enumerate() {
745            if let (Some(cursor), Some(view)) = (cursor, view) {
746                let (d, t, _) = cursor.get_with(view);
747                heap.push(Reverse((t, d, i)));
748            }
749        }
750
751        while let Some(Reverse((time, data, first))) = heap.pop() {
752            let mut diff = Diff::ZERO;
753            let mut crossed = false;
754            let mut next = Some(first);
755            while let Some(i) = next {
756                let view = views[i].as_ref().expect("live cursors have a view");
757                let cursor = cursors[i].take().expect("heap entries are live cursors");
758                diff += cursor.get_with(view).2;
759                cursors[i] = cursor.step();
760                if let Some(cursor) = &cursors[i] {
761                    let chunk = chunks[i].as_ref().expect("live cursors have a chunk");
762                    if cursor.reads_from(chunk) {
763                        let (d, t, _) = cursor.get_with(view);
764                        heap.push(Reverse((t, d, i)));
765                    } else {
766                        crossed = true;
767                    }
768                }
769
770                // Every cursor at the same update contributes to it before it is pushed, so a
771                // round only ends between distinct updates.
772                next = match heap.peek() {
773                    Some(Reverse((t, d, _))) if *t == time && *d == data => {
774                        heap.pop().map(|Reverse((_, _, j))| j)
775                    }
776                    _ => None,
777                };
778            }
779
780            if diff != Diff::ZERO {
781                merged.push_ref((data, time, diff));
782            }
783            if crossed {
784                break;
785            }
786        }
787    }
788
789    merged.finish()
790}
791
792impl<D: Data> Drop for CorrectionV2<D> {
793    fn drop(&mut self) {
794        for bucket in self.chain.buckets() {
795            bucket
796                .chains
797                .iter()
798                .for_each(|c| self.account_chain_dropped(c));
799        }
800        self.pending_low
801            .iter()
802            .for_each(|c| self.account_chain_dropped(c));
803        if !self.emitted.is_empty() {
804            self.account_chain_dropped(&self.emitted);
805        }
806        self.update_metrics_inner(Default::default(), 0);
807    }
808}
809
810/// Running totals over the buffer's contents, plus the introspection logging that reports them.
811///
812/// Shared by the buffer and by every structure that holds its chains, so the totals are
813/// maintained where chains are minted and retired rather than by a walk over the buffer.
814///
815/// Correctness rests on a discipline the code already keeps for the introspection logging: a
816/// chain the buffer holds has been announced exactly once through [`Accounting::chain_created`],
817/// and is retired exactly once, through [`Accounting::chain_dropped`], before it is dropped or
818/// consumed. Chains are immutable once announced, so the totals a chain contributes at
819/// retirement are the ones it contributed at announcement. `metrics_totals_match_walk` checks
820/// the result against a walk.
821#[derive(Clone, Debug)]
822struct Accounting {
823    /// The totals, shared with every clone.
824    totals: Arc<Totals>,
825    /// Introspection logging, absent when the sink is not logging.
826    logging: Option<ChannelLogging>,
827}
828
829/// The running totals behind an [`Accounting`].
830///
831/// Atomics because the buffer must stay `Send`, not because the counters are contended: the sink
832/// touches its buffer from one thread at a time, which is why `Relaxed` suffices.
833#[derive(Debug, Default)]
834struct Totals {
835    /// Number of updates the buffer holds.
836    records: AtomicUsize,
837    /// Serialized size of what the buffer holds, in bytes.
838    size: AtomicUsize,
839    /// Number of allocations holding it: one per chunk, plus the staging vector while it has one.
840    allocations: AtomicUsize,
841}
842
843impl Accounting {
844    /// Construct accounting that reports to the given logging, if any.
845    fn new(logging: Option<ChannelLogging>) -> Self {
846        Self {
847            totals: Default::default(),
848            logging,
849        }
850    }
851
852    /// Return the current totals as `(records, size metrics)`.
853    ///
854    /// The reported capacity equals the size: chunk bodies are exactly sized when minted, and the
855    /// staging vector's slack is bounded by one chunk's worth of updates and not tracked.
856    fn totals(&self) -> (usize, SizeMetrics) {
857        let size = self.totals.size.load(atomic::Ordering::Relaxed);
858        let metrics = SizeMetrics {
859            size,
860            capacity: size,
861            allocations: self.totals.allocations.load(atomic::Ordering::Relaxed),
862        };
863        (self.totals.records.load(atomic::Ordering::Relaxed), metrics)
864    }
865
866    /// Account for a chain the buffer now holds.
867    fn chain_created<D: Data>(&self, chain: &Chain<D>) {
868        let relaxed = atomic::Ordering::Relaxed;
869        self.totals.records.fetch_add(chain.update_count, relaxed);
870        self.totals.size.fetch_add(chain.size, relaxed);
871        self.totals
872            .allocations
873            .fetch_add(chain.chunks.len(), relaxed);
874        if let Some(logging) = &self.logging {
875            logging.chain_created(chain.update_count);
876        }
877    }
878
879    /// Account for a chain the buffer no longer holds.
880    fn chain_dropped<D: Data>(&self, chain: &Chain<D>) {
881        sub_checked(&self.totals.records, chain.update_count);
882        sub_checked(&self.totals.size, chain.size);
883        sub_checked(&self.totals.allocations, chain.chunks.len());
884        if let Some(logging) = &self.logging {
885            logging.chain_dropped(chain.update_count);
886        }
887    }
888
889    /// Announce the staging area as an empty chain, which is how its population is reported.
890    fn stage_created(&self) {
891        if let Some(logging) = &self.logging {
892            logging.chain_created(0);
893        }
894    }
895
896    /// Retire the staging area, which holds `records` updates, `bytes` bytes, and `allocations`
897    /// allocations.
898    ///
899    /// Balances [`Accounting::stage_created`]: the chain the stage is reported as is dropped at
900    /// its current length.
901    fn stage_dropped(&self, records: usize, bytes: usize, allocations: usize) {
902        sub_checked(&self.totals.records, records);
903        sub_checked(&self.totals.size, bytes);
904        sub_checked(&self.totals.allocations, allocations);
905        if let Some(logging) = &self.logging {
906            logging.chain_dropped(records);
907        }
908    }
909
910    /// Account for the staging area changing by `records` updates, `bytes` bytes, and
911    /// `allocations` allocations.
912    fn stage_diff(&self, records: isize, bytes: isize, allocations: isize) {
913        add_signed(&self.totals.records, records);
914        add_signed(&self.totals.size, bytes);
915        add_signed(&self.totals.allocations, allocations);
916
917        // The stage is reported as a chain that is dropped and re-created at its new length.
918        let Some(logging) = &self.logging else { return };
919        if records > 0 {
920            logging.chain_created(usize::try_from(records).expect("positive"));
921            logging.chain_dropped(0);
922        } else if records < 0 {
923            logging.chain_created(0);
924            logging.chain_dropped(usize::try_from(-records).expect("positive"));
925        }
926    }
927
928    /// Report the change from `old` to `new` size metrics to introspection.
929    fn report_size_metrics(&self, new: SizeMetrics, old: SizeMetrics) {
930        let Some(logging) = &self.logging else { return };
931        let i = |x: usize| isize::try_from(x).expect("must fit");
932        logging.report_size_diff(i(new.size) - i(old.size));
933        logging.report_capacity_diff(i(new.capacity) - i(old.capacity));
934        logging.report_allocations_diff(i(new.allocations) - i(old.allocations));
935    }
936}
937
938/// Add a signed `diff` to `counter`.
939///
940/// # Panics
941///
942/// Panics if the result would be negative, which means a retirement was never matched by an
943/// announcement.
944fn add_signed(counter: &AtomicUsize, diff: isize) {
945    let relaxed = atomic::Ordering::Relaxed;
946    if diff >= 0 {
947        counter.fetch_add(usize::try_from(diff).expect("non-negative"), relaxed);
948    } else {
949        sub_checked(counter, usize::try_from(-diff).expect("positive"));
950    }
951}
952
953/// Subtract `sub` from `counter`.
954///
955/// # Panics
956///
957/// Panics if the result would be negative, which means a retirement was never matched by an
958/// announcement.
959fn sub_checked(counter: &AtomicUsize, sub: usize) {
960    let prev = counter.fetch_sub(sub, atomic::Ordering::Relaxed);
961    assert!(prev >= sub, "total retired below zero");
962}
963
964/// A bucket of `Chain`s, for use in a [`BucketChain`].
965///
966/// All chains are individually sorted by (time, data) and consolidated, but updates can appear in
967/// multiple chains, so consumers must merge the chains to obtain consolidated updates.
968struct ChainBucket<D: Data> {
969    /// The contained chains.
970    ///
971    /// Maintained with the chain invariant on pushes; splits can leave it violated until the next
972    /// push restores it.
973    chains: Vec<Chain<D>>,
974    /// The size factor of subsequent chains required by the chain invariant.
975    chain_proportionality: f64,
976    /// Running totals and introspection logging.
977    accounting: Accounting,
978}
979
980impl<D: Data> fmt::Debug for ChainBucket<D> {
981    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
982        f.debug_struct("ChainBucket")
983            .field("chains", &self.chains)
984            .finish_non_exhaustive()
985    }
986}
987
988impl<D: Data> ChainBucket<D> {
989    /// Construct a new, empty `ChainBucket`.
990    fn new(chain_proportionality: f64, accounting: Accounting) -> Self {
991        Self {
992            chains: Vec::new(),
993            chain_proportionality,
994            accounting,
995        }
996    }
997
998    /// Push a chain onto the bucket, restoring the chain invariant.
999    fn push_chain(&mut self, chain: Chain<D>) {
1000        if chain.is_empty() {
1001            return;
1002        }
1003        self.accounting.chain_created(&chain);
1004        self.chains.push(chain);
1005
1006        // Restore the chain invariant.
1007        let prop = self.chain_proportionality;
1008        let merge_needed = |chains: &[Chain<_>]| match chains {
1009            [.., prev, last] => {
1010                let last_len = f64::cast_lossy(last.update_count);
1011                let prev_len = f64::cast_lossy(prev.update_count);
1012                last_len * prop > prev_len
1013            }
1014            _ => false,
1015        };
1016
1017        while merge_needed(&self.chains) {
1018            let a = self.chains.pop().unwrap();
1019            let b = self.chains.pop().unwrap();
1020            self.accounting.chain_dropped(&a);
1021            self.accounting.chain_dropped(&b);
1022
1023            // A merged chain waits about `chain_proportionality` times longer for its next merge
1024            // than its inputs did, so it is one generation deeper.
1025            let depth = a.depth.max(b.depth).saturating_add(1);
1026            let cursors = [a, b].into_iter().filter_map(Chain::into_cursor).collect();
1027            let merged = merge_cursors(cursors, depth);
1028            if !merged.is_empty() {
1029                self.accounting.chain_created(&merged);
1030                self.chains.push(merged);
1031            }
1032        }
1033    }
1034
1035    /// Convert the bucket into its contained chains.
1036    fn into_chains(self) -> Vec<Chain<D>> {
1037        self.chains
1038    }
1039}
1040
1041impl<D: Data> Bucket for ChainBucket<D> {
1042    type Timestamp = Timestamp;
1043
1044    fn split(self, timestamp: &Self::Timestamp, fuel: &mut i64) -> (Self, Self) {
1045        let mut lower = Self::new(self.chain_proportionality, self.accounting.clone());
1046        let mut upper = Self::new(self.chain_proportionality, self.accounting.clone());
1047
1048        for chain in self.chains {
1049            // Whole chunks are reused; at most one chunk straddling the timestamp is copied per
1050            // chain. Account fuel at chunk granularity.
1051            *fuel = fuel.saturating_sub(i64::try_from(chain.chunks.len()).expect("must fit"));
1052
1053            self.accounting.chain_dropped(&chain);
1054            let (lo, hi) = chain.split_at_time(*timestamp);
1055            for (part, target) in [(lo, &mut lower), (hi, &mut upper)] {
1056                if !part.is_empty() {
1057                    target.accounting.chain_created(&part);
1058                    target.chains.push(part);
1059                }
1060            }
1061        }
1062
1063        (lower, upper)
1064    }
1065}
1066
1067/// A chain of [`Chunk`]s containing updates.
1068///
1069/// All updates in a chain are sorted by (time, data) and consolidated.
1070///
1071/// Note that, in contrast to [`Chunk`]s, chains can be empty. Though we generally try to avoid
1072/// keeping around empty chains.
1073#[derive(Debug)]
1074struct Chain<D: Data> {
1075    /// The contained chunks.
1076    chunks: Vec<Chunk<D>>,
1077    /// The number of updates contained in all chunks.
1078    update_count: usize,
1079    /// The serialized size of all chunks, in bytes.
1080    ///
1081    /// Maintained as chunks are pushed, so metrics never walk the chunks.
1082    size: usize,
1083    /// The deepest generation among the chunks, see [`Chunk::depth`].
1084    depth: u8,
1085}
1086
1087impl<D: Data> Chain<D> {
1088    /// Construct an empty chain.
1089    fn new() -> Self {
1090        Self {
1091            chunks: Default::default(),
1092            update_count: 0,
1093            size: 0,
1094            depth: 0,
1095        }
1096    }
1097
1098    /// Return whether the chain is empty.
1099    fn is_empty(&self) -> bool {
1100        self.chunks.is_empty()
1101    }
1102
1103    /// Push a chunk onto the chain.
1104    ///
1105    /// All updates in the chunk must sort after all updates already in the chain, in
1106    /// (time, data)-order, to ensure the chain remains sorted.
1107    fn push_chunk(&mut self, chunk: Chunk<D>) {
1108        mz_ore::soft_assert_no_log!(self.can_accept_chunk(&chunk));
1109
1110        self.update_count += chunk.len();
1111        self.size += chunk.size();
1112        self.depth = self.depth.max(chunk.depth);
1113        self.chunks.push(chunk);
1114    }
1115
1116    /// Return whether the chain can accept the given chunk at its end while preserving
1117    /// (time, data)-order.
1118    ///
1119    /// NOTE: The cached boundary times settle every case but a tie. On a tie the boundary updates
1120    /// themselves are compared, which reads both chunks. The reads are scoped, so neither chunk
1121    /// stays materialized, but a spilled chunk costs a copy out of the pool. The only caller is the
1122    /// soft assertion in [`Chain::push_chunk`], and soft assertions are live in any build started
1123    /// with `MZ_SOFT_ASSERTIONS` set, so this cost is not confined to debug builds. Ties are
1124    /// reached whenever a run of updates at a single timestamp spans a chunk boundary, which
1125    /// [`ChunkBuilder`] produces for any such run larger than its byte limit.
1126    fn can_accept_chunk(&self, chunk: &Chunk<D>) -> bool {
1127        match self.chunks.last() {
1128            None => true,
1129            Some(last) => match last.last_time().cmp(&chunk.first_time()) {
1130                Ordering::Less => true,
1131                Ordering::Greater => false,
1132                Ordering::Equal => last.with_view(|last_view| {
1133                    chunk.with_view(|view| {
1134                        let time = chunk.first_time();
1135                        let (dc, _, _) = last_view.get(last.len() - 1);
1136                        let (d, _, _) = view.get(0);
1137                        refs_cmp::<D>((time, dc), (time, d)).is_lt()
1138                    })
1139                }),
1140            },
1141        }
1142    }
1143
1144    /// Convert the chain into a cursor over the contained updates.
1145    fn into_cursor(self) -> Option<Cursor<D>> {
1146        let chunks = self.chunks.into_iter().map(Rc::new).collect();
1147        Cursor::new(chunks)
1148    }
1149
1150    /// Return an iterator over the contained updates.
1151    ///
1152    /// Reads each chunk scoped, so a spilled chunk stays spillable after the read. This is the
1153    /// read behind [`CorrectionV2::updates_before`], after which `emitted` rests until the next
1154    /// consolidation merges it: a batch as large as the snapshot after hydration would otherwise
1155    /// sit on the heap for that whole time. Each chunk's updates are made owned in one pass, so
1156    /// the iterator holds at most one chunk's worth of owned updates.
1157    fn iter(&self) -> impl Iterator<Item = (D, Timestamp, Diff)> + '_ {
1158        self.chunks.iter().flat_map(|c| {
1159            c.with_view(|view| {
1160                (0..c.len())
1161                    .map(|i| {
1162                        let (d, t, r) = view.get(i);
1163                        (D::into_owned(d), t, r)
1164                    })
1165                    .collect::<Vec<_>>()
1166            })
1167        })
1168    }
1169
1170    /// Count the distinct times of updates at times before `time`, up to the given cap.
1171    ///
1172    /// The scan uses one binary search per distinct time, so its cost is bounded by
1173    /// O(cap log chunks).
1174    fn distinct_times_before(&self, time: Timestamp, cap: usize) -> usize {
1175        let mut count = 0;
1176        let mut chunk_idx = 0;
1177        let mut offset = 0;
1178        while count < cap && chunk_idx < self.chunks.len() {
1179            let chunk = &self.chunks[chunk_idx];
1180            let current = chunk.index(offset).1;
1181            if current >= time {
1182                break;
1183            }
1184            count += 1;
1185            // Skip to the first update at a time greater than `current`.
1186            match chunk.find_time_greater_than(current) {
1187                Some(idx) => offset = idx,
1188                None => {
1189                    // All later updates at `current` are in subsequent chunks.
1190                    chunk_idx += 1;
1191                    offset = 0;
1192                    while chunk_idx < self.chunks.len() {
1193                        match self.chunks[chunk_idx].find_time_greater_than(current) {
1194                            Some(idx) => {
1195                                offset = idx;
1196                                break;
1197                            }
1198                            None => chunk_idx += 1,
1199                        }
1200                    }
1201                }
1202            }
1203        }
1204        count
1205    }
1206
1207    /// Split the chain at the given time.
1208    ///
1209    /// Returns two chains, the first containing all updates at times < `time`, the second
1210    /// containing all updates at times >= `time`. Chunks fully on either side of `time` are
1211    /// reused; only a chunk straddling `time` is copied.
1212    fn split_at_time(mut self, time: Timestamp) -> (Self, Self) {
1213        let mut lower = Self::new();
1214        let mut upper = Self::new();
1215
1216        let Some(skip_ts) = time.step_back() else {
1217            // Nothing sorts before `time`.
1218            return (lower, self);
1219        };
1220
1221        for chunk in self.chunks.drain(..) {
1222            // Route whole chunks by cached boundary times, so a chunk that lands entirely on one
1223            // side is moved without paging it in. Only a straddling chunk is materialized here.
1224            // With soft assertions on, `push_chunk` can still page in a chunk whose boundary time
1225            // ties the chain's last one, see `Chain::can_accept_chunk`.
1226            if chunk.last_time() < time {
1227                lower.push_chunk(chunk);
1228            } else if chunk.first_time() >= time {
1229                upper.push_chunk(chunk);
1230            } else {
1231                // The chunk straddles `time`; copy its two halves.
1232                let idx = chunk
1233                    .find_time_greater_than(skip_ts)
1234                    .expect("straddles time");
1235                let view = chunk.view();
1236                let mut builder = ChainBuilder::at_depth(chunk.depth);
1237                for i in 0..idx {
1238                    builder.push_ref(view.get(i));
1239                }
1240                for part in builder.finish().chunks {
1241                    lower.push_chunk(part);
1242                }
1243                let mut builder = ChainBuilder::at_depth(chunk.depth);
1244                for i in idx..chunk.len() {
1245                    builder.push_ref(view.get(i));
1246                }
1247                for part in builder.finish().chunks {
1248                    upper.push_chunk(part);
1249                }
1250            }
1251        }
1252
1253        (lower, upper)
1254    }
1255}
1256
1257/// A builder that constructs a [`Chain`] from a stream of updates.
1258///
1259/// Wraps a [`ChunkBuilder`] and drains its minted chunks into a [`Chain`]. Pushed updates must
1260/// arrive in (time, data) sorted order.
1261struct ChainBuilder<D: Data> {
1262    builder: ChunkBuilder<D>,
1263    chain: Chain<D>,
1264}
1265
1266impl<D: Data> Default for ChainBuilder<D> {
1267    fn default() -> Self {
1268        Self::at_depth(0)
1269    }
1270}
1271
1272impl<D: Data> ChainBuilder<D> {
1273    /// A builder whose chunks are minted at the given generational depth.
1274    fn at_depth(depth: u8) -> Self {
1275        Self {
1276            builder: ChunkBuilder::at_depth(depth),
1277            chain: Chain::new(),
1278        }
1279    }
1280
1281    /// Push a reference-form update into the builder.
1282    fn push_ref(&mut self, update: Ref<'_, (D, Timestamp, Diff)>) {
1283        if let Some(chunk) = self.builder.push(update) {
1284            self.chain.push_chunk(chunk);
1285        }
1286    }
1287
1288    /// Push an owned-form update into the builder.
1289    fn push_owned(&mut self, update: &(D, Timestamp, Diff)) {
1290        if let Some(chunk) = self.builder.push(update) {
1291            self.chain.push_chunk(chunk);
1292        }
1293    }
1294
1295    /// Push the updates produced by a cursor into the builder.
1296    fn push_cursor(&mut self, cursor: Cursor<D>) {
1297        let mut rest = Some(cursor);
1298        // One borrow per chunk: see `Chunk::view` for why this must not move into the inner loop.
1299        while let Some(cursor) = rest.take() {
1300            let chunk = cursor.chunk_handle();
1301            let view = chunk.view();
1302            rest = Some(cursor);
1303
1304            while let Some(cursor) = rest.as_ref() {
1305                if !cursor.reads_from(&chunk) {
1306                    break;
1307                }
1308                self.push_ref(cursor.get_with(&view));
1309                rest = rest.take().expect("checked above").step();
1310            }
1311        }
1312    }
1313
1314    /// Finish building, returning the assembled [`Chain`].
1315    fn finish(self) -> Chain<D> {
1316        let Self { builder, mut chain } = self;
1317        if let Some(chunk) = builder.finish() {
1318            chain.push_chunk(chunk);
1319        }
1320        chain
1321    }
1322}
1323
1324impl<D: Data> Extend<(D, Timestamp, Diff)> for ChainBuilder<D> {
1325    fn extend<I: IntoIterator<Item = (D, Timestamp, Diff)>>(&mut self, iter: I) {
1326        for update in iter {
1327            self.push_owned(&update);
1328        }
1329    }
1330}
1331
1332/// A cursor over updates in a chain.
1333///
1334/// A cursor provides two guarantees:
1335///  * Produced updates are ordered and consolidated.
1336///  * A cursor always yields at least one update.
1337///
1338/// The second guarantee is enforced through the type system: Every method that steps a cursor
1339/// forward consumes `self` and returns an `Option<Cursor>` that's `None` if the operation stepped
1340/// over the last update.
1341///
1342/// A cursor holds on to `Rc<Chunk>`s, allowing multiple cursors to produce updates from the same
1343/// chunks concurrently. As soon as a cursor is done producing updates from a [`Chunk`] it drops
1344/// its reference. Once the last cursor is done with a [`Chunk`] its memory can be reclaimed.
1345#[derive(Clone, Debug)]
1346struct Cursor<D: Data> {
1347    /// The chunks from which updates can still be produced.
1348    chunks: VecDeque<Rc<Chunk<D>>>,
1349    /// The current offset into `chunks.front()`.
1350    chunk_offset: usize,
1351    /// An optional limit for the number of updates the cursor will produce.
1352    limit: Option<usize>,
1353    /// An optional overwrite for the timestamp of produced updates.
1354    overwrite_ts: Option<Timestamp>,
1355}
1356
1357impl<D: Data> Cursor<D> {
1358    /// Construct a cursor over a list of chunks.
1359    ///
1360    /// Returns `None` if `chunks` is empty.
1361    fn new(chunks: VecDeque<Rc<Chunk<D>>>) -> Option<Self> {
1362        if chunks.is_empty() {
1363            return None;
1364        }
1365
1366        Some(Self {
1367            chunks,
1368            chunk_offset: 0,
1369            limit: None,
1370            overwrite_ts: None,
1371        })
1372    }
1373
1374    /// Set a limit for the number of updates this cursor will produce.
1375    ///
1376    /// # Panics
1377    ///
1378    /// Panics if there is already a limit lower than the new one.
1379    fn set_limit(mut self, limit: usize) -> Option<Self> {
1380        assert!(self.limit.is_none_or(|l| l >= limit));
1381
1382        if limit == 0 {
1383            return None;
1384        }
1385
1386        // Release chunks made unreachable by the limit.
1387        let mut count = 0;
1388        let mut idx = 0;
1389        let mut offset = self.chunk_offset;
1390        while idx < self.chunks.len() && count < limit {
1391            let chunk = &self.chunks[idx];
1392            count += chunk.len() - offset;
1393            idx += 1;
1394            offset = 0;
1395        }
1396        self.chunks.truncate(idx);
1397
1398        if count > limit {
1399            self.limit = Some(limit);
1400        }
1401
1402        Some(self)
1403    }
1404
1405    /// Get a reference to the current update.
1406    ///
1407    /// Single-access only. A loop over a cursor must hoist [`Cursor::chunk_handle`]'s
1408    /// [`Chunk::view`] and read through [`Cursor::get_with`] instead.
1409    fn get(&self) -> Ref<'_, (D, Timestamp, Diff)> {
1410        let chunk = self.get_chunk();
1411        let (d, t, r) = chunk.index(self.chunk_offset);
1412        let t = self.overwrite_ts.unwrap_or(t);
1413        (d, t, r)
1414    }
1415
1416    /// Get a reference to the current update, reading through an already-borrowed view.
1417    ///
1418    /// # Panics
1419    ///
1420    /// Panics if `view` is not a view of the cursor's current chunk. Guard loops with
1421    /// [`Cursor::reads_from`], which is how a caller learns the cursor has crossed into the next
1422    /// chunk and the view must be refreshed.
1423    fn get_with<'a>(&self, view: &ChunkView<'a, D>) -> Ref<'a, (D, Timestamp, Diff)> {
1424        debug_assert_eq!(view.len(), self.get_chunk().len(), "view of another chunk");
1425        let (d, t, r) = view.get(self.chunk_offset);
1426        let t = self.overwrite_ts.unwrap_or(t);
1427        (d, t, r)
1428    }
1429
1430    /// A shared handle on the chunk the cursor currently reads from.
1431    ///
1432    /// Held by callers that hoist a [`Chunk::view`], so the view's borrow outlives the cursor
1433    /// steps taken against it.
1434    fn chunk_handle(&self) -> Rc<Chunk<D>> {
1435        Rc::clone(&self.chunks[0])
1436    }
1437
1438    /// Whether the cursor still reads from `chunk`.
1439    fn reads_from(&self, chunk: &Rc<Chunk<D>>) -> bool {
1440        Rc::ptr_eq(&self.chunks[0], chunk)
1441    }
1442
1443    /// Get a reference to the current chunk.
1444    fn get_chunk(&self) -> &Chunk<D> {
1445        &self.chunks[0]
1446    }
1447
1448    /// Step to the next update.
1449    ///
1450    /// Returns the stepped cursor, or `None` if the step was over the last update.
1451    fn step(mut self) -> Option<Self> {
1452        if self.chunk_offset == self.get_chunk().len() - 1 {
1453            return self.skip_chunk().map(|(c, _)| c);
1454        }
1455
1456        self.chunk_offset += 1;
1457
1458        if let Some(limit) = &mut self.limit {
1459            *limit -= 1;
1460            if *limit == 0 {
1461                return None;
1462            }
1463        }
1464
1465        Some(self)
1466    }
1467
1468    /// Skip the remainder of the current chunk.
1469    ///
1470    /// Returns the forwarded cursor and the number of updates skipped, or `None` if no chunks are
1471    /// left after the skip.
1472    fn skip_chunk(mut self) -> Option<(Self, usize)> {
1473        let chunk = self.chunks.pop_front().expect("cursor invariant");
1474
1475        if self.chunks.is_empty() {
1476            return None;
1477        }
1478
1479        let skipped = chunk.len() - self.chunk_offset;
1480        self.chunk_offset = 0;
1481
1482        if let Some(limit) = &mut self.limit {
1483            if skipped >= *limit {
1484                return None;
1485            }
1486            *limit -= skipped;
1487        }
1488
1489        Some((self, skipped))
1490    }
1491
1492    /// Skip all updates with times <= the given time.
1493    ///
1494    /// Returns the forwarded cursor and the number of updates skipped, or `None` if no updates are
1495    /// left after the skip.
1496    fn skip_time(mut self, time: Timestamp) -> Option<(Self, usize)> {
1497        if self.overwrite_ts.is_some_and(|ts| ts <= time) {
1498            return None;
1499        } else if self.get().1 > time {
1500            return Some((self, 0));
1501        }
1502
1503        let mut skipped = 0;
1504
1505        let new_offset = loop {
1506            let chunk = self.get_chunk();
1507            if let Some(index) = chunk.find_time_greater_than(time) {
1508                break index;
1509            }
1510
1511            let (cursor, count) = self.skip_chunk()?;
1512            self = cursor;
1513            skipped += count;
1514        };
1515
1516        skipped += new_offset - self.chunk_offset;
1517        self.chunk_offset = new_offset;
1518
1519        Some((self, skipped))
1520    }
1521
1522    /// Advance all updates in this cursor by the given `since_ts`.
1523    ///
1524    /// Returns a list of cursors, each of which yields ordered and consolidated updates that have
1525    /// been advanced by `since_ts`.
1526    fn advance_by(mut self, since_ts: Timestamp) -> Vec<Self> {
1527        // If the cursor has an `overwrite_ts`, all its updates are at the same time already. We
1528        // only need to advance the `overwrite_ts` by the `since_ts`.
1529        if let Some(ts) = self.overwrite_ts {
1530            if ts < since_ts {
1531                self.overwrite_ts = Some(since_ts);
1532            }
1533            return vec![self];
1534        }
1535
1536        // Otherwise we need to split the cursor so that each new cursor only yields runs of
1537        // updates that are correctly (time, data)-ordered when advanced by `since_ts`. We achieve
1538        // this by splitting the cursor at each time <= `since_ts`.
1539        let mut splits = Vec::new();
1540        let mut remaining = Some(self);
1541
1542        while let Some(cursor) = remaining.take() {
1543            let (_, time, _) = cursor.get();
1544            if time >= since_ts {
1545                splits.push(cursor);
1546                break;
1547            }
1548
1549            let mut current = cursor.clone();
1550            if let Some((cursor, skipped)) = cursor.skip_time(time) {
1551                remaining = Some(cursor);
1552                current = current.set_limit(skipped).expect("skipped at least 1");
1553            }
1554            current.overwrite_ts = Some(since_ts);
1555            splits.push(current);
1556        }
1557
1558        splits
1559    }
1560
1561    /// Drain the cursor into a [`Chain`].
1562    ///
1563    /// This reuses the underlying chunks if possible, and writes new ones at `depth` otherwise.
1564    /// Reused chunks keep their own depth.
1565    fn into_chain(self, depth: u8) -> Chain<D> {
1566        match self.try_unwrap(depth) {
1567            Ok(chain) => chain,
1568            Err((_, cursor)) => {
1569                let mut builder = ChainBuilder::at_depth(depth);
1570                builder.push_cursor(cursor);
1571                builder.finish()
1572            }
1573        }
1574    }
1575
1576    /// Attempt to unwrap the cursor into a [`Chain`].
1577    ///
1578    /// This operation efficiently reuses chunks by directly inserting them into the output chain
1579    /// where possible.
1580    ///
1581    /// An unwrap is only successful if the cursor's `limit` and `overwrite_ts` are both `None` and
1582    /// the cursor has unique references to its chunks. If the unwrap fails, this method returns an
1583    /// `Err` containing the cursor in an unchanged state, allowing the caller to convert it into a
1584    /// chain by copying chunks rather than reusing them.
1585    fn try_unwrap(self, depth: u8) -> Result<Chain<D>, (&'static str, Self)> {
1586        if self.limit.is_some() {
1587            return Err(("cursor with limit", self));
1588        }
1589        if self.overwrite_ts.is_some() {
1590            return Err(("cursor with overwrite_ts", self));
1591        }
1592        if self.chunks.iter().any(|c| Rc::strong_count(c) != 1) {
1593            return Err(("cursor on shared chunks", self));
1594        }
1595
1596        let mut builder = ChainBuilder::at_depth(depth);
1597        let mut remaining = Some(self);
1598
1599        // We might be partway through the first chunk, in which case we can't reuse it but need to
1600        // allocate a new one to contain only the updates the cursor can still yield.
1601        while let Some(cursor) = remaining.take() {
1602            if cursor.chunk_offset == 0 {
1603                remaining = Some(cursor);
1604                break;
1605            }
1606            let update = cursor.get();
1607            builder.push_ref(update);
1608            remaining = cursor.step();
1609        }
1610
1611        let mut chain = builder.finish();
1612        if let Some(cursor) = remaining {
1613            for chunk in cursor.chunks {
1614                let chunk = Rc::into_inner(chunk).expect("checked above");
1615                chain.push_chunk(chunk);
1616            }
1617        }
1618
1619        Ok(chain)
1620    }
1621}
1622
1623/// A non-empty chunk of updates, backed by a columnar region.
1624///
1625/// All updates in a chunk are sorted by (time, data) and consolidated.
1626///
1627/// Chunks are immutable once created. They are produced by [`ChunkBuilder`], which mints a
1628/// new chunk whenever its in-progress columnar container reaches a fixed serialized byte
1629/// boundary (~2 MiB, matching the ship granularity used elsewhere in the codebase), so each
1630/// chunk corresponds to a single, predictably sized allocation.
1631struct Chunk<D: Data> {
1632    /// The body in the buffer pool, which may spill it.
1633    ///
1634    /// Empty when the body stays on the heap, which is what
1635    /// [`mz_timely_util::columnar::chunk::try_spill_ref`] decides, and emptied again by
1636    /// [`Chunk::body`]: a materialized chunk holds its body on the heap for the rest of its
1637    /// life, so keeping the pool's copy alive would only double the footprint and hold budget
1638    /// the pool could give to a chunk that still needs it. [`Chunk::with_view`] reads a copy and
1639    /// leaves the body here.
1640    ///
1641    /// A `Mutex` (not `RefCell`) keeps the chunk `Sync`: cursors hold chunks behind a shared
1642    /// `Rc`, and the iterator returned by [`CorrectionV2::updates_before`] borrows them across
1643    /// the persist writer's `await`, so `&Chunk` must be `Send`. The lock is taken at
1644    /// materialization and for scoped reads, and is otherwise uncontended (the sink runs
1645    /// single-threaded per worker).
1646    pooled: Mutex<Option<ChunkHandle>>,
1647    /// The materialized form, populated lazily by [`Chunk::body`] on first access.
1648    ///
1649    /// An `OnceLock` for the same `Sync` reason as `pooled`. Once set the slot is never
1650    /// cleared, so its address is stable and [`Chunk::index`] can hand out `Ref<'_>` borrows
1651    /// tied to `&self`. The allocation is freed when the chunk drops, which bounds resident
1652    /// memory to the chunks under an active merge front.
1653    resident: OnceLock<ColumnBody<(D, Timestamp, Diff)>>,
1654    /// Number of updates, cached so `len` and chain bookkeeping never page the chunk in.
1655    len: usize,
1656    /// Serialized size of the body in words, cached so size accounting never pages it in.
1657    body_words: usize,
1658    /// Time of the first update, cached so boundary checks (`split_at_time`, `can_accept`) route
1659    /// a resting chunk without materializing it.
1660    first_time: Timestamp,
1661    /// Time of the last update, cached likewise.
1662    last_time: Timestamp,
1663    /// The generational depth: 0 for chunks built from staged updates or written by a read,
1664    /// one more than the deepest input for chunks written by a bucket's chain merge. The pool
1665    /// treats deeper chunks as colder, and compresses them past the floor set by
1666    /// [`chunk::set_compress_min_depth`].
1667    depth: u8,
1668}
1669
1670impl<D: Data> fmt::Debug for Chunk<D> {
1671    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1672        write!(f, "Chunk(<{}>)", self.len())
1673    }
1674}
1675
1676impl<D: Data> Chunk<D> {
1677    /// Mint a chunk from the given non-empty body, emptying it.
1678    ///
1679    /// Reads the metadata a resting chunk must answer (length, size, boundary times) while the
1680    /// body is still in hand, then offers it to the buffer pool. A body the pool takes is
1681    /// encoded straight into its slot, which leaves the typed allocation behind for the caller
1682    /// to refill, and materializes lazily on first read. A body the pool declines is
1683    /// serialized out of the typed allocation instead, because a resident chunk has to own it.
1684    ///
1685    /// # Panics
1686    ///
1687    /// Panics if the body is empty. Chunks are non-empty by construction; [`ChunkBuilder`]
1688    /// only ever mints from a populated body.
1689    fn mint(body: &mut ColumnBody<(D, Timestamp, Diff)>, depth: u8) -> Self {
1690        let (len, first_time, last_time) = {
1691            let borrowed = body.borrow();
1692            let len = borrowed.len();
1693            assert!(len > 0, "chunks are non-empty");
1694            (len, borrowed.get(0).1, borrowed.get(len - 1).1)
1695        };
1696        let body_words = body.length_in_bytes() / std::mem::size_of::<u64>();
1697
1698        // TODO: chunks resting in far-future buckets go untouched for longer than their merge
1699        // depth suggests, and could spill deeper still.
1700        let (pooled, resident) = match chunk::try_spill_ref(body, depth) {
1701            Some(handle) => {
1702                body.clear();
1703                (Mutex::new(Some(handle)), OnceLock::new())
1704            }
1705            None => {
1706                // Serialize into an exactly sized vector, which leaves the typed allocation
1707                // in place for the caller and keeps its spare capacity out of a body that is
1708                // held for the rest of the chunk's life.
1709                let mut words = Vec::with_capacity(body_words);
1710                indexed::encode(&mut words, &body.borrow());
1711                body.clear();
1712                let cell = OnceLock::new();
1713                if cell.set(ColumnBody::Words(words)).is_err() {
1714                    unreachable!("cell is fresh");
1715                }
1716                (Mutex::new(None), cell)
1717            }
1718        };
1719        Self {
1720            pooled,
1721            body_words,
1722            resident,
1723            len,
1724            first_time,
1725            last_time,
1726            depth,
1727        }
1728    }
1729
1730    /// Materialize the chunk's body, taking it out of the pool on first access.
1731    ///
1732    /// The returned reference is valid for as long as `&self`: the `OnceLock` slot is never
1733    /// cleared once populated, so its contents have a stable address. Taking the body frees the
1734    /// pool's copy, so the chunk holds exactly one.
1735    ///
1736    /// For reads by merges and splits, whose output replaces the chunk. A read after which the
1737    /// chunk survives uses [`Chunk::with_view`], which keeps the body spillable.
1738    fn body(&self) -> &ColumnBody<(D, Timestamp, Diff)> {
1739        self.resident.get_or_init(|| {
1740            let handle = self
1741                .pooled
1742                .lock()
1743                .expect("pool handle mutex poisoned")
1744                .take()
1745                .expect("a chunk the pool declined is materialized at construction");
1746            let mut words = Vec::new();
1747            handle.take(&mut words);
1748            ColumnBody::Words(words)
1749        })
1750    }
1751
1752    /// Call `f` with a view of the chunk's body, leaving the body where it is.
1753    ///
1754    /// A materialized chunk is read in place. Otherwise the body is copied out of the pool for
1755    /// the call and the copy dropped after it, so the pool keeps its slot and may still spill
1756    /// the body.
1757    fn with_view<R>(&self, f: impl FnOnce(ChunkView<'_, D>) -> R) -> R {
1758        if let Some(body) = self.resident.get() {
1759            return f(body.borrow());
1760        }
1761        let mut words = Vec::new();
1762        {
1763            let pooled = self.pooled.lock().expect("pool handle mutex poisoned");
1764            match pooled.as_ref() {
1765                Some(handle) => handle.read_into(&mut words),
1766                // Materialized between the check above and the lock.
1767                None => {
1768                    drop(pooled);
1769                    let body = self.resident.get().expect("a taken body is materialized");
1770                    return f(body.borrow());
1771                }
1772            }
1773        }
1774        let body = ColumnBody::<(D, Timestamp, Diff)>::Words(words);
1775        f(body.borrow())
1776    }
1777
1778    /// Return the number of updates in the chunk.
1779    fn len(&self) -> usize {
1780        self.len
1781    }
1782
1783    /// Borrow the chunk's body, paging it in if necessary.
1784    ///
1785    /// Any caller that touches more than one update must hoist this out of its loop and index the
1786    /// returned view. `ColumnBody::borrow` on a serialized body rebuilds the struct-of-arrays
1787    /// view from the serialized header on every call, so borrowing per element pays that decode
1788    /// per element.
1789    fn view(&self) -> ChunkView<'_, D> {
1790        self.body().borrow()
1791    }
1792
1793    /// Return the update at the given index, paging the chunk in if necessary.
1794    ///
1795    /// Single-access only. Indexing a hoisted [`Chunk::view`] is the loop form.
1796    ///
1797    /// # Panics
1798    ///
1799    /// Panics if the given index is not populated.
1800    fn index(&self, idx: usize) -> Ref<'_, (D, Timestamp, Diff)> {
1801        self.view().get(idx)
1802    }
1803
1804    /// Return the time of the first update, without materializing the chunk.
1805    fn first_time(&self) -> Timestamp {
1806        self.first_time
1807    }
1808
1809    /// Return the time of the last update, without materializing the chunk.
1810    fn last_time(&self) -> Timestamp {
1811        self.last_time
1812    }
1813
1814    /// Return the index of the first update at a time greater than `time`, or `None` if no such
1815    /// update exists.
1816    ///
1817    /// The early-out uses the cached last time, so a chunk whose updates are all at or before
1818    /// `time` is skipped without paging it in.
1819    fn find_time_greater_than(&self, time: Timestamp) -> Option<usize> {
1820        if self.last_time <= time {
1821            return None;
1822        }
1823
1824        let view = self.view();
1825        let mut lower = 0;
1826        let mut upper = self.len;
1827        while lower < upper {
1828            let idx = (lower + upper) / 2;
1829            if view.get(idx).1 > time {
1830                upper = idx;
1831            } else {
1832                lower = idx + 1;
1833            }
1834        }
1835
1836        Some(lower)
1837    }
1838
1839    /// Return the serialized size of the chunk's body in bytes, for use in metrics.
1840    ///
1841    /// The chunk holds exactly one copy of its body, in the pool or on the heap, so this is its
1842    /// size either way. A scoped read's copy lives only for the read. It is not a resident-byte
1843    /// number: what the pool has evicted is the pool's business, and it publishes
1844    /// `resident_bytes` itself (`mz_ore::pool::PoolStats`).
1845    fn size(&self) -> usize {
1846        self.body_words * std::mem::size_of::<u64>()
1847    }
1848}
1849
1850/// Builder that mints fixed-size [`Chunk`]s from a stream of updates.
1851///
1852/// Updates accumulate in one body, which is minted into a chunk as soon as it reports itself at
1853/// capacity, so every chunk is a single, predictably sized body. Minting a spilled body leaves
1854/// the typed allocation in place, so a builder that mints many chunks grows one allocation once.
1855struct ChunkBuilder<D: Data> {
1856    /// The updates pushed since the last mint.
1857    current: ColumnBody<(D, Timestamp, Diff)>,
1858    /// The generational depth minted chunks carry.
1859    depth: u8,
1860}
1861
1862impl<D: Data> ChunkBuilder<D> {
1863    /// A builder whose chunks are minted at the given generational depth.
1864    fn at_depth(depth: u8) -> Self {
1865        Self {
1866            current: Default::default(),
1867            depth,
1868        }
1869    }
1870
1871    /// Push an update, returning a chunk if the push completed one.
1872    ///
1873    /// Accepts whatever [`ColumnBody`]'s [`PushInto`] impl accepts, both the
1874    /// `Ref<'_, (D, T, R)>` refs produced by cursors and `&(D, T, R)` references to owned
1875    /// tuples drained from the staging buffer.
1876    ///
1877    /// [`PushInto`]: timely::container::PushInto
1878    #[inline]
1879    fn push<T>(&mut self, item: T) -> Option<Chunk<D>>
1880    where
1881        ColumnBody<(D, Timestamp, Diff)>: PushInto<T>,
1882    {
1883        PushInto::push_into(&mut self.current, item);
1884        // The ship test walks the body's slice lengths, so it costs per push, not per byte.
1885        self.current
1886            .at_capacity()
1887            .then(|| Chunk::mint(&mut self.current, self.depth))
1888    }
1889
1890    /// Mint the updates pushed since the last mint, if any.
1891    fn finish(mut self) -> Option<Chunk<D>> {
1892        (!self.current.is_empty()).then(|| Chunk::mint(&mut self.current, self.depth))
1893    }
1894}
1895
1896/// A buffer for staging updates before they are inserted into the sorted chains.
1897#[derive(Debug)]
1898struct Stage<D> {
1899    /// The contained updates.
1900    ///
1901    /// Grows on demand rather than being allocated at the ship size, so a sink that never
1902    /// stages that much never holds it. One of these exists per sink and worker, which is why
1903    /// the eager allocation is worth avoiding even though a staging area is small.
1904    data: Vec<(D, Timestamp, Diff)>,
1905    /// Bytes held by `data`: the tuples plus what they own on the heap.
1906    bytes: usize,
1907    /// How many bytes to accumulate before shipping a batch, from
1908    /// `compute_correction_v2_chunk_size`.
1909    ///
1910    /// Shipping less often costs staging memory and saves inserts: it is the number of chains
1911    /// minted per update, and every chain minted is a chain some later read has to merge.
1912    /// Counted in bytes with the heap included, so the staging memory is bounded by this
1913    /// setting regardless of row width, and a wide row fills the stage in fewer updates.
1914    chunk_size: usize,
1915    /// Running totals and introspection logging.
1916    ///
1917    /// We want to report the number of records in the stage. To do so, we pretend that the stage
1918    /// is a chain, and every time the number of updates inside changes, the chain gets dropped and
1919    /// re-created.
1920    accounting: Accounting,
1921}
1922
1923impl<D: Data> Stage<D> {
1924    fn new(accounting: Accounting, chunk_size: usize) -> Self {
1925        // For logging, we pretend the stage consists of a single chain.
1926        accounting.stage_created();
1927
1928        Self {
1929            data: Vec::new(),
1930            bytes: 0,
1931            chunk_size,
1932            accounting,
1933        }
1934    }
1935
1936    /// Bytes an update occupies: its datum, heap included, plus time and diff.
1937    fn update_bytes(update: &(D, Timestamp, Diff)) -> usize {
1938        update.0.data_bytes() + std::mem::size_of::<(Timestamp, Diff)>()
1939    }
1940
1941    /// Insert a batch of updates, possibly producing a batch of sorted, consolidated updates
1942    /// ready to be stored.
1943    ///
1944    /// Everything staged ships together, the new batch included, once the staged bytes reach
1945    /// `chunk_size`.
1946    fn insert(
1947        &mut self,
1948        updates: &mut Vec<(D, Timestamp, Diff)>,
1949    ) -> Option<Vec<(D, Timestamp, Diff)>> {
1950        if updates.is_empty() {
1951            return None;
1952        }
1953
1954        let before = self.snapshot();
1955
1956        self.bytes += updates.iter().map(Self::update_bytes).sum::<usize>();
1957        self.data.append(updates);
1958
1959        let maybe_ready = if self.bytes >= self.chunk_size {
1960            let mut ready = std::mem::take(&mut self.data);
1961            self.bytes = 0;
1962            consolidate(&mut ready);
1963            Some(ready)
1964        } else {
1965            None
1966        };
1967
1968        self.account_since(before);
1969
1970        maybe_ready
1971    }
1972
1973    /// Flush all currently staged updates, returning them sorted and consolidated.
1974    fn flush(&mut self) -> Option<Vec<(D, Timestamp, Diff)>> {
1975        let before = self.snapshot();
1976
1977        consolidate(&mut self.data);
1978        self.bytes = 0;
1979        let data = (!self.data.is_empty()).then(|| std::mem::take(&mut self.data));
1980
1981        self.account_since(before);
1982        data
1983    }
1984
1985    /// Advance the times of staged updates by the given `since`.
1986    fn advance_times(&mut self, since: &Antichain<Timestamp>) {
1987        let Some(since_ts) = since.as_option() else {
1988            // If the since is the empty frontier, discard all updates.
1989            let before = self.snapshot();
1990            self.data.clear();
1991            self.bytes = 0;
1992            self.account_since(before);
1993            return;
1994        };
1995
1996        for (_, time, _) in &mut self.data {
1997            *time = std::cmp::max(*time, *since_ts);
1998        }
1999    }
2000
2001    /// The staged length, bytes, and allocation count, taken before a mutation for
2002    /// [`Stage::account_since`].
2003    fn snapshot(&self) -> (isize, isize, isize) {
2004        let i = |x: usize| isize::try_from(x).expect("must fit");
2005        let allocated = isize::from(self.data.capacity() > 0);
2006        (i(self.data.len()), i(self.bytes), allocated)
2007    }
2008
2009    /// Account for the change to the stage since `before`, a [`Stage::snapshot`].
2010    fn account_since(&self, before: (isize, isize, isize)) {
2011        let (len, bytes, allocations) = self.snapshot();
2012        self.accounting
2013            .stage_diff(len - before.0, bytes - before.1, allocations - before.2);
2014    }
2015}
2016
2017impl<D> Drop for Stage<D> {
2018    fn drop(&mut self) {
2019        let allocated = usize::from(self.data.capacity() > 0);
2020        self.accounting
2021            .stage_dropped(self.data.len(), self.bytes, allocated);
2022    }
2023}
2024
2025/// Sort and consolidate the given list of updates.
2026///
2027/// This function is the same as [`differential_dataflow::consolidation::consolidate_updates`],
2028/// except that it sorts updates by (time, data) instead of (data, time).
2029fn consolidate<D: Data>(updates: &mut Vec<(D, Timestamp, Diff)>) {
2030    if updates.len() <= 1 {
2031        return;
2032    }
2033
2034    let diff = |update: &(_, _, Diff)| update.2;
2035
2036    updates.sort_unstable_by(|(d1, t1, _), (d2, t2, _)| (t1, d1).cmp(&(t2, d2)));
2037
2038    let mut offset = 0;
2039    let mut accum = diff(&updates[0]);
2040
2041    for idx in 1..updates.len() {
2042        let this = &updates[idx];
2043        let prev = &updates[idx - 1];
2044        if this.0 == prev.0 && this.1 == prev.1 {
2045            accum += diff(&updates[idx]);
2046        } else {
2047            if accum != Diff::ZERO {
2048                updates.swap(offset, idx - 1);
2049                updates[offset].2 = accum;
2050                offset += 1;
2051            }
2052            accum = diff(&updates[idx]);
2053        }
2054    }
2055
2056    if accum != Diff::ZERO {
2057        let len = updates.len();
2058        updates.swap(offset, len - 1);
2059        updates[offset].2 = accum;
2060        offset += 1;
2061    }
2062
2063    updates.truncate(offset);
2064}
2065
2066/// Compare two `(time, data)` pairs of columnar refs that have unrelated input lifetimes.
2067///
2068/// `<D::Container as Borrow>::Ref<'a>` is an associated-type projection through a trait, so
2069/// the compiler treats it as invariant in `'a` and won't auto-shorten the inputs by variance.
2070/// We instead explicitly reborrow both to a fresh, local lifetime `'x` via
2071/// [`Columnar::reborrow`] before letting the inner `cmp` pick up the `for<'a> Ref<'a>: Ord`
2072/// bound on [`Data`].
2073#[inline]
2074fn refs_cmp<D: Data>(a: (Timestamp, Ref<'_, D>), b: (Timestamp, Ref<'_, D>)) -> Ordering {
2075    #[inline]
2076    fn cmp<'x, D: Data>(a: (Timestamp, Ref<'x, D>), b: (Timestamp, Ref<'x, D>)) -> Ordering {
2077        a.cmp(&b)
2078    }
2079    cmp::<D>((a.0, D::reborrow(a.1)), (b.0, D::reborrow(b.1)))
2080}
2081
2082#[cfg(test)]
2083mod tests {
2084    use mz_ore::metrics::MetricsRegistry;
2085    use mz_persist_client::cfg::PersistConfig;
2086    use mz_persist_client::metrics::Metrics;
2087    use mz_repr::{Diff, Timestamp};
2088
2089    use super::*;
2090    use crate::sink::correction::CorrectionV1;
2091
2092    #[mz_ore::test]
2093    fn chain_builder_update_count_matches_items() {
2094        let mut builder = ChainBuilder::<i64>::default();
2095        for i in 0..10_u64 {
2096            let d = i64::try_from(i).expect("fits");
2097            builder.push_owned(&(d, Timestamp::new(i), Diff::ONE));
2098        }
2099        let chain = builder.finish();
2100        assert_eq!(chain.update_count, chain.iter().count());
2101    }
2102
2103    /// Push enough updates to cross at least one `mint()` boundary, forcing the
2104    /// `Align` encode -> `from_bytes` decode roundtrip (the spilling path this data
2105    /// structure exists to support), and assert `iter()` roundtrips values, order,
2106    /// and diffs across the spill boundary.
2107    #[mz_ore::test]
2108    #[cfg_attr(miri, ignore)] // too slow: crossing the ~2 MiB mint boundary needs ~200k updates
2109    fn chain_builder_roundtrips_across_mint_boundary() {
2110        // A single `mint()` fires near the ~2 MiB (`SHIP_WORDS`) serialized boundary. With
2111        // three 8-byte columns per update that's tens of thousands of updates; pushing 200k
2112        // comfortably forces multiple mints.
2113        let count = 200_000_u64;
2114
2115        let mut builder = ChainBuilder::<i64>::default();
2116        for i in 0..count {
2117            let d = i64::try_from(i).expect("fits");
2118            builder.push_owned(&(d, Timestamp::new(i), Diff::ONE));
2119        }
2120        let chain = builder.finish();
2121
2122        // Crossing the mint boundary must have produced more than one chunk; otherwise the spill
2123        // path (each minted chunk is paged out and read back through the pager) wouldn't be
2124        // exercised. The chunk payload itself is now behind the pager (see [`Chunk`]), so we
2125        // assert on chunk count rather than inspecting the column variant directly.
2126        assert!(
2127            chain.chunks.len() > 1,
2128            "expected multiple minted chunks, got {} chunk(s): {:?}",
2129            chain.chunks.len(),
2130            chain.chunks,
2131        );
2132
2133        // `iter()` must roundtrip every update, in order, with correct diffs.
2134        assert_eq!(chain.update_count, usize::try_from(count).expect("fits"));
2135        let mut expected = 0_u64;
2136        for (d, t, r) in chain.iter() {
2137            assert_eq!(d, i64::try_from(expected).expect("fits"));
2138            assert_eq!(t, Timestamp::new(expected));
2139            assert_eq!(r, Diff::ONE);
2140            expected += 1;
2141        }
2142        assert_eq!(expected, count);
2143    }
2144
2145    impl DataBytes for String {
2146        fn data_bytes(&self) -> usize {
2147            std::mem::size_of::<Self>() + self.len()
2148        }
2149    }
2150
2151    impl DataBytes for i64 {
2152        fn data_bytes(&self) -> usize {
2153            std::mem::size_of::<Self>()
2154        }
2155    }
2156
2157    /// A k-way merge over chains of several chunks each consolidates across all of them, with
2158    /// its rounds ending wherever any input crosses a chunk boundary.
2159    #[mz_ore::test]
2160    #[cfg_attr(miri, ignore)] // slow under Miri
2161    fn merge_many_consolidates_across_chunks() {
2162        // Large enough that the sparsest chain still spans several chunks.
2163        let count = 600_000_i64;
2164        // Chain `k` holds the updates of `inputs[k]`, sorted by (time, data).
2165        let inputs: [Vec<(i64, Timestamp, Diff)>; 3] = [
2166            (0..count)
2167                .map(|d| (d, Timestamp::new(0), Diff::ONE))
2168                .collect(),
2169            (0..count)
2170                .filter(|d| d % 2 == 0)
2171                .map(|d| (d, Timestamp::new(0), Diff::ONE))
2172                .collect(),
2173            (0..count)
2174                .filter(|d| d % 3 == 0)
2175                .map(|d| (d, Timestamp::new(0), -Diff::ONE))
2176                .collect(),
2177        ];
2178
2179        let mut expected = std::collections::BTreeMap::new();
2180        for (d, _, r) in inputs.iter().flatten() {
2181            *expected.entry(*d).or_insert(Diff::ZERO) += *r;
2182        }
2183        expected.retain(|_, r| *r != Diff::ZERO);
2184
2185        let cursors = inputs
2186            .iter()
2187            .map(|updates| {
2188                let mut builder = ChainBuilder::<i64>::default();
2189                for update in updates {
2190                    builder.push_owned(update);
2191                }
2192                let chain = builder.finish();
2193                assert!(chain.chunks.len() > 1, "expected multiple minted chunks");
2194                chain.into_cursor().expect("non-empty")
2195            })
2196            .collect();
2197        let merged = merge_many(cursors, 0);
2198
2199        let actual: Vec<_> = merged.iter().map(|(d, _, r)| (d, r)).collect();
2200        let expected: Vec<_> = expected.into_iter().collect();
2201        assert_eq!(actual, expected);
2202    }
2203
2204    fn sink_metrics() -> SinkMetrics {
2205        let registry = MetricsRegistry::new();
2206        let metrics = Metrics::new(&PersistConfig::new_for_tests(), &registry);
2207        metrics.sink.clone()
2208    }
2209
2210    /// Insert single updates of the given payload width until the stage ships, returning how
2211    /// many updates that took.
2212    fn updates_until_ship(width: usize, chunk_size: usize) -> usize {
2213        let metrics = sink_metrics();
2214        let mut v2 = CorrectionV2::<String>::new(
2215            metrics.clone(),
2216            metrics.for_worker(0),
2217            None,
2218            3.0,
2219            chunk_size,
2220        );
2221        for i in 0..chunk_size {
2222            let datum = format!("{i:0width$}");
2223            v2.insert(&mut vec![(datum, Timestamp::from(0), Diff::ONE)]);
2224            if v2.stage.data.is_empty() {
2225                assert_eq!(v2.stage.bytes, 0, "a shipped stage holds no bytes");
2226                return i + 1;
2227            }
2228        }
2229        panic!("the stage never shipped");
2230    }
2231
2232    /// The stage ships by bytes, heap included, so wider rows ship after fewer updates.
2233    #[mz_ore::test]
2234    #[cfg_attr(miri, ignore)] // slow under Miri
2235    fn stage_ships_by_bytes() {
2236        let chunk_size = 16 * 1024;
2237        let narrow = updates_until_ship(8, chunk_size);
2238        let wide = updates_until_ship(1024, chunk_size);
2239
2240        let per_update = |width: usize| {
2241            std::mem::size_of::<String>() + width + std::mem::size_of::<(Timestamp, Diff)>()
2242        };
2243        assert_eq!(narrow, chunk_size.div_ceil(per_update(8)));
2244        assert_eq!(wide, chunk_size.div_ceil(per_update(1024)));
2245        assert!(wide < narrow);
2246    }
2247
2248    /// Run the same stepwise-drain workload through `CorrectionV1` and `CorrectionV2` and assert
2249    /// that they emit the same updates at every step.
2250    ///
2251    /// Models the `write_batches` operator catching up through many distinct timestamps: the
2252    /// desired input runs ahead, batches are written one timestamp at a time, and written updates
2253    /// come back negated through the persist feedback.
2254    #[mz_ore::test]
2255    // Columnation regions are not Stacked Borrows compliant: later pushes invalidate the
2256    // provenance of previously stored items under Miri.
2257    #[cfg_attr(miri, ignore)]
2258    fn equivalence_with_v1() {
2259        let sink_metrics = sink_metrics();
2260
2261        let mut v1 =
2262            CorrectionV1::<String>::new(sink_metrics.clone(), sink_metrics.for_worker(0), 1);
2263        let mut v2 = CorrectionV2::<String>::new(
2264            sink_metrics.clone(),
2265            sink_metrics.for_worker(0),
2266            None,
2267            3.0,
2268            8 * 1024,
2269        );
2270
2271        let num_ts = 50;
2272        let keys = 4;
2273
2274        // Upsert-style input: every timestamp updates each key, retracting the previous value.
2275        let batch = |t: u64| -> Vec<(String, Timestamp, Diff)> {
2276            (0..keys)
2277                .flat_map(|k| {
2278                    let addition = (format!("{k}-{t}"), Timestamp::from(t), Diff::ONE);
2279                    let retraction = t
2280                        .checked_sub(1)
2281                        .map(|p| (format!("{k}-{p}"), Timestamp::from(t), -Diff::ONE));
2282                    std::iter::once(addition).chain(retraction)
2283                })
2284                .collect()
2285        };
2286
2287        // Pre-fill both with all batches, like a catch-up where the input runs ahead.
2288        for t in 0..num_ts {
2289            v1.insert(&mut batch(t));
2290            v2.insert(&mut batch(t));
2291        }
2292
2293        // Drain stepwise, with persist feedback, comparing emissions.
2294        for t in 0..num_ts {
2295            let upper = Antichain::from_elem(Timestamp::from(t + 1));
2296
2297            let mut out1: Vec<_> = v1.updates_before(&upper).collect();
2298            let mut out2: Vec<_> = v2.updates_before(&upper).collect();
2299            out1.sort();
2300            out2.sort();
2301            assert_eq!(out1, out2, "diverged at t={t}");
2302
2303            v1.insert_negated(&mut out1.clone());
2304            v2.insert_negated(&mut out2);
2305            v1.advance_since(upper.clone());
2306            v2.advance_since(upper);
2307        }
2308
2309        // Compare the final state at the since.
2310        let upper = Antichain::from_elem(Timestamp::from(num_ts + 1));
2311        v1.consolidate_at_since();
2312        v2.consolidate_at_since();
2313        let mut out1: Vec<_> = v1.updates_before(&upper).collect();
2314        let mut out2: Vec<_> = v2.updates_before(&upper).collect();
2315        out1.sort();
2316        out2.sort();
2317        assert_eq!(out1, out2);
2318    }
2319
2320    /// A since jump across many distinct buffered timestamps must collapse them onto the since.
2321    #[mz_ore::test]
2322    // Columnation regions are not Stacked Borrows compliant: later pushes invalidate the
2323    // provenance of previously stored items under Miri.
2324    #[cfg_attr(miri, ignore)]
2325    fn since_jump() {
2326        let sink_metrics = sink_metrics();
2327        let mut v2 = CorrectionV2::<String>::new(
2328            sink_metrics.clone(),
2329            sink_metrics.for_worker(0),
2330            None,
2331            3.0,
2332            8 * 1024,
2333        );
2334
2335        let num_ts = 100;
2336        for t in 0..num_ts {
2337            v2.insert(&mut vec![
2338                (format!("a-{t}"), Timestamp::from(t), Diff::ONE),
2339                (format!("a-{t}"), Timestamp::from(t), -Diff::ONE),
2340                (format!("b-{t}"), Timestamp::from(t), Diff::ONE),
2341            ]);
2342        }
2343
2344        v2.advance_since(Antichain::from_elem(Timestamp::from(num_ts)));
2345        v2.consolidate_at_since();
2346
2347        let upper = Antichain::from_elem(Timestamp::from(num_ts + 1));
2348        let out: Vec<_> = v2.updates_before(&upper).collect();
2349        assert_eq!(out.len(), usize::try_from(num_ts).unwrap());
2350        assert!(
2351            out.iter()
2352                .all(|(_, t, r)| *t == Timestamp::from(num_ts) && *r == Diff::ONE)
2353        );
2354    }
2355
2356    /// The maintained record, size, and allocation totals must equal a walk over the chunks.
2357    ///
2358    /// Guards the delta accounting in [`Chain::push_chunk`] and [`Stage::account_since`]: a
2359    /// site that grows a chain or the stage without adjusting the totals shows up here.
2360    #[mz_ore::test]
2361    #[cfg_attr(miri, ignore)] // slow under Miri, and the bookkeeping it checks has no unsafe code
2362    fn metrics_totals_match_walk() {
2363        let sink_metrics = sink_metrics();
2364        let mut v2 = CorrectionV2::<String>::new(
2365            sink_metrics.clone(),
2366            sink_metrics.for_worker(0),
2367            None,
2368            3.0,
2369            8 * 1024,
2370        );
2371
2372        let num_ts = 200_u64;
2373        for t in 0..num_ts {
2374            v2.insert(&mut vec![
2375                (format!("a-{t}"), Timestamp::from(t), Diff::ONE),
2376                (format!("b-{t}"), Timestamp::from(t), Diff::ONE),
2377            ]);
2378        }
2379
2380        // Drain part of the buffer, so chains sit in `emitted` and `pending_low` as well as in
2381        // the bucket chain.
2382        let upper = Antichain::from_elem(Timestamp::from(num_ts / 2));
2383        let drained = v2.updates_before(&upper).count();
2384        assert_eq!(drained, usize::try_from(num_ts).expect("fits"));
2385
2386        // Advancing the since and consolidating there exercises the peel, split, and merge paths,
2387        // each of which retires and mints chains.
2388        v2.advance_since(Antichain::from_elem(Timestamp::from(num_ts / 2)));
2389        v2.consolidate_at_since();
2390        v2.insert(&mut vec![(
2391            "late".to_owned(),
2392            Timestamp::from(num_ts),
2393            Diff::ONE,
2394        )]);
2395
2396        let mut chains: Vec<&Chain<String>> = vec![&v2.emitted];
2397        chains.extend(&v2.pending_low);
2398        for bucket in v2.chain.buckets() {
2399            chains.extend(&bucket.chains);
2400        }
2401
2402        let mut size = v2.stage.bytes;
2403        let mut records = v2.stage.data.len();
2404        let mut allocations = usize::from(v2.stage.data.capacity() > 0);
2405        for chain in chains {
2406            size += chain.chunks.iter().map(Chunk::size).sum::<usize>();
2407            records += chain.chunks.iter().map(Chunk::len).sum::<usize>();
2408            allocations += chain.chunks.len();
2409        }
2410
2411        assert!(size > 0, "workload must leave chunks behind");
2412        assert_eq!(v2.prev_size.size, size);
2413        assert_eq!(v2.prev_size.capacity, size);
2414        assert_eq!(v2.prev_size.allocations, allocations);
2415        assert_eq!(v2.prev_update_count, records);
2416    }
2417
2418    #[mz_ore::test]
2419    #[cfg_attr(miri, ignore)] // slow under Miri
2420    fn bucket_merges_deepen_and_reads_reset_depth() {
2421        let sink_metrics = sink_metrics();
2422        let mut v2 = CorrectionV2::<String>::new(
2423            sink_metrics.clone(),
2424            sink_metrics.for_worker(0),
2425            None,
2426            3.0,
2427            8 * 1024,
2428        );
2429
2430        // Many stage flushes at one time land in one bucket, whose chain merges deepen.
2431        let mut i = 0;
2432        let deepest = loop {
2433            v2.insert(&mut vec![(
2434                format!("{i:08}"),
2435                Timestamp::from(0),
2436                Diff::ONE,
2437            )]);
2438            i += 1;
2439            let chains: Vec<&Chain<String>> = v2.chain.buckets().flat_map(|b| &b.chains).collect();
2440            for chain in &chains {
2441                let max = chain.chunks.iter().map(|c| c.depth).max().unwrap_or(0);
2442                assert_eq!(chain.depth, max, "a chain's depth is its deepest chunk's");
2443            }
2444            let deepest = chains.iter().map(|c| c.depth).max().unwrap_or(0);
2445            if deepest >= 2 && chains.len() >= 2 {
2446                break deepest;
2447            }
2448            assert!(i < 1_000_000, "bucket merges never deepened");
2449        };
2450        assert!(deepest >= 2);
2451
2452        let upper = Antichain::from_elem(Timestamp::from(1));
2453        let read = v2.updates_before(&upper).count();
2454        assert_eq!(read, i);
2455        assert!(
2456            v2.emitted.chunks.iter().all(|c| c.depth == 0),
2457            "a read writes the youngest generation",
2458        );
2459    }
2460
2461    /// Every chain announced to introspection is retired by the time the buffer drops, the
2462    /// stage's included, whether or not the stage still holds updates.
2463    #[mz_ore::test]
2464    #[cfg_attr(miri, ignore)] // slow under Miri
2465    fn logged_chains_balance_on_drop() {
2466        use crate::sink::correction::LoggingEvent;
2467
2468        for staged in [false, true] {
2469            let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2470            let sink_metrics = sink_metrics();
2471            let mut v2 = CorrectionV2::<String>::new(
2472                sink_metrics.clone(),
2473                sink_metrics.for_worker(0),
2474                Some(ChannelLogging::new(tx)),
2475                3.0,
2476                // One byte ships every insert, one mebibyte ships none of them.
2477                if staged { 1 << 20 } else { 1 },
2478            );
2479            for t in 0..100_u64 {
2480                v2.insert(&mut vec![(format!("{t}"), Timestamp::from(t), Diff::ONE)]);
2481            }
2482            assert_eq!(v2.stage.data.is_empty(), !staged);
2483            drop(v2);
2484
2485            let (mut batches, mut records) = (0_isize, 0_isize);
2486            while let Ok(event) = rx.try_recv() {
2487                match event {
2488                    LoggingEvent::ChainCreated(n) => {
2489                        batches += 1;
2490                        records += isize::try_from(n).expect("fits");
2491                    }
2492                    LoggingEvent::ChainDropped(n) => {
2493                        batches -= 1;
2494                        records -= isize::try_from(n).expect("fits");
2495                    }
2496                    _ => {}
2497                }
2498            }
2499            assert_eq!((batches, records), (0, 0), "staged={staged}");
2500        }
2501    }
2502
2503    /// Reads must not observe updates at or beyond their `upper`, even when the `upper` is not
2504    /// beyond the `since`.
2505    #[mz_ore::test]
2506    // Columnation regions are not Stacked Borrows compliant: later pushes invalidate the
2507    // provenance of previously stored items under Miri.
2508    #[cfg_attr(miri, ignore)]
2509    fn upper_not_beyond_since() {
2510        let sink_metrics = sink_metrics();
2511        let mut v2 = CorrectionV2::<String>::new(
2512            sink_metrics.clone(),
2513            sink_metrics.for_worker(0),
2514            None,
2515            3.0,
2516            8 * 1024,
2517        );
2518
2519        v2.insert(&mut vec![(
2520            "a".to_owned(),
2521            Timestamp::from(5_u64),
2522            Diff::ONE,
2523        )]);
2524        v2.advance_since(Antichain::from_elem(Timestamp::from(10_u64)));
2525
2526        // The update logically lives at time 10 now, so a read before 7 must be empty.
2527        let upper = Antichain::from_elem(Timestamp::from(7_u64));
2528        assert_eq!(v2.updates_before(&upper).count(), 0);
2529
2530        // A read before 11 must emit it, advanced to the since.
2531        let upper = Antichain::from_elem(Timestamp::from(11_u64));
2532        let out: Vec<_> = v2.updates_before(&upper).collect();
2533        assert_eq!(
2534            out,
2535            vec![("a".to_owned(), Timestamp::from(10_u64), Diff::ONE)]
2536        );
2537    }
2538
2539    /// One pool shared by every test in the module. A pool reserves a large slab of address
2540    /// space, so one per test exhausts the VM map under parallel test threads.
2541    fn test_pool() -> mz_ore::pool::Pool {
2542        static POOL: std::sync::OnceLock<mz_ore::pool::Pool> = std::sync::OnceLock::new();
2543        POOL.get_or_init(|| mz_ore::pool::Pool::new().expect("pool reservation"))
2544            .clone()
2545    }
2546
2547    /// Route this thread's chunk spills through the test pool for the duration of `f`.
2548    ///
2549    /// Without an override no pool is installed and every chunk stays resident, so the tests
2550    /// would not exercise [`Chunk::body`]'s read-back at all. The override is thread-scoped,
2551    /// so concurrently running tests do not race on it.
2552    fn with_spill_pool<R>(f: impl FnOnce() -> R) -> R {
2553        chunk::set_spill_override(Some(test_pool()));
2554        let result = f();
2555        chunk::set_spill_override(None);
2556        result
2557    }
2558
2559    /// Assert every chunk in the chain went to the pool, so a read must fetch it back.
2560    fn assert_all_spilled<D: Data>(chain: &Chain<D>) {
2561        assert!(
2562            chain
2563                .chunks
2564                .iter()
2565                .all(|c| c.pooled.lock().expect("not poisoned").is_some()),
2566            "expected every chunk to spill, got {:?}",
2567            chain.chunks,
2568        );
2569    }
2570
2571    /// Build a chain crossing the mint boundary while every chunk spills, then assert `iter()`
2572    /// (the read path behind `updates_before`) reads each chunk back, roundtrips values, order,
2573    /// and diffs, and leaves every body in the pool.
2574    #[mz_ore::test]
2575    #[cfg_attr(miri, ignore)] // the pool's mapped regions are unsupported under miri
2576    fn iter_roundtrips_through_pool() {
2577        let count = 200_000_u64;
2578        with_spill_pool(|| {
2579            let mut builder = ChainBuilder::<i64>::default();
2580            for i in 0..count {
2581                let d = i64::try_from(i).expect("fits");
2582                builder.push_owned(&(d, Timestamp::new(i), Diff::ONE));
2583            }
2584            let chain = builder.finish();
2585            assert!(chain.chunks.len() > 1, "expected multiple minted chunks");
2586            assert_eq!(chain.update_count, usize::try_from(count).expect("fits"));
2587            assert_all_spilled(&chain);
2588
2589            let mut expected = 0_u64;
2590            for (d, t, r) in chain.iter() {
2591                assert_eq!(d, i64::try_from(expected).expect("fits"));
2592                assert_eq!(t, Timestamp::new(expected));
2593                assert_eq!(r, Diff::ONE);
2594                expected += 1;
2595            }
2596            assert_eq!(expected, count);
2597            // The read is scoped: every body is still in the pool and none was materialized.
2598            assert_all_spilled(&chain);
2599            assert!(chain.chunks.iter().all(|c| c.resident.get().is_none()));
2600        });
2601    }
2602
2603    /// Drive a [`Cursor`] over a spilled, multi-chunk chain to completion (the access pattern
2604    /// merges use). Each step reads the front chunk back via [`Chunk::body`]; assert the
2605    /// cursor yields every update in order.
2606    #[mz_ore::test]
2607    #[cfg_attr(miri, ignore)] // the pool's mapped regions are unsupported under miri
2608    fn cursor_steps_through_pool() {
2609        let count = 200_000_u64;
2610        with_spill_pool(|| {
2611            let mut builder = ChainBuilder::<i64>::default();
2612            for i in 0..count {
2613                let d = i64::try_from(i).expect("fits");
2614                builder.push_owned(&(d, Timestamp::new(i), Diff::ONE));
2615            }
2616            let chain = builder.finish();
2617            assert!(chain.chunks.len() > 1, "expected multiple minted chunks");
2618            assert_all_spilled(&chain);
2619
2620            let mut rest = chain.into_cursor();
2621            let mut expected = 0_u64;
2622            while let Some(cursor) = rest.take() {
2623                let (d, t, r) = cursor.get();
2624                assert_eq!(i64::into_owned(d), i64::try_from(expected).expect("fits"));
2625                assert_eq!(t, Timestamp::new(expected));
2626                assert_eq!(r, Diff::ONE);
2627                expected += 1;
2628                rest = cursor.step();
2629            }
2630            assert_eq!(expected, count);
2631        });
2632    }
2633}