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(), ®istry);
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}