Skip to main content

mz_repr/row/
encode.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//! A permanent storage encoding for rows.
11//!
12//! See row.proto for details.
13
14use std::fmt::Debug;
15use std::ops::AddAssign;
16use std::sync::Arc;
17
18use anyhow::anyhow;
19use arrow::array::{
20    Array, ArrayBuilder, ArrayRef, BinaryArray, BinaryBuilder, BooleanArray, BooleanBufferBuilder,
21    BooleanBuilder, FixedSizeBinaryArray, FixedSizeBinaryBuilder, Float32Array, Float32Builder,
22    Float64Array, Float64Builder, Int16Array, Int16Builder, Int32Array, Int32Builder, Int64Array,
23    Int64Builder, ListArray, ListBuilder, MapArray, StringArray, StringBuilder, StructArray,
24    UInt8Array, UInt8Builder, UInt16Array, UInt16Builder, UInt32Array, UInt32Builder, UInt64Array,
25    UInt64Builder, make_array,
26};
27use arrow::buffer::{BooleanBuffer, Buffer, NullBuffer, OffsetBuffer, ScalarBuffer};
28use arrow::datatypes::{DataType, Field, Fields, ToByteSlice};
29use bytes::{BufMut, Bytes};
30use chrono::Timelike;
31use dec::{Context, Decimal, OrderedDecimal};
32use itertools::{EitherOrBoth, Itertools};
33use mz_ore::assert_none;
34use mz_ore::cast::CastFrom;
35use mz_persist_types::Codec;
36use mz_persist_types::arrow::ArrayOrd;
37use mz_persist_types::columnar::{ColumnDecoder, ColumnEncoder, FixedSizeCodec, Schema};
38use mz_persist_types::stats::{
39    ColumnNullStats, ColumnStatKinds, ColumnarStats, ColumnarStatsBuilder, FixedSizeBytesStatsKind,
40    OptionStats, PrimitiveStats, StructStats,
41};
42use mz_proto::chrono::ProtoNaiveTime;
43use mz_proto::{ProtoType, RustType, TryFromProtoError};
44use prost::Message;
45use uuid::Uuid;
46
47use crate::adt::array::{ArrayDimension, PackedArrayDimension};
48use crate::adt::date::Date;
49use crate::adt::datetime::PackedNaiveTime;
50use crate::adt::interval::PackedInterval;
51use crate::adt::jsonb::{JsonbPacker, JsonbRef};
52use crate::adt::mz_acl_item::{PackedAclItem, PackedMzAclItem};
53use crate::adt::numeric::{NUMERIC_DATUM_MAX_PRECISION, Numeric, PackedNumeric};
54use crate::adt::range::{Range, RangeInner, RangeLowerBound, RangeUpperBound};
55use crate::adt::timestamp::{CheckedTimestamp, PackedNaiveDateTime};
56use crate::row::proto_datum::DatumType;
57use crate::row::{
58    ProtoArray, ProtoArrayDimension, ProtoDatum, ProtoDatumOther, ProtoDict, ProtoDictElement,
59    ProtoNumeric, ProtoRange, ProtoRangeInner, ProtoRow,
60};
61use crate::stats::{fixed_stats_from_column, numeric_stats_from_column, stats_for_json};
62use crate::{Datum, ProtoRelationDesc, RelationDesc, Row, RowPacker, SqlScalarType, Timestamp};
63
64// TODO(parkmycar): Benchmark the difference between `FixedSizeBinaryArray` and `BinaryArray`.
65//
66// `FixedSizeBinaryArray`s push empty bytes when a value is null which for larger binary types
67// could result in poor performance.
68#[allow(clippy::as_conversions)]
69mod fixed_binary_sizes {
70    use super::*;
71
72    pub const TIME_FIXED_BYTES: i32 = PackedNaiveTime::SIZE as i32;
73    pub const TIMESTAMP_FIXED_BYTES: i32 = PackedNaiveDateTime::SIZE as i32;
74    pub const INTERVAL_FIXED_BYTES: i32 = PackedInterval::SIZE as i32;
75    pub const ACL_ITEM_FIXED_BYTES: i32 = PackedAclItem::SIZE as i32;
76    pub const _MZ_ACL_ITEM_FIXED_BYTES: i32 = PackedMzAclItem::SIZE as i32;
77    pub const ARRAY_DIMENSION_FIXED_BYTES: i32 = PackedArrayDimension::SIZE as i32;
78
79    pub const UUID_FIXED_BYTES: i32 = 16;
80    static_assertions::const_assert_eq!(UUID_FIXED_BYTES as usize, std::mem::size_of::<Uuid>());
81}
82use fixed_binary_sizes::*;
83
84/// Returns true iff the ordering of the "raw" and Persist-encoded versions of this columm would match:
85/// ie. `sort(encode(column)) == encode(sort(column))`. This encoding has been designed so that this
86/// is true for many types.
87pub fn preserves_order(scalar_type: &SqlScalarType) -> bool {
88    match scalar_type {
89        // These types have short, fixed-length encodings that are designed to sort identically.
90        SqlScalarType::Bool
91        | SqlScalarType::Int16
92        | SqlScalarType::Int32
93        | SqlScalarType::Int64
94        | SqlScalarType::UInt16
95        | SqlScalarType::UInt32
96        | SqlScalarType::UInt64
97        | SqlScalarType::Date
98        | SqlScalarType::Time
99        | SqlScalarType::Timestamp { .. }
100        | SqlScalarType::TimestampTz { .. }
101        | SqlScalarType::Interval
102        | SqlScalarType::Bytes
103        | SqlScalarType::String
104        | SqlScalarType::Uuid
105        | SqlScalarType::MzTimestamp
106        | SqlScalarType::MzAclItem
107        | SqlScalarType::AclItem => true,
108        // We sort records lexicographically; a record has a meaningful sort if all its fields do.
109        SqlScalarType::Record { fields, .. } => fields
110            .iter()
111            .all(|(_, field_type)| preserves_order(&field_type.scalar_type)),
112        // Our floating-point encoding preserves order generally, but differs when comparing
113        // -0 and 0. Opt these out for now.
114        SqlScalarType::Float32 | SqlScalarType::Float64 => false,
115        // Numeric is sensitive to similar ordering issues as floating point numbers, and requires
116        // some special handling we don't have yet.
117        SqlScalarType::Numeric { .. } => false,
118        // For all other types: either the encoding is known to not preserve ordering, or we
119        // don't yet care to make strong guarantees one way or the other.
120        SqlScalarType::PgLegacyChar
121        | SqlScalarType::PgLegacyName
122        | SqlScalarType::Char { .. }
123        | SqlScalarType::VarChar { .. }
124        | SqlScalarType::Jsonb
125        | SqlScalarType::Array(_)
126        | SqlScalarType::List { .. }
127        | SqlScalarType::Oid
128        | SqlScalarType::Map { .. }
129        | SqlScalarType::RegProc
130        | SqlScalarType::RegType
131        | SqlScalarType::RegClass
132        | SqlScalarType::Int2Vector
133        | SqlScalarType::Range { .. } => false,
134    }
135}
136
137/// An encoder for a column of [`Datum`]s.
138#[derive(Debug)]
139struct DatumEncoder {
140    nullable: bool,
141    encoder: DatumColumnEncoder,
142}
143
144impl DatumEncoder {
145    fn goodbytes(&self) -> usize {
146        self.encoder.goodbytes()
147    }
148
149    fn push(&mut self, datum: Datum) {
150        assert!(
151            !datum.is_null() || self.nullable,
152            "tried pushing Null into non-nullable column"
153        );
154        self.encoder.push(datum);
155    }
156
157    fn push_invalid(&mut self) {
158        self.encoder.push_invalid();
159    }
160
161    fn finish(self) -> ArrayRef {
162        self.encoder.finish()
163    }
164}
165
166/// An encoder for a single column of [`Datum`]s. To encode an entire row see
167/// [`RowColumnarEncoder`].
168///
169/// Note: We specifically structure the encoder as an enum instead of using trait objects because
170/// Datum encoding is an extremely hot path and downcasting objects is relatively slow.
171#[derive(Debug)]
172enum DatumColumnEncoder {
173    Bool(BooleanBuilder),
174    U8(UInt8Builder),
175    U16(UInt16Builder),
176    U32(UInt32Builder),
177    U64(UInt64Builder),
178    I16(Int16Builder),
179    I32(Int32Builder),
180    I64(Int64Builder),
181    F32(Float32Builder),
182    F64(Float64Builder),
183    Numeric {
184        /// The raw bytes so we can losslessly roundtrip Numerics.
185        binary_values: BinaryBuilder,
186        /// Also maintain a float64 approximation for sorting.
187        approx_values: Float64Builder,
188        /// Re-usable `libdecimal` context for conversions.
189        numeric_context: Context<Numeric>,
190    },
191    Bytes(BinaryBuilder),
192    String(StringBuilder),
193    Date(Int32Builder),
194    Time(FixedSizeBinaryBuilder),
195    Timestamp(FixedSizeBinaryBuilder),
196    TimestampTz(FixedSizeBinaryBuilder),
197    MzTimestamp(UInt64Builder),
198    Interval(FixedSizeBinaryBuilder),
199    Uuid(FixedSizeBinaryBuilder),
200    AclItem(FixedSizeBinaryBuilder),
201    MzAclItem(BinaryBuilder),
202    Range(BinaryBuilder),
203    /// Hand rolled "StringBuilder" that reduces the number of copies required
204    /// to serialize JSON.
205    ///
206    /// An alternative would be to use [`StringBuilder`] but that requires
207    /// serializing to an intermediary string, and then copying that
208    /// intermediary string into an underlying buffer.
209    Jsonb {
210        /// Monotonically increasing offsets of each encoded segment.
211        offsets: Vec<i32>,
212        /// Buffer that contains UTF-8 encoded JSON.
213        buf: Vec<u8>,
214        /// Null entries, if any.
215        nulls: Option<BooleanBufferBuilder>,
216    },
217    Array {
218        /// Binary encoded `ArrayDimension`s.
219        dims: ListBuilder<FixedSizeBinaryBuilder>,
220        /// Lengths of each `Array` in this column.
221        val_lengths: Vec<usize>,
222        /// Contiguous array of underlying data.
223        vals: Box<DatumColumnEncoder>,
224        /// Null entires, if any.
225        nulls: Option<BooleanBufferBuilder>,
226    },
227    List {
228        /// Lengths of each `List` in this column.
229        lengths: Vec<usize>,
230        /// Contiguous array of underlying data.
231        values: Box<DatumColumnEncoder>,
232        /// Null entires, if any.
233        nulls: Option<BooleanBufferBuilder>,
234    },
235    Map {
236        /// Lengths of each `Map` in this column
237        lengths: Vec<usize>,
238        /// Contiguous array of key data.
239        keys: StringBuilder,
240        /// Contiguous array of val data.
241        vals: Box<DatumColumnEncoder>,
242        /// Null entires, if any.
243        nulls: Option<BooleanBufferBuilder>,
244    },
245    Record {
246        /// Columns in the record.
247        fields: Vec<DatumEncoder>,
248        /// Null entries, if any.
249        nulls: Option<BooleanBufferBuilder>,
250        /// Number of values we've pushed into this builder thus far.
251        length: usize,
252    },
253    /// Special encoder for a [`SqlScalarType::Record`] that has no inner fields.
254    ///
255    /// We have a special case for this scenario because Arrow does not allow a
256    /// [`StructArray`] (what normally use to encod a `Record`) with no fields.
257    RecordEmpty(BooleanBuilder),
258}
259
260impl DatumColumnEncoder {
261    fn goodbytes(&self) -> usize {
262        match self {
263            DatumColumnEncoder::Bool(a) => a.len(),
264            DatumColumnEncoder::U8(a) => a.values_slice().to_byte_slice().len(),
265            DatumColumnEncoder::U16(a) => a.values_slice().to_byte_slice().len(),
266            DatumColumnEncoder::U32(a) => a.values_slice().to_byte_slice().len(),
267            DatumColumnEncoder::U64(a) => a.values_slice().to_byte_slice().len(),
268            DatumColumnEncoder::I16(a) => a.values_slice().to_byte_slice().len(),
269            DatumColumnEncoder::I32(a) => a.values_slice().to_byte_slice().len(),
270            DatumColumnEncoder::I64(a) => a.values_slice().to_byte_slice().len(),
271            DatumColumnEncoder::F32(a) => a.values_slice().to_byte_slice().len(),
272            DatumColumnEncoder::F64(a) => a.values_slice().to_byte_slice().len(),
273            DatumColumnEncoder::Numeric {
274                binary_values,
275                approx_values,
276                ..
277            } => {
278                binary_values.values_slice().len()
279                    + approx_values.values_slice().to_byte_slice().len()
280            }
281            DatumColumnEncoder::Bytes(a) => a.values_slice().len(),
282            DatumColumnEncoder::String(a) => a.values_slice().len(),
283            DatumColumnEncoder::Date(a) => a.values_slice().to_byte_slice().len(),
284            DatumColumnEncoder::Time(a) => a.len() * PackedNaiveTime::SIZE,
285            DatumColumnEncoder::Timestamp(a) => a.len() * PackedNaiveDateTime::SIZE,
286            DatumColumnEncoder::TimestampTz(a) => a.len() * PackedNaiveDateTime::SIZE,
287            DatumColumnEncoder::MzTimestamp(a) => a.values_slice().to_byte_slice().len(),
288            DatumColumnEncoder::Interval(a) => a.len() * PackedInterval::SIZE,
289            DatumColumnEncoder::Uuid(a) => a.len() * size_of::<Uuid>(),
290            DatumColumnEncoder::AclItem(a) => a.len() * PackedAclItem::SIZE,
291            DatumColumnEncoder::MzAclItem(a) => a.values_slice().len(),
292            DatumColumnEncoder::Range(a) => a.values_slice().len(),
293            DatumColumnEncoder::Jsonb { buf, .. } => buf.len(),
294            DatumColumnEncoder::Array { dims, vals, .. } => {
295                dims.len() * PackedArrayDimension::SIZE + vals.goodbytes()
296            }
297            DatumColumnEncoder::List { values, .. } => values.goodbytes(),
298            DatumColumnEncoder::Map { keys, vals, .. } => {
299                keys.values_slice().len() + vals.goodbytes()
300            }
301            DatumColumnEncoder::Record { fields, .. } => fields.iter().map(|f| f.goodbytes()).sum(),
302            DatumColumnEncoder::RecordEmpty(a) => a.len(),
303        }
304    }
305
306    fn push<'e, 'd>(&'e mut self, datum: Datum<'d>) {
307        match (self, datum) {
308            (DatumColumnEncoder::Bool(bool_builder), Datum::True) => {
309                bool_builder.append_value(true)
310            }
311            (DatumColumnEncoder::Bool(bool_builder), Datum::False) => {
312                bool_builder.append_value(false)
313            }
314            (DatumColumnEncoder::U8(builder), Datum::UInt8(val)) => builder.append_value(val),
315            (DatumColumnEncoder::U16(builder), Datum::UInt16(val)) => builder.append_value(val),
316            (DatumColumnEncoder::U32(builder), Datum::UInt32(val)) => builder.append_value(val),
317            (DatumColumnEncoder::U64(builder), Datum::UInt64(val)) => builder.append_value(val),
318            (DatumColumnEncoder::I16(builder), Datum::Int16(val)) => builder.append_value(val),
319            (DatumColumnEncoder::I32(builder), Datum::Int32(val)) => builder.append_value(val),
320            (DatumColumnEncoder::I64(builder), Datum::Int64(val)) => builder.append_value(val),
321            (DatumColumnEncoder::F32(builder), Datum::Float32(val)) => builder.append_value(*val),
322            (DatumColumnEncoder::F64(builder), Datum::Float64(val)) => builder.append_value(*val),
323            (
324                DatumColumnEncoder::Numeric {
325                    approx_values,
326                    binary_values,
327                    numeric_context,
328                },
329                Datum::Numeric(val),
330            ) => {
331                let float_approx = numeric_context.try_into_f64(val.0).unwrap_or_else(|_| {
332                    numeric_context.clear_status();
333                    if val.0.is_negative() {
334                        f64::NEG_INFINITY
335                    } else {
336                        f64::INFINITY
337                    }
338                });
339                let packed = PackedNumeric::from_value(val.0);
340
341                approx_values.append_value(float_approx);
342                binary_values.append_value(packed.as_bytes());
343            }
344            (DatumColumnEncoder::String(builder), Datum::String(val)) => builder.append_value(val),
345            (DatumColumnEncoder::Bytes(builder), Datum::Bytes(val)) => builder.append_value(val),
346            (DatumColumnEncoder::Date(builder), Datum::Date(val)) => {
347                builder.append_value(val.pg_epoch_days())
348            }
349            (DatumColumnEncoder::Time(builder), Datum::Time(val)) => {
350                let packed = PackedNaiveTime::from_value(val);
351                builder
352                    .append_value(packed.as_bytes())
353                    .expect("known correct size");
354            }
355            (DatumColumnEncoder::Timestamp(builder), Datum::Timestamp(val)) => {
356                let packed = PackedNaiveDateTime::from_value(val.to_naive());
357                builder
358                    .append_value(packed.as_bytes())
359                    .expect("known correct size");
360            }
361            (DatumColumnEncoder::TimestampTz(builder), Datum::TimestampTz(val)) => {
362                let packed = PackedNaiveDateTime::from_value(val.to_naive());
363                builder
364                    .append_value(packed.as_bytes())
365                    .expect("known correct size");
366            }
367            (DatumColumnEncoder::MzTimestamp(builder), Datum::MzTimestamp(val)) => {
368                builder.append_value(val.into());
369            }
370            (DatumColumnEncoder::Interval(builder), Datum::Interval(val)) => {
371                let packed = PackedInterval::from_value(val);
372                builder
373                    .append_value(packed.as_bytes())
374                    .expect("known correct size");
375            }
376            (DatumColumnEncoder::Uuid(builder), Datum::Uuid(val)) => builder
377                .append_value(val.as_bytes())
378                .expect("known correct size"),
379            (DatumColumnEncoder::AclItem(builder), Datum::AclItem(val)) => {
380                let packed = PackedAclItem::from_value(val);
381                builder
382                    .append_value(packed.as_bytes())
383                    .expect("known correct size");
384            }
385            (DatumColumnEncoder::MzAclItem(builder), Datum::MzAclItem(val)) => {
386                let packed = PackedMzAclItem::from_value(val);
387                builder.append_value(packed.as_bytes());
388            }
389            (DatumColumnEncoder::Range(builder), d @ Datum::Range(_)) => {
390                let proto = ProtoDatum::from(d);
391                let bytes = proto.encode_to_vec();
392                builder.append_value(&bytes);
393            }
394            (
395                DatumColumnEncoder::Jsonb {
396                    offsets,
397                    buf,
398                    nulls,
399                },
400                d @ Datum::JsonNull
401                | d @ Datum::True
402                | d @ Datum::False
403                | d @ Datum::Numeric(_)
404                | d @ Datum::String(_)
405                | d @ Datum::List(_)
406                | d @ Datum::Map(_),
407            ) => {
408                // TODO(parkmycar): Why do we need to re-borrow here?
409                let mut buf = buf;
410                let json = JsonbRef::from_datum(d);
411
412                // Serialize our JSON.
413                json.to_writer(&mut buf)
414                    .expect("failed to serialize Datum to jsonb");
415                let offset: i32 = buf.len().try_into().expect("wrote more than 4GB of JSON");
416                offsets.push(offset);
417
418                if let Some(nulls) = nulls {
419                    nulls.append(true);
420                }
421            }
422            (
423                DatumColumnEncoder::Array {
424                    dims,
425                    val_lengths,
426                    vals,
427                    nulls,
428                },
429                Datum::Array(array),
430            ) => {
431                // Store our array dimensions.
432                for dimension in array.dims() {
433                    let packed = PackedArrayDimension::from_value(dimension);
434                    dims.values()
435                        .append_value(packed.as_bytes())
436                        .expect("known correct size");
437                }
438                dims.append(true);
439
440                // Store the values of the array.
441                let mut count = 0;
442                for datum in array.elements() {
443                    count += 1;
444                    vals.push(datum);
445                }
446                val_lengths.push(count);
447
448                if let Some(nulls) = nulls {
449                    nulls.append(true);
450                }
451            }
452            (
453                DatumColumnEncoder::List {
454                    lengths,
455                    values,
456                    nulls,
457                },
458                Datum::List(list),
459            ) => {
460                let mut count = 0;
461                for datum in list {
462                    count += 1;
463                    values.push(datum);
464                }
465                lengths.push(count);
466
467                if let Some(nulls) = nulls {
468                    nulls.append(true);
469                }
470            }
471            (
472                DatumColumnEncoder::Map {
473                    lengths,
474                    keys,
475                    vals,
476                    nulls,
477                },
478                Datum::Map(map),
479            ) => {
480                let mut count = 0;
481                for (key, datum) in &map {
482                    count += 1;
483                    keys.append_value(key);
484                    vals.push(datum);
485                }
486                lengths.push(count);
487
488                if let Some(nulls) = nulls {
489                    nulls.append(true);
490                }
491            }
492            (
493                DatumColumnEncoder::Record {
494                    fields,
495                    nulls,
496                    length,
497                },
498                Datum::List(records),
499            ) => {
500                let mut count = 0;
501                // `zip_eq` will panic if the number of records != number of fields.
502                for (datum, encoder) in records.into_iter().zip_eq(fields.iter_mut()) {
503                    count += 1;
504                    encoder.push(datum);
505                }
506                assert_eq!(count, fields.len());
507
508                length.add_assign(1);
509                if let Some(nulls) = nulls.as_mut() {
510                    nulls.append(true);
511                }
512            }
513            (DatumColumnEncoder::RecordEmpty(builder), Datum::List(records)) => {
514                assert_none!(records.into_iter().next());
515                builder.append_value(true);
516            }
517            (encoder, Datum::Null) => encoder.push_invalid(),
518            (encoder, datum) => panic!("can't encode {datum:?} into {encoder:?}"),
519        }
520    }
521
522    fn push_invalid(&mut self) {
523        match self {
524            DatumColumnEncoder::Bool(builder) => builder.append_null(),
525            DatumColumnEncoder::U8(builder) => builder.append_null(),
526            DatumColumnEncoder::U16(builder) => builder.append_null(),
527            DatumColumnEncoder::U32(builder) => builder.append_null(),
528            DatumColumnEncoder::U64(builder) => builder.append_null(),
529            DatumColumnEncoder::I16(builder) => builder.append_null(),
530            DatumColumnEncoder::I32(builder) => builder.append_null(),
531            DatumColumnEncoder::I64(builder) => builder.append_null(),
532            DatumColumnEncoder::F32(builder) => builder.append_null(),
533            DatumColumnEncoder::F64(builder) => builder.append_null(),
534            DatumColumnEncoder::Numeric {
535                approx_values,
536                binary_values,
537                numeric_context: _,
538            } => {
539                approx_values.append_null();
540                binary_values.append_null();
541            }
542            DatumColumnEncoder::String(builder) => builder.append_null(),
543            DatumColumnEncoder::Bytes(builder) => builder.append_null(),
544            DatumColumnEncoder::Date(builder) => builder.append_null(),
545            DatumColumnEncoder::Time(builder) => builder.append_null(),
546            DatumColumnEncoder::Timestamp(builder) => builder.append_null(),
547            DatumColumnEncoder::TimestampTz(builder) => builder.append_null(),
548            DatumColumnEncoder::MzTimestamp(builder) => builder.append_null(),
549            DatumColumnEncoder::Interval(builder) => builder.append_null(),
550            DatumColumnEncoder::Uuid(builder) => builder.append_null(),
551            DatumColumnEncoder::AclItem(builder) => builder.append_null(),
552            DatumColumnEncoder::MzAclItem(builder) => builder.append_null(),
553            DatumColumnEncoder::Range(builder) => builder.append_null(),
554            DatumColumnEncoder::Jsonb {
555                offsets,
556                buf: _,
557                nulls,
558            } => {
559                let nulls = nulls.get_or_insert_with(|| {
560                    let mut buf = BooleanBufferBuilder::new(offsets.len());
561                    // The offsets buffer has one more value than there are elements.
562                    buf.append_n(offsets.len() - 1, true);
563                    buf
564                });
565
566                offsets.push(offsets.last().copied().unwrap_or(0));
567                nulls.append(false);
568            }
569            DatumColumnEncoder::Array {
570                dims,
571                val_lengths,
572                vals: _,
573                nulls,
574            } => {
575                let nulls = nulls.get_or_insert_with(|| {
576                    let mut buf = BooleanBufferBuilder::new(dims.len() + 1);
577                    buf.append_n(dims.len(), true);
578                    buf
579                });
580                dims.append_null();
581
582                val_lengths.push(0);
583                nulls.append(false);
584            }
585            DatumColumnEncoder::List {
586                lengths,
587                values: _,
588                nulls,
589            } => {
590                let nulls = nulls.get_or_insert_with(|| {
591                    let mut buf = BooleanBufferBuilder::new(lengths.len() + 1);
592                    buf.append_n(lengths.len(), true);
593                    buf
594                });
595
596                lengths.push(0);
597                nulls.append(false);
598            }
599            DatumColumnEncoder::Map {
600                lengths,
601                keys: _,
602                vals: _,
603                nulls,
604            } => {
605                let nulls = nulls.get_or_insert_with(|| {
606                    let mut buf = BooleanBufferBuilder::new(lengths.len() + 1);
607                    buf.append_n(lengths.len(), true);
608                    buf
609                });
610
611                lengths.push(0);
612                nulls.append(false);
613            }
614            DatumColumnEncoder::Record {
615                fields,
616                nulls,
617                length,
618            } => {
619                let nulls = nulls.get_or_insert_with(|| {
620                    let mut buf = BooleanBufferBuilder::new(*length + 1);
621                    buf.append_n(*length, true);
622                    buf
623                });
624                nulls.append(false);
625                length.add_assign(1);
626
627                for field in fields {
628                    field.push_invalid();
629                }
630            }
631            DatumColumnEncoder::RecordEmpty(builder) => builder.append_null(),
632        }
633    }
634
635    fn finish(self) -> ArrayRef {
636        match self {
637            DatumColumnEncoder::Bool(mut builder) => {
638                let array = builder.finish();
639                Arc::new(array)
640            }
641            DatumColumnEncoder::U8(mut builder) => {
642                let array = builder.finish();
643                Arc::new(array)
644            }
645            DatumColumnEncoder::U16(mut builder) => {
646                let array = builder.finish();
647                Arc::new(array)
648            }
649            DatumColumnEncoder::U32(mut builder) => {
650                let array = builder.finish();
651                Arc::new(array)
652            }
653            DatumColumnEncoder::U64(mut builder) => {
654                let array = builder.finish();
655                Arc::new(array)
656            }
657            DatumColumnEncoder::I16(mut builder) => {
658                let array = builder.finish();
659                Arc::new(array)
660            }
661            DatumColumnEncoder::I32(mut builder) => {
662                let array = builder.finish();
663                Arc::new(array)
664            }
665            DatumColumnEncoder::I64(mut builder) => {
666                let array = builder.finish();
667                Arc::new(array)
668            }
669            DatumColumnEncoder::F32(mut builder) => {
670                let array = builder.finish();
671                Arc::new(array)
672            }
673            DatumColumnEncoder::F64(mut builder) => {
674                let array = builder.finish();
675                Arc::new(array)
676            }
677            DatumColumnEncoder::Numeric {
678                mut approx_values,
679                mut binary_values,
680                numeric_context: _,
681            } => {
682                let approx_array = approx_values.finish();
683                let binary_array = binary_values.finish();
684
685                assert_eq!(approx_array.len(), binary_array.len());
686                // This is O(n) so we only enable it for debug assertions.
687                debug_assert_eq!(approx_array.logical_nulls(), binary_array.logical_nulls());
688
689                let fields = Fields::from(vec![
690                    Field::new("approx", approx_array.data_type().clone(), true),
691                    Field::new("binary", binary_array.data_type().clone(), true),
692                ]);
693                let nulls = approx_array.logical_nulls();
694                let array = StructArray::new(
695                    fields,
696                    vec![Arc::new(approx_array), Arc::new(binary_array)],
697                    nulls,
698                );
699                Arc::new(array)
700            }
701            DatumColumnEncoder::String(mut builder) => {
702                let array = builder.finish();
703                Arc::new(array)
704            }
705            DatumColumnEncoder::Bytes(mut builder) => {
706                let array = builder.finish();
707                Arc::new(array)
708            }
709            DatumColumnEncoder::Date(mut builder) => {
710                let array = builder.finish();
711                Arc::new(array)
712            }
713            DatumColumnEncoder::Time(mut builder) => {
714                let array = builder.finish();
715                Arc::new(array)
716            }
717            DatumColumnEncoder::Timestamp(mut builder) => {
718                let array = builder.finish();
719                Arc::new(array)
720            }
721            DatumColumnEncoder::TimestampTz(mut builder) => {
722                let array = builder.finish();
723                Arc::new(array)
724            }
725            DatumColumnEncoder::MzTimestamp(mut builder) => {
726                let array = builder.finish();
727                Arc::new(array)
728            }
729            DatumColumnEncoder::Interval(mut builder) => {
730                let array = builder.finish();
731                Arc::new(array)
732            }
733            DatumColumnEncoder::Uuid(mut builder) => {
734                let array = builder.finish();
735                Arc::new(array)
736            }
737            DatumColumnEncoder::AclItem(mut builder) => Arc::new(builder.finish()),
738            DatumColumnEncoder::MzAclItem(mut builder) => Arc::new(builder.finish()),
739            DatumColumnEncoder::Range(mut builder) => Arc::new(builder.finish()),
740            DatumColumnEncoder::Jsonb {
741                offsets,
742                buf,
743                mut nulls,
744            } => {
745                let values = Buffer::from_vec(buf);
746                let offsets = OffsetBuffer::new(ScalarBuffer::from(offsets));
747                let nulls = nulls.as_mut().map(|n| NullBuffer::from(n.finish()));
748                let array = StringArray::new(offsets, values, nulls);
749                Arc::new(array)
750            }
751            DatumColumnEncoder::Array {
752                mut dims,
753                val_lengths,
754                vals,
755                mut nulls,
756            } => {
757                let nulls = nulls.as_mut().map(|n| NullBuffer::from(n.finish()));
758                let vals = vals.finish();
759
760                // Note: Values in an Array can always be Null, regardless of whether or not the
761                // column is nullable.
762                let field = Field::new_list_field(vals.data_type().clone(), true);
763                let val_offsets = OffsetBuffer::from_lengths(val_lengths);
764                let values =
765                    ListArray::new(Arc::new(field), val_offsets, Arc::new(vals), nulls.clone());
766
767                let dims = dims.finish();
768                assert_eq!(values.len(), dims.len());
769
770                // Note: The inner arrays are always nullable, and we let the higher-level array
771                // drive nullability for the entire column.
772                let fields = Fields::from(vec![
773                    Field::new("dims", dims.data_type().clone(), true),
774                    Field::new("vals", values.data_type().clone(), true),
775                ]);
776                let array = StructArray::new(fields, vec![Arc::new(dims), Arc::new(values)], nulls);
777
778                Arc::new(array)
779            }
780            DatumColumnEncoder::List {
781                lengths,
782                values,
783                mut nulls,
784            } => {
785                let values = values.finish();
786
787                // Note: Values in an Array can always be Null, regardless of whether or not the
788                // column is nullable.
789                let field = Field::new_list_field(values.data_type().clone(), true);
790                let offsets = OffsetBuffer::<i32>::from_lengths(lengths.iter().copied());
791                let nulls = nulls.as_mut().map(|n| NullBuffer::from(n.finish()));
792
793                let array = ListArray::new(Arc::new(field), offsets, values, nulls);
794                Arc::new(array)
795            }
796            DatumColumnEncoder::Map {
797                lengths,
798                mut keys,
799                vals,
800                mut nulls,
801            } => {
802                let keys = keys.finish();
803                let vals = vals.finish();
804
805                let offsets = OffsetBuffer::<i32>::from_lengths(lengths.iter().copied());
806                let nulls = nulls.as_mut().map(|n| NullBuffer::from(n.finish()));
807
808                // Note: Values in an Map can always be Null, regardless of whether or not the
809                // column is nullable, but Keys cannot.
810                assert_none!(keys.logical_nulls());
811                let key_field = Arc::new(Field::new("key", keys.data_type().clone(), false));
812                let val_field = Arc::new(Field::new("val", vals.data_type().clone(), true));
813                let fields = Fields::from(vec![Arc::clone(&key_field), Arc::clone(&val_field)]);
814                let entries = StructArray::new(fields, vec![Arc::new(keys), vals], None);
815
816                // Note: DatumMap is always sorted, and Arrow enforces that the inner 'map_entries'
817                // array can never be null.
818                let field = Field::new("map_entries", entries.data_type().clone(), false);
819                let array = ListArray::new(Arc::new(field), offsets, Arc::new(entries), nulls);
820                Arc::new(array)
821            }
822            DatumColumnEncoder::Record {
823                fields,
824                mut nulls,
825                length: _,
826            } => {
827                let (fields, arrays): (Vec<_>, Vec<_>) = fields
828                    .into_iter()
829                    .enumerate()
830                    .map(|(tag, encoder)| {
831                        // Note: We mark all columns as nullable at the Arrow/Parquet level because
832                        // it has a negligible performance difference, but it protects us from
833                        // unintended nullability changes in the columns of SQL objects.
834                        //
835                        // See: <https://github.com/MaterializeInc/database-issues/issues/2488>
836                        let nullable = true;
837                        let array = encoder.finish();
838                        let field =
839                            Field::new(tag.to_string(), array.data_type().clone(), nullable);
840                        (field, array)
841                    })
842                    .unzip();
843                let nulls = nulls.as_mut().map(|n| NullBuffer::from(n.finish()));
844
845                let array = StructArray::new(Fields::from(fields), arrays, nulls);
846                Arc::new(array)
847            }
848            DatumColumnEncoder::RecordEmpty(mut builder) => Arc::new(builder.finish()),
849        }
850    }
851}
852
853/// A decoder for a column of [`Datum`]s.
854///
855/// Note: We specifically structure the decoder as an enum instead of using trait objects because
856/// Datum decoding is an extremely hot path and downcasting objects is relatively slow.
857#[derive(Debug)]
858enum DatumColumnDecoder {
859    Bool(BooleanArray),
860    U8(UInt8Array),
861    U16(UInt16Array),
862    U32(UInt32Array),
863    U64(UInt64Array),
864    I16(Int16Array),
865    I32(Int32Array),
866    I64(Int64Array),
867    F32(Float32Array),
868    F64(Float64Array),
869    Numeric(BinaryArray),
870    Bytes(BinaryArray),
871    String(StringArray),
872    Date(Int32Array),
873    Time(FixedSizeBinaryArray),
874    Timestamp(FixedSizeBinaryArray),
875    TimestampTz(FixedSizeBinaryArray),
876    MzTimestamp(UInt64Array),
877    Interval(FixedSizeBinaryArray),
878    Uuid(FixedSizeBinaryArray),
879    Json(StringArray),
880    Array {
881        dim_offsets: OffsetBuffer<i32>,
882        dims: FixedSizeBinaryArray,
883        val_offsets: OffsetBuffer<i32>,
884        vals: Box<DatumColumnDecoder>,
885        nulls: Option<NullBuffer>,
886    },
887    List {
888        offsets: OffsetBuffer<i32>,
889        values: Box<DatumColumnDecoder>,
890        nulls: Option<NullBuffer>,
891    },
892    Map {
893        offsets: OffsetBuffer<i32>,
894        keys: StringArray,
895        vals: Box<DatumColumnDecoder>,
896        nulls: Option<NullBuffer>,
897    },
898    RecordEmpty(BooleanArray),
899    Record {
900        fields: Vec<Box<DatumColumnDecoder>>,
901        nulls: Option<NullBuffer>,
902    },
903    Range(BinaryArray),
904    MzAclItem(BinaryArray),
905    AclItem(FixedSizeBinaryArray),
906}
907
908impl DatumColumnDecoder {
909    fn get<'a>(&'a self, idx: usize, packer: &'a mut RowPacker) {
910        let datum = match self {
911            DatumColumnDecoder::Bool(array) => array
912                .is_valid(idx)
913                .then(|| array.value(idx))
914                .map(|x| if x { Datum::True } else { Datum::False }),
915            DatumColumnDecoder::U8(array) => array
916                .is_valid(idx)
917                .then(|| array.value(idx))
918                .map(Datum::UInt8),
919            DatumColumnDecoder::U16(array) => array
920                .is_valid(idx)
921                .then(|| array.value(idx))
922                .map(Datum::UInt16),
923            DatumColumnDecoder::U32(array) => array
924                .is_valid(idx)
925                .then(|| array.value(idx))
926                .map(Datum::UInt32),
927            DatumColumnDecoder::U64(array) => array
928                .is_valid(idx)
929                .then(|| array.value(idx))
930                .map(Datum::UInt64),
931            DatumColumnDecoder::I16(array) => array
932                .is_valid(idx)
933                .then(|| array.value(idx))
934                .map(Datum::Int16),
935            DatumColumnDecoder::I32(array) => array
936                .is_valid(idx)
937                .then(|| array.value(idx))
938                .map(Datum::Int32),
939            DatumColumnDecoder::I64(array) => array
940                .is_valid(idx)
941                .then(|| array.value(idx))
942                .map(Datum::Int64),
943            DatumColumnDecoder::F32(array) => array
944                .is_valid(idx)
945                .then(|| array.value(idx))
946                .map(|x| Datum::Float32(ordered_float::OrderedFloat(x))),
947            DatumColumnDecoder::F64(array) => array
948                .is_valid(idx)
949                .then(|| array.value(idx))
950                .map(|x| Datum::Float64(ordered_float::OrderedFloat(x))),
951            DatumColumnDecoder::Numeric(array) => array.is_valid(idx).then(|| {
952                let val = array.value(idx);
953                let val = PackedNumeric::from_bytes(val)
954                    .expect("failed to roundtrip Numeric")
955                    .into_value();
956                Datum::Numeric(OrderedDecimal(val))
957            }),
958            DatumColumnDecoder::String(array) => array
959                .is_valid(idx)
960                .then(|| array.value(idx))
961                .map(Datum::String),
962            DatumColumnDecoder::Bytes(array) => array
963                .is_valid(idx)
964                .then(|| array.value(idx))
965                .map(Datum::Bytes),
966            DatumColumnDecoder::Date(array) => {
967                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
968                    let date = Date::from_pg_epoch(x).expect("failed to roundtrip");
969                    Datum::Date(date)
970                })
971            }
972            DatumColumnDecoder::Time(array) => {
973                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
974                    let packed = PackedNaiveTime::from_bytes(x).expect("failed to roundtrip time");
975                    Datum::Time(packed.into_value())
976                })
977            }
978            DatumColumnDecoder::Timestamp(array) => {
979                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
980                    let packed = PackedNaiveDateTime::from_bytes(x)
981                        .expect("failed to roundtrip PackedNaiveDateTime");
982                    let timestamp = CheckedTimestamp::from_timestamplike(packed.into_value())
983                        .expect("failed to roundtrip timestamp");
984                    Datum::Timestamp(timestamp)
985                })
986            }
987            DatumColumnDecoder::TimestampTz(array) => {
988                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
989                    let packed = PackedNaiveDateTime::from_bytes(x)
990                        .expect("failed to roundtrip PackedNaiveDateTime");
991                    let timestamp =
992                        CheckedTimestamp::from_timestamplike(packed.into_value().and_utc())
993                            .expect("failed to roundtrip timestamp");
994                    Datum::TimestampTz(timestamp)
995                })
996            }
997            DatumColumnDecoder::MzTimestamp(array) => array
998                .is_valid(idx)
999                .then(|| array.value(idx))
1000                .map(|x| Datum::MzTimestamp(Timestamp::from(x))),
1001            DatumColumnDecoder::Interval(array) => {
1002                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
1003                    let packed =
1004                        PackedInterval::from_bytes(x).expect("failed to roundtrip interval");
1005                    Datum::Interval(packed.into_value())
1006                })
1007            }
1008            DatumColumnDecoder::Uuid(array) => {
1009                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
1010                    let uuid = Uuid::from_slice(x).expect("failed to roundtrip uuid");
1011                    Datum::Uuid(uuid)
1012                })
1013            }
1014            DatumColumnDecoder::AclItem(array) => {
1015                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
1016                    let packed =
1017                        PackedAclItem::from_bytes(x).expect("failed to roundtrip MzAclItem");
1018                    Datum::AclItem(packed.into_value())
1019                })
1020            }
1021            DatumColumnDecoder::MzAclItem(array) => {
1022                array.is_valid(idx).then(|| array.value(idx)).map(|x| {
1023                    let packed =
1024                        PackedMzAclItem::from_bytes(x).expect("failed to roundtrip MzAclItem");
1025                    Datum::MzAclItem(packed.into_value())
1026                })
1027            }
1028            DatumColumnDecoder::Range(array) => {
1029                let Some(val) = array.is_valid(idx).then(|| array.value(idx)) else {
1030                    packer.push(Datum::Null);
1031                    return;
1032                };
1033
1034                let proto = ProtoDatum::decode(val).expect("failed to roundtrip Range");
1035                packer
1036                    .try_push_proto(&proto)
1037                    .expect("failed to pack ProtoRange");
1038
1039                // Return early because we've already packed the necessary Datums.
1040                return;
1041            }
1042            DatumColumnDecoder::Json(array) => {
1043                let Some(val) = array.is_valid(idx).then(|| array.value(idx)) else {
1044                    packer.push(Datum::Null);
1045                    return;
1046                };
1047                JsonbPacker::new(packer)
1048                    .pack_str(val)
1049                    .expect("failed to roundtrip JSON");
1050
1051                // Return early because we've already packed the necessary Datums.
1052                return;
1053            }
1054            DatumColumnDecoder::Array {
1055                dim_offsets,
1056                dims,
1057                val_offsets,
1058                vals,
1059                nulls,
1060            } => {
1061                let is_valid = nulls.as_ref().map(|n| n.is_valid(idx)).unwrap_or(true);
1062                if !is_valid {
1063                    packer.push(Datum::Null);
1064                    return;
1065                }
1066
1067                let start: usize = dim_offsets[idx]
1068                    .try_into()
1069                    .expect("unexpected negative offset");
1070                let end: usize = dim_offsets[idx + 1]
1071                    .try_into()
1072                    .expect("unexpected negative offset");
1073                let dimensions = (start..end).map(|idx| {
1074                    PackedArrayDimension::from_bytes(dims.value(idx))
1075                        .expect("failed to roundtrip ArrayDimension")
1076                        .into_value()
1077                });
1078
1079                let start: usize = val_offsets[idx]
1080                    .try_into()
1081                    .expect("unexpected negative offset");
1082                let end: usize = val_offsets[idx + 1]
1083                    .try_into()
1084                    .expect("unexpected negative offset");
1085                packer
1086                    .push_array_with_row_major(dimensions, |packer| {
1087                        for x in start..end {
1088                            vals.get(x, packer);
1089                        }
1090                        // Return the numer of Datums we just packed.
1091                        end - start
1092                    })
1093                    .expect("failed to pack Array");
1094
1095                // Return early because we've already packed the necessary Datums.
1096                return;
1097            }
1098            DatumColumnDecoder::List {
1099                offsets,
1100                values,
1101                nulls,
1102            } => {
1103                let is_valid = nulls.as_ref().map(|n| n.is_valid(idx)).unwrap_or(true);
1104                if !is_valid {
1105                    packer.push(Datum::Null);
1106                    return;
1107                }
1108
1109                let start: usize = offsets[idx].try_into().expect("unexpected negative offset");
1110                let end: usize = offsets[idx + 1]
1111                    .try_into()
1112                    .expect("unexpected negative offset");
1113
1114                packer.push_list_with(|packer| {
1115                    for idx in start..end {
1116                        values.get(idx, packer)
1117                    }
1118                });
1119
1120                // Return early because we've already packed the necessary Datums.
1121                return;
1122            }
1123            DatumColumnDecoder::Map {
1124                offsets,
1125                keys,
1126                vals,
1127                nulls,
1128            } => {
1129                let is_valid = nulls.as_ref().map(|n| n.is_valid(idx)).unwrap_or(true);
1130                if !is_valid {
1131                    packer.push(Datum::Null);
1132                    return;
1133                }
1134
1135                let start: usize = offsets[idx].try_into().expect("unexpected negative offset");
1136                let end: usize = offsets[idx + 1]
1137                    .try_into()
1138                    .expect("unexpected negative offset");
1139
1140                packer.push_dict_with(|packer| {
1141                    for idx in start..end {
1142                        packer.push(Datum::String(keys.value(idx)));
1143                        vals.get(idx, packer);
1144                    }
1145                });
1146
1147                // Return early because we've already packed the necessary Datums.
1148                return;
1149            }
1150            DatumColumnDecoder::RecordEmpty(array) => array.is_valid(idx).then(Datum::empty_list),
1151            DatumColumnDecoder::Record { fields, nulls } => {
1152                let is_valid = nulls.as_ref().map(|n| n.is_valid(idx)).unwrap_or(true);
1153                if !is_valid {
1154                    packer.push(Datum::Null);
1155                    return;
1156                }
1157
1158                // let mut datums = Vec::with_capacity(fields.len());
1159                packer.push_list_with(|packer| {
1160                    for field in fields {
1161                        field.get(idx, packer);
1162                    }
1163                });
1164
1165                // Return early because we've already packed the necessary Datums.
1166                return;
1167            }
1168        };
1169
1170        match datum {
1171            Some(d) => packer.push(d),
1172            None => packer.push(Datum::Null),
1173        }
1174    }
1175
1176    fn stats(&self) -> ColumnStatKinds {
1177        match self {
1178            DatumColumnDecoder::Bool(a) => PrimitiveStats::<bool>::from_column(a).into(),
1179            DatumColumnDecoder::U8(a) => PrimitiveStats::<u8>::from_column(a).into(),
1180            DatumColumnDecoder::U16(a) => PrimitiveStats::<u16>::from_column(a).into(),
1181            DatumColumnDecoder::U32(a) => PrimitiveStats::<u32>::from_column(a).into(),
1182            DatumColumnDecoder::U64(a) => PrimitiveStats::<u64>::from_column(a).into(),
1183            DatumColumnDecoder::I16(a) => PrimitiveStats::<i16>::from_column(a).into(),
1184            DatumColumnDecoder::I32(a) => PrimitiveStats::<i32>::from_column(a).into(),
1185            DatumColumnDecoder::I64(a) => PrimitiveStats::<i64>::from_column(a).into(),
1186            DatumColumnDecoder::F32(a) => PrimitiveStats::<f32>::from_column(a).into(),
1187            DatumColumnDecoder::F64(a) => PrimitiveStats::<f64>::from_column(a).into(),
1188            DatumColumnDecoder::Numeric(a) => numeric_stats_from_column(a),
1189            DatumColumnDecoder::String(a) => PrimitiveStats::<String>::from_column(a).into(),
1190            DatumColumnDecoder::Bytes(a) => PrimitiveStats::<Vec<u8>>::from_column(a).into(),
1191            DatumColumnDecoder::Date(a) => PrimitiveStats::<i32>::from_column(a).into(),
1192            DatumColumnDecoder::Time(a) => {
1193                fixed_stats_from_column(a, FixedSizeBytesStatsKind::PackedTime)
1194            }
1195            DatumColumnDecoder::Timestamp(a) => {
1196                fixed_stats_from_column(a, FixedSizeBytesStatsKind::PackedDateTime)
1197            }
1198            DatumColumnDecoder::TimestampTz(a) => {
1199                fixed_stats_from_column(a, FixedSizeBytesStatsKind::PackedDateTime)
1200            }
1201            DatumColumnDecoder::MzTimestamp(a) => PrimitiveStats::<u64>::from_column(a).into(),
1202            DatumColumnDecoder::Interval(a) => {
1203                fixed_stats_from_column(a, FixedSizeBytesStatsKind::PackedInterval)
1204            }
1205            DatumColumnDecoder::Uuid(a) => {
1206                fixed_stats_from_column(a, FixedSizeBytesStatsKind::Uuid)
1207            }
1208            DatumColumnDecoder::AclItem(_)
1209            | DatumColumnDecoder::MzAclItem(_)
1210            | DatumColumnDecoder::Range(_) => ColumnStatKinds::None,
1211            DatumColumnDecoder::Json(a) => stats_for_json(a.iter()).values,
1212            DatumColumnDecoder::Array { .. }
1213            | DatumColumnDecoder::List { .. }
1214            | DatumColumnDecoder::Map { .. }
1215            | DatumColumnDecoder::Record { .. }
1216            | DatumColumnDecoder::RecordEmpty(_) => ColumnStatKinds::None,
1217        }
1218    }
1219
1220    fn goodbytes(&self) -> usize {
1221        match self {
1222            DatumColumnDecoder::Bool(a) => ArrayOrd::Bool(a.clone()).goodbytes(),
1223            DatumColumnDecoder::U8(a) => ArrayOrd::UInt8(a.clone()).goodbytes(),
1224            DatumColumnDecoder::U16(a) => ArrayOrd::UInt16(a.clone()).goodbytes(),
1225            DatumColumnDecoder::U32(a) => ArrayOrd::UInt32(a.clone()).goodbytes(),
1226            DatumColumnDecoder::U64(a) => ArrayOrd::UInt64(a.clone()).goodbytes(),
1227            DatumColumnDecoder::I16(a) => ArrayOrd::Int16(a.clone()).goodbytes(),
1228            DatumColumnDecoder::I32(a) => ArrayOrd::Int32(a.clone()).goodbytes(),
1229            DatumColumnDecoder::I64(a) => ArrayOrd::Int64(a.clone()).goodbytes(),
1230            DatumColumnDecoder::F32(a) => ArrayOrd::Float32(a.clone()).goodbytes(),
1231            DatumColumnDecoder::F64(a) => ArrayOrd::Float64(a.clone()).goodbytes(),
1232            DatumColumnDecoder::Numeric(a) => ArrayOrd::Binary(a.clone()).goodbytes(),
1233            DatumColumnDecoder::String(a) => ArrayOrd::String(a.clone()).goodbytes(),
1234            DatumColumnDecoder::Bytes(a) => ArrayOrd::Binary(a.clone()).goodbytes(),
1235            DatumColumnDecoder::Date(a) => ArrayOrd::Int32(a.clone()).goodbytes(),
1236            DatumColumnDecoder::Time(a) => ArrayOrd::FixedSizeBinary(a.clone()).goodbytes(),
1237            DatumColumnDecoder::Timestamp(a) => ArrayOrd::FixedSizeBinary(a.clone()).goodbytes(),
1238            DatumColumnDecoder::TimestampTz(a) => ArrayOrd::FixedSizeBinary(a.clone()).goodbytes(),
1239            DatumColumnDecoder::MzTimestamp(a) => ArrayOrd::UInt64(a.clone()).goodbytes(),
1240            DatumColumnDecoder::Interval(a) => ArrayOrd::FixedSizeBinary(a.clone()).goodbytes(),
1241            DatumColumnDecoder::Uuid(a) => ArrayOrd::FixedSizeBinary(a.clone()).goodbytes(),
1242            DatumColumnDecoder::AclItem(a) => ArrayOrd::FixedSizeBinary(a.clone()).goodbytes(),
1243            DatumColumnDecoder::MzAclItem(a) => ArrayOrd::Binary(a.clone()).goodbytes(),
1244            DatumColumnDecoder::Range(a) => ArrayOrd::Binary(a.clone()).goodbytes(),
1245            DatumColumnDecoder::Json(a) => ArrayOrd::String(a.clone()).goodbytes(),
1246            DatumColumnDecoder::Array { dims, vals, .. } => {
1247                (dims.len() * PackedArrayDimension::SIZE) + vals.goodbytes()
1248            }
1249            DatumColumnDecoder::List { values, .. } => values.goodbytes(),
1250            DatumColumnDecoder::Map { keys, vals, .. } => {
1251                ArrayOrd::String(keys.clone()).goodbytes() + vals.goodbytes()
1252            }
1253            DatumColumnDecoder::Record { fields, .. } => fields.iter().map(|f| f.goodbytes()).sum(),
1254            DatumColumnDecoder::RecordEmpty(a) => ArrayOrd::Bool(a.clone()).goodbytes(),
1255        }
1256    }
1257}
1258
1259impl Schema<Row> for RelationDesc {
1260    type ArrowColumn = arrow::array::StructArray;
1261    type Statistics = OptionStats<StructStats>;
1262
1263    type Decoder = RowColumnarDecoder;
1264    type Encoder = RowColumnarEncoder;
1265
1266    fn decoder(&self, col: Self::ArrowColumn) -> Result<Self::Decoder, anyhow::Error> {
1267        RowColumnarDecoder::new(col, self)
1268    }
1269
1270    fn encoder(&self) -> Result<Self::Encoder, anyhow::Error> {
1271        RowColumnarEncoder::new(self)
1272            .ok_or_else(|| anyhow::anyhow!("Cannot encode a RelationDesc with no columns"))
1273    }
1274}
1275
1276/// A [`ColumnDecoder`] for a [`Row`].
1277#[derive(Debug)]
1278pub struct RowColumnarDecoder {
1279    /// The length of all columns in this decoder; matching all child arrays and the null array
1280    /// if present.
1281    len: usize,
1282    /// Field-specific information: the user-readable field name, the null count (or None if the
1283    /// column is non-nullable), and the decoder which wraps the column-specific array.
1284    decoders: Vec<(Arc<str>, Option<usize>, DatumColumnDecoder)>,
1285    /// The null buffer for this row, if present. (At time of writing, all rows are assumed to be
1286    /// logically nullable.)
1287    nullability: Option<NullBuffer>,
1288}
1289
1290/// Merge the provided null buffer with the existing array's null buffer, if any.
1291fn mask_nulls(column: &ArrayRef, null_mask: Option<&NullBuffer>) -> ArrayRef {
1292    if null_mask.is_none() {
1293        Arc::clone(column)
1294    } else {
1295        // We calculate stats on the nested arrays, so make sure we don't count entries
1296        // that are masked off at a higher level.
1297        let nulls = NullBuffer::union(null_mask, column.nulls());
1298        let data = column
1299            .to_data()
1300            .into_builder()
1301            .nulls(nulls)
1302            .build()
1303            .expect("changed only null mask");
1304        make_array(data)
1305    }
1306}
1307
1308impl RowColumnarDecoder {
1309    /// Creates a [`RowColumnarDecoder`] that decodes from the provided [`StructArray`].
1310    ///
1311    /// Returns an error if the schema of the [`StructArray`] does not match
1312    /// the provided [`RelationDesc`].
1313    pub fn new(col: StructArray, desc: &RelationDesc) -> Result<Self, anyhow::Error> {
1314        let inner_columns = col.columns();
1315        let desc_columns = desc.typ().columns();
1316
1317        if desc_columns.len() > inner_columns.len() {
1318            anyhow::bail!(
1319                "provided array has too few columns! {desc_columns:?} > {inner_columns:?}"
1320            );
1321        }
1322
1323        // For performance reasons we downcast just a single time.
1324        let mut decoders = Vec::with_capacity(desc_columns.len());
1325
1326        let null_mask = col.nulls();
1327
1328        // The columns of the `StructArray` are named with their column index.
1329        for (col_idx, col_name, col_type) in desc.iter_all() {
1330            let field_name = col_idx.to_stable_name();
1331            let column = col.column_by_name(&field_name).ok_or_else(|| {
1332                anyhow::anyhow!(
1333                    "StructArray did not contain column name {field_name}, found {:?}",
1334                    col.column_names()
1335                )
1336            })?;
1337            let column = mask_nulls(column, null_mask);
1338            let null_count = col_type.nullable.then(|| column.null_count());
1339            let decoder = array_to_decoder(&column, &col_type.scalar_type)?;
1340            decoders.push((col_name.as_str().into(), null_count, decoder));
1341        }
1342
1343        Ok(RowColumnarDecoder {
1344            len: col.len(),
1345            decoders,
1346            nullability: col.logical_nulls(),
1347        })
1348    }
1349
1350    // Returns the number of null entries in this array of Row structs. This
1351    // will be 0 when `Row` is encoded directly, but could be non-zero when it's
1352    // used inside `SourceDataEncoder`.
1353    pub fn null_count(&self) -> usize {
1354        self.nullability.as_ref().map_or(0, |n| n.null_count())
1355    }
1356}
1357
1358impl ColumnDecoder<Row> for RowColumnarDecoder {
1359    fn decode(&self, idx: usize, val: &mut Row) {
1360        let mut packer = val.packer();
1361
1362        for (_, _, decoder) in &self.decoders {
1363            decoder.get(idx, &mut packer);
1364        }
1365    }
1366
1367    fn is_null(&self, idx: usize) -> bool {
1368        let Some(nullability) = self.nullability.as_ref() else {
1369            return false;
1370        };
1371        nullability.is_null(idx)
1372    }
1373
1374    fn goodbytes(&self) -> usize {
1375        let decoders_size: usize = self
1376            .decoders
1377            .iter()
1378            .map(|(_name, _null_count, decoder)| decoder.goodbytes())
1379            .sum();
1380
1381        decoders_size
1382            + self
1383                .nullability
1384                .as_ref()
1385                .map(|nulls| nulls.inner().inner().len())
1386                .unwrap_or(0)
1387    }
1388
1389    fn stats(&self) -> StructStats {
1390        StructStats {
1391            len: self.len,
1392            cols: self
1393                .decoders
1394                .iter()
1395                .map(|(name, null_count, decoder)| {
1396                    let name = name.to_string();
1397                    let stats = ColumnarStats {
1398                        nulls: null_count.map(|count| ColumnNullStats { count }),
1399                        values: decoder.stats(),
1400                    };
1401                    (name, stats)
1402                })
1403                .collect(),
1404        }
1405    }
1406}
1407
1408/// A [`ColumnEncoder`] for a [`Row`].
1409#[derive(Debug)]
1410pub struct RowColumnarEncoder {
1411    encoders: Vec<DatumEncoder>,
1412    // TODO(parkmycar): Replace the `usize` with a `ColumnIdx` type.
1413    col_names: Vec<(usize, Arc<str>)>,
1414    // TODO(parkmycar): Optionally omit this.
1415    nullability: BooleanBufferBuilder,
1416}
1417
1418impl RowColumnarEncoder {
1419    /// Creates a [`RowColumnarEncoder`] for the provided [`RelationDesc`].
1420    ///
1421    /// Returns `None` if the provided [`RelationDesc`] has no columns.
1422    ///
1423    /// # Note
1424    /// Internally we represent a [`Row`] as a [`StructArray`] which is
1425    /// required to have at least one field. Instead of handling this case by
1426    /// adding some special "internal" column we let a higher level encoder
1427    /// (e.g. `SourceDataColumnarEncoder`) handle this case.
1428    pub fn new(desc: &RelationDesc) -> Option<Self> {
1429        if desc.typ().columns().is_empty() {
1430            return None;
1431        }
1432
1433        let (col_names, encoders): (Vec<_>, Vec<_>) = desc
1434            .iter_all()
1435            .map(|(col_idx, col_name, col_type)| {
1436                let encoder = scalar_type_to_encoder(&col_type.scalar_type)
1437                    .expect("failed to create encoder");
1438                let encoder = DatumEncoder {
1439                    nullable: col_type.nullable,
1440                    encoder,
1441                };
1442
1443                // We name the Fields in Parquet with the column index, but for
1444                // backwards compat use the column name for stats.
1445                //
1446                // NOTE: name-keyed stats are sound only while durable
1447                // relations never carry duplicate column names (the planner
1448                // enforces this) and a dropped column's name can never be
1449                // reused by a later version (persist rejects schema
1450                // migrations containing drops). Filter pushdown consults
1451                // these stats by name; violating either invariant attaches
1452                // one column's stats to another and yields wrong results.
1453                let name = (col_idx.to_raw(), col_name.as_str().into());
1454
1455                (name, encoder)
1456            })
1457            .unzip();
1458
1459        Some(RowColumnarEncoder {
1460            encoders,
1461            col_names,
1462            nullability: BooleanBufferBuilder::new(100),
1463        })
1464    }
1465}
1466
1467impl ColumnEncoder<Row> for RowColumnarEncoder {
1468    type FinishedColumn = StructArray;
1469
1470    fn goodbytes(&self) -> usize {
1471        self.encoders.iter().map(|e| e.goodbytes()).sum()
1472    }
1473
1474    fn append(&mut self, val: &Row) {
1475        let mut num_datums = 0;
1476        for (datum, encoder) in val.iter().zip_eq(self.encoders.iter_mut()) {
1477            encoder.push(datum);
1478            num_datums += 1;
1479        }
1480        assert_eq!(
1481            num_datums,
1482            self.encoders.len(),
1483            "tried to encode {val:?}, but only have {:?}",
1484            self.encoders
1485        );
1486
1487        self.nullability.append(true);
1488    }
1489
1490    fn append_null(&mut self) {
1491        for encoder in self.encoders.iter_mut() {
1492            encoder.push_invalid();
1493        }
1494        self.nullability.append(false);
1495    }
1496
1497    fn finish(self) -> Self::FinishedColumn {
1498        let RowColumnarEncoder {
1499            encoders,
1500            col_names,
1501            nullability,
1502            ..
1503        } = self;
1504
1505        let (arrays, fields): (Vec<_>, Vec<_>) = col_names
1506            .iter()
1507            .zip_eq(encoders)
1508            .map(|((col_idx, _col_name), encoder)| {
1509                // Note: We mark all columns as nullable at the Arrow/Parquet level because it has
1510                // a negligible performance difference, but it protects us from unintended
1511                // nullability changes in the columns of SQL objects.
1512                //
1513                // See: <https://github.com/MaterializeInc/database-issues/issues/2488>
1514                let nullable = true;
1515                let array = encoder.finish();
1516                let field = Field::new(col_idx.to_string(), array.data_type().clone(), nullable);
1517
1518                (array, field)
1519            })
1520            .multiunzip();
1521
1522        let null_buffer = NullBuffer::from(BooleanBuffer::from(nullability));
1523
1524        let array = StructArray::new(Fields::from(fields), arrays, Some(null_buffer));
1525
1526        array
1527    }
1528}
1529
1530/// Small helper method to make downcasting an [`Array`] return an error.
1531///
1532/// Note: it is _super_ important that we downcast as few times as possible. Datum encoding is a
1533/// very hot path and downcasting is relatively slow
1534#[inline]
1535fn downcast_array<T: 'static>(array: &Arc<dyn Array>) -> Result<&T, anyhow::Error> {
1536    array
1537        .as_any()
1538        .downcast_ref::<T>()
1539        .ok_or_else(|| anyhow!("expected {}, found {array:?}", std::any::type_name::<T>()))
1540}
1541
1542/// Small helper function to downcast from an array to a [`DatumColumnDecoder`].
1543///
1544/// Note: it is _super_ important that we downcast as few times as possible. Datum encoding is a
1545/// very hot path and downcasting is relatively slow
1546fn array_to_decoder(
1547    array: &Arc<dyn Array>,
1548    col_ty: &SqlScalarType,
1549) -> Result<DatumColumnDecoder, anyhow::Error> {
1550    let decoder = match (array.data_type(), col_ty) {
1551        (DataType::Boolean, SqlScalarType::Bool) => {
1552            let array = downcast_array::<BooleanArray>(array)?;
1553            DatumColumnDecoder::Bool(array.clone())
1554        }
1555        (DataType::UInt8, SqlScalarType::PgLegacyChar) => {
1556            let array = downcast_array::<UInt8Array>(array)?;
1557            DatumColumnDecoder::U8(array.clone())
1558        }
1559        (DataType::UInt16, SqlScalarType::UInt16) => {
1560            let array = downcast_array::<UInt16Array>(array)?;
1561            DatumColumnDecoder::U16(array.clone())
1562        }
1563        (
1564            DataType::UInt32,
1565            SqlScalarType::UInt32
1566            | SqlScalarType::Oid
1567            | SqlScalarType::RegClass
1568            | SqlScalarType::RegProc
1569            | SqlScalarType::RegType,
1570        ) => {
1571            let array = downcast_array::<UInt32Array>(array)?;
1572            DatumColumnDecoder::U32(array.clone())
1573        }
1574        (DataType::UInt64, SqlScalarType::UInt64) => {
1575            let array = downcast_array::<UInt64Array>(array)?;
1576            DatumColumnDecoder::U64(array.clone())
1577        }
1578        (DataType::Int16, SqlScalarType::Int16) => {
1579            let array = downcast_array::<Int16Array>(array)?;
1580            DatumColumnDecoder::I16(array.clone())
1581        }
1582        (DataType::Int32, SqlScalarType::Int32) => {
1583            let array = downcast_array::<Int32Array>(array)?;
1584            DatumColumnDecoder::I32(array.clone())
1585        }
1586        (DataType::Int64, SqlScalarType::Int64) => {
1587            let array = downcast_array::<Int64Array>(array)?;
1588            DatumColumnDecoder::I64(array.clone())
1589        }
1590        (DataType::Float32, SqlScalarType::Float32) => {
1591            let array = downcast_array::<Float32Array>(array)?;
1592            DatumColumnDecoder::F32(array.clone())
1593        }
1594        (DataType::Float64, SqlScalarType::Float64) => {
1595            let array = downcast_array::<Float64Array>(array)?;
1596            DatumColumnDecoder::F64(array.clone())
1597        }
1598        (DataType::Struct(_), SqlScalarType::Numeric { .. }) => {
1599            let array = downcast_array::<StructArray>(array)?;
1600            // Note: We only use the approx column for sorting, and ignore it
1601            // when decoding.
1602            let binary_values = array
1603                .column_by_name("binary")
1604                .expect("missing binary column");
1605
1606            let array = downcast_array::<BinaryArray>(binary_values)?;
1607            DatumColumnDecoder::Numeric(array.clone())
1608        }
1609        (
1610            DataType::Utf8,
1611            SqlScalarType::String
1612            | SqlScalarType::PgLegacyName
1613            | SqlScalarType::Char { .. }
1614            | SqlScalarType::VarChar { .. },
1615        ) => {
1616            let array = downcast_array::<StringArray>(array)?;
1617            DatumColumnDecoder::String(array.clone())
1618        }
1619        (DataType::Binary, SqlScalarType::Bytes) => {
1620            let array = downcast_array::<BinaryArray>(array)?;
1621            DatumColumnDecoder::Bytes(array.clone())
1622        }
1623        (DataType::Int32, SqlScalarType::Date) => {
1624            let array = downcast_array::<Int32Array>(array)?;
1625            DatumColumnDecoder::Date(array.clone())
1626        }
1627        (DataType::FixedSizeBinary(TIME_FIXED_BYTES), SqlScalarType::Time) => {
1628            let array = downcast_array::<FixedSizeBinaryArray>(array)?;
1629            DatumColumnDecoder::Time(array.clone())
1630        }
1631        (DataType::FixedSizeBinary(TIMESTAMP_FIXED_BYTES), SqlScalarType::Timestamp { .. }) => {
1632            let array = downcast_array::<FixedSizeBinaryArray>(array)?;
1633            DatumColumnDecoder::Timestamp(array.clone())
1634        }
1635        (DataType::FixedSizeBinary(TIMESTAMP_FIXED_BYTES), SqlScalarType::TimestampTz { .. }) => {
1636            let array = downcast_array::<FixedSizeBinaryArray>(array)?;
1637            DatumColumnDecoder::TimestampTz(array.clone())
1638        }
1639        (DataType::UInt64, SqlScalarType::MzTimestamp) => {
1640            let array = downcast_array::<UInt64Array>(array)?;
1641            DatumColumnDecoder::MzTimestamp(array.clone())
1642        }
1643        (DataType::FixedSizeBinary(INTERVAL_FIXED_BYTES), SqlScalarType::Interval) => {
1644            let array = downcast_array::<FixedSizeBinaryArray>(array)?;
1645            DatumColumnDecoder::Interval(array.clone())
1646        }
1647        (DataType::FixedSizeBinary(UUID_FIXED_BYTES), SqlScalarType::Uuid) => {
1648            let array = downcast_array::<FixedSizeBinaryArray>(array)?;
1649            DatumColumnDecoder::Uuid(array.clone())
1650        }
1651        (DataType::FixedSizeBinary(ACL_ITEM_FIXED_BYTES), SqlScalarType::AclItem) => {
1652            let array = downcast_array::<FixedSizeBinaryArray>(array)?;
1653            DatumColumnDecoder::AclItem(array.clone())
1654        }
1655        (DataType::Binary, SqlScalarType::MzAclItem) => {
1656            let array = downcast_array::<BinaryArray>(array)?;
1657            DatumColumnDecoder::MzAclItem(array.clone())
1658        }
1659        (DataType::Binary, SqlScalarType::Range { .. }) => {
1660            let array = downcast_array::<BinaryArray>(array)?;
1661            DatumColumnDecoder::Range(array.clone())
1662        }
1663        (DataType::Utf8, SqlScalarType::Jsonb) => {
1664            let array = downcast_array::<StringArray>(array)?;
1665            DatumColumnDecoder::Json(array.clone())
1666        }
1667        (DataType::Struct(_), s @ SqlScalarType::Array(_) | s @ SqlScalarType::Int2Vector) => {
1668            let element_type = match s {
1669                SqlScalarType::Array(inner) => inner,
1670                SqlScalarType::Int2Vector => &SqlScalarType::Int16,
1671                _ => unreachable!("checked above"),
1672            };
1673
1674            let array = downcast_array::<StructArray>(array)?;
1675            let nulls = array.nulls().cloned();
1676
1677            let dims = array
1678                .column_by_name("dims")
1679                .expect("missing dimensions column");
1680            let dims = downcast_array::<ListArray>(dims).cloned()?;
1681            let dim_offsets = dims.offsets().clone();
1682            let dims = downcast_array::<FixedSizeBinaryArray>(dims.values()).cloned()?;
1683
1684            let vals = array.column_by_name("vals").expect("missing values column");
1685            let vals = downcast_array::<ListArray>(vals)?;
1686            let val_offsets = vals.offsets().clone();
1687            let vals = array_to_decoder(vals.values(), element_type)?;
1688
1689            DatumColumnDecoder::Array {
1690                dim_offsets,
1691                dims,
1692                val_offsets,
1693                vals: Box::new(vals),
1694                nulls,
1695            }
1696        }
1697        (DataType::List(_), SqlScalarType::List { element_type, .. }) => {
1698            let array = downcast_array::<ListArray>(array)?;
1699            let inner_decoder = array_to_decoder(array.values(), &*element_type)?;
1700            DatumColumnDecoder::List {
1701                offsets: array.offsets().clone(),
1702                values: Box::new(inner_decoder),
1703                nulls: array.nulls().cloned(),
1704            }
1705        }
1706        (DataType::Map(_, true), SqlScalarType::Map { value_type, .. }) => {
1707            let array = downcast_array::<MapArray>(array)?;
1708            let keys = downcast_array::<StringArray>(array.keys())?;
1709            let vals = array_to_decoder(array.values(), value_type)?;
1710            DatumColumnDecoder::Map {
1711                offsets: array.offsets().clone(),
1712                keys: keys.clone(),
1713                vals: Box::new(vals),
1714                nulls: array.nulls().cloned(),
1715            }
1716        }
1717        (DataType::List(_), SqlScalarType::Map { value_type, .. }) => {
1718            let array: &ListArray = downcast_array(array)?;
1719            let entries: &StructArray = downcast_array(array.values())?;
1720            let [keys, values]: &[ArrayRef; 2] = entries.columns().try_into()?;
1721            let keys: &StringArray = downcast_array(keys)?;
1722            let vals: DatumColumnDecoder = array_to_decoder(values, value_type)?;
1723            DatumColumnDecoder::Map {
1724                offsets: array.offsets().clone(),
1725                keys: keys.clone(),
1726                vals: Box::new(vals),
1727                nulls: array.nulls().cloned(),
1728            }
1729        }
1730        (DataType::Boolean, SqlScalarType::Record { fields, .. }) if fields.is_empty() => {
1731            let empty_record_array = downcast_array::<BooleanArray>(array)?;
1732            DatumColumnDecoder::RecordEmpty(empty_record_array.clone())
1733        }
1734        (DataType::Struct(_), SqlScalarType::Record { fields, .. }) => {
1735            let record_array = downcast_array::<StructArray>(array)?;
1736            let null_mask = record_array.nulls();
1737            let mut decoders = Vec::with_capacity(fields.len());
1738            for (tag, (_name, col_type)) in fields.iter().enumerate() {
1739                let inner_array = record_array
1740                    .column_by_name(&tag.to_string())
1741                    .ok_or_else(|| anyhow::anyhow!("no column named '{tag}'"))?;
1742                let inner_array = mask_nulls(inner_array, null_mask);
1743                let inner_decoder = array_to_decoder(&inner_array, &col_type.scalar_type)?;
1744
1745                decoders.push(Box::new(inner_decoder));
1746            }
1747
1748            DatumColumnDecoder::Record {
1749                fields: decoders,
1750                nulls: record_array.nulls().cloned(),
1751            }
1752        }
1753        (x, y) => {
1754            let msg = format!("can't decode column of {x:?} for scalar type {y:?}");
1755            mz_ore::soft_panic_or_log!("{msg}");
1756            anyhow::bail!("{msg}");
1757        }
1758    };
1759
1760    Ok(decoder)
1761}
1762
1763/// Small helper function to create a [`DatumColumnEncoder`] from a [`SqlScalarType`]
1764fn scalar_type_to_encoder(col_ty: &SqlScalarType) -> Result<DatumColumnEncoder, anyhow::Error> {
1765    let encoder = match &col_ty {
1766        SqlScalarType::Bool => DatumColumnEncoder::Bool(BooleanBuilder::new()),
1767        SqlScalarType::PgLegacyChar => DatumColumnEncoder::U8(UInt8Builder::new()),
1768        SqlScalarType::UInt16 => DatumColumnEncoder::U16(UInt16Builder::new()),
1769        SqlScalarType::UInt32
1770        | SqlScalarType::Oid
1771        | SqlScalarType::RegClass
1772        | SqlScalarType::RegProc
1773        | SqlScalarType::RegType => DatumColumnEncoder::U32(UInt32Builder::new()),
1774        SqlScalarType::UInt64 => DatumColumnEncoder::U64(UInt64Builder::new()),
1775        SqlScalarType::Int16 => DatumColumnEncoder::I16(Int16Builder::new()),
1776        SqlScalarType::Int32 => DatumColumnEncoder::I32(Int32Builder::new()),
1777        SqlScalarType::Int64 => DatumColumnEncoder::I64(Int64Builder::new()),
1778        SqlScalarType::Float32 => DatumColumnEncoder::F32(Float32Builder::new()),
1779        SqlScalarType::Float64 => DatumColumnEncoder::F64(Float64Builder::new()),
1780        SqlScalarType::Numeric { .. } => DatumColumnEncoder::Numeric {
1781            approx_values: Float64Builder::new(),
1782            binary_values: BinaryBuilder::new(),
1783            numeric_context: crate::adt::numeric::cx_datum().clone(),
1784        },
1785        SqlScalarType::String
1786        | SqlScalarType::PgLegacyName
1787        | SqlScalarType::Char { .. }
1788        | SqlScalarType::VarChar { .. } => DatumColumnEncoder::String(StringBuilder::new()),
1789        SqlScalarType::Bytes => DatumColumnEncoder::Bytes(BinaryBuilder::new()),
1790        SqlScalarType::Date => DatumColumnEncoder::Date(Int32Builder::new()),
1791        SqlScalarType::Time => {
1792            DatumColumnEncoder::Time(FixedSizeBinaryBuilder::new(TIME_FIXED_BYTES))
1793        }
1794        SqlScalarType::Timestamp { .. } => {
1795            DatumColumnEncoder::Timestamp(FixedSizeBinaryBuilder::new(TIMESTAMP_FIXED_BYTES))
1796        }
1797        SqlScalarType::TimestampTz { .. } => {
1798            DatumColumnEncoder::TimestampTz(FixedSizeBinaryBuilder::new(TIMESTAMP_FIXED_BYTES))
1799        }
1800        SqlScalarType::MzTimestamp => DatumColumnEncoder::MzTimestamp(UInt64Builder::new()),
1801        SqlScalarType::Interval => {
1802            DatumColumnEncoder::Interval(FixedSizeBinaryBuilder::new(INTERVAL_FIXED_BYTES))
1803        }
1804        SqlScalarType::Uuid => {
1805            DatumColumnEncoder::Uuid(FixedSizeBinaryBuilder::new(UUID_FIXED_BYTES))
1806        }
1807        SqlScalarType::AclItem => {
1808            DatumColumnEncoder::AclItem(FixedSizeBinaryBuilder::new(ACL_ITEM_FIXED_BYTES))
1809        }
1810        SqlScalarType::MzAclItem => DatumColumnEncoder::MzAclItem(BinaryBuilder::new()),
1811        SqlScalarType::Range { .. } => DatumColumnEncoder::Range(BinaryBuilder::new()),
1812        SqlScalarType::Jsonb => DatumColumnEncoder::Jsonb {
1813            offsets: vec![0],
1814            buf: Vec::new(),
1815            nulls: None,
1816        },
1817        s @ SqlScalarType::Array(_) | s @ SqlScalarType::Int2Vector => {
1818            let element_type = match s {
1819                SqlScalarType::Array(inner) => inner,
1820                SqlScalarType::Int2Vector => &SqlScalarType::Int16,
1821                _ => unreachable!("checked above"),
1822            };
1823            let inner = scalar_type_to_encoder(element_type)?;
1824            DatumColumnEncoder::Array {
1825                dims: ListBuilder::new(FixedSizeBinaryBuilder::new(ARRAY_DIMENSION_FIXED_BYTES)),
1826                val_lengths: Vec::new(),
1827                vals: Box::new(inner),
1828                nulls: None,
1829            }
1830        }
1831        SqlScalarType::List { element_type, .. } => {
1832            let inner = scalar_type_to_encoder(&*element_type)?;
1833            DatumColumnEncoder::List {
1834                lengths: Vec::new(),
1835                values: Box::new(inner),
1836                nulls: None,
1837            }
1838        }
1839        SqlScalarType::Map { value_type, .. } => {
1840            let inner = scalar_type_to_encoder(&*value_type)?;
1841            DatumColumnEncoder::Map {
1842                lengths: Vec::new(),
1843                keys: StringBuilder::new(),
1844                vals: Box::new(inner),
1845                nulls: None,
1846            }
1847        }
1848        SqlScalarType::Record { fields, .. } if fields.is_empty() => {
1849            DatumColumnEncoder::RecordEmpty(BooleanBuilder::new())
1850        }
1851        SqlScalarType::Record { fields, .. } => {
1852            let encoders = fields
1853                .iter()
1854                .map(|(_name, ty)| {
1855                    scalar_type_to_encoder(&ty.scalar_type).map(|e| DatumEncoder {
1856                        nullable: ty.nullable,
1857                        encoder: e,
1858                    })
1859                })
1860                .collect::<Result<_, _>>()?;
1861
1862            DatumColumnEncoder::Record {
1863                fields: encoders,
1864                nulls: None,
1865                length: 0,
1866            }
1867        }
1868    };
1869    Ok(encoder)
1870}
1871
1872impl Codec for Row {
1873    type Storage = ProtoRow;
1874    type Schema = RelationDesc;
1875
1876    fn codec_name() -> String {
1877        "protobuf[Row]".into()
1878    }
1879
1880    /// Encodes a row into the permanent storage format.
1881    ///
1882    /// This perfectly round-trips through [Row::decode]. It's guaranteed to be
1883    /// readable by future versions of Materialize through v(TODO: Figure out
1884    /// our policy).
1885    fn encode<B>(&self, buf: &mut B)
1886    where
1887        B: BufMut,
1888    {
1889        self.into_proto()
1890            .encode(buf)
1891            .expect("no required fields means no initialization errors");
1892    }
1893
1894    /// Decodes a row from the permanent storage format.
1895    ///
1896    /// This perfectly round-trips through [Row::encode]. It can read rows
1897    /// encoded by historical versions of Materialize back to v(TODO: Figure out
1898    /// our policy).
1899    fn decode(buf: &[u8], schema: &RelationDesc) -> Result<Row, String> {
1900        // NB: We could easily implement this directly instead of via
1901        // `decode_from`, but do this so that we get maximal coverage of the
1902        // more complicated codepath.
1903        //
1904        // The length of the encoded ProtoRow (i.e. `buf.len()`) doesn't perfect
1905        // predict the length of the resulting Row, but it's definitely
1906        // correlated, so probably a decent estimate.
1907        let mut row = Row::with_capacity(buf.len());
1908        <Self as Codec>::decode_from(&mut row, buf, &mut None, schema)?;
1909        Ok(row)
1910    }
1911
1912    fn decode_from<'a>(
1913        &mut self,
1914        buf: &'a [u8],
1915        storage: &mut Option<ProtoRow>,
1916        schema: &RelationDesc,
1917    ) -> Result<(), String> {
1918        let mut proto = storage.take().unwrap_or_default();
1919        proto.clear();
1920        proto.merge(buf).map_err(|err| err.to_string())?;
1921        let ret = self.decode_from_proto(&proto, schema);
1922        storage.replace(proto);
1923        ret
1924    }
1925
1926    fn validate(row: &Self, desc: &Self::Schema) -> Result<(), String> {
1927        for x in Itertools::zip_longest(desc.iter_types(), row.iter()) {
1928            match x {
1929                EitherOrBoth::Both(typ, datum) if datum.is_instance_of_sql(typ) => continue,
1930                _ => return Err(format!("row {:?} did not match desc {:?}", row, desc)),
1931            };
1932        }
1933        Ok(())
1934    }
1935
1936    fn encode_schema(schema: &Self::Schema) -> Bytes {
1937        schema.into_proto().encode_to_vec().into()
1938    }
1939
1940    fn decode_schema(buf: &Bytes) -> Self::Schema {
1941        let proto = ProtoRelationDesc::decode(buf.as_ref()).expect("valid schema");
1942        proto.into_rust().expect("valid schema")
1943    }
1944}
1945
1946impl<'a> From<Datum<'a>> for ProtoDatum {
1947    fn from(x: Datum<'a>) -> Self {
1948        let datum_type = match x {
1949            Datum::False => DatumType::Other(ProtoDatumOther::False.into()),
1950            Datum::True => DatumType::Other(ProtoDatumOther::True.into()),
1951            Datum::Int16(x) => DatumType::Int16(x.into()),
1952            Datum::Int32(x) => DatumType::Int32(x),
1953            Datum::UInt8(x) => DatumType::Uint8(x.into()),
1954            Datum::UInt16(x) => DatumType::Uint16(x.into()),
1955            Datum::UInt32(x) => DatumType::Uint32(x),
1956            Datum::UInt64(x) => DatumType::Uint64(x),
1957            Datum::Int64(x) => DatumType::Int64(x),
1958            Datum::Float32(x) => DatumType::Float32(x.into_inner()),
1959            Datum::Float64(x) => DatumType::Float64(x.into_inner()),
1960            Datum::Date(x) => DatumType::Date(x.into_proto()),
1961            Datum::Time(x) => DatumType::Time(ProtoNaiveTime {
1962                secs: x.num_seconds_from_midnight(),
1963                frac: x.nanosecond(),
1964            }),
1965            Datum::Timestamp(x) => DatumType::Timestamp(x.into_proto()),
1966            Datum::TimestampTz(x) => DatumType::TimestampTz(x.into_proto()),
1967            Datum::Interval(x) => DatumType::Interval(x.into_proto()),
1968            Datum::Bytes(x) => DatumType::Bytes(Bytes::copy_from_slice(x)),
1969            Datum::String(x) => DatumType::String(x.to_owned()),
1970            Datum::Array(x) => DatumType::Array(ProtoArray {
1971                elements: Some(ProtoRow {
1972                    datums: x.elements().iter().map(|x| x.into()).collect(),
1973                }),
1974                dims: x
1975                    .dims()
1976                    .into_iter()
1977                    .map(|x| ProtoArrayDimension {
1978                        lower_bound: i64::cast_from(x.lower_bound),
1979                        length: u64::cast_from(x.length),
1980                    })
1981                    .collect(),
1982            }),
1983            Datum::List(x) => DatumType::List(ProtoRow {
1984                datums: x.iter().map(|x| x.into()).collect(),
1985            }),
1986            Datum::Map(x) => DatumType::Dict(ProtoDict {
1987                elements: x
1988                    .iter()
1989                    .map(|(k, v)| ProtoDictElement {
1990                        key: k.to_owned(),
1991                        val: Some(v.into()),
1992                    })
1993                    .collect(),
1994            }),
1995            Datum::Numeric(x) => {
1996                // TODO: Do we need this defensive clone?
1997                let mut x = x.0.clone();
1998                if let Some((bcd, scale)) = x.to_packed_bcd() {
1999                    DatumType::Numeric(ProtoNumeric { bcd, scale })
2000                } else if x.is_nan() {
2001                    DatumType::Other(ProtoDatumOther::NumericNaN.into())
2002                } else if x.is_infinite() {
2003                    if x.is_negative() {
2004                        DatumType::Other(ProtoDatumOther::NumericNegInf.into())
2005                    } else {
2006                        DatumType::Other(ProtoDatumOther::NumericPosInf.into())
2007                    }
2008                } else if x.is_special() {
2009                    panic!("internal error: unhandled special numeric value: {}", x);
2010                } else {
2011                    panic!(
2012                        "internal error: to_packed_bcd returned None for non-special value: {}",
2013                        x
2014                    )
2015                }
2016            }
2017            Datum::JsonNull => DatumType::Other(ProtoDatumOther::JsonNull.into()),
2018            Datum::Uuid(x) => DatumType::Uuid(x.as_bytes().to_vec()),
2019            Datum::MzTimestamp(x) => DatumType::MzTimestamp(x.into()),
2020            Datum::Dummy => DatumType::Other(ProtoDatumOther::Dummy.into()),
2021            Datum::Null => DatumType::Other(ProtoDatumOther::Null.into()),
2022            Datum::Range(super::Range { inner }) => DatumType::Range(Box::new(ProtoRange {
2023                inner: inner.map(|RangeInner { lower, upper }| {
2024                    Box::new(ProtoRangeInner {
2025                        lower_inclusive: lower.inclusive,
2026                        lower: lower.bound.map(|bound| Box::new(bound.datum().into())),
2027                        upper_inclusive: upper.inclusive,
2028                        upper: upper.bound.map(|bound| Box::new(bound.datum().into())),
2029                    })
2030                }),
2031            })),
2032            Datum::MzAclItem(x) => DatumType::MzAclItem(x.into_proto()),
2033            Datum::AclItem(x) => DatumType::AclItem(x.into_proto()),
2034        };
2035        ProtoDatum {
2036            datum_type: Some(datum_type),
2037        }
2038    }
2039}
2040
2041impl RowPacker<'_> {
2042    pub(crate) fn try_push_proto(&mut self, x: &ProtoDatum) -> Result<(), String> {
2043        match &x.datum_type {
2044            Some(DatumType::Other(o)) => match ProtoDatumOther::try_from(*o) {
2045                Ok(ProtoDatumOther::Unknown) => return Err("unknown datum type".into()),
2046                Ok(ProtoDatumOther::Null) => self.push(Datum::Null),
2047                Ok(ProtoDatumOther::False) => self.push(Datum::False),
2048                Ok(ProtoDatumOther::True) => self.push(Datum::True),
2049                Ok(ProtoDatumOther::JsonNull) => self.push(Datum::JsonNull),
2050                Ok(ProtoDatumOther::Dummy) => {
2051                    // We plan to remove the `Dummy` variant soon (materialize#17099). To prepare for that, we
2052                    // emit a log to Sentry here, to notify us of any instances that might have
2053                    // been made durable.
2054                    #[cfg(feature = "tracing")]
2055                    tracing::error!("protobuf decoding found Dummy datum");
2056                    self.push(Datum::Dummy);
2057                }
2058                Ok(ProtoDatumOther::NumericPosInf) => self.push(Datum::from(Numeric::infinity())),
2059                Ok(ProtoDatumOther::NumericNegInf) => self.push(Datum::from(-Numeric::infinity())),
2060                Ok(ProtoDatumOther::NumericNaN) => self.push(Datum::from(Numeric::nan())),
2061                Err(_) => return Err(format!("unknown datum type: {}", o)),
2062            },
2063            Some(DatumType::Int16(x)) => {
2064                let x = i16::try_from(*x)
2065                    .map_err(|_| format!("int16 field stored with out of range value: {}", *x))?;
2066                self.push(Datum::Int16(x))
2067            }
2068            Some(DatumType::Int32(x)) => self.push(Datum::Int32(*x)),
2069            Some(DatumType::Int64(x)) => self.push(Datum::Int64(*x)),
2070            Some(DatumType::Uint8(x)) => {
2071                let x = u8::try_from(*x)
2072                    .map_err(|_| format!("uint8 field stored with out of range value: {}", *x))?;
2073                self.push(Datum::UInt8(x))
2074            }
2075            Some(DatumType::Uint16(x)) => {
2076                let x = u16::try_from(*x)
2077                    .map_err(|_| format!("uint16 field stored with out of range value: {}", *x))?;
2078                self.push(Datum::UInt16(x))
2079            }
2080            Some(DatumType::Uint32(x)) => self.push(Datum::UInt32(*x)),
2081            Some(DatumType::Uint64(x)) => self.push(Datum::UInt64(*x)),
2082            Some(DatumType::Float32(x)) => self.push(Datum::Float32((*x).into())),
2083            Some(DatumType::Float64(x)) => self.push(Datum::Float64((*x).into())),
2084            Some(DatumType::Bytes(x)) => self.push(Datum::Bytes(x)),
2085            Some(DatumType::String(x)) => self.push(Datum::String(x)),
2086            Some(DatumType::Uuid(x)) => {
2087                // Uuid internally has a [u8; 16] so we'll have to do at least
2088                // one copy, but there's currently an additional one when the
2089                // Vec is created. Perhaps the protobuf Bytes support will let
2090                // us fix one of them.
2091                let u = Uuid::from_slice(x).map_err(|err| err.to_string())?;
2092                self.push(Datum::Uuid(u));
2093            }
2094            Some(DatumType::Date(x)) => self.push(Datum::Date(x.clone().into_rust()?)),
2095            Some(DatumType::Time(x)) => self.push(Datum::Time(x.clone().into_rust()?)),
2096            Some(DatumType::Timestamp(x)) => self.push(Datum::Timestamp(x.clone().into_rust()?)),
2097            Some(DatumType::TimestampTz(x)) => {
2098                self.push(Datum::TimestampTz(x.clone().into_rust()?))
2099            }
2100            Some(DatumType::Interval(x)) => self.push(Datum::Interval(
2101                x.clone()
2102                    .into_rust()
2103                    .map_err(|e: TryFromProtoError| e.to_string())?,
2104            )),
2105            Some(DatumType::List(x)) => self.push_list_with(|row| -> Result<(), String> {
2106                for d in x.datums.iter() {
2107                    row.try_push_proto(d)?;
2108                }
2109                Ok(())
2110            })?,
2111            Some(DatumType::Array(x)) => {
2112                let dims = x
2113                    .dims
2114                    .iter()
2115                    .map(|x| ArrayDimension {
2116                        lower_bound: isize::cast_from(x.lower_bound),
2117                        length: usize::cast_from(x.length),
2118                    })
2119                    .collect::<Vec<_>>();
2120                match x.elements.as_ref() {
2121                    None => self.try_push_array(&dims, [].iter()),
2122                    Some(elements) => {
2123                        // TODO: Could we avoid this Row alloc if we made a
2124                        // push_array_with?
2125                        let elements_row = Row::try_from(elements)?;
2126                        self.try_push_array(&dims, elements_row.iter())
2127                    }
2128                }
2129                .map_err(|err| err.to_string())?
2130            }
2131            Some(DatumType::Dict(x)) => self.push_dict_with(|row| -> Result<(), String> {
2132                let mut prev_key: Option<&str> = None;
2133                for e in x.elements.iter() {
2134                    // Map keys must be unique and strictly ascending; iterating a
2135                    // map that violates this trips a debug_assert. A crafted
2136                    // proto can, so reject it as a decode error here instead.
2137                    if let Some(prev) = prev_key
2138                        && e.key.as_str() <= prev
2139                    {
2140                        return Err(format!(
2141                            "dict keys must be unique and in ascending order, \
2142                             but {:?} came after {:?}",
2143                            e.key, prev,
2144                        ));
2145                    }
2146                    prev_key = Some(e.key.as_str());
2147                    row.push(Datum::from(e.key.as_str()));
2148                    let val = e
2149                        .val
2150                        .as_ref()
2151                        .ok_or_else(|| format!("missing val for key: {}", e.key))?;
2152                    row.try_push_proto(val)?;
2153                }
2154                Ok(())
2155            })?,
2156            Some(DatumType::Numeric(x)) => {
2157                // Reminder that special values like NaN, PosInf, and NegInf are
2158                // represented as variants of ProtoDatumOther.
2159                //
2160                // `decPackedToNumber` (called via `Decimal::from_packed_bcd`)
2161                // doesn't bounds-check its input: it segfaults on an empty bcd,
2162                // and it writes one BCD digit per input nibble into `Numeric`'s
2163                // fixed-size coefficient, overrunning it for a longer bcd. Both
2164                // are reachable from untrusted proto bytes, so bound the length
2165                // before descending into the FFI. `to_packed_bcd` never emits
2166                // more than MAX_BCD_LEN bytes: one nibble per digit plus a sign
2167                // nibble, rounded up to whole bytes.
2168                const MAX_BCD_LEN: usize =
2169                    mz_ore::cast::u8_to_usize(NUMERIC_DATUM_MAX_PRECISION) / 2 + 1;
2170                if x.bcd.is_empty() || x.bcd.len() > MAX_BCD_LEN {
2171                    return Err(format!(
2172                        "ProtoNumeric.bcd length {} not in 1..={}",
2173                        x.bcd.len(),
2174                        MAX_BCD_LEN,
2175                    ));
2176                }
2177                let n = Decimal::from_packed_bcd(&x.bcd, x.scale).map_err(|err| err.to_string())?;
2178                self.push(Datum::from(n))
2179            }
2180            Some(DatumType::MzTimestamp(x)) => self.push(Datum::MzTimestamp((*x).into())),
2181            Some(DatumType::Range(inner)) => {
2182                let ProtoRange { inner } = &**inner;
2183                match inner {
2184                    None => self.push_range(Range { inner: None }).unwrap(),
2185                    Some(inner) => {
2186                        let ProtoRangeInner {
2187                            lower_inclusive,
2188                            lower,
2189                            upper_inclusive,
2190                            upper,
2191                        } = &**inner;
2192
2193                        // Range bounds must not be `Datum::Null`, because
2194                        // `push_range_with` panics on that invariant. Reject
2195                        // untrusted proto bytes that would push a `Null` bound
2196                        // before calling it.
2197                        let is_null_proto = |d: &ProtoDatum| {
2198                            matches!(
2199                                d.datum_type,
2200                                Some(DatumType::Other(o))
2201                                    if ProtoDatumOther::try_from(o)
2202                                        == Ok(ProtoDatumOther::Null)
2203                            )
2204                        };
2205                        if lower.as_deref().is_some_and(is_null_proto)
2206                            || upper.as_deref().is_some_and(is_null_proto)
2207                        {
2208                            return Err("range bound cannot be Null".into());
2209                        }
2210
2211                        self.push_range_with(
2212                            RangeLowerBound {
2213                                inclusive: *lower_inclusive,
2214                                bound: lower
2215                                    .as_ref()
2216                                    .map(|d| |row: &mut RowPacker| row.try_push_proto(&*d)),
2217                            },
2218                            RangeUpperBound {
2219                                inclusive: *upper_inclusive,
2220                                bound: upper
2221                                    .as_ref()
2222                                    .map(|d| |row: &mut RowPacker| row.try_push_proto(&*d)),
2223                            },
2224                        )
2225                        .map_err(|err| err.to_string())?;
2226                    }
2227                }
2228            }
2229            Some(DatumType::MzAclItem(x)) => self.push(Datum::MzAclItem(x.clone().into_rust()?)),
2230            Some(DatumType::AclItem(x)) => self.push(Datum::AclItem(x.clone().into_rust()?)),
2231            None => return Err("unknown datum type".into()),
2232        };
2233        Ok(())
2234    }
2235}
2236
2237/// TODO: remove this in favor of [`RustType::from_proto`].
2238impl TryFrom<&ProtoRow> for Row {
2239    type Error = String;
2240
2241    fn try_from(x: &ProtoRow) -> Result<Self, Self::Error> {
2242        // TODO: Try to pre-size this.
2243        // see https://github.com/MaterializeInc/database-issues/issues/3640
2244        let mut row = Row::default();
2245        let mut packer = row.packer();
2246        for d in x.datums.iter() {
2247            packer.try_push_proto(d)?;
2248        }
2249        Ok(row)
2250    }
2251}
2252
2253impl RustType<ProtoRow> for Row {
2254    fn into_proto(&self) -> ProtoRow {
2255        let datums = self.iter().map(|x| x.into()).collect();
2256        ProtoRow { datums }
2257    }
2258
2259    fn from_proto(proto: ProtoRow) -> Result<Self, TryFromProtoError> {
2260        // TODO: Try to pre-size this.
2261        // see https://github.com/MaterializeInc/database-issues/issues/3640
2262        let mut row = Row::default();
2263        let mut packer = row.packer();
2264        for d in proto.datums.iter() {
2265            packer
2266                .try_push_proto(d)
2267                .map_err(TryFromProtoError::RowConversionError)?;
2268        }
2269        Ok(row)
2270    }
2271}
2272
2273#[cfg(test)]
2274mod tests {
2275    use std::collections::BTreeSet;
2276
2277    use arrow::array::{ArrayData, make_array};
2278    use arrow::compute::SortOptions;
2279    use arrow::datatypes::ArrowNativeType;
2280    use arrow::row::SortField;
2281    use chrono::{DateTime, NaiveDate, NaiveTime, Utc};
2282    use mz_ore::assert_err;
2283    use mz_ore::collections::CollectionExt;
2284    use mz_persist::indexed::columnar::arrow::realloc_array;
2285    use mz_persist::metrics::ColumnarMetrics;
2286    use mz_persist_types::Codec;
2287    use mz_persist_types::arrow::{ArrayBound, ArrayOrd};
2288    use mz_persist_types::columnar::{codec_to_schema, schema_to_codec};
2289    use mz_proto::{ProtoType, RustType};
2290    use proptest::prelude::*;
2291    use proptest::strategy::Strategy;
2292    use uuid::Uuid;
2293
2294    use super::*;
2295    use crate::adt::array::ArrayDimension;
2296    use crate::adt::interval::Interval;
2297    use crate::adt::numeric::Numeric;
2298    use crate::adt::timestamp::CheckedTimestamp;
2299    use crate::relation::arb_relation_desc;
2300    use crate::{ColumnName, RowArena, SqlColumnType, arb_datum_for_column, arb_row_for_relation};
2301    use crate::{Datum, RelationDesc, Row, SqlScalarType};
2302
2303    #[mz_ore::test]
2304    fn proto_row_invalid_range_is_error() {
2305        // A ProtoRow with a range whose bounds have inconsistent datum kinds
2306        // (or a null/extra bound) must decode to an error, not panic. The range
2307        // packer used to `assert!` these invariants. Regression for the
2308        // row_proto_roundtrip cargo-fuzz finding.
2309        use prost::Message;
2310        let bytes: &[u8] = &[
2311            0x0a, 0x03, 0xaa, 0x01, 0x00, 0x0a, 0x03, 0xaa, 0x01, 0x00, 0x0a, 0x03, 0xa2, 0x01,
2312            0x00, 0x0a, 0x03, 0xaa, 0x01, 0x00, 0x0a, 0x20, 0xfa, 0x01, 0x1d, 0x1d, 0x9f, 0x00,
2313            0x00, 0x00, 0xaa, 0x01, 0x00, 0x0a, 0x13, 0xf8, 0x01, 0x08, 0xaa, 0x0a, 0x03, 0xba,
2314            0x01, 0x00, 0x22, 0x03, 0xba, 0x01, 0x00, 0x12, 0x03, 0xaa, 0x01, 0x00, 0x0a, 0x03,
2315            0xaa, 0x01, 0x00,
2316        ];
2317        let proto = ProtoRow::decode(bytes).expect("crash input decodes as a proto");
2318        let result: Result<Row, _> = proto.into_rust();
2319        assert_err!(result);
2320    }
2321
2322    #[mz_ore::test]
2323    fn proto_row_unordered_dict_keys_is_error() {
2324        // A ProtoRow with a dict whose keys are duplicated or not in ascending
2325        // order must decode to an error, not panic. Iterating such a map trips
2326        // a debug_assert. Regression for the row_proto_roundtrip cargo-fuzz
2327        // finding.
2328        use prost::Message;
2329        let bytes: &[u8] = &[
2330            0x0a, 0x32, 0x18, 0x4e, 0x18, 0x18, 0x68, 0x4e, 0xe8, 0x68, 0x57, 0xba, 0x01, 0x0a,
2331            0x0a, 0x08, 0x60, 0xff, 0xff, 0x10, 0x12, 0x02, 0x10, 0x10, 0x99, 0x68, 0x0a, 0x18,
2332            0x18, 0x4e, 0x18, 0x18, 0x68, 0x4e, 0xe8, 0x5b, 0x18, 0x68, 0x57, 0xba, 0x01, 0x0a,
2333            0x0a, 0x08, 0x60, 0xff, 0xff, 0x10, 0x12, 0x02, 0x18, 0x10,
2334        ];
2335        let proto = ProtoRow::decode(bytes).expect("crash input decodes as a proto");
2336        let result: Result<Row, _> = proto.into_rust();
2337        assert_err!(result);
2338    }
2339
2340    #[mz_ore::test]
2341    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
2342    fn proto_row_oversized_numeric_bcd_is_error() {
2343        // A `ProtoNumeric` whose `bcd` holds more digits than `Numeric`'s
2344        // coefficient can hold must decode to an error. `decPackedToNumber`
2345        // writes one BCD digit per input nibble into the fixed-size
2346        // coefficient without bounds-checking it, so a long `bcd` writes past
2347        // the end of the `Decimal` and eventually segfaults.
2348        for len in [21, 64, 4096] {
2349            let mut bcd = vec![0x99u8; len];
2350            // DECPPLUS, so the sign nibble check passes and the digit copy runs.
2351            *bcd.last_mut().unwrap() = 0x9c;
2352            let proto = ProtoRow {
2353                datums: vec![ProtoDatum {
2354                    datum_type: Some(DatumType::Numeric(ProtoNumeric { bcd, scale: 0 })),
2355                }],
2356            };
2357            assert_err!(Row::try_from(&proto));
2358        }
2359
2360        // The bound must not reject anything we encode: a numeric at max
2361        // precision packs into exactly MAX_BCD_LEN bytes.
2362        let mut cx = crate::adt::numeric::cx_datum();
2363        let digits = "9".repeat(usize::from(NUMERIC_DATUM_MAX_PRECISION));
2364        let n: Numeric = cx.parse(digits.as_str()).unwrap();
2365        let row = Row::pack_slice(&[Datum::Numeric(OrderedDecimal(n))]);
2366        assert_eq!(Row::try_from(&row.into_proto()).unwrap(), row);
2367    }
2368
2369    #[track_caller]
2370    fn roundtrip_datum<'a>(
2371        ty: SqlColumnType,
2372        datum: impl Iterator<Item = Datum<'a>>,
2373        metrics: &ColumnarMetrics,
2374    ) {
2375        let desc = RelationDesc::builder().with_column("a", ty).finish();
2376        let rows = datum.map(|d| Row::pack_slice(&[d])).collect();
2377        roundtrip_rows(&desc, rows, metrics)
2378    }
2379
2380    #[track_caller]
2381    fn roundtrip_rows(desc: &RelationDesc, rows: Vec<Row>, metrics: &ColumnarMetrics) {
2382        let mut encoder = <RelationDesc as Schema<Row>>::encoder(desc).unwrap();
2383        for row in &rows {
2384            encoder.append(row);
2385        }
2386        let col = encoder.finish();
2387
2388        // Exercise reallocating columns with lgalloc.
2389        let col = realloc_array(&col, metrics);
2390        // Exercise our ProtoArray format.
2391        {
2392            let proto = col.to_data().into_proto();
2393            let bytes = proto.encode_to_vec();
2394            let proto = mz_persist_types::arrow::ProtoArrayData::decode(&bytes[..]).unwrap();
2395            let array_data: ArrayData = proto.into_rust().unwrap();
2396
2397            let col_rnd = StructArray::from(array_data.clone());
2398            assert_eq!(col, col_rnd);
2399
2400            let col_dyn = arrow::array::make_array(array_data);
2401            let col_dyn = col_dyn.as_any().downcast_ref::<StructArray>().unwrap();
2402            assert_eq!(&col, col_dyn);
2403        }
2404
2405        let decoder = <RelationDesc as Schema<Row>>::decoder(desc, col.clone()).unwrap();
2406        let stats = decoder.stats();
2407
2408        // Collect all of our lower and upper bounds.
2409        let arena = RowArena::new();
2410        let (stats, stat_nulls): (Vec<_>, Vec<_>) = desc
2411            .iter()
2412            .map(|(name, ty)| {
2413                let col_stats = stats.cols.get(name.as_str()).unwrap();
2414                let lower_upper =
2415                    crate::stats::col_values(&ty.scalar_type, &col_stats.values, &arena);
2416                let null_count = col_stats.nulls.map_or(0, |n| n.count);
2417
2418                (lower_upper, null_count)
2419            })
2420            .unzip();
2421        // Track how many nulls we saw for each column so we can assert stats match.
2422        let mut actual_nulls = vec![0usize; stats.len()];
2423
2424        let mut rnd_row = Row::default();
2425        for (idx, og_row) in rows.iter().enumerate() {
2426            decoder.decode(idx, &mut rnd_row);
2427            assert_eq!(og_row, &rnd_row);
2428
2429            // Check for each Datum in each Row that we're within our stats bounds.
2430            for (c_idx, (rnd_datum, ty)) in rnd_row.iter().zip_eq(desc.typ().columns()).enumerate()
2431            {
2432                let lower_upper = stats[c_idx];
2433
2434                // Assert our stat bounds are correct.
2435                if rnd_datum.is_null() {
2436                    actual_nulls[c_idx] += 1;
2437                } else if let Some((lower, upper)) = lower_upper {
2438                    assert!(rnd_datum >= lower, "{rnd_datum:?} is not >= {lower:?}");
2439                    assert!(rnd_datum <= upper, "{rnd_datum:?} is not <= {upper:?}");
2440                } else {
2441                    match &ty.scalar_type {
2442                        // JSON stats are handled separately.
2443                        SqlScalarType::Jsonb => (),
2444                        // We don't collect stats for these types.
2445                        SqlScalarType::AclItem
2446                        | SqlScalarType::MzAclItem
2447                        | SqlScalarType::Range { .. }
2448                        | SqlScalarType::Array(_)
2449                        | SqlScalarType::Map { .. }
2450                        | SqlScalarType::List { .. }
2451                        | SqlScalarType::Record { .. }
2452                        | SqlScalarType::Int2Vector => (),
2453                        other => panic!("should have collected stats for {other:?}"),
2454                    }
2455                }
2456            }
2457        }
2458
2459        // Validate that the null counts in our stats matched the actual counts.
2460        for (col_idx, (stats_count, actual_count)) in
2461            stat_nulls.iter().zip_eq(actual_nulls.iter()).enumerate()
2462        {
2463            assert_eq!(
2464                stats_count, actual_count,
2465                "column {col_idx} has incorrect number of nulls!"
2466            );
2467        }
2468
2469        // Validate that we can convert losslessly to codec and back
2470        let codec = schema_to_codec::<Row>(desc, &col).unwrap();
2471        let col2 = codec_to_schema::<Row>(desc, &codec).unwrap();
2472        assert_eq!(col2.as_ref(), &col);
2473
2474        // Validate that we only generate supported array types
2475        let converter = arrow::row::RowConverter::new(vec![SortField::new_with_options(
2476            col.data_type().clone(),
2477            SortOptions {
2478                descending: false,
2479                nulls_first: false,
2480            },
2481        )])
2482        .expect("sortable");
2483        let rows = converter
2484            .convert_columns(&[Arc::new(col.clone())])
2485            .expect("convertible");
2486        let mut row_vec = rows.iter().collect::<Vec<_>>();
2487        row_vec.sort();
2488        let row_col = converter
2489            .convert_rows(row_vec)
2490            .expect("convertible")
2491            .into_element();
2492        assert_eq!(row_col.len(), col.len());
2493
2494        let ord = ArrayOrd::new(&col);
2495        let mut indices = (0..u64::usize_as(col.len())).collect::<Vec<_>>();
2496        indices.sort_by_key(|i| ord.at(i.as_usize()));
2497        let indices = UInt64Array::from(indices);
2498        let ord_col = ::arrow::compute::take(&col, &indices, None).expect("takeable");
2499        assert_eq!(row_col.as_ref(), ord_col.as_ref());
2500
2501        // Check that our order matches the datum-native order when `preserves_order` is true.
2502        let ordered_prefix_len = desc
2503            .iter()
2504            .take_while(|(_, c)| preserves_order(&c.scalar_type))
2505            .count();
2506        let decoder = <RelationDesc as Schema<Row>>::decoder_any(desc, ord_col.as_ref()).unwrap();
2507        let rows = (0..ord_col.len()).map(|i| {
2508            let mut row = Row::default();
2509            decoder.decode(i, &mut row);
2510            row
2511        });
2512        for (a, b) in rows.tuple_windows() {
2513            let a_prefix = a.iter().take(ordered_prefix_len);
2514            let b_prefix = b.iter().take(ordered_prefix_len);
2515            assert!(
2516                a_prefix.cmp(b_prefix).is_le(),
2517                "ordering should be consistent on preserves_order columns: {:#?}\n{:?}\n{:?}",
2518                desc.iter().take(ordered_prefix_len).collect_vec(),
2519                a.iter().take(ordered_prefix_len).collect_vec(),
2520                b.iter().take(ordered_prefix_len).collect_vec()
2521            );
2522        }
2523
2524        // Check that our size estimates are consistent.
2525        assert_eq!(
2526            ord.goodbytes(),
2527            (0..col.len()).map(|i| ord.at(i).goodbytes()).sum::<usize>(),
2528            "total size should match the sum of the sizes at each index"
2529        );
2530
2531        // Check that our lower bounds work as expected.
2532        if !ord_col.is_empty() {
2533            let min_idx = indices.values()[0].as_usize();
2534            let lower_bound = ArrayBound::new(ord_col, min_idx);
2535            let max_encoded_len = 1000;
2536            if let Some(proto) = lower_bound.to_proto_lower(max_encoded_len) {
2537                assert!(
2538                    proto.encoded_len() <= max_encoded_len,
2539                    "should respect the max len"
2540                );
2541                let array_data = proto.into_rust().expect("valid array");
2542                let new_lower_bound = ArrayBound::new(make_array(array_data), 0);
2543                assert!(
2544                    new_lower_bound.get() <= lower_bound.get(),
2545                    "proto-roundtripped bound should be <= the original"
2546                );
2547            }
2548        }
2549    }
2550
2551    #[mz_ore::test]
2552    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
2553    fn proptest_datums() {
2554        let strat = any::<SqlColumnType>().prop_flat_map(|ty| {
2555            proptest::collection::vec(arb_datum_for_column(ty.clone()), 0..16)
2556                .prop_map(move |d| (ty.clone(), d))
2557        });
2558        let metrics = ColumnarMetrics::disconnected();
2559
2560        proptest!(|((ty, datums) in strat)| {
2561            roundtrip_datum(ty.clone(), datums.iter().map(Datum::from), &metrics);
2562        })
2563    }
2564
2565    #[mz_ore::test]
2566    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
2567    fn proptest_non_empty_relation_descs() {
2568        let strat = arb_relation_desc(1..8).prop_flat_map(|desc| {
2569            proptest::collection::vec(arb_row_for_relation(&desc), 0..12)
2570                .prop_map(move |rows| (desc.clone(), rows))
2571        });
2572        let metrics = ColumnarMetrics::disconnected();
2573
2574        proptest!(|((desc, rows) in strat)| {
2575            roundtrip_rows(&desc, rows, &metrics)
2576        })
2577    }
2578
2579    #[mz_ore::test]
2580    fn empty_relation_desc_returns_error() {
2581        let empty_desc = RelationDesc::empty();
2582        let result = <RelationDesc as Schema<Row>>::encoder(&empty_desc);
2583        assert_err!(result);
2584    }
2585
2586    #[mz_ore::test]
2587    fn smoketest_collections() {
2588        let mut row = Row::default();
2589        let mut packer = row.packer();
2590        let metrics = ColumnarMetrics::disconnected();
2591
2592        packer
2593            .try_push_array(
2594                &[ArrayDimension {
2595                    lower_bound: 0,
2596                    length: 3,
2597                }],
2598                [Datum::UInt32(4), Datum::UInt32(5), Datum::UInt32(6)],
2599            )
2600            .unwrap();
2601
2602        let array = row.unpack_first();
2603        roundtrip_datum(
2604            SqlScalarType::Array(Box::new(SqlScalarType::UInt32)).nullable(true),
2605            [array].into_iter(),
2606            &metrics,
2607        );
2608    }
2609
2610    #[mz_ore::test]
2611    fn smoketest_row() {
2612        let desc = RelationDesc::builder()
2613            .with_column("a", SqlScalarType::Int64.nullable(true))
2614            .with_column("b", SqlScalarType::String.nullable(true))
2615            .with_column("c", SqlScalarType::Bool.nullable(true))
2616            .with_column(
2617                "d",
2618                SqlScalarType::List {
2619                    element_type: Box::new(SqlScalarType::UInt32),
2620                    custom_id: None,
2621                }
2622                .nullable(true),
2623            )
2624            .with_column(
2625                "e",
2626                SqlScalarType::Map {
2627                    value_type: Box::new(SqlScalarType::Int16),
2628                    custom_id: None,
2629                }
2630                .nullable(true),
2631            )
2632            .finish();
2633        let mut encoder = <RelationDesc as Schema<Row>>::encoder(&desc).unwrap();
2634
2635        let mut og_row = Row::default();
2636        {
2637            let mut packer = og_row.packer();
2638            packer.push(Datum::Int64(100));
2639            packer.push(Datum::String("hello world"));
2640            packer.push(Datum::True);
2641            packer.push_list([Datum::UInt32(1), Datum::UInt32(2), Datum::UInt32(3)]);
2642            packer.push_dict([("bar", Datum::Int16(9)), ("foo", Datum::Int16(3))]);
2643        }
2644        let mut og_row_2 = Row::default();
2645        {
2646            let mut packer = og_row_2.packer();
2647            packer.push(Datum::Null);
2648            packer.push(Datum::Null);
2649            packer.push(Datum::Null);
2650            packer.push(Datum::Null);
2651            packer.push(Datum::Null);
2652        }
2653
2654        encoder.append(&og_row);
2655        encoder.append(&og_row_2);
2656        let col = encoder.finish();
2657
2658        let decoder = <RelationDesc as Schema<Row>>::decoder(&desc, col).unwrap();
2659
2660        let mut rnd_row = Row::default();
2661        decoder.decode(0, &mut rnd_row);
2662        assert_eq!(og_row, rnd_row);
2663
2664        let mut rnd_row = Row::default();
2665        decoder.decode(1, &mut rnd_row);
2666        assert_eq!(og_row_2, rnd_row);
2667    }
2668
2669    #[mz_ore::test]
2670    fn test_nested_list() {
2671        let desc = RelationDesc::builder()
2672            .with_column(
2673                "a",
2674                SqlScalarType::List {
2675                    element_type: Box::new(SqlScalarType::List {
2676                        element_type: Box::new(SqlScalarType::Int64),
2677                        custom_id: None,
2678                    }),
2679                    custom_id: None,
2680                }
2681                .nullable(false),
2682            )
2683            .finish();
2684        let mut encoder = <RelationDesc as Schema<Row>>::encoder(&desc).unwrap();
2685
2686        let mut og_row = Row::default();
2687        {
2688            let mut packer = og_row.packer();
2689            packer.push_list_with(|inner| {
2690                inner.push_list([Datum::Int64(1), Datum::Int64(2)]);
2691                inner.push_list([Datum::Int64(5)]);
2692                inner.push_list([Datum::Int64(9), Datum::Int64(99), Datum::Int64(999)]);
2693            });
2694        }
2695
2696        encoder.append(&og_row);
2697        let col = encoder.finish();
2698
2699        let decoder = <RelationDesc as Schema<Row>>::decoder(&desc, col).unwrap();
2700        let mut rnd_row = Row::default();
2701        decoder.decode(0, &mut rnd_row);
2702
2703        assert_eq!(og_row, rnd_row);
2704    }
2705
2706    #[mz_ore::test]
2707    fn test_record() {
2708        let desc = RelationDesc::builder()
2709            .with_column(
2710                "a",
2711                SqlScalarType::Record {
2712                    fields: [
2713                        (
2714                            ColumnName::from("foo"),
2715                            SqlScalarType::Int64.nullable(false),
2716                        ),
2717                        (
2718                            ColumnName::from("bar"),
2719                            SqlScalarType::String.nullable(true),
2720                        ),
2721                        (
2722                            ColumnName::from("baz"),
2723                            SqlScalarType::List {
2724                                element_type: Box::new(SqlScalarType::UInt32),
2725                                custom_id: None,
2726                            }
2727                            .nullable(false),
2728                        ),
2729                    ]
2730                    .into(),
2731                    custom_id: None,
2732                }
2733                .nullable(true),
2734            )
2735            .finish();
2736        let mut encoder = <RelationDesc as Schema<Row>>::encoder(&desc).unwrap();
2737
2738        let mut og_row = Row::default();
2739        {
2740            let mut packer = og_row.packer();
2741            packer.push_list_with(|inner| {
2742                inner.push(Datum::Int64(42));
2743                inner.push(Datum::Null);
2744                inner.push_list([Datum::UInt32(1), Datum::UInt32(2), Datum::UInt32(3)]);
2745            });
2746        }
2747        let null_row = Row::pack_slice(&[Datum::Null]);
2748
2749        encoder.append(&og_row);
2750        encoder.append(&null_row);
2751        let col = encoder.finish();
2752
2753        let decoder = <RelationDesc as Schema<Row>>::decoder(&desc, col).unwrap();
2754        let mut rnd_row = Row::default();
2755
2756        decoder.decode(0, &mut rnd_row);
2757        assert_eq!(og_row, rnd_row);
2758
2759        rnd_row.packer();
2760        decoder.decode(1, &mut rnd_row);
2761        assert_eq!(null_row, rnd_row);
2762    }
2763
2764    #[mz_ore::test]
2765    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decNumberFromInt32` on OS `linux`
2766    fn roundtrip() {
2767        let mut row = Row::default();
2768        let mut packer = row.packer();
2769        packer.extend([
2770            Datum::False,
2771            Datum::True,
2772            Datum::Int16(1),
2773            Datum::Int32(2),
2774            Datum::Int64(3),
2775            Datum::Float32(4f32.into()),
2776            Datum::Float64(5f64.into()),
2777            Datum::Date(
2778                NaiveDate::from_ymd_opt(6, 7, 8)
2779                    .unwrap()
2780                    .try_into()
2781                    .unwrap(),
2782            ),
2783            Datum::Time(NaiveTime::from_hms_opt(9, 10, 11).unwrap()),
2784            Datum::Timestamp(
2785                CheckedTimestamp::from_timestamplike(
2786                    NaiveDate::from_ymd_opt(12, 13 % 12, 14)
2787                        .unwrap()
2788                        .and_time(NaiveTime::from_hms_opt(15, 16, 17).unwrap()),
2789                )
2790                .unwrap(),
2791            ),
2792            Datum::TimestampTz(
2793                CheckedTimestamp::from_timestamplike(DateTime::from_naive_utc_and_offset(
2794                    NaiveDate::from_ymd_opt(18, 19 % 12, 20)
2795                        .unwrap()
2796                        .and_time(NaiveTime::from_hms_opt(21, 22, 23).unwrap()),
2797                    Utc,
2798                ))
2799                .unwrap(),
2800            ),
2801            Datum::Interval(Interval {
2802                months: 24,
2803                days: 42,
2804                micros: 25,
2805            }),
2806            Datum::Bytes(&[26, 27]),
2807            Datum::String("28"),
2808            Datum::from(Numeric::from(29)),
2809            Datum::from(Numeric::infinity()),
2810            Datum::from(-Numeric::infinity()),
2811            Datum::from(Numeric::nan()),
2812            Datum::JsonNull,
2813            Datum::Uuid(Uuid::from_u128(30)),
2814            Datum::Dummy,
2815            Datum::Null,
2816        ]);
2817        packer
2818            .try_push_array(
2819                &[ArrayDimension {
2820                    lower_bound: 2,
2821                    length: 2,
2822                }],
2823                vec![Datum::Int32(31), Datum::Int32(32)],
2824            )
2825            .expect("valid array");
2826        packer.push_list_with(|packer| {
2827            packer.push(Datum::String("33"));
2828            packer.push_list_with(|packer| {
2829                packer.push(Datum::String("34"));
2830                packer.push(Datum::String("35"));
2831            });
2832            packer.push(Datum::String("36"));
2833            packer.push(Datum::String("37"));
2834        });
2835        packer.push_dict_with(|row| {
2836            // Add a bunch of data to the hash to ensure we don't get a
2837            // HashMap's random iteration anywhere in the encode/decode path.
2838            let mut i = 38;
2839            for _ in 0..20 {
2840                row.push(Datum::String(&i.to_string()));
2841                row.push(Datum::Int32(i + 1));
2842                i += 2;
2843            }
2844        });
2845
2846        let mut desc = RelationDesc::builder();
2847        for (idx, _) in row.iter().enumerate() {
2848            // HACK(parkmycar): We don't currently validate the types of the `RelationDesc` are
2849            // correct, just the number of columns. So we can fill in any type here.
2850            desc = desc.with_column(idx.to_string(), SqlScalarType::Int32.nullable(true));
2851        }
2852        let desc = desc.finish();
2853
2854        let encoded = row.encode_to_vec();
2855        assert_eq!(Row::decode(&encoded, &desc), Ok(row));
2856    }
2857
2858    #[mz_ore::test]
2859    fn smoketest_projection() {
2860        let desc = RelationDesc::builder()
2861            .with_column("a", SqlScalarType::Int64.nullable(true))
2862            .with_column("b", SqlScalarType::String.nullable(true))
2863            .with_column("c", SqlScalarType::Bool.nullable(true))
2864            .finish();
2865        let mut encoder = <RelationDesc as Schema<Row>>::encoder(&desc).unwrap();
2866
2867        let mut og_row = Row::default();
2868        {
2869            let mut packer = og_row.packer();
2870            packer.push(Datum::Int64(100));
2871            packer.push(Datum::String("hello world"));
2872            packer.push(Datum::True);
2873        }
2874        let mut og_row_2 = Row::default();
2875        {
2876            let mut packer = og_row_2.packer();
2877            packer.push(Datum::Null);
2878            packer.push(Datum::Null);
2879            packer.push(Datum::Null);
2880        }
2881
2882        encoder.append(&og_row);
2883        encoder.append(&og_row_2);
2884        let col = encoder.finish();
2885
2886        let projected_desc = desc.apply_demand(&BTreeSet::from([0, 2]));
2887
2888        let decoder = <RelationDesc as Schema<Row>>::decoder(&projected_desc, col).unwrap();
2889
2890        let mut rnd_row = Row::default();
2891        decoder.decode(0, &mut rnd_row);
2892        let expected_row = Row::pack_slice(&[Datum::Int64(100), Datum::True]);
2893        assert_eq!(expected_row, rnd_row);
2894
2895        let mut rnd_row = Row::default();
2896        decoder.decode(1, &mut rnd_row);
2897        let expected_row = Row::pack_slice(&[Datum::Null, Datum::Null]);
2898        assert_eq!(expected_row, rnd_row);
2899    }
2900}