Skip to main content

mz_row_spine/
lib.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//! Types and traits in support of containers for row-encoded byte slices.
11//!
12//! This includes the vanilla `bytes_container` that holds byte slices in contiguous
13//! allocations, as well as a `dictionary` encoding wrapper that is able to rewrite
14//! the byte slices to use spare tags in each column to reference common values.
15
16pub use self::arc_batch::{ArcBatch, ArcBuilder};
17pub use self::dictionary::DatumContainer;
18pub use self::dictionary::DatumSeq;
19pub use self::dictionary::builders::RowRowColPagedState;
20pub use self::offset_opt::OffsetOptimized;
21pub use self::spines::{
22    ArcOrdKeyBuilder, ArcOrdKeySpine, ArcOrdValBuilder, ArcOrdValSpine, FundedValRowSpine,
23    RowBatcher, RowBuilder, RowRowBatcher, RowRowBuilder, RowRowColPagedBuilder, RowRowSpine,
24    RowSpine, RowValBatcher, RowValBuilder, RowValSpine, ValRowBatcher, ValRowBuilder,
25    ValRowColPagedBuilder, ValRowSpine,
26};
27
28mod arc_batch;
29
30use differential_dataflow::trace::implementations::OffsetList;
31
32/// Enable per-column dictionary compression in row containers.
33pub static DICTIONARY_COMPRESSION: std::sync::atomic::AtomicBool =
34    std::sync::atomic::AtomicBool::new(false);
35
36/// Spines specialized to contain `Row` types in keys and values.
37mod spines {
38    use columnation::Columnation;
39    use differential_dataflow::trace::implementations::BatchContainer;
40    use differential_dataflow::trace::implementations::Layout;
41    use differential_dataflow::trace::implementations::Update;
42    use differential_dataflow::trace::implementations::Vector;
43    use differential_dataflow::trace::implementations::merge_batcher::MergeBatcher;
44    use differential_dataflow::trace::implementations::ord_neu::{
45        OrdKeyBatch, OrdKeyBuilder, OrdValBatch, OrdValBuilder,
46    };
47    use differential_dataflow::trace::implementations::spine_fueled::Spine;
48    use mz_repr::Row;
49    use mz_timely_util::columnation::{ColInternalMerger, ColumnationStack};
50
51    use crate::arc_batch::{ArcBatch, ArcBuilder};
52    use crate::{DatumContainer, OffsetOptimized};
53
54    /// Batcher matching `mz_compute::typedefs::KeyValBatcher`, redeclared
55    /// locally so this crate does not need to depend on `mz_compute`.
56    type KeyValBatcher<K, V, T, D> = MergeBatcher<ColInternalMerger<(K, V), T, D>>;
57    type KeyBatcher<K, T, D> = KeyValBatcher<K, (), T, D>;
58
59    pub type RowRowSpine<T, R> = Spine<ArcBatch<OrdValBatch<RowRowLayout<((Row, Row), T, R)>>>>;
60    pub type RowRowBatcher<T, R> = KeyValBatcher<Row, Row, T, R>;
61    pub type RowRowBuilder<T, R> = ArcBuilder<crate::dictionary::builders::RowRowBuilder<T, R>>;
62
63    /// `RowRowBuilder` variant that consumes [`ColumnBody`] chunks. Pairs with
64    /// any batcher whose chains are bodies, spillable ([`AccountedChunkBatcher`]
65    /// behind an [`UnchunkBuilder`]) or resident ([`Col2ValColBatcher`]) alike,
66    /// so the `Paged` in the name records where it started rather than a
67    /// restriction. Installs a dictionary codec at seal time, gathering
68    /// statistics from the sealed chain, so columnar arrangements compress on
69    /// the same footing as the columnation-fed [`RowRowBuilder`].
70    ///
71    /// [`AccountedChunkBatcher`]: mz_timely_util::columnar::chunk::AccountedChunkBatcher
72    /// [`Col2ValColBatcher`]: mz_timely_util::columnar::Col2ValColBatcher
73    /// [`ColumnBody`]: mz_timely_util::columnar::body::ColumnBody
74    /// [`UnchunkBuilder`]: mz_timely_util::columnar::chunk::UnchunkBuilder
75    pub type RowRowColPagedBuilder<T, R> =
76        ArcBuilder<crate::dictionary::builders::RowRowColPagedBuilder<T, R>>;
77
78    pub type RowValSpine<V, T, R> = Spine<ArcBatch<OrdValBatch<RowValLayout<((Row, V), T, R)>>>>;
79    pub type RowValBatcher<V, T, R> = KeyValBatcher<Row, V, T, R>;
80    pub type RowValBuilder<V, T, R> =
81        ArcBuilder<crate::dictionary::builders::RowValBuilder<V, T, R>>;
82
83    /// Key-only `Row` spine. `DC` is the diff container, see `RowLayout`.
84    pub type RowSpine<T, R, DC = ColumnationStack<R>> =
85        Spine<ArcBatch<OrdKeyBatch<RowLayout<((Row, ()), T, R), DC>>>>;
86    pub type RowBatcher<T, R> = KeyBatcher<Row, T, R>;
87    pub type RowBuilder<T, R, DC = ColumnationStack<R>> =
88        ArcBuilder<crate::dictionary::builders::RowBuilder<T, R, DC>>;
89
90    pub type ValRowSpine<K, T, R> = Spine<ArcBatch<OrdValBatch<ValRowLayout<((K, Row), T, R)>>>>;
91    /// A `ValRowSpine` whose optional exertion is funded by its input.
92    pub type FundedValRowSpine<K, T, R> =
93        mz_timely_util::funded_spine::Spine<ArcBatch<OrdValBatch<ValRowLayout<((K, Row), T, R)>>>>;
94    pub type ValRowBatcher<K, T, R> = KeyValBatcher<K, Row, T, R>;
95    pub type ValRowBuilder<K, T, R> =
96        ArcBuilder<crate::dictionary::builders::ValRowBuilder<K, T, R>>;
97
98    /// `ValRowBuilder` variant that consumes [`Column`] chunks. Pairs with
99    /// `Col2ValPagedBatcher<K, Row, T, R>` for the spillable arrange path where
100    /// keys are arbitrary `Columnar` values (e.g. `UpsertKey`) and values are
101    /// packed `Row` bytes. Installs a dictionary codec on the value container at
102    /// seal time, gathering statistics from the sealed `Column` chain; keys are
103    /// not `Row`-shaped and so are left uncompressed.
104    ///
105    /// [`Column`]: mz_timely_util::columnar::Column
106    pub type ValRowColPagedBuilder<K, T, R> =
107        ArcBuilder<crate::dictionary::builders::ValRowColPagedBuilder<K, T, R>>;
108
109    /// A generic `Arc`-backed key/value spine, for callers outside `mz_compute` that need an
110    /// arrangement over non-`Row`-specialized types. The `Arc` handle rides on the local
111    /// [`ArcBatch`] newtype, so no differential-side `Arc` batch impls are required.
112    pub type ArcOrdValSpine<K, V, T, R> = Spine<ArcBatch<OrdValBatch<Vector<((K, V), T, R)>>>>;
113    /// Generic `Arc`-backed key-only spine. See [`ArcOrdValSpine`].
114    pub type ArcOrdKeySpine<K, T, R> = Spine<ArcBatch<OrdKeyBatch<Vector<((K, ()), T, R)>>>>;
115    /// Builder pairing with [`ArcOrdValSpine`].
116    pub type ArcOrdValBuilder<K, V, T, R> =
117        ArcBuilder<OrdValBuilder<Vector<((K, V), T, R)>, Vec<((K, V), T, R)>>>;
118    /// Builder pairing with [`ArcOrdKeySpine`].
119    pub type ArcOrdKeyBuilder<K, T, R> =
120        ArcBuilder<OrdKeyBuilder<Vector<((K, ()), T, R)>, Vec<((K, ()), T, R)>>>;
121
122    /// A layout based on timely stacks
123    pub struct RowRowLayout<U: Update<Key = Row, Val = Row>> {
124        phantom: std::marker::PhantomData<U>,
125    }
126    pub struct RowValLayout<U: Update<Key = Row>> {
127        phantom: std::marker::PhantomData<U>,
128    }
129    /// Layout for key-only `Row` updates. `DC` is the diff `BatchContainer`, a
130    /// columnation stack by default.
131    pub struct RowLayout<U, DC = ColumnationStack<<U as Update>::Diff>>
132    where
133        U: Update<Key = Row, Val = ()>,
134    {
135        phantom: std::marker::PhantomData<(U, DC)>,
136    }
137    /// Mirror of [`RowValLayout`] with the roles swapped: arbitrary `Columnation`
138    /// keys with `Row` values stored as packed bytes in a [`DatumContainer`].
139    pub struct ValRowLayout<U: Update<Val = Row>> {
140        phantom: std::marker::PhantomData<U>,
141    }
142
143    impl<U: Update<Key = Row, Val = Row>> Layout for RowRowLayout<U>
144    where
145        U::Time: Columnation,
146        U::Diff: Columnation,
147    {
148        type KeyContainer = DatumContainer;
149        type ValContainer = DatumContainer;
150        type TimeContainer = ColumnationStack<U::Time>;
151        type DiffContainer = ColumnationStack<U::Diff>;
152        type OffsetContainer = OffsetOptimized;
153    }
154    impl<U: Update<Key = Row>> Layout for RowValLayout<U>
155    where
156        U::Val: Columnation,
157        U::Time: Columnation,
158        U::Diff: Columnation,
159    {
160        type KeyContainer = DatumContainer;
161        type ValContainer = ColumnationStack<U::Val>;
162        type TimeContainer = ColumnationStack<U::Time>;
163        type DiffContainer = ColumnationStack<U::Diff>;
164        type OffsetContainer = OffsetOptimized;
165    }
166    impl<U: Update<Key = Row, Val = ()>, DC> Layout for RowLayout<U, DC>
167    where
168        U::Time: Columnation,
169        DC: BatchContainer<Owned = U::Diff>,
170    {
171        type KeyContainer = DatumContainer;
172        type ValContainer = ColumnationStack<()>;
173        type TimeContainer = ColumnationStack<U::Time>;
174        type DiffContainer = DC;
175        type OffsetContainer = OffsetOptimized;
176    }
177    impl<U: Update<Val = Row>> Layout for ValRowLayout<U>
178    where
179        U::Key: Columnation,
180        U::Time: Columnation,
181        U::Diff: Columnation,
182    {
183        type KeyContainer = ColumnationStack<U::Key>;
184        type ValContainer = DatumContainer;
185        type TimeContainer = ColumnationStack<U::Time>;
186        type DiffContainer = ColumnationStack<U::Diff>;
187        type OffsetContainer = OffsetOptimized;
188    }
189}
190
191#[cfg(test)]
192mod tests {
193    use crate::DatumContainer;
194    use crate::spines::{RowLayout, RowRowLayout, RowValLayout};
195    use differential_dataflow::trace::implementations::BatchContainer;
196    use differential_dataflow::trace::implementations::ord_neu::{OrdKeyBatch, OrdValBatch};
197    use mz_repr::adt::date::Date;
198    use mz_repr::adt::interval::Interval;
199    use mz_repr::{Datum, Diff, Row, SqlScalarType, Timestamp};
200    use mz_timely_util::columnation::ColumnationStack;
201
202    fn assert_send_sync<T: Send + Sync>() {}
203
204    /// The batch types backing our spines must stay `Send + Sync`, so that batches
205    /// can be shared across threads (for example behind an `Arc`) to serve reads
206    /// from outside the worker that maintains the trace. This holds because the
207    /// backing containers bottom out in `Vec`s, lgalloc regions, and `CompactBytes`,
208    /// all of which are thread-safe.
209    #[mz_ore::test]
210    fn batches_are_send_sync() {
211        assert_send_sync::<OrdValBatch<RowRowLayout<((Row, Row), Timestamp, Diff)>>>();
212        assert_send_sync::<OrdValBatch<RowValLayout<((Row, Row), Timestamp, Diff)>>>();
213        assert_send_sync::<OrdKeyBatch<RowLayout<((Row, ()), Timestamp, Diff)>>>();
214        assert_send_sync::<ColumnationStack<((Row, Row), Timestamp, Diff)>>();
215    }
216
217    /// A chain of chunks seals with the chain's dictionary codecs installed,
218    /// so its batch compresses like one sealed from columns directly. A seal
219    /// that skips the install stores raw bytes and, per
220    /// `DatumContainer::promote_stats_to_codec`, drags every merge it joins
221    /// down with it.
222    #[mz_ore::test]
223    #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
224    fn chunked_seal_installs_codecs() {
225        use std::sync::atomic::Ordering;
226
227        use differential_dataflow::trace::implementations::ord_neu::OrdValBatch;
228        use differential_dataflow::trace::{Builder, Description};
229        use mz_timely_util::columnar::body::ColumnBody;
230        use mz_timely_util::columnar::chunk::{ColumnChunk, UnchunkBuilder};
231        use timely::container::PushInto;
232        use timely::progress::{Antichain, Timestamp as _};
233
234        use crate::ArcBatch;
235
236        // Gate the dictionary path on. Safe for other tests: the flag only
237        // controls whether codecs are built, never decode results.
238        crate::DICTIONARY_COMPRESSION.store(true, Ordering::Relaxed);
239
240        // Low-cardinality rows, well under `STATS_THRESHOLD` (64Ki pushes), so
241        // a codec can only come from the seal-time install.
242        let time = Timestamp::minimum();
243        let updates: Vec<((Row, Row), Timestamp, i64)> = (0..2_000i64)
244            .map(|i| {
245                (
246                    (
247                        Row::pack_slice(&[Datum::Int64(i)]),
248                        Row::pack_slice(&[Datum::String("a repeated string value")]),
249                    ),
250                    time,
251                    1i64,
252                )
253            })
254            .collect();
255
256        // Cut the sorted run the way a merge batcher's chain is cut.
257        let columns = || -> Vec<ColumnBody<((Row, Row), Timestamp, i64)>> {
258            updates
259                .chunks(250)
260                .map(|part| {
261                    let mut column: ColumnBody<((Row, Row), Timestamp, i64)> = Default::default();
262                    for update in part {
263                        column.push_into(update);
264                    }
265                    column
266                })
267                .collect()
268        };
269        let description = || {
270            Description::new(
271                Antichain::from_elem(time),
272                Antichain::new(),
273                Antichain::from_elem(time),
274            )
275        };
276
277        type Paged = crate::RowRowColPagedBuilder<Timestamp, i64>;
278        type Chunked = UnchunkBuilder<Paged, (Row, Row), Timestamp, i64>;
279
280        let mut chain = columns();
281        let from_columns = <Paged as Builder>::seal(&mut chain, description());
282
283        let mut chain: Vec<ColumnChunk<(Row, Row), Timestamp, i64>> =
284            columns().into_iter().map(ColumnChunk::from_body).collect();
285        let from_chunks = <Chunked as Builder>::seal(&mut chain, description());
286
287        assert!(
288            from_chunks.0.storage.keys.has_codec() && from_chunks.0.storage.vals.vals.has_codec(),
289            "the sealed batch carries the chain's codecs"
290        );
291
292        // Same updates and same cuts, so the chunked seal has to reach the
293        // encoded size the column seal does.
294        let heap = |c: &DatumContainer| {
295            let mut size = 0;
296            c.heap_size(|_, cap| size += cap);
297            size
298        };
299        let encoded = |b: &ArcBatch<OrdValBatch<RowRowLayout<((Row, Row), Timestamp, i64)>>>| {
300            heap(&b.0.storage.keys) + heap(&b.0.storage.vals.vals)
301        };
302
303        assert!(
304            encoded(&from_chunks) <= encoded(&from_columns),
305            "chunked seal should compress like the column seal: chunks={} columns={}",
306            encoded(&from_chunks),
307            encoded(&from_columns),
308        );
309    }
310
311    #[mz_ore::test]
312    #[cfg_attr(miri, ignore)] // unsupported operation: integer-to-pointer casts and `ptr::with_exposed_provenance` are not supported
313    fn test_round_trip() {
314        fn round_trip(datums: Vec<Datum>) {
315            let row = Row::pack(datums.clone());
316
317            let mut container = DatumContainer::with_capacity(row.byte_len());
318            container.push_own(&row);
319
320            // When run under miri this catches undefined bytes written to data
321            // eg by calling push_copy! on a type which contains undefined padding values
322            println!("{:?}", container.index(0).iter.data);
323
324            let datums2 = container.index(0).collect::<Vec<_>>();
325            assert_eq!(datums, datums2);
326        }
327
328        round_trip(vec![]);
329        round_trip(
330            SqlScalarType::enumerate()
331                .iter()
332                .flat_map(|r#type| r#type.interesting_datums())
333                .collect(),
334        );
335        round_trip(vec![
336            Datum::Null,
337            Datum::Null,
338            Datum::False,
339            Datum::True,
340            Datum::Int16(-21),
341            Datum::Int32(-42),
342            Datum::Int64(-2_147_483_648 - 42),
343            Datum::UInt8(0),
344            Datum::UInt8(1),
345            Datum::UInt16(0),
346            Datum::UInt16(1),
347            Datum::UInt16(1 << 8),
348            Datum::UInt32(0),
349            Datum::UInt32(1),
350            Datum::UInt32(1 << 8),
351            Datum::UInt32(1 << 16),
352            Datum::UInt32(1 << 24),
353            Datum::UInt64(0),
354            Datum::UInt64(1),
355            Datum::UInt64(1 << 8),
356            Datum::UInt64(1 << 16),
357            Datum::UInt64(1 << 24),
358            Datum::UInt64(1 << 32),
359            Datum::UInt64(1 << 40),
360            Datum::UInt64(1 << 48),
361            Datum::UInt64(1 << 56),
362            Datum::Date(Date::from_pg_epoch(365 * 45 + 21).unwrap()),
363            Datum::Interval(Interval {
364                months: 312,
365                ..Default::default()
366            }),
367            Datum::Interval(Interval::new(0, 0, 1_012_312)),
368            Datum::Bytes(&[]),
369            Datum::Bytes(&[0, 2, 1, 255]),
370            Datum::String(""),
371            Datum::String("العَرَبِيَّة"),
372        ]);
373    }
374
375    /// Exercises the *compressed* encode→decode paths, which the dyncfg-gated
376    /// `test_round_trip` never reaches (it installs no codec). We drive the codec
377    /// directly: observe a sample, build a codec via both `new_from([c1, c2])`
378    /// (the merge path) and `new_safe` (the safe-tag path), then round-trip every
379    /// row through it. We additionally assert the dictionary actually engaged, so
380    /// the test keeps covering the compressed branch rather than silently
381    /// degrading to raw fall-through.
382    #[mz_ore::test]
383    #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
384    fn test_codec_round_trip() {
385        use crate::row_codec::ColumnsCodec;
386
387        // Rows with a small set of repeated, multi-byte string values, so the
388        // dictionary installs entries (MisraGries keeps values with len > 1 and
389        // count > 1). Mixing in an integer column exercises the raw fall-through
390        // (and thus the new soundness `debug_assert`) alongside dictionary hits.
391        let values = ["apple", "banana", "cherry"];
392        let rows: Vec<Row> = (0..3_000)
393            .map(|i| {
394                Row::pack_slice(&[
395                    Datum::String(values[i % values.len()]),
396                    Datum::Int64(i64::try_from(i).unwrap()),
397                    Datum::String(values[(i / 7) % values.len()]),
398                ])
399            })
400            .collect();
401
402        // Accumulate statistics in two independent observers, so the merge in
403        // `new_from([&stats1, &stats2])` is actually exercised.
404        let mut stats1 = ColumnsCodec::default();
405        let mut stats2 = ColumnsCodec::default();
406        let mut scratch = Vec::new();
407        for (i, row) in rows.iter().enumerate() {
408            scratch.clear();
409            let stats = if i % 2 == 0 { &mut stats1 } else { &mut stats2 };
410            stats.encode(ColumnsCodec::borrow_row(row), &mut scratch);
411        }
412
413        let merged = ColumnsCodec::new_from([&stats1, &stats2]);
414        let safe = stats1.new_safe();
415        for mut codec in [merged, safe] {
416            let mut compressed_any = false;
417            for row in &rows {
418                let mut buf = Vec::new();
419                codec.encode(ColumnsCodec::borrow_row(row), &mut buf);
420
421                let decoded = codec.decode(&buf).collect::<Vec<_>>();
422                let expected = ColumnsCodec::borrow_row(row).collect::<Vec<_>>();
423                assert_eq!(decoded, expected, "round-trip mismatch for {row:?}");
424
425                compressed_any |= buf.len() < row.data().len();
426            }
427            assert!(
428                compressed_any,
429                "dictionary never engaged; test no longer covers the compressed path",
430            );
431        }
432    }
433
434    /// Regression test for a dictionary-codec soundness bug in the safe-install
435    /// path (`new_safe`), reachable with the paged batcher enabled.
436    ///
437    /// A from-scratch container stores its pre-install rows *raw* while gathering
438    /// statistics, then installs a *safe* codec via `new_safe`. `new_safe` used to
439    /// discard the first-byte bitmap gathered over those raw rows. That bitmap is
440    /// soundness-critical: a later `new_from` merge consults it to decide which
441    /// one-byte tags are free to hand out as dictionary keys. With the bitmap
442    /// dropped, the merge could assign a dictionary tag equal to a raw datum's
443    /// first byte, after which `decode` resolves that literal datum to the
444    /// dictionary entry — returning the wrong value.
445    ///
446    /// We drive the lifecycle directly: observe short strings (first byte
447    /// `StringTiny`) into the pre-install statistics, install a safe codec, then
448    /// feed it many distinct *long* strings (first byte `StringShort`)
449    /// post-install so the merge has heavy hitters to compress. Merging via
450    /// `new_from` and re-encoding the short strings then exercises the raw
451    /// fall-through whose first byte the merge must not have claimed as a tag.
452    /// Before the fix the `StringTiny` tag was handed out and the round-trip
453    /// produced a long string (and tripped `encode`'s soundness `debug_assert`).
454    #[mz_ore::test]
455    #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
456    fn test_safe_codec_merge_bitmap_carryover() {
457        use crate::row_codec::ColumnsCodec;
458
459        // Short strings: length < 256, so they encode with the `StringTiny` tag.
460        // Unique, so MisraGries never makes them dictionary entries; they always
461        // fall through raw, exposing their first byte.
462        let short_rows: Vec<Row> = (0..256)
463            .map(|i| Row::pack_slice(&[Datum::String(&format!("s{i}"))]))
464            .collect();
465        // Long strings: length >= 256, so they encode with the `StringShort` tag —
466        // a *different* first byte than the short strings. Distinct values, each
467        // repeated, so the post-install codec accrues many heavy hitters and the
468        // merge assigns dictionary tags across the low byte range, reaching the
469        // short strings' `StringTiny` tag unless the bitmap reserves it.
470        let long_values: Vec<String> = (0..64).map(|i| format!("{i:0>300}")).collect();
471
472        // Pre-install statistics observe only the short strings' first bytes.
473        let mut stats = ColumnsCodec::default();
474        let mut scratch = Vec::new();
475        for row in &short_rows {
476            scratch.clear();
477            stats.encode(ColumnsCodec::borrow_row(row), &mut scratch);
478        }
479
480        // Install a safe codec, then feed it the long strings post-install so it
481        // accrues heavy hitters (and observes only the `StringShort` first byte).
482        let mut safe = stats.new_safe();
483        for _ in 0..8 {
484            for v in &long_values {
485                let row = Row::pack_slice(&[Datum::String(v)]);
486                scratch.clear();
487                safe.encode(ColumnsCodec::borrow_row(&row), &mut scratch);
488            }
489        }
490
491        // Merge, then round-trip the short strings. With the bitmap carried over,
492        // no dictionary tag collides with the short strings' first byte; without
493        // it, one does.
494        let mut merged = ColumnsCodec::new_from([&safe]);
495        for row in &short_rows {
496            let mut buf = Vec::new();
497            merged.encode(ColumnsCodec::borrow_row(row), &mut buf);
498            let decoded = merged.decode(&buf).collect::<Vec<_>>();
499            let expected = ColumnsCodec::borrow_row(row).collect::<Vec<_>>();
500            assert_eq!(decoded, expected, "round-trip mismatch for {row:?}");
501        }
502    }
503
504    /// Confirms the structural assumption underpinning `SAFE_TAG_BASE`: every
505    /// datum the row format produces encodes with a first byte strictly less
506    /// than `SAFE_TAG_BASE`. If `mz_repr` ever introduces a tag that crosses
507    /// the boundary, `DictionaryCodec::new_safe` would assign a dictionary tag
508    /// that collides with a literal datum first-byte, breaking decoding.
509    #[mz_ore::test]
510    fn test_safe_tag_base() {
511        use crate::row_codec::SAFE_TAG_BASE;
512        let check = |datum: Datum| {
513            let row = Row::pack_slice(&[datum]);
514            let data = row.data();
515            assert!(!data.is_empty(), "empty encoding for {datum:?}");
516            assert!(
517                data[0] < SAFE_TAG_BASE,
518                "datum {datum:?} encodes with first byte {} >= SAFE_TAG_BASE ({}); \
519                 a new row tag has crossed the safe boundary",
520                data[0],
521                SAFE_TAG_BASE,
522            );
523        };
524        for ty in SqlScalarType::enumerate().iter() {
525            for datum in ty.interesting_datums() {
526                check(datum);
527            }
528        }
529    }
530
531    /// A batch built via the builder's `push`/`done` path (as the `reduce` operator
532    /// does) that stays under `STATS_THRESHOLD` never installs a codec at build time.
533    /// `done` now promotes the gathered statistics into the codec slot, so the batch
534    /// carries a codec + heavy-hitter summary and does not poison a later merge.
535    ///
536    /// This drives that container lifecycle directly: gather raw (well under the
537    /// threshold), promote at "done", then merge two such containers the way a spine
538    /// compaction does. With promotion the merge takes the `new_from` path and
539    /// compresses; without it both inputs are codec-less and the merge stays raw.
540    /// Every merged row must still round-trip.
541    #[mz_ore::test]
542    #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
543    fn push_done_promotion_avoids_merge_poison() {
544        use std::sync::atomic::Ordering;
545        use timely::container::PushInto;
546
547        // Gate the dictionary path on. Safe for other tests: the flag only controls
548        // whether `DatumContainer` gathers stats; it never changes decode results.
549        crate::DICTIONARY_COMPRESSION.store(true, Ordering::Relaxed);
550
551        // Low-cardinality rows, well under `STATS_THRESHOLD` (64Ki): a repeated
552        // multi-byte string the dictionary compresses, plus an integer column that
553        // exercises raw fall-through.
554        let rows: Vec<Row> = (0..2_000i64)
555            .map(|i| {
556                Row::pack_slice(&[
557                    Datum::Int64(i % 8),
558                    Datum::String("a repeated string value"),
559                ])
560            })
561            .collect();
562
563        // Build a container the way the push/done path does: gather raw without ever
564        // crossing `STATS_THRESHOLD`, optionally promoting at "done".
565        let build = |promote: bool| {
566            let mut c = DatumContainer::with_capacity(rows.len());
567            for row in &rows {
568                c.push_into(row);
569            }
570            if promote {
571                c.promote_stats_to_codec();
572            }
573            c
574        };
575
576        // Merge two containers as a spine compaction does: allocate via
577        // `merge_capacity`, then copy every row through.
578        let merge = |a: &DatumContainer, b: &DatumContainer| {
579            let mut m = DatumContainer::merge_capacity(a, b);
580            for i in 0..a.len() {
581                m.push_into(a.index(i));
582            }
583            for i in 0..b.len() {
584                m.push_into(b.index(i));
585            }
586            m
587        };
588
589        let heap = |c: &DatumContainer| {
590            let mut size = 0;
591            c.heap_size(|_, cap| size += cap);
592            size
593        };
594
595        // Codec-less inputs (no promotion): the merge cannot `new_from` and stays raw.
596        let poisoned = merge(&build(false), &build(false));
597        // Promoted inputs carry a codec + summary: the merge `new_from`s and compresses.
598        let compressed = merge(&build(true), &build(true));
599
600        // Round-trip: every merged row decodes back to the corresponding input row
601        // (the merge here concatenates a's rows then b's rows, no consolidation).
602        assert_eq!(compressed.len(), rows.len() * 2);
603        for i in 0..compressed.len() {
604            let got = compressed.index(i).collect::<Vec<_>>();
605            let want = rows[i % rows.len()].iter().collect::<Vec<_>>();
606            assert_eq!(got, want, "merged row {i} round-trips");
607        }
608
609        // The promoted merge must actually compress relative to the poisoned one,
610        // confirming promotion carried a usable summary into `new_from`.
611        assert!(
612            heap(&compressed) < heap(&poisoned),
613            "promotion should let the merge compress: compressed={} poisoned={}",
614            heap(&compressed),
615            heap(&poisoned),
616        );
617    }
618}
619
620/// A `[u8]`-specialized container.
621mod bytes_container {
622
623    use differential_dataflow::trace::implementations::BatchContainer;
624    use timely::container::PushInto;
625
626    use mz_ore::region::Region;
627
628    /// A slice container with four bytes overhead per slice.
629    pub struct BytesContainer {
630        /// Total length of `batches`, maintained because recomputation is expensive.
631        length: usize,
632        batches: Vec<BytesBatch>,
633    }
634
635    impl BytesContainer {
636        /// Visit contained allocations to determine their size and capacity.
637        #[inline]
638        pub fn heap_size(&self, mut callback: impl FnMut(usize, usize)) {
639            // Calculate heap size for local, stash, and stash entries
640            callback(
641                self.batches.len() * std::mem::size_of::<BytesBatch>(),
642                self.batches.capacity() * std::mem::size_of::<BytesBatch>(),
643            );
644            for batch in self.batches.iter() {
645                batch.offsets.heap_size(&mut callback);
646                callback(batch.storage.len(), batch.storage.capacity());
647            }
648        }
649    }
650
651    impl BatchContainer for BytesContainer {
652        type Owned = Vec<u8>;
653        type ReadItem<'a> = &'a [u8];
654
655        #[inline]
656        fn into_owned<'a>(item: Self::ReadItem<'a>) -> Self::Owned {
657            item.to_vec()
658        }
659
660        #[inline]
661        fn clone_onto<'a>(item: Self::ReadItem<'a>, other: &mut Self::Owned) {
662            other.clear();
663            other.extend_from_slice(item);
664        }
665
666        #[inline(always)]
667        fn push_ref(&mut self, item: Self::ReadItem<'_>) {
668            self.push_into(item);
669        }
670
671        #[inline(always)]
672        fn push_own(&mut self, item: &Self::Owned) {
673            self.push_into(item.as_slice())
674        }
675
676        fn clear(&mut self) {
677            self.batches.clear();
678            self.batches.push(BytesBatch::with_capacities(0, 0));
679            self.length = 0;
680        }
681
682        fn with_capacity(size: usize) -> Self {
683            Self {
684                length: 0,
685                batches: vec![BytesBatch::with_capacities(size, size)],
686            }
687        }
688
689        fn merge_capacity(cont1: &Self, cont2: &Self) -> Self {
690            let mut item_cap = 1;
691            let mut byte_cap = 0;
692            for batch in cont1.batches.iter() {
693                item_cap += batch.offsets.len() - 1;
694                byte_cap += batch.storage.len();
695            }
696            for batch in cont2.batches.iter() {
697                item_cap += batch.offsets.len() - 1;
698                byte_cap += batch.storage.len();
699            }
700            Self {
701                length: 0,
702                batches: vec![BytesBatch::with_capacities(item_cap, byte_cap)],
703            }
704        }
705
706        #[inline(always)]
707        fn reborrow<'b, 'a: 'b>(item: Self::ReadItem<'a>) -> Self::ReadItem<'b> {
708            item
709        }
710
711        #[inline]
712        fn index(&self, mut index: usize) -> Self::ReadItem<'_> {
713            for batch in self.batches.iter() {
714                if index < batch.len() {
715                    return batch.index(index);
716                }
717                index -= batch.len();
718            }
719            panic!("Index out of bounds");
720        }
721
722        #[inline(always)]
723        fn len(&self) -> usize {
724            self.length
725        }
726    }
727
728    impl PushInto<&[u8]> for BytesContainer {
729        #[inline]
730        fn push_into(&mut self, item: &[u8]) {
731            self.length += 1;
732            if let Some(batch) = self.batches.last_mut() {
733                let success = batch.try_push(item);
734                if !success {
735                    // double the lengths from `batch`.
736                    let item_cap = 2 * batch.offsets.len();
737                    let byte_cap = std::cmp::max(2 * batch.storage.capacity(), item.len());
738                    let mut new_batch = BytesBatch::with_capacities(item_cap, byte_cap);
739                    assert!(new_batch.try_push(item));
740                    self.batches.push(new_batch);
741                }
742            }
743        }
744    }
745
746    /// A batch of slice storage.
747    ///
748    /// The backing storage for this batch will not be resized.
749    pub struct BytesBatch {
750        offsets: crate::OffsetOptimized,
751        storage: Region<u8>,
752        len: usize,
753    }
754
755    impl BytesBatch {
756        /// Either accepts the slice and returns true,
757        /// or does not and returns false.
758        fn try_push(&mut self, slice: &[u8]) -> bool {
759            if self.storage.len() + slice.len() <= self.storage.capacity() {
760                self.storage.extend_from_slice(slice);
761                self.offsets.push_into(self.storage.len());
762                self.len += 1;
763                true
764            } else {
765                false
766            }
767        }
768        #[inline]
769        fn index(&self, index: usize) -> &[u8] {
770            let lower = self.offsets.index(index);
771            let upper = self.offsets.index(index + 1);
772            &self.storage[lower..upper]
773        }
774        #[inline(always)]
775        fn len(&self) -> usize {
776            mz_ore::soft_assert_eq_no_log!(self.len, self.offsets.len() - 1);
777            self.len
778        }
779
780        fn with_capacities(item_cap: usize, byte_cap: usize) -> Self {
781            // TODO: be wary of `byte_cap` greater than 2^32.
782            let mut offsets = crate::OffsetOptimized::with_capacity(item_cap + 1);
783            offsets.push_into(0);
784            Self {
785                offsets,
786                storage: Region::new_auto(byte_cap.next_power_of_two()),
787                len: 0,
788            }
789        }
790    }
791}
792
793mod offset_opt {
794    use differential_dataflow::trace::implementations::BatchContainer;
795    use differential_dataflow::trace::implementations::OffsetList;
796    use timely::container::PushInto;
797
798    enum OffsetStride {
799        Empty,
800        Zero,
801        Striding(usize, usize),
802        Saturated(usize, usize, usize),
803    }
804
805    impl OffsetStride {
806        /// Accepts or rejects a newly pushed element.
807        #[inline]
808        fn push(&mut self, item: usize) -> bool {
809            match self {
810                OffsetStride::Empty => {
811                    if item == 0 {
812                        *self = OffsetStride::Zero;
813                        true
814                    } else {
815                        false
816                    }
817                }
818                OffsetStride::Zero => {
819                    *self = OffsetStride::Striding(item, 2);
820                    true
821                }
822                OffsetStride::Striding(stride, count) => {
823                    if item == *stride * *count {
824                        *count += 1;
825                        true
826                    } else if item == *stride * (*count - 1) {
827                        *self = OffsetStride::Saturated(*stride, *count, 1);
828                        true
829                    } else {
830                        false
831                    }
832                }
833                OffsetStride::Saturated(stride, count, reps) => {
834                    if item == *stride * (*count - 1) {
835                        *reps += 1;
836                        true
837                    } else {
838                        false
839                    }
840                }
841            }
842        }
843
844        #[inline]
845        fn index(&self, index: usize) -> usize {
846            match self {
847                OffsetStride::Empty => {
848                    panic!("Empty OffsetStride")
849                }
850                OffsetStride::Zero => 0,
851                OffsetStride::Striding(stride, _steps) => *stride * index,
852                OffsetStride::Saturated(stride, steps, _reps) => {
853                    if index < *steps {
854                        *stride * index
855                    } else {
856                        *stride * (*steps - 1)
857                    }
858                }
859            }
860        }
861
862        #[inline]
863        fn len(&self) -> usize {
864            match self {
865                OffsetStride::Empty => 0,
866                OffsetStride::Zero => 1,
867                OffsetStride::Striding(_stride, steps) => *steps,
868                OffsetStride::Saturated(_stride, steps, reps) => *steps + *reps,
869            }
870        }
871    }
872
873    pub struct OffsetOptimized {
874        strided: OffsetStride,
875        spilled: OffsetList,
876    }
877
878    impl BatchContainer for OffsetOptimized {
879        type Owned = usize;
880        type ReadItem<'a> = usize;
881
882        #[inline]
883        fn into_owned<'a>(item: Self::ReadItem<'a>) -> Self::Owned {
884            item
885        }
886
887        #[inline]
888        fn push_ref(&mut self, item: Self::ReadItem<'_>) {
889            self.push_into(item)
890        }
891
892        #[inline]
893        fn push_own(&mut self, item: &Self::Owned) {
894            self.push_into(*item)
895        }
896
897        fn clear(&mut self) {
898            self.strided = OffsetStride::Empty;
899            self.spilled.clear();
900        }
901
902        fn with_capacity(_size: usize) -> Self {
903            Self {
904                strided: OffsetStride::Empty,
905                spilled: OffsetList::with_capacity(0),
906            }
907        }
908
909        fn merge_capacity(_cont1: &Self, _cont2: &Self) -> Self {
910            Self {
911                strided: OffsetStride::Empty,
912                spilled: OffsetList::with_capacity(0),
913            }
914        }
915
916        #[inline]
917        fn reborrow<'b, 'a: 'b>(item: Self::ReadItem<'a>) -> Self::ReadItem<'b> {
918            item
919        }
920
921        #[inline]
922        fn index(&self, index: usize) -> Self::ReadItem<'_> {
923            if index < self.strided.len() {
924                self.strided.index(index)
925            } else {
926                self.spilled.index(index - self.strided.len())
927            }
928        }
929
930        #[inline]
931        fn len(&self) -> usize {
932            self.strided.len() + self.spilled.len()
933        }
934    }
935
936    impl PushInto<usize> for OffsetOptimized {
937        #[inline]
938        fn push_into(&mut self, item: usize) {
939            if !self.spilled.is_empty() {
940                self.spilled.push(item);
941            } else {
942                let inserted = self.strided.push(item);
943                if !inserted {
944                    self.spilled.push(item);
945                }
946            }
947        }
948    }
949
950    impl OffsetOptimized {
951        pub fn heap_size(&self, callback: impl FnMut(usize, usize)) {
952            crate::offset_list_size(&self.spilled, callback);
953        }
954    }
955}
956
957/// Helper to compute the size of an [`OffsetList`] in memory.
958#[inline]
959pub(crate) fn offset_list_size(data: &OffsetList, mut callback: impl FnMut(usize, usize)) {
960    // Private `vec_size` because we should only use it where data isn't region-allocated.
961    // `T: Copy` makes sure the implementation is correct even if types change!
962    #[inline(always)]
963    fn vec_size<T: Copy>(data: &Vec<T>, mut callback: impl FnMut(usize, usize)) {
964        let size_of_t = std::mem::size_of::<T>();
965        callback(data.len() * size_of_t, data.capacity() * size_of_t);
966    }
967
968    vec_size(&data.smol, &mut callback);
969    vec_size(&data.chonk, callback);
970}
971
972/// A `Row`-specialized container using dictionary compression.
973///
974/// The approach is to establish for each column lists of common values, and to use "unoccupied"
975/// tags in the row encoding (e.g. where we would indicate types) to replace these common values.
976/// This substitution is opt-in, in that we don't need to do it, and in particular do not do it
977/// while we are collecting preliminary information about common values, and then start to use it
978/// once we believe we have enough information. Once we have started to use the substitutions we
979/// cannot change the meaning of a reserved byte pattern, for the container we are populating.
980///
981/// Each from-scratch container observes `STATS_THRESHOLD` records before establishing a mapping
982/// from spare tags to common values. Containers that are formed from merging other containers
983/// use those input containers' common values to populate a codec and use it immediately.
984///
985/// The dictionary behavior is controlled by the `DICTIONARY_COMPRESSION` flag, which if disabled
986/// prevents the construction of codecs, which when absent simply cause the wrapper to behave as
987/// a no-op that fails to use any spare tags for common values. The flag is set once, when a
988/// replica is created (from compute's `InstanceConfig::arrangement_dictionary_compression`, itself
989/// captured from the `enable_arrangement_dictionary_compression_alpha` dyncfg at that moment), and is
990/// not changed for the life of the process; flipping the dyncfg only affects replicas created
991/// afterwards. Even with the flag fixed, a single replica can hold a mix of compressed and
992/// uncompressed containers — e.g. containers that never observed enough records to install a
993/// codec, or that were merged from uncompressed inputs.
994mod dictionary {
995
996    use differential_dataflow::trace::implementations::BatchContainer;
997
998    use mz_repr::{Row, RowRef};
999
1000    use super::row_codec::{ColumnsCodec, ColumnsIter};
1001
1002    /// Wrapper types that exist to support the creation of dictionary codecs.
1003    ///
1004    /// These types interpose at the seal() call, to traverse the data that is being sealed and
1005    /// then construct codecs that are used to encode the row-shaped keys and values. There are
1006    /// several variants, corresponding to the RowRow, RowVal, and Row-only spine types.
1007    pub mod builders {
1008
1009        use columnar::{Columnar, Index};
1010        use columnation::Columnation;
1011        use differential_dataflow::difference::Semigroup;
1012        use differential_dataflow::lattice::Lattice;
1013        use differential_dataflow::trace::Builder;
1014        use differential_dataflow::trace::Description;
1015        use differential_dataflow::trace::implementations::BatchContainer;
1016        use differential_dataflow::trace::implementations::ord_neu::{OrdKeyBatch, OrdKeyBuilder};
1017        use differential_dataflow::trace::implementations::ord_neu::{OrdValBatch, OrdValBuilder};
1018        use mz_timely_util::columnar::Column;
1019        use mz_timely_util::columnar::body::ColumnBody;
1020        use mz_timely_util::columnar::chunk::ChainState;
1021        use mz_timely_util::columnation::ColumnationStack as TimelyStack;
1022        use timely::progress::Timestamp;
1023
1024        use mz_repr::{Row, RowRef};
1025
1026        use super::super::row_codec::ColumnsCodec;
1027        use super::{DatumContainer, DatumSeq};
1028        use crate::DICTIONARY_COMPRESSION;
1029        use crate::spines::{RowLayout, RowRowLayout, RowValLayout, ValRowLayout};
1030
1031        /// Gather encoding statistics across `rows` and produce a codec from them.
1032        ///
1033        /// Accepts anything that borrows as a [`RowRef`], so it serves both the
1034        /// columnation-fed builders (which yield `&Row`) and the paged builders
1035        /// (which yield `&RowRef` straight out of a [`Column`] chunk).
1036        ///
1037        /// Returns `None` when dictionary compression is disabled.
1038        fn build_codec<'a, B>(rows: impl IntoIterator<Item = &'a B>) -> Option<ColumnsCodec>
1039        where
1040            B: std::borrow::Borrow<RowRef> + ?Sized + 'a,
1041        {
1042            if !DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed) {
1043                return None;
1044            }
1045            let mut stats = ColumnsCodec::default();
1046            observe_rows(&mut stats, rows);
1047            Some(ColumnsCodec::new_from([&stats]))
1048        }
1049
1050        /// Fold `rows` into `stats`, the accumulator a codec is built from.
1051        fn observe_rows<'a, B>(stats: &mut ColumnsCodec, rows: impl IntoIterator<Item = &'a B>)
1052        where
1053            B: std::borrow::Borrow<RowRef> + ?Sized + 'a,
1054        {
1055            for row in rows {
1056                let row = row.borrow();
1057                if !row.is_empty() {
1058                    // Gather stats only; the encoded output would be thrown away here, so
1059                    // `observe` skips the per-value lookup and the throwaway-buffer memcpy
1060                    // that `encode` would do (see `ColumnsCodec::observe`).
1061                    stats.observe(DatumSeq::borrow_as(row).bytes_iter());
1062                }
1063            }
1064        }
1065
1066        pub struct RowRowBuilder<
1067            T: Lattice + Timestamp + Columnation,
1068            R: Ord + Semigroup + Columnation + 'static,
1069        > {
1070            inner: OrdValBuilder<RowRowLayout<((Row, Row), T, R)>, TimelyStack<((Row, Row), T, R)>>,
1071        }
1072
1073        impl<T: Lattice + Timestamp + Columnation, R: Ord + Semigroup + Columnation + 'static>
1074            Builder for RowRowBuilder<T, R>
1075        {
1076            type Input = TimelyStack<((Row, Row), T, R)>;
1077            type Time = T;
1078            type Output = OrdValBatch<RowRowLayout<((Row, Row), T, R)>>;
1079
1080            fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1081                Self {
1082                    inner: Builder::with_capacity(keys, vals, upds),
1083                }
1084            }
1085            fn push(&mut self, chunk: &mut Self::Input) {
1086                self.inner.push(chunk)
1087            }
1088            fn done(self, description: Description<Self::Time>) -> Self::Output {
1089                // The push/done build path (e.g. the `reduce` operator, which builds
1090                // batches with `Builder::new()` + `push` + `done` rather than `seal`)
1091                // never runs `seal`'s codec install. Install a codec here from the
1092                // statistics gathered during `push`, mirroring `seal` — but without
1093                // building a dictionary or re-encoding the rows; see
1094                // `DatumContainer::promote_stats_to_codec` for why a codec-less batch
1095                // must be avoided even though its rows stay raw.
1096                let mut inner = self.inner;
1097                inner.result.keys.promote_stats_to_codec();
1098                inner.result.vals.vals.promote_stats_to_codec();
1099                inner.done(description)
1100            }
1101            fn seal(
1102                chain: &mut Vec<Self::Input>,
1103                description: Description<Self::Time>,
1104            ) -> Self::Output {
1105                let key_codec = build_codec(
1106                    chain
1107                        .iter()
1108                        .flat_map(|link| link.iter().map(|((k, _), _, _)| k)),
1109                );
1110                let val_codec = build_codec(
1111                    chain
1112                        .iter()
1113                        .flat_map(|link| link.iter().map(|((_, v), _, _)| v)),
1114                );
1115
1116                use differential_dataflow::trace::implementations::BuilderInput;
1117
1118                let (keys, vals, upds) = <Self::Input as BuilderInput<
1119                    DatumContainer,
1120                    DatumContainer,
1121                >>::key_val_upd_counts(&chain[..]);
1122                let mut builder = Self::with_capacity(keys, vals, upds);
1123                // The seal path installs a codec directly, so the per-container stats
1124                // gatherer (which `with_capacity` may have allocated) is dead weight and
1125                // would contradict the `stats: None once codec installed` invariant.
1126                builder.inner.result.keys.codec = key_codec;
1127                builder.inner.result.keys.stats = None;
1128                builder.inner.result.vals.vals.codec = val_codec;
1129                builder.inner.result.vals.vals.stats = None;
1130
1131                for mut chunk in chain.drain(..) {
1132                    builder.push(&mut chunk);
1133                }
1134
1135                builder.done(description)
1136            }
1137        }
1138
1139        pub struct RowValBuilder<
1140            V: Ord + Clone + Columnation + 'static,
1141            T: Lattice + Timestamp + Columnation,
1142            R: Ord + Semigroup + Columnation + 'static,
1143        > {
1144            inner: OrdValBuilder<RowValLayout<((Row, V), T, R)>, TimelyStack<((Row, V), T, R)>>,
1145        }
1146
1147        impl<
1148            V: Ord + Clone + Columnation,
1149            T: Lattice + Timestamp + Columnation,
1150            R: Ord + Semigroup + Columnation + 'static,
1151        > Builder for RowValBuilder<V, T, R>
1152        {
1153            type Input = TimelyStack<((Row, V), T, R)>;
1154            type Time = T;
1155            type Output = OrdValBatch<RowValLayout<((Row, V), T, R)>>;
1156
1157            fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1158                Self {
1159                    inner: Builder::with_capacity(keys, vals, upds),
1160                }
1161            }
1162            fn push(&mut self, chunk: &mut Self::Input) {
1163                self.inner.push(chunk)
1164            }
1165            fn done(self, description: Description<Self::Time>) -> Self::Output {
1166                // See `RowRowBuilder::done`: install a codec on the `Row`-shaped key
1167                // container for the push/done (e.g. `reduce`) path that skips `seal`.
1168                let mut inner = self.inner;
1169                inner.result.keys.promote_stats_to_codec();
1170                inner.done(description)
1171            }
1172            fn seal(
1173                chain: &mut Vec<Self::Input>,
1174                description: Description<Self::Time>,
1175            ) -> Self::Output {
1176                let key_codec = build_codec(
1177                    chain
1178                        .iter()
1179                        .flat_map(|link| link.iter().map(|((k, _), _, _)| k)),
1180                );
1181
1182                use differential_dataflow::trace::implementations::BuilderInput;
1183
1184                let (keys, vals, upds) = <Self::Input as BuilderInput<
1185                    DatumContainer,
1186                    TimelyStack<V>,
1187                >>::key_val_upd_counts(&chain[..]);
1188                let mut builder = Self::with_capacity(keys, vals, upds);
1189                // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1190                builder.inner.result.keys.codec = key_codec;
1191                builder.inner.result.keys.stats = None;
1192
1193                for mut chunk in chain.drain(..) {
1194                    builder.push(&mut chunk);
1195                }
1196
1197                builder.done(description)
1198            }
1199        }
1200
1201        pub struct RowBuilder<
1202            T: Lattice + Timestamp + Columnation,
1203            R: Ord + Semigroup + Columnation + 'static,
1204            DC: BatchContainer<Owned = R> = TimelyStack<R>,
1205        > {
1206            inner: OrdKeyBuilder<RowLayout<((Row, ()), T, R), DC>, TimelyStack<((Row, ()), T, R)>>,
1207        }
1208
1209        impl<T, R, DC> Builder for RowBuilder<T, R, DC>
1210        where
1211            T: Lattice + Timestamp + Columnation,
1212            R: Ord + Semigroup + Columnation + 'static,
1213            DC: BatchContainer<Owned = R>,
1214        {
1215            type Input = TimelyStack<((Row, ()), T, R)>;
1216            type Time = T;
1217            type Output = OrdKeyBatch<RowLayout<((Row, ()), T, R), DC>>;
1218
1219            fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1220                Self {
1221                    inner: Builder::with_capacity(keys, vals, upds),
1222                }
1223            }
1224            fn push(&mut self, chunk: &mut Self::Input) {
1225                self.inner.push(chunk)
1226            }
1227            fn done(self, description: Description<Self::Time>) -> Self::Output {
1228                // See `RowRowBuilder::done`: install a codec on the `Row`-shaped key
1229                // container for the push/done (e.g. `reduce`) path that skips `seal`.
1230                let mut inner = self.inner;
1231                inner.result.keys.promote_stats_to_codec();
1232                inner.done(description)
1233            }
1234            fn seal(
1235                chain: &mut Vec<Self::Input>,
1236                description: Description<Self::Time>,
1237            ) -> Self::Output {
1238                let key_codec = build_codec(
1239                    chain
1240                        .iter()
1241                        .flat_map(|link| link.iter().map(|((k, _), _, _)| k)),
1242                );
1243
1244                use differential_dataflow::trace::implementations::BuilderInput;
1245
1246                let (keys, vals, upds) = <Self::Input as BuilderInput<
1247                    DatumContainer,
1248                    TimelyStack<()>,
1249                >>::key_val_upd_counts(&chain[..]);
1250                let mut builder = Self::with_capacity(keys, vals, upds);
1251                // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1252                builder.inner.result.keys.codec = key_codec;
1253                builder.inner.result.keys.stats = None;
1254
1255                for mut chunk in chain.drain(..) {
1256                    builder.push(&mut chunk);
1257                }
1258
1259                builder.done(description)
1260            }
1261        }
1262
1263        /// Mirror of [`RowValBuilder`] with the roles swapped: arbitrary keys and
1264        /// `Row` *values*, so the dictionary codec is built for and installed on the
1265        /// value container.
1266        pub struct ValRowBuilder<
1267            K: Ord + Clone + Columnation + 'static,
1268            T: Lattice + Timestamp + Columnation,
1269            R: Ord + Semigroup + Columnation + 'static,
1270        > {
1271            inner: OrdValBuilder<ValRowLayout<((K, Row), T, R)>, TimelyStack<((K, Row), T, R)>>,
1272        }
1273
1274        impl<
1275            K: Ord + Clone + Columnation,
1276            T: Lattice + Timestamp + Columnation,
1277            R: Ord + Semigroup + Columnation + 'static,
1278        > Builder for ValRowBuilder<K, T, R>
1279        {
1280            type Input = TimelyStack<((K, Row), T, R)>;
1281            type Time = T;
1282            type Output = OrdValBatch<ValRowLayout<((K, Row), T, R)>>;
1283
1284            fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1285                Self {
1286                    inner: Builder::with_capacity(keys, vals, upds),
1287                }
1288            }
1289            fn push(&mut self, chunk: &mut Self::Input) {
1290                self.inner.push(chunk)
1291            }
1292            fn done(self, description: Description<Self::Time>) -> Self::Output {
1293                // See `RowRowBuilder::done`: install a codec on the `Row`-shaped value
1294                // container for the push/done (e.g. `reduce`) path that skips `seal`.
1295                let mut inner = self.inner;
1296                inner.result.vals.vals.promote_stats_to_codec();
1297                inner.done(description)
1298            }
1299            fn seal(
1300                chain: &mut Vec<Self::Input>,
1301                description: Description<Self::Time>,
1302            ) -> Self::Output {
1303                let val_codec = build_codec(
1304                    chain
1305                        .iter()
1306                        .flat_map(|link| link.iter().map(|((_, v), _, _)| v)),
1307                );
1308
1309                use differential_dataflow::trace::implementations::BuilderInput;
1310
1311                let (keys, vals, upds) = <Self::Input as BuilderInput<
1312                    TimelyStack<K>,
1313                    DatumContainer,
1314                >>::key_val_upd_counts(&chain[..]);
1315                let mut builder = Self::with_capacity(keys, vals, upds);
1316                // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1317                builder.inner.result.vals.vals.codec = val_codec;
1318                builder.inner.result.vals.vals.stats = None;
1319
1320                for mut chunk in chain.drain(..) {
1321                    builder.push(&mut chunk);
1322                }
1323
1324                builder.done(description)
1325            }
1326        }
1327
1328        /// Counterpart of [`RowRowBuilder`] that consumes [`ColumnBody`] chunks
1329        /// instead of columnation stacks, whether or not the batcher that
1330        /// produced them spills. Mirrors `RowRowBuilder::seal`:
1331        /// it gathers key and value statistics from the sealed chain and
1332        /// installs codecs directly, then drops the per-container stats gatherer.
1333        pub struct RowRowColPagedBuilder<
1334            T: Lattice + Timestamp + Columnation + Columnar,
1335            R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1336        > {
1337            inner: OrdValBuilder<RowRowLayout<((Row, Row), T, R)>, ColumnBody<((Row, Row), T, R)>>,
1338        }
1339
1340        impl<
1341            T: Lattice + Timestamp + Columnation + Columnar,
1342            R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1343        > Builder for RowRowColPagedBuilder<T, R>
1344        {
1345            type Input = ColumnBody<((Row, Row), T, R)>;
1346            type Time = T;
1347            type Output = OrdValBatch<RowRowLayout<((Row, Row), T, R)>>;
1348
1349            fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1350                Self {
1351                    inner: Builder::with_capacity(keys, vals, upds),
1352                }
1353            }
1354            fn push(&mut self, chunk: &mut Self::Input) {
1355                self.inner.push(chunk)
1356            }
1357            fn done(self, description: Description<Self::Time>) -> Self::Output {
1358                // See `RowRowBuilder::done`: a builder driven by `push` and `done`
1359                // alone gets its codec from the statistics gathered during the
1360                // pushes, since a codec-less batch poisons the merges it joins.
1361                // Installing a codec here is a no-op once one is installed.
1362                let mut inner = self.inner;
1363                inner.result.keys.promote_stats_to_codec();
1364                inner.result.vals.vals.promote_stats_to_codec();
1365                inner.done(description)
1366            }
1367            fn seal(
1368                chain: &mut Vec<Self::Input>,
1369                description: Description<Self::Time>,
1370            ) -> Self::Output {
1371                // `into_index_iter` yields the value column's `Row`s as `&RowRef`,
1372                // which `build_codec` consumes directly.
1373                let key_codec = build_codec(
1374                    chain
1375                        .iter()
1376                        .flat_map(|c| c.borrow().into_index_iter().map(|((k, _), _, _)| k)),
1377                );
1378                let val_codec = build_codec(
1379                    chain
1380                        .iter()
1381                        .flat_map(|c| c.borrow().into_index_iter().map(|((_, v), _, _)| v)),
1382                );
1383
1384                use differential_dataflow::trace::implementations::BuilderInput;
1385
1386                let (keys, vals, upds) = <Self::Input as BuilderInput<
1387                    DatumContainer,
1388                    DatumContainer,
1389                >>::key_val_upd_counts(&chain[..]);
1390                let mut builder = Self::with_capacity(keys, vals, upds);
1391                // See `RowRowBuilder::seal`: install the codecs and drop the
1392                // now-redundant per-container stats gatherer.
1393                builder.inner.result.keys.codec = key_codec;
1394                builder.inner.result.keys.stats = None;
1395                builder.inner.result.vals.vals.codec = val_codec;
1396                builder.inner.result.vals.vals.stats = None;
1397
1398                for mut chunk in chain.drain(..) {
1399                    builder.push(&mut chunk);
1400                }
1401
1402                builder.done(description)
1403            }
1404        }
1405
1406        /// The chain-wide state of [`RowRowColPagedBuilder`]: the statistics its
1407        /// key and value codecs are built from, and the container sizes the
1408        /// chain implies.
1409        ///
1410        /// [`RowRowColPagedBuilder`]: crate::RowRowColPagedBuilder
1411        #[derive(Default)]
1412        pub struct RowRowColPagedState {
1413            keys: ColumnsCodec,
1414            vals: ColumnsCodec,
1415            /// Whether `keys` and `vals` hold statistics from any body.
1416            gathered: bool,
1417            counts: (usize, usize, usize),
1418        }
1419
1420        impl<
1421            T: Lattice + Timestamp + Columnation + Columnar,
1422            R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1423        > ChainState for RowRowColPagedBuilder<T, R>
1424        {
1425            type State = RowRowColPagedState;
1426
1427            fn wants_bodies() -> bool {
1428                DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed)
1429            }
1430
1431            fn observe(state: &mut Self::State, input: &Self::Input) {
1432                if Self::wants_bodies() {
1433                    // `into_index_iter` yields the key and value columns' `Row`s
1434                    // as `&RowRef`, which `observe_rows` consumes directly.
1435                    observe_rows(
1436                        &mut state.keys,
1437                        input.borrow().into_index_iter().map(|((k, _), _, _)| k),
1438                    );
1439                    observe_rows(
1440                        &mut state.vals,
1441                        input.borrow().into_index_iter().map(|((_, v), _, _)| v),
1442                    );
1443                    state.gathered = true;
1444                }
1445
1446                use differential_dataflow::trace::implementations::BuilderInput;
1447                let (keys, vals, upds) = <Self::Input as BuilderInput<
1448                    DatumContainer,
1449                    DatumContainer,
1450                >>::key_val_upd_counts(
1451                    std::slice::from_ref(input)
1452                );
1453                state.counts.0 += keys;
1454                state.counts.1 += vals;
1455                state.counts.2 += upds;
1456            }
1457
1458            fn observe_records(state: &mut Self::State, records: usize) {
1459                // Key and value counts need the bodies. Sizing the update
1460                // containers alone is what a record count buys.
1461                state.counts.2 += records;
1462            }
1463
1464            fn from_state(state: Self::State) -> Self {
1465                let mut builder =
1466                    Self::with_capacity(state.counts.0, state.counts.1, state.counts.2);
1467                // Install from the statistics that were gathered, not from the
1468                // compression flag: the flag is process-global and a worker can
1469                // flip it mid-seal, and a codec built from no observations would
1470                // carry an empty summary and displace the gatherer that `done`
1471                // would otherwise promote.
1472                if state.gathered {
1473                    // See `RowRowBuilder::seal`: install the codecs the chain's
1474                    // statistics imply and drop the now-redundant per-container
1475                    // stats gatherer.
1476                    builder.inner.result.keys.codec = Some(ColumnsCodec::new_from([&state.keys]));
1477                    builder.inner.result.keys.stats = None;
1478                    builder.inner.result.vals.vals.codec =
1479                        Some(ColumnsCodec::new_from([&state.vals]));
1480                    builder.inner.result.vals.vals.stats = None;
1481                }
1482                builder
1483            }
1484        }
1485
1486        /// Paged counterpart of [`ValRowBuilder`] that consumes [`Column`]
1487        /// chunks. Keys are arbitrary `Columnar` values (not `Row`-shaped) and
1488        /// stay uncompressed; only the value container receives a codec.
1489        pub struct ValRowColPagedBuilder<
1490            K: Ord + Clone + Columnation + Columnar + 'static,
1491            T: Lattice + Timestamp + Columnation + Columnar,
1492            R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1493        > {
1494            inner: OrdValBuilder<ValRowLayout<((K, Row), T, R)>, Column<((K, Row), T, R)>>,
1495        }
1496
1497        impl<
1498            K: Ord + Clone + Columnation + Columnar + 'static,
1499            T: Lattice + Timestamp + Columnation + Columnar,
1500            R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1501        > Builder for ValRowColPagedBuilder<K, T, R>
1502        where
1503            for<'a> columnar::Ref<'a, K>: Copy + Ord,
1504            for<'a, 'b> &'a K: PartialEq<columnar::Ref<'b, K>>,
1505            for<'a> TimelyStack<K>: timely::container::PushInto<columnar::Ref<'a, K>>,
1506        {
1507            type Input = Column<((K, Row), T, R)>;
1508            type Time = T;
1509            type Output = OrdValBatch<ValRowLayout<((K, Row), T, R)>>;
1510
1511            fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1512                Self {
1513                    inner: Builder::with_capacity(keys, vals, upds),
1514                }
1515            }
1516            fn push(&mut self, chunk: &mut Self::Input) {
1517                self.inner.push(chunk)
1518            }
1519            fn done(self, description: Description<Self::Time>) -> Self::Output {
1520                self.inner.done(description)
1521            }
1522            fn seal(
1523                chain: &mut Vec<Self::Input>,
1524                description: Description<Self::Time>,
1525            ) -> Self::Output {
1526                let val_codec = build_codec(
1527                    chain
1528                        .iter()
1529                        .flat_map(|c| c.borrow().into_index_iter().map(|((_, v), _, _)| v)),
1530                );
1531
1532                use differential_dataflow::trace::implementations::BuilderInput;
1533
1534                let (keys, vals, upds) = <Self::Input as BuilderInput<
1535                    TimelyStack<K>,
1536                    DatumContainer,
1537                >>::key_val_upd_counts(&chain[..]);
1538                let mut builder = Self::with_capacity(keys, vals, upds);
1539                // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1540                builder.inner.result.vals.vals.codec = val_codec;
1541                builder.inner.result.vals.vals.stats = None;
1542
1543                for mut chunk in chain.drain(..) {
1544                    builder.push(&mut chunk);
1545                }
1546
1547                builder.done(description)
1548            }
1549        }
1550    }
1551
1552    pub struct DatumContainer {
1553        /// Encoder/decoder used to translate between row bytes and the stored bytes.
1554        /// `None` until enough pushes have been observed (or if compression is disabled).
1555        codec: Option<ColumnsCodec>,
1556        /// The stored, possibly-encoded, row bytes.
1557        inner: super::bytes_container::BytesContainer,
1558        /// Staging buffer for ingested `Row` types.
1559        staging: Vec<u8>,
1560        /// Statistics gatherer, used to build a safe codec after enough pushes.
1561        /// `None` once the codec has been installed or if compression is disabled.
1562        stats: Option<ColumnsCodec>,
1563    }
1564
1565    impl BatchContainer for DatumContainer {
1566        type Owned = Row;
1567        type ReadItem<'a> = DatumSeq<'a>;
1568
1569        fn with_capacity(size: usize) -> Self {
1570            let stats = if crate::DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed)
1571            {
1572                Some(Default::default())
1573            } else {
1574                None
1575            };
1576
1577            Self {
1578                codec: None,
1579                inner: BatchContainer::with_capacity(size),
1580                staging: Vec::new(),
1581                stats,
1582            }
1583        }
1584        fn merge_capacity(cont1: &Self, cont2: &Self) -> Self {
1585            // We only build a merged codec when *both* inputs carry one. A codec is
1586            // sound only for the data whose tag usage it observed, so we cannot reuse
1587            // one side's codec to decode the other side's rows. When exactly one side
1588            // is compressed we conservatively produce an uncompressed container rather
1589            // than risk a tag collision; the merged container re-gathers stats and may
1590            // install a fresh codec later via the `STATS_THRESHOLD` path.
1591            let codec = match (&cont1.codec, &cont2.codec) {
1592                (Some(c1), Some(c2)) => Some(ColumnsCodec::new_from([c1, c2])),
1593                _ => None,
1594            };
1595
1596            Self {
1597                codec,
1598                inner: BatchContainer::merge_capacity(&cont1.inner, &cont2.inner),
1599                staging: Vec::new(),
1600                stats: None,
1601            }
1602        }
1603        #[inline]
1604        fn index(&self, index: usize) -> Self::ReadItem<'_> {
1605            let data = self.inner.index(index);
1606            let iter = if let Some(codec) = &self.codec {
1607                codec.decode(data)
1608            } else {
1609                // Safety: without a codec we only push rows or datumseqs into `self.inner`.
1610                // Each retrieved byte slice should be row-encoded data, as long as we have
1611                // not unset the codec in the interim.
1612                unsafe { ColumnsIter::without_codec(data) }
1613            };
1614            DatumSeq { iter }
1615        }
1616        #[inline(always)]
1617        fn len(&self) -> usize {
1618            self.inner.len()
1619        }
1620
1621        #[inline(always)]
1622        fn reborrow<'b, 'a: 'b>(item: Self::ReadItem<'a>) -> Self::ReadItem<'b> {
1623            item
1624        }
1625
1626        #[inline(always)]
1627        fn into_owned<'a>(item: Self::ReadItem<'a>) -> Self::Owned {
1628            // Fast path: unencoded data is already row-formatted bytes.
1629            if item.iter.index.is_none() {
1630                // SAFETY: `iter.data` is raw row-encoded bytes when there is no codec.
1631                return unsafe { Row::from_bytes_unchecked(item.iter.data) };
1632            }
1633            Row::pack(item)
1634        }
1635
1636        #[inline(always)]
1637        fn clone_onto<'a>(item: Self::ReadItem<'a>, other: &mut Self::Owned) {
1638            // Fast path: unencoded data is already row-formatted bytes.
1639            if item.iter.index.is_none() {
1640                let mut packer = other.packer();
1641                // SAFETY: `iter.data` is raw row-encoded bytes when there is no codec.
1642                unsafe { packer.extend_by_slice_unchecked(item.iter.data) };
1643                return;
1644            }
1645            other.packer().extend(item);
1646        }
1647
1648        #[inline(always)]
1649        fn push_ref(&mut self, item: Self::ReadItem<'_>) {
1650            // Fast path: both sides unencoded — push raw bytes directly.
1651            if self.codec.is_none() && self.stats.is_none() && item.iter.index.is_none() {
1652                self.inner.push_ref(item.iter.data);
1653                return;
1654            }
1655            self.push_into(item);
1656        }
1657
1658        #[inline(always)]
1659        fn push_own(&mut self, item: &Self::Owned) {
1660            // Fast path: container is unencoded — push raw row bytes directly.
1661            if self.codec.is_none() && self.stats.is_none() {
1662                self.inner.push_ref(item.data());
1663                return;
1664            }
1665            self.push_into(item);
1666        }
1667
1668        #[inline(always)]
1669        fn clear(&mut self) {
1670            self.inner.clear();
1671            self.staging.clear();
1672            // Reset to the same state as a fresh `with_capacity`: drop any installed
1673            // codec and restore stats gathering (if compression is enabled). Keeping a
1674            // now-empty codec would leave `codec.is_some()`, which permanently routes
1675            // pushes down the encode path with an empty dictionary and prevents the
1676            // `STATS_THRESHOLD` install logic from ever re-engaging compression.
1677            self.codec = None;
1678            self.stats = if crate::DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed)
1679            {
1680                Some(Default::default())
1681            } else {
1682                None
1683            };
1684        }
1685    }
1686
1687    impl DatumContainer {
1688        /// Visit contained allocations to determine their size and capacity.
1689        #[inline]
1690        pub fn heap_size(&self, mut callback: impl FnMut(usize, usize)) {
1691            self.inner.heap_size(&mut callback);
1692            // The staging buffer and the (possibly absent) codec and stats gatherer all
1693            // hold heap allocations that the bare `inner` accounting misses.
1694            callback(self.staging.len(), self.staging.capacity());
1695            if let Some(codec) = &self.codec {
1696                codec.heap_size(&mut callback);
1697            }
1698            if let Some(stats) = &self.stats {
1699                stats.heap_size(&mut callback);
1700            }
1701        }
1702
1703        /// Promote a gathered-but-uninstalled statistics summary into the codec slot.
1704        ///
1705        /// A container filled via the builder's `push`/`done` path — as the `reduce`
1706        /// operator does, building batches with `Builder::new()` + `push` + `done`
1707        /// rather than `seal` — gathers statistics on every push but never reaches
1708        /// `seal`'s codec install, and only crosses the mid-formation
1709        /// `STATS_THRESHOLD` install if it grows past it. A smaller such container
1710        /// would otherwise be finalized with no codec at all, even with the flag on.
1711        ///
1712        /// That is a problem not because this batch needs compressing — its rows are
1713        /// already stored raw and we deliberately do *not* re-encode them here — but
1714        /// because a codec-less batch poisons future merges: [`Self::merge_capacity`]
1715        /// keys off the presence of a codec, so a codec-less input forces the merged
1716        /// container onto the uncompressed path. Moving the gathered statistics into
1717        /// the codec slot leaves the batch carrying a codec whose retained heavy-hitter
1718        /// summary a later merge can rebuild from via `ColumnsCodec::new_from`, while
1719        /// installing no dictionary: the empty `decode` map resolves every stored
1720        /// (raw) column through the literal-datum fall-through, so reads stay correct.
1721        ///
1722        /// We move the summary as-is rather than building a dictionary via `new_safe`
1723        /// / `new_from` (which reset the summary): unlike `seal` and the mid-formation
1724        /// install, `done` has no further rows to re-observe, so a reset summary would
1725        /// leave the eventual merge nothing to rebuild from.
1726        pub(crate) fn promote_stats_to_codec(&mut self) {
1727            if self.codec.is_none() {
1728                self.codec = self.stats.take();
1729            }
1730        }
1731
1732        /// Whether a codec is installed. A container without one stores raw
1733        /// bytes and blocks compression in every merge it takes part in.
1734        #[cfg(test)]
1735        pub(crate) fn has_codec(&self) -> bool {
1736            self.codec.is_some()
1737        }
1738    }
1739
1740    use timely::container::PushInto;
1741    impl PushInto<Row> for DatumContainer {
1742        #[inline(always)]
1743        fn push_into(&mut self, item: Row) {
1744            self.push_into(&item);
1745        }
1746    }
1747
1748    impl PushInto<&Row> for DatumContainer {
1749        #[inline(always)]
1750        fn push_into(&mut self, item: &Row) {
1751            self.push_into(DatumSeq::borrow_as(item));
1752        }
1753    }
1754
1755    impl PushInto<&RowRef> for DatumContainer {
1756        #[inline(always)]
1757        fn push_into(&mut self, item: &RowRef) {
1758            self.push_into(DatumSeq::borrow_as(item));
1759        }
1760    }
1761
1762    /// Number of pushes a from-scratch container observes before it turns its
1763    /// gathered stats into a safe codec.
1764    ///
1765    /// A safe codec has at most `256 - SAFE_TAG_BASE` (= 134) dictionary slots per
1766    /// column, so we only need to identify ~134 genuinely-popular values. The
1767    /// `MisraGries` summary retains up to `2 * k` (= 1024) distinct candidates
1768    /// between tidies and reduces to `k` (= 512), comfortably more than 134, so the
1769    /// threshold just needs to be large enough that heavy hitters accumulate counts
1770    /// well above 1 before we freeze the codec. 64Ki pushes gives that headroom while
1771    /// keeping the pre-codec (uncompressed) window short.
1772    const STATS_THRESHOLD: usize = 64 * 1024;
1773
1774    impl PushInto<DatumSeq<'_>> for DatumContainer {
1775        #[inline]
1776        fn push_into(&mut self, item: DatumSeq<'_>) {
1777            // Fast path: container and item are both unencoded.
1778            // This is the hot path when dictionary compression is disabled.
1779            if self.codec.is_none() && self.stats.is_none() && item.iter.index.is_none() {
1780                self.inner.push_ref(item.iter.data);
1781                return;
1782            }
1783
1784            // Check if we've gathered enough stats to install a safe codec.
1785            if self.codec.is_none() && self.stats.is_some() && self.inner.len() >= STATS_THRESHOLD {
1786                let stats = self.stats.take().unwrap();
1787                self.codec = Some(stats.new_safe());
1788            }
1789
1790            if let Some(codec) = &mut self.codec {
1791                // Encode using the installed codec.
1792                codec.encode(item.bytes_iter(), &mut self.staging);
1793            } else if let Some(stats) = &mut self.stats {
1794                // Stats-gathering phase: feed the statistics but store raw bytes.
1795                // `observe` updates the heavy-hitter/tag summaries without encoding, so
1796                // we copy each row exactly once (below) instead of also encoding it into
1797                // a buffer we would immediately discard.
1798                stats.observe(item.bytes_iter());
1799                for slice in item.bytes_iter() {
1800                    self.staging.extend_from_slice(slice);
1801                }
1802            } else {
1803                // No codec, no stats: raw copy.
1804                for slice in item.bytes_iter() {
1805                    self.staging.extend_from_slice(slice);
1806                }
1807            }
1808            self.inner.push_ref(&self.staging[..]);
1809            self.staging.clear();
1810        }
1811    }
1812
1813    use mz_repr::{Datum, read_datum};
1814
1815    /// A reference that can be resolved to a sequence of `Datum`s.
1816    ///
1817    /// This type must "compare" as if decoded to a `Row`, which means it needs to track
1818    /// various nuances of `Row::cmp`, which at the moment is first by length, and then by
1819    /// the raw binary slice backing the row. Neither of those are explicit in this struct.
1820    /// We will need to produce them in order to perform comparisons.
1821    #[derive(Debug)]
1822    pub struct DatumSeq<'a> {
1823        pub iter: ColumnsIter<'a>,
1824    }
1825
1826    impl<'a> DatumSeq<'a> {
1827        #[inline(always)]
1828        fn borrow_as(other: &'a RowRef) -> Self {
1829            Self {
1830                iter: ColumnsCodec::borrow_row(other),
1831            }
1832        }
1833
1834        /// Borrow a `Row` as a `DatumSeq` so that it can be used to seek into a
1835        /// trace whose key/value container is a [`DatumContainer`].
1836        #[inline]
1837        pub fn from_row(row: &'a Row) -> Self {
1838            Self::borrow_as(row)
1839        }
1840
1841        #[inline]
1842        pub fn to_row(&self) -> Row {
1843            // Fast path: unencoded data is already row-formatted bytes.
1844            if self.iter.index.is_none() {
1845                return unsafe { Row::from_bytes_unchecked(self.iter.data) };
1846            }
1847            Row::pack(*self)
1848        }
1849    }
1850
1851    impl<'a> Copy for DatumSeq<'a> {}
1852    impl<'a> Clone for DatumSeq<'a> {
1853        #[inline(always)]
1854        fn clone(&self) -> Self {
1855            *self
1856        }
1857    }
1858
1859    use std::cmp::Ordering;
1860    impl<'a, 'b> PartialEq<DatumSeq<'a>> for DatumSeq<'b> {
1861        #[inline(always)]
1862        fn eq(&self, other: &DatumSeq<'a>) -> bool {
1863            // Fast path: both sides are unencoded raw row bytes.
1864            if self.iter.index.is_none() && other.iter.index.is_none() {
1865                return self.iter.data == other.iter.data;
1866            }
1867            Iterator::eq(self.iter, other.iter)
1868        }
1869    }
1870    impl<'a> Eq for DatumSeq<'a> {}
1871    impl<'a, 'b> PartialOrd<DatumSeq<'a>> for DatumSeq<'b> {
1872        #[inline(always)]
1873        fn partial_cmp(&self, other: &DatumSeq<'a>) -> Option<Ordering> {
1874            // Fast path: both sides are unencoded raw row bytes.
1875            if self.iter.index.is_none() && other.iter.index.is_none() {
1876                let left = self.iter.data;
1877                let right = other.iter.data;
1878                return Some(match left.len().cmp(&right.len()) {
1879                    Ordering::Equal => left.cmp(right),
1880                    other => other,
1881                });
1882            }
1883            // Slow path: at least one side is dictionary-encoded.
1884            // Fused length + lexicographic comparison in a single pass per side.
1885            // Row ordering is: shorter < longer; equal lengths compared lexicographically.
1886            //
1887            // We compare byte-by-byte (via `flatten`) rather than slice-by-slice on
1888            // purpose: a dictionary tag expands to a multi-byte value on one side while
1889            // the other side may store those same bytes raw, so the per-column slice
1890            // boundaries do not line up between the two iterators. Decoding to a flat
1891            // byte stream is the only representation in which both sides are directly
1892            // comparable. This path is cold — it only runs when at least one operand is
1893            // dictionary-encoded; the common unencoded case is handled by the fast path
1894            // above with a single slice comparison.
1895            let mut left = self.iter.flatten();
1896            let mut right = other.iter.flatten();
1897            let mut first_diff = Ordering::Equal;
1898            loop {
1899                match (left.next(), right.next()) {
1900                    (Some(l), Some(r)) => {
1901                        if first_diff == Ordering::Equal {
1902                            first_diff = l.cmp(r);
1903                        }
1904                    }
1905                    // Left exhausted first: left is shorter, so Less.
1906                    (None, Some(_)) => return Some(Ordering::Less),
1907                    // Right exhausted first: right is shorter, so Greater.
1908                    (Some(_), None) => return Some(Ordering::Greater),
1909                    // Same length: use first lexicographic difference.
1910                    (None, None) => return Some(first_diff),
1911                }
1912            }
1913        }
1914    }
1915    impl<'a> Ord for DatumSeq<'a> {
1916        #[inline(always)]
1917        fn cmp(&self, other: &Self) -> Ordering {
1918            self.partial_cmp(other).unwrap()
1919        }
1920    }
1921
1922    impl<'a> PartialEq<&'a Row> for DatumSeq<'a> {
1923        #[inline(always)]
1924        fn eq(&self, other: &&'a Row) -> bool {
1925            self.eq(&Self::borrow_as(*other))
1926        }
1927    }
1928
1929    // Lifetimes decoupled (`'b` independent of `'a`): the arrange machinery
1930    // requires `for<'b> DatumSeq<'a>: PartialEq<&'b RowRef>`, i.e. a fixed
1931    // `DatumSeq` must compare against a `&RowRef` of any lifetime.
1932    impl<'a, 'b> PartialEq<&'b RowRef> for DatumSeq<'a> {
1933        #[inline(always)]
1934        fn eq(&self, other: &&'b RowRef) -> bool {
1935            self.eq(&DatumSeq::borrow_as(*other))
1936        }
1937    }
1938
1939    impl<'a> DatumSeq<'a> {
1940        #[inline(always)]
1941        pub fn bytes_iter(self) -> ColumnsIter<'a> {
1942            self.iter
1943        }
1944    }
1945
1946    impl<'a> Iterator for DatumSeq<'a> {
1947        type Item = Datum<'a>;
1948        #[inline(always)]
1949        fn next(&mut self) -> Option<Self::Item> {
1950            // Delegate to `ColumnsIter`, which handles both the codec and no-codec
1951            // cases. The no-codec scan hot path is served directly by `extend_datums`
1952            // (which decodes without going through this iterator), so the only callers
1953            // left here are the codec-encoded `extend_datums`/`to_row` paths and tests;
1954            // none warrant a dedicated no-codec fast path.
1955            self.iter
1956                .next()
1957                .map(|mut bytes| unsafe { read_datum(&mut bytes) })
1958        }
1959    }
1960
1961    use mz_repr::RowArena;
1962    use mz_repr::fixed_length::ExtendDatums;
1963    impl<'long> ExtendDatums for DatumSeq<'long> {
1964        #[inline]
1965        fn extend_datums<'a>(
1966            &'a self,
1967            _arena: &'a RowArena,
1968            target: &mut Vec<Datum<'a>>,
1969            max: Option<usize>,
1970        ) {
1971            // Branch on codec presence ONCE per row rather than once per datum.
1972            // With no codec (the common, feature-off case) push raw datums in a
1973            // tight loop, matching the pre-dictionary path; with a codec, fall
1974            // back to the per-column iterator. This keeps the codec check out of
1975            // the per-datum loop — the source of the feature-off scan overhead.
1976            if self.iter.index.is_none() {
1977                let mut data = self.iter.data;
1978                match max {
1979                    Some(max) => {
1980                        let mut n = 0;
1981                        while n < max && !data.is_empty() {
1982                            target.push(unsafe { read_datum(&mut data) });
1983                            n += 1;
1984                        }
1985                    }
1986                    None => {
1987                        while !data.is_empty() {
1988                            target.push(unsafe { read_datum(&mut data) });
1989                        }
1990                    }
1991                }
1992            } else {
1993                match max {
1994                    Some(max) => target.extend((*self).take(max)),
1995                    None => target.extend(*self),
1996                }
1997            }
1998        }
1999    }
2000}
2001
2002/// Traits abstracting the processes of encoding and decoding row-encoded byte sequences.
2003///
2004/// It is unsafe to use these types to encode byte sequences that are not row-encoded,
2005/// as they are parsed out of contiguous `[u8]` slices using `mz_repr::read_datum`.
2006mod row_codec {
2007
2008    pub use self::misra_gries::MisraGries;
2009    pub use columns::{ColumnsCodec, ColumnsIter};
2010    pub use dictionary::DictionaryCodec;
2011    #[cfg(test)]
2012    pub use dictionary::SAFE_TAG_BASE;
2013
2014    // Deterministic hasher state for the codecs' hash maps: a fixed-seed
2015    // `ahash::RandomState` shared with `mz_timely_util`'s consolidation hasher, so
2016    // the heavy-hitter summaries — and therefore which values each codec compresses
2017    // — are identical across runs and replicas, as the old `BTreeMap` backing was.
2018    use mz_timely_util::hash::fixed_state;
2019
2020    // The codecs encode and decode `[u8]` data specific to the `[Row]` encoding. They
2021    // soundly decode data they themselves encoded from valid `[Row]` data, but may be
2022    // unsound if asked to decode data that was not row-encoded, or was encoded with a
2023    // different codec. `ColumnsCodec` (a per-column wrapper around `DictionaryCodec`) is
2024    // the only codec the spine instantiates; the methods are inherent rather than behind
2025    // a `Codec` trait because nothing ever dispatches over codecs generically.
2026
2027    mod columns {
2028
2029        use mz_repr::{RowRef, read_datum};
2030
2031        use super::DictionaryCodec;
2032
2033        /// Independently encodes each column.
2034        #[derive(Default, Debug)]
2035        pub struct ColumnsCodec {
2036            columns: Vec<DictionaryCodec>,
2037        }
2038
2039        impl ColumnsCodec {
2040            /// Decode a row-encoded byte slice into per-column byte slices.
2041            pub(crate) fn decode<'a>(&'a self, bytes: &'a [u8]) -> ColumnsIter<'a> {
2042                ColumnsIter {
2043                    index: Some(self),
2044                    column: 0,
2045                    data: bytes,
2046                }
2047            }
2048            /// Encode a sequence of column byte slices, updating per-column statistics.
2049            pub(crate) fn encode<'a, I>(&mut self, iter: I, output: &mut Vec<u8>)
2050            where
2051                I: IntoIterator<Item = &'a [u8]>,
2052            {
2053                for (index, bytes) in iter.into_iter().enumerate() {
2054                    if self.columns.len() <= index {
2055                        self.columns.push(Default::default());
2056                    }
2057                    self.columns[index].encode(std::iter::once(bytes), output);
2058                }
2059            }
2060
2061            /// Construct a codec valid for the union of the supplied codecs' data.
2062            pub(crate) fn new_from<'a>(stats: impl IntoIterator<Item = &'a Self>) -> Self {
2063                // An empty `stats` iterator yields a zero-column codec, which encodes and
2064                // decodes nothing; callers merging no inputs get an inert (but sound) codec.
2065                let stats = stats.into_iter().collect::<Vec<_>>();
2066                let cols = stats.iter().map(|s| s.columns.len()).max().unwrap_or(0);
2067                let mut columns = Vec::with_capacity(cols);
2068                let default: DictionaryCodec = Default::default();
2069                for index in 0..cols {
2070                    columns.push(DictionaryCodec::new_from(
2071                        stats
2072                            .iter()
2073                            .map(|s| s.columns.get(index).unwrap_or(&default)),
2074                    ));
2075                }
2076                Self { columns }
2077            }
2078
2079            /// Reveal a row's bytes for fast-path comparison, with no codec to consult.
2080            #[inline(always)]
2081            pub(crate) fn borrow_row(row: &RowRef) -> ColumnsIter<'_> {
2082                ColumnsIter {
2083                    index: None,
2084                    column: 0,
2085                    data: row.data(),
2086                }
2087            }
2088        }
2089
2090        impl ColumnsCodec {
2091            /// Visit contained allocations to determine their size and capacity.
2092            pub(crate) fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2093                let elem = std::mem::size_of::<DictionaryCodec>();
2094                callback(self.columns.len() * elem, self.columns.capacity() * elem);
2095                for column in &self.columns {
2096                    column.heap_size(callback);
2097                }
2098            }
2099        }
2100
2101        impl ColumnsCodec {
2102            /// Record a row's column values in the statistics without encoding.
2103            ///
2104            /// Used during the stats-gathering phase, where we want the heavy-hitter
2105            /// and tag-usage information but store the row raw, so encoding into a
2106            /// throwaway buffer would be pure waste.
2107            #[inline]
2108            pub(crate) fn observe<'a, I>(&mut self, iter: I)
2109            where
2110                I: IntoIterator<Item = &'a [u8]>,
2111            {
2112                for (index, bytes) in iter.into_iter().enumerate() {
2113                    if self.columns.len() <= index {
2114                        self.columns.push(Default::default());
2115                    }
2116                    self.columns[index].observe(bytes);
2117                }
2118            }
2119        }
2120
2121        impl ColumnsCodec {
2122            /// Construct a codec using only structurally safe tags.
2123            ///
2124            /// Consumes `self`: this is only ever called on stats that have just been
2125            /// `take`n out of a container and are about to be discarded, so we move the
2126            /// per-column `MisraGries` summaries through rather than cloning them.
2127            pub(crate) fn new_safe(self) -> Self {
2128                let columns = self
2129                    .columns
2130                    .into_iter()
2131                    .map(DictionaryCodec::new_safe)
2132                    .collect();
2133                Self { columns }
2134            }
2135        }
2136
2137        #[derive(Debug, Copy, Clone)]
2138        pub struct ColumnsIter<'a> {
2139            // `None` when iterating an owned row directly, with no codec to consult.
2140            pub index: Option<&'a ColumnsCodec>,
2141            pub column: usize,
2142            pub data: &'a [u8],
2143        }
2144
2145        impl<'a> Iterator for ColumnsIter<'a> {
2146            type Item = &'a [u8];
2147            #[inline(always)]
2148            fn next(&mut self) -> Option<Self::Item> {
2149                if self.data.is_empty() {
2150                    None
2151                } else if let Some(bytes) = self
2152                    .index
2153                    .as_ref()
2154                    .and_then(|i| i.columns.get(self.column))
2155                    .and_then(|i| i.decode.get(self.data[0].into()))
2156                {
2157                    self.data = &self.data[1..];
2158                    self.column += 1;
2159                    Some(bytes)
2160                } else {
2161                    let mut data = self.data;
2162                    let data_len = data.len();
2163                    unsafe {
2164                        read_datum(&mut data);
2165                    }
2166                    let (prev, next) = self.data.split_at(data_len - data.len());
2167                    self.data = next;
2168                    self.column += 1;
2169                    Some(prev)
2170                }
2171            }
2172        }
2173
2174        impl<'a> ColumnsIter<'a> {
2175            /// Create a column iterator without a codec.
2176            ///
2177            /// This requires the data to be row-formatted, and it will be erroneous otherwise.
2178            #[inline(always)]
2179            pub unsafe fn without_codec(data: &'a [u8]) -> Self {
2180                Self {
2181                    index: None,
2182                    column: 0,
2183                    data,
2184                }
2185            }
2186        }
2187    }
2188
2189    /// A dictionary encoding codec for `[Row]` data.
2190    ///
2191    /// The dictionary harvests unused tags within each column and uses them to
2192    /// represent popular values within that column. There are two mechanisms it
2193    /// uses to accomplish this:
2194    ///
2195    /// 1. Statically free tags: `SAFE_TAG_BASE` is taken as an exclusive upper bound
2196    ///    on the tags that will be used by `[Row]`, and tags greater or equal to this
2197    ///    value are always safe to use.
2198    /// 2. Dynamically free tags: having seen an entire collection, we can use any
2199    ///    tag not otherwise used by the collection, as it would not be ambiguous.
2200    ///
2201    /// It goes without saying that if either of these approaches are incorrect,
2202    /// there are calamitous unsoundness implications.
2203    mod dictionary {
2204        // The `encode` map is a pure value->tag lookup table (never iterated for logic),
2205        // so `mz_ore::collections::HashMap`'s order-hiding would suffice — but it offers
2206        // no fixed-seed constructor, and we want the same deterministic hasher as the
2207        // summary above. `heap_size`'s `keys()` walk is an order-insensitive sum.
2208        #![allow(clippy::disallowed_types)]
2209
2210        use std::collections::HashMap;
2211
2212        use super::fixed_state;
2213        pub use super::{BytesMap, MisraGries};
2214
2215        /// First byte value that is structurally unused by the datum encoding.
2216        /// All byte values >= this are safe to use as dictionary tags without
2217        /// observing the data, since no datum's first byte can have this value.
2218        ///
2219        /// `mz_repr`'s `Row` `Tag` enum currently has 84 variants (discriminants
2220        /// 0..=83), so the truly tight bound is 84. We deliberately pick a larger,
2221        /// round-ish constant to leave headroom for new tags without having to also
2222        /// bump the safe set, and the `test_safe_tag_base` test pins the real
2223        /// invariant: every datum the row format produces must encode with a first
2224        /// byte strictly less than this value. If a future tag crosses the boundary
2225        /// that test fails loudly rather than silently corrupting decoding.
2226        pub const SAFE_TAG_BASE: u8 = 122;
2227
2228        /// Per-column dictionary codec. Encodes column byte slices, replacing popular
2229        /// values with spare tags; decoding is performed by `ColumnsIter` reading the
2230        /// `decode` map directly.
2231        #[derive(Default, Debug)]
2232        pub struct DictionaryCodec {
2233            // Looked up once per value on the encode path; mostly misses (only popular
2234            // values compress), so a hash map beats a `BTreeMap`'s byte-slice walk. The
2235            // map is only ever read via `get` — never iterated — so its hasher seed has
2236            // no observable effect; the populated maps are built with `fixed_state` in
2237            // `new_from`/`new_safe` for consistency, while the derived-`Default` (stats
2238            // accumulator) variant stays empty and is never consulted.
2239            encode: HashMap<Vec<u8>, u8, ahash::RandomState>,
2240            pub decode: BytesMap,
2241            stats: (MisraGries<Vec<u8>>, [u64; 4]),
2242        }
2243
2244        impl DictionaryCodec {
2245            /// Encode a sequence of byte slices.
2246            ///
2247            /// Encoding also records statistics about the structure of the input.
2248            ///
2249            /// Decoding has no symmetric method here: a column's bytes are decoded by
2250            /// `ColumnsIter`, which consults the `decode` map directly.
2251            pub(super) fn encode<'a, I>(&mut self, iter: I, output: &mut Vec<u8>)
2252            where
2253                I: IntoIterator<Item = &'a [u8]>,
2254            {
2255                for bytes in iter.into_iter() {
2256                    mz_ore::soft_assert_no_log!(
2257                        !bytes.is_empty(),
2258                        "row encoding never yields empty column slices",
2259                    );
2260                    // If we have an index referencing `bytes`, use the index key.
2261                    if let Some(b) = self.encode.get(bytes) {
2262                        output.push(*b);
2263                    } else {
2264                        // Raw fall-through. Soundness rests on `bytes[0]` never being a
2265                        // tag we hand out as a dictionary key: `new_from`/`new_safe` only
2266                        // assign dictionary tags from first-byte values that were never
2267                        // observed (or are `>= SAFE_TAG_BASE`, which no datum first-byte
2268                        // can equal). If a literal datum's first byte collided with a
2269                        // dictionary tag, `decode` would resolve it to the dictionary
2270                        // entry instead of reading the datum. This `debug_assert` makes
2271                        // the load-bearing "no later first-byte outside the observed
2272                        // union" invariant self-checking.
2273                        mz_ore::soft_assert_no_log!(
2274                            self.decode.get(bytes[0].into()).is_none(),
2275                            "raw datum first-byte {} collides with a dictionary tag; \
2276                             decode would be ambiguous",
2277                            bytes[0],
2278                        );
2279                        output.extend(bytes);
2280                    }
2281                    self.observe(bytes);
2282                }
2283            }
2284
2285            /// Construct a new encoder from supplied statistics.
2286            pub(super) fn new_from<'a>(stats: impl IntoIterator<Item = &'a Self>) -> Self {
2287                // Collect most popular bytes from combined containers.
2288                let mut mg = MisraGries::default();
2289                let mut tags: [u64; 4] = [0; 4];
2290                for stat in stats.into_iter() {
2291                    for (thing, count) in stat.stats.0.clone().done() {
2292                        mg.update(thing, count);
2293                    }
2294                    tags[0] |= stat.stats.1[0];
2295                    tags[1] |= stat.stats.1[1];
2296                    tags[2] |= stat.stats.1[2];
2297                    tags[3] |= stat.stats.1[3];
2298                }
2299                let mut mg = mg
2300                    .done()
2301                    .into_iter()
2302                    .filter(|(next_bytes, count)| next_bytes.len() > 1 && count > &1);
2303                // Establish encoding and decoding rules.
2304                let mut encode = HashMap::with_hasher(fixed_state());
2305                let mut decode = BytesMap::default();
2306                for tag in 0..=255 {
2307                    let tag_idx: usize = (tag % 4).into();
2308                    let shift = tag >> 2;
2309                    if (tags[tag_idx] >> shift) & 0x01 != 0 {
2310                        // Tag is used by a literal datum first-byte; reserve the slot.
2311                        decode.push(None);
2312                    } else if let Some((next_bytes, _count)) = mg.next() {
2313                        decode.push(Some(&next_bytes[..]));
2314                        encode.insert(next_bytes, tag);
2315                    } else {
2316                        // Unused tag, but the heavy-hitter supply is exhausted. We must
2317                        // still push a slot so that `decode`'s index stays aligned with
2318                        // the tag value: every iteration pushes exactly once, keeping the
2319                        // map length 256 and `decode.get(tag)` addressable by tag.
2320                        decode.push(None);
2321                    }
2322                }
2323
2324                Self {
2325                    encode,
2326                    decode,
2327                    stats: (MisraGries::default(), [0u64; 4]),
2328                }
2329            }
2330        }
2331
2332        impl DictionaryCodec {
2333            /// Visit contained allocations to determine their size and capacity.
2334            ///
2335            /// The `encode` table is approximated as one logical entry's worth of bytes
2336            /// per element for size and its reserved `capacity()` for capacity; the
2337            /// dominant terms (the owned key bytes and the `decode` map's byte arena)
2338            /// are accounted exactly.
2339            pub fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2340                let entry = std::mem::size_of::<(Vec<u8>, u8)>();
2341                callback(self.encode.len() * entry, self.encode.capacity() * entry);
2342                for key in self.encode.keys() {
2343                    callback(key.len(), key.capacity());
2344                }
2345                self.decode.heap_size(callback);
2346                self.stats.0.heap_size(callback);
2347            }
2348
2349            /// Record a single column value in this codec's statistics without
2350            /// producing any encoded output.
2351            ///
2352            /// Statistics come in two decoupled parts, with very different costs and
2353            /// purposes:
2354            ///
2355            /// 1. The tag bitmap (`stats.1`) records which first-byte values have been
2356            ///    observed. It is cheap (four `u64` ORs) and *soundness critical*:
2357            ///    `new_from`'s dynamic-tag path only hands out tags that this bitmap
2358            ///    reports as unused, so it must stay accurate for the entire life of the
2359            ///    codec, including on the hot encode path.
2360            /// 2. The MisraGries summary (`stats.0`) tracks heavy hitters and only
2361            ///    affects *which* values a future codec compresses, never correctness.
2362            ///    It is the expensive part (a `BTreeMap` insert per column per row). We
2363            ///    keep feeding it after install, on the hot encode path, on purpose: a
2364            ///    later merge rebuilds the merged codec from these summaries via
2365            ///    `new_from`. If we froze the summary at install time, then as the
2366            ///    collection evolves — records cancel under consolidation, the popular
2367            ///    set drifts — the codec could never reclaim slots for newly-popular
2368            ///    values and would eventually be left compressing values that no longer
2369            ///    occur, ceasing to compress the ones that do.
2370            #[inline]
2371            pub fn observe(&mut self, bytes: &[u8]) {
2372                mz_ore::soft_assert_no_log!(
2373                    !bytes.is_empty(),
2374                    "row encoding never yields empty column slices",
2375                );
2376                let tag = bytes[0];
2377                let tag_idx: usize = (tag % 4).into();
2378                self.stats.1[tag_idx] |= 1 << (tag >> 2);
2379                self.stats.0.insert_ref(bytes);
2380            }
2381
2382            /// Construct a codec using only structurally safe tags (>= SAFE_TAG_BASE).
2383            /// These tags never collide with datum first-bytes, so the codec can be
2384            /// installed without observing all data first.
2385            pub(super) fn new_safe(stats: Self) -> Self {
2386                // The container stores its pre-install rows raw, so the first-byte
2387                // bitmap (`stats.1`) gathered while observing them must carry over to
2388                // the installed codec. The bitmap is soundness-critical: a later
2389                // `new_from` merge consults it to decide which one-byte tags are free
2390                // to hand out as dictionary keys. If we dropped it here, the merge
2391                // could assign a dictionary tag equal to a pre-install datum's first
2392                // byte, after which `decode` would resolve that literal datum to the
2393                // dictionary entry. The MisraGries summary (`stats.0`), by contrast,
2394                // is consumed below to seed the dictionary and is reset, since the
2395                // installed codec re-accumulates it from rows it sees post-install.
2396                let (mg, observed_tags) = stats.stats;
2397                let mut mg = mg
2398                    .done()
2399                    .into_iter()
2400                    .filter(|(next_bytes, count)| next_bytes.len() > 1 && count > &1);
2401                let mut encode = HashMap::with_hasher(fixed_state());
2402                let mut decode = BytesMap::default();
2403                // Fill slots 0..SAFE_TAG_BASE with None (reserved for datum tags).
2404                for _ in 0..SAFE_TAG_BASE {
2405                    decode.push(None);
2406                }
2407                // Assign dictionary entries to safe tags.
2408                for tag in SAFE_TAG_BASE..=255 {
2409                    if let Some((next_bytes, _count)) = mg.next() {
2410                        decode.push(Some(&next_bytes[..]));
2411                        encode.insert(next_bytes, tag);
2412                    }
2413                }
2414                Self {
2415                    encode,
2416                    decode,
2417                    stats: (MisraGries::default(), observed_tags),
2418                }
2419            }
2420        }
2421    }
2422
2423    /// A map from `0 .. something` to `Option<&[u8]>`.
2424    ///
2425    /// Non-empty slices are pushed in order, and can be retrieved by index.
2426    /// Pushing an empty slice is equivalent to pushing `None`.
2427    #[derive(Debug)]
2428    pub struct BytesMap {
2429        offsets: Vec<usize>,
2430        bytes: Vec<u8>,
2431    }
2432    impl Default for BytesMap {
2433        #[inline(always)]
2434        fn default() -> Self {
2435            Self {
2436                offsets: vec![0],
2437                bytes: Vec::new(),
2438            }
2439        }
2440    }
2441    impl BytesMap {
2442        #[inline]
2443        fn push(&mut self, input: Option<&[u8]>) {
2444            if let Some(bytes) = input {
2445                self.bytes.extend(bytes);
2446            }
2447            self.offsets.push(self.bytes.len());
2448        }
2449        /// Visit contained allocations to determine their size and capacity.
2450        fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2451            let off = std::mem::size_of::<usize>();
2452            callback(self.offsets.len() * off, self.offsets.capacity() * off);
2453            callback(self.bytes.len(), self.bytes.capacity());
2454        }
2455        #[inline]
2456        fn get(&self, index: usize) -> Option<&[u8]> {
2457            if index < self.offsets.len() - 1 {
2458                let lower = self.offsets[index];
2459                let upper = self.offsets[index + 1];
2460                if lower < upper {
2461                    Some(&self.bytes[lower..upper])
2462                } else {
2463                    None
2464                }
2465            } else {
2466                None
2467            }
2468        }
2469    }
2470
2471    mod misra_gries {
2472        // The summary must iterate its entries (to extract heavy hitters in `done`, to
2473        // `tidy`, and to size itself), which `mz_ore::collections::HashMap` deliberately
2474        // forbids. We instead get determinism from the fixed-seed hasher (`fixed_state`)
2475        // plus the total-order sort in `done`; `tidy`/`heap_size` are order-insensitive.
2476        #![allow(clippy::disallowed_types)]
2477
2478        use std::collections::HashMap;
2479        use std::hash::Hash;
2480
2481        use super::fixed_state;
2482
2483        /// Maintains a summary of "heavy hitters" in a presented collection of items.
2484        ///
2485        /// Uses a hash map internally so that repeated observations of the same
2486        /// element only allocate once (on first sighting), and so the per-element
2487        /// `insert_ref` is an O(1) hash rather than an O(log n) walk of byte-slice
2488        /// comparisons. This is the hot path: one lookup per column per row, fed both
2489        /// while gathering stats and on the steady-state encode path. The hasher is
2490        /// fixed-seed (see [`fixed_state`]) so the summary — and thus which values a
2491        /// codec compresses — stays deterministic across runs and replicas.
2492        ///
2493        /// Tidy is performed when the number of *distinct* elements exceeds `2 * k`,
2494        /// reducing to at most `k` entries.
2495        #[derive(Clone, Debug)]
2496        pub struct MisraGries<T: Ord + Hash> {
2497            inner: HashMap<T, usize, ahash::RandomState>,
2498            k: usize,
2499        }
2500
2501        impl<T: Ord + Hash> Default for MisraGries<T> {
2502            #[inline(always)]
2503            fn default() -> Self {
2504                Self {
2505                    inner: HashMap::with_hasher(fixed_state()),
2506                    k: 512,
2507                }
2508            }
2509        }
2510
2511        impl<T: Ord + Hash> MisraGries<T> {
2512            /// Inserts an additional element to the summary.
2513            #[inline(always)]
2514            pub fn insert(&mut self, element: T) {
2515                self.update(element, 1);
2516            }
2517            /// Inserts multiple copies of an element to the summary.
2518            #[inline]
2519            pub fn update(&mut self, element: T, count: usize) {
2520                *self.inner.entry(element).or_insert(0) += count;
2521                if self.inner.len() > 2 * self.k {
2522                    self.tidy();
2523                }
2524            }
2525
2526            /// Completes the summary, and extracts the items and their counts.
2527            pub fn done(self) -> Vec<(T, usize)> {
2528                let mut result: Vec<_> = self.inner.into_iter().collect();
2529                // Descending count, ties broken by key, so the values a codec selects
2530                // are deterministic regardless of hash-map iteration order.
2531                result.sort_by(|x, y| y.1.cmp(&x.1).then_with(|| x.0.cmp(&y.0)));
2532                result
2533            }
2534
2535            /// Reduces the summary down to at most `k` distinct items by
2536            /// subtracting the (k+1)-th largest count from all entries and
2537            /// discarding those that drop to zero or below.
2538            fn tidy(&mut self) {
2539                let mut counts: Vec<usize> = self.inner.values().copied().collect();
2540                counts.sort_unstable_by(|a, b| b.cmp(a));
2541                // The (k+1)-th largest count, or 0 if fewer than k+1 entries.
2542                let sub_weight = counts.get(self.k).copied().unwrap_or(0);
2543                if sub_weight > 0 {
2544                    self.inner.retain(|_, count| {
2545                        *count = count.saturating_sub(sub_weight);
2546                        *count > 0
2547                    });
2548                }
2549            }
2550        }
2551
2552        impl MisraGries<Vec<u8>> {
2553            /// Visit contained allocations to determine their size and capacity.
2554            ///
2555            /// The hash table is approximated as one logical entry per element for
2556            /// size and its reserved `capacity()` for capacity; the owned key bytes
2557            /// are accounted exactly.
2558            pub fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2559                let entry = std::mem::size_of::<(Vec<u8>, usize)>();
2560                callback(self.inner.len() * entry, self.inner.capacity() * entry);
2561                for key in self.inner.keys() {
2562                    callback(key.len(), key.capacity());
2563                }
2564            }
2565
2566            /// Insert a borrowed byte slice, only allocating if the key is new.
2567            #[inline]
2568            pub fn insert_ref(&mut self, element: &[u8]) {
2569                if let Some(count) = self.inner.get_mut(element) {
2570                    *count += 1;
2571                } else {
2572                    self.insert(element.to_owned());
2573                }
2574            }
2575        }
2576
2577        impl<T: Ord + Hash> std::ops::AddAssign for MisraGries<T> {
2578            fn add_assign(&mut self, rhs: Self) {
2579                for (element, count) in rhs.done() {
2580                    self.update(element, count);
2581                }
2582            }
2583        }
2584    }
2585}