Skip to main content

mz_storage_types/
sources.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Types and traits related to the introduction of changing collections into `dataflow`.
11
12use std::collections::BTreeMap;
13use std::fmt::Debug;
14use std::hash::Hash;
15use std::ops::{Add, AddAssign, Deref, DerefMut};
16use std::str::FromStr;
17use std::sync::Arc;
18use std::time::Duration;
19
20use arrow::array::{Array, ArrayRef, BinaryArray, BinaryBuilder, NullArray, StructArray};
21use arrow::datatypes::{Field, Fields};
22use bytes::{BufMut, Bytes};
23use columnation::Columnation;
24use itertools::EitherOrBoth::Both;
25use itertools::Itertools;
26use kafka::KafkaSourceExportDetails;
27use load_generator::{LoadGeneratorOutput, LoadGeneratorSourceExportDetails};
28use mz_ore::assert_none;
29use mz_persist_types::Codec;
30use mz_persist_types::arrow::ArrayOrd;
31use mz_persist_types::columnar::{ColumnDecoder, ColumnEncoder, Schema};
32use mz_persist_types::stats::{
33    ColumnNullStats, ColumnStatKinds, ColumnarStats, ColumnarStatsBuilder, PrimitiveStats,
34    StructStats,
35};
36use mz_proto::{IntoRustIfSome, ProtoType, RustType, TryFromProtoError};
37#[cfg(any(test, feature = "proptest"))]
38use mz_repr::arb_row_for_relation;
39use mz_repr::{
40    CatalogItemId, Datum, GlobalId, ProtoRelationDesc, ProtoRow, RelationDesc, Row,
41    RowColumnarDecoder, RowColumnarEncoder,
42};
43use mz_sql_parser::ast::{Ident, IdentError, UnresolvedItemName};
44#[cfg(any(test, feature = "proptest"))]
45use proptest::prelude::any;
46#[cfg(any(test, feature = "proptest"))]
47use proptest::strategy::Strategy;
48use prost::Message;
49use serde::{Deserialize, Serialize};
50use timely::order::{PartialOrder, TotalOrder};
51use timely::progress::timestamp::Refines;
52use timely::progress::{PathSummary, Timestamp};
53
54use crate::AlterCompatible;
55use crate::connections::inline::{
56    ConnectionAccess, ConnectionResolver, InlinedConnection, IntoInlineConnection,
57    ReferencedConnection,
58};
59use crate::controller::AlterError;
60use crate::errors::{DataflowError, ProtoDataflowError};
61use crate::instances::StorageInstanceId;
62use crate::sources::sql_server::SqlServerSourceExportDetails;
63
64pub mod casts;
65pub mod encoding;
66pub mod envelope;
67pub mod kafka;
68pub mod load_generator;
69pub mod mysql;
70pub mod postgres;
71pub mod sql_server;
72
73pub use crate::sources::envelope::SourceEnvelope;
74pub use crate::sources::kafka::KafkaSourceConnection;
75pub use crate::sources::load_generator::LoadGeneratorSourceConnection;
76pub use crate::sources::mysql::{MySqlSourceConnection, MySqlSourceExportDetails};
77pub use crate::sources::postgres::{PostgresSourceConnection, PostgresSourceExportDetails};
78pub use crate::sources::sql_server::{SqlServerSourceConnection, SqlServerSourceExtras};
79
80include!(concat!(env!("OUT_DIR"), "/mz_storage_types.sources.rs"));
81
82/// A description of a source ingestion
83#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq)]
84pub struct IngestionDescription<S: 'static = (), C: ConnectionAccess = InlinedConnection> {
85    /// The source description.
86    pub desc: SourceDesc<C>,
87    /// Collections to be exported by this ingestion.
88    ///
89    /// # Notes
90    /// - For multi-output sources:
91    ///     - Add exports by adding a new [`SourceExport`].
92    ///     - Remove exports by removing the [`SourceExport`].
93    ///
94    ///   Re-rendering/executing the source after making these modifications
95    ///   adds and drops the subsource, respectively.
96    /// - For old-syntax sources this field includes the primary source's ID,
97    ///   which might need to be filtered out to understand which exports are
98    ///   explicit ingestion exports. New-syntax sources (with source tables)
99    ///   list only their exports here.
100    /// - This field does _not_ include the remap collection, which is tracked
101    ///   in its own field.
102    pub source_exports: BTreeMap<GlobalId, SourceExport<S>>,
103    /// The ID of the instance in which to install the source.
104    pub instance_id: StorageInstanceId,
105    /// The ID of this ingestion's remap/progress collection.
106    pub remap_collection_id: GlobalId,
107    /// The storage metadata for the remap/progress collection
108    pub remap_metadata: S,
109}
110
111impl IngestionDescription {
112    pub fn new(
113        desc: SourceDesc,
114        instance_id: StorageInstanceId,
115        remap_collection_id: GlobalId,
116    ) -> Self {
117        Self {
118            desc,
119            remap_metadata: (),
120            source_exports: BTreeMap::new(),
121            instance_id,
122            remap_collection_id,
123        }
124    }
125}
126
127impl<S> IngestionDescription<S> {
128    /// Return an iterator over the `GlobalId`s of `self`'s collections.
129    /// This will contain ids for the remap collection, subsources,
130    /// tables for this source, and the primary collection ID, even if
131    /// no data will be exported to the primary collection.
132    pub fn collection_ids(&self) -> impl Iterator<Item = GlobalId> + '_ {
133        // Expand self so that any new fields added generate a compiler error to
134        // increase the likelihood of developers seeing this function.
135        let IngestionDescription {
136            desc: _,
137            remap_metadata: _,
138            source_exports,
139            instance_id: _,
140            remap_collection_id,
141        } = &self;
142
143        source_exports
144            .keys()
145            .copied()
146            .chain(std::iter::once(*remap_collection_id))
147    }
148}
149
150impl<S: Debug + Eq + PartialEq + AlterCompatible> AlterCompatible for IngestionDescription<S> {
151    fn alter_compatible(
152        &self,
153        id: GlobalId,
154        other: &IngestionDescription<S>,
155    ) -> Result<(), AlterError> {
156        if self == other {
157            return Ok(());
158        }
159        let IngestionDescription {
160            desc,
161            remap_metadata,
162            source_exports,
163            instance_id,
164            remap_collection_id,
165        } = self;
166
167        let compatibility_checks = [
168            (desc.alter_compatible(id, &other.desc).is_ok(), "desc"),
169            (remap_metadata == &other.remap_metadata, "remap_metadata"),
170            (
171                source_exports
172                    .iter()
173                    .merge_join_by(&other.source_exports, |(l_key, _), (r_key, _)| {
174                        l_key.cmp(r_key)
175                    })
176                    .all(|r| match r {
177                        Both(
178                            (
179                                _,
180                                SourceExport {
181                                    storage_metadata: l_metadata,
182                                    details: l_details,
183                                    data_config: l_data_config,
184                                },
185                            ),
186                            (
187                                _,
188                                SourceExport {
189                                    storage_metadata: r_metadata,
190                                    details: r_details,
191                                    data_config: r_data_config,
192                                },
193                            ),
194                        ) => {
195                            l_metadata.alter_compatible(id, r_metadata).is_ok()
196                                && l_details.alter_compatible(id, r_details).is_ok()
197                                && l_data_config.alter_compatible(id, r_data_config).is_ok()
198                        }
199                        _ => true,
200                    }),
201                "source_exports",
202            ),
203            (instance_id == &other.instance_id, "instance_id"),
204            (
205                remap_collection_id == &other.remap_collection_id,
206                "remap_collection_id",
207            ),
208        ];
209        for (compatible, field) in compatibility_checks {
210            if !compatible {
211                tracing::warn!(
212                    "IngestionDescription incompatible at {field}:\nself:\n{:#?}\n\nother\n{:#?}",
213                    self,
214                    other
215                );
216
217                return Err(AlterError { id });
218            }
219        }
220
221        Ok(())
222    }
223}
224
225impl<R: ConnectionResolver> IntoInlineConnection<IngestionDescription, R>
226    for IngestionDescription<(), ReferencedConnection>
227{
228    fn into_inline_connection(self, r: R) -> IngestionDescription {
229        let IngestionDescription {
230            desc,
231            remap_metadata,
232            source_exports,
233            instance_id,
234            remap_collection_id,
235        } = self;
236
237        IngestionDescription {
238            desc: desc.into_inline_connection(r),
239            remap_metadata,
240            source_exports,
241            instance_id,
242            remap_collection_id,
243        }
244    }
245}
246
247#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq)]
248pub struct SourceExport<S = (), C: ConnectionAccess = InlinedConnection> {
249    /// The collection metadata needed to write the exported data
250    pub storage_metadata: S,
251    /// Details necessary for the source to export data to this export's collection.
252    pub details: SourceExportDetails,
253    /// Config necessary to handle (e.g. decode and envelope) the data for this export.
254    pub data_config: SourceExportDataConfig<C>,
255}
256
257pub trait SourceTimestamp:
258    Timestamp + Columnation + Refines<()> + std::fmt::Display + Sync
259{
260    fn encode_row(&self) -> Row;
261    fn decode_row(row: &Row) -> Self;
262}
263
264impl SourceTimestamp for MzOffset {
265    fn encode_row(&self) -> Row {
266        Row::pack([Datum::UInt64(self.offset)])
267    }
268
269    fn decode_row(row: &Row) -> Self {
270        let mut datums = row.iter();
271        match (datums.next(), datums.next()) {
272            (Some(Datum::UInt64(offset)), None) => MzOffset::from(offset),
273            _ => panic!("invalid row {row:?}"),
274        }
275    }
276}
277
278/// Universal language for describing message positions in Materialize, in a source independent
279/// way. Individual sources like Kafka or File sources should explicitly implement their own offset
280/// type that converts to/From MzOffsets. A 0-MzOffset denotes an empty stream.
281#[derive(
282    Copy,
283    Clone,
284    Default,
285    Debug,
286    PartialEq,
287    PartialOrd,
288    Eq,
289    Ord,
290    Hash,
291    Serialize,
292    Deserialize
293)]
294pub struct MzOffset {
295    pub offset: u64,
296}
297
298impl differential_dataflow::difference::Semigroup for MzOffset {
299    fn plus_equals(&mut self, rhs: &Self) {
300        self.offset.plus_equals(&rhs.offset)
301    }
302}
303
304impl differential_dataflow::difference::IsZero for MzOffset {
305    fn is_zero(&self) -> bool {
306        self.offset.is_zero()
307    }
308}
309
310impl mz_persist_types::Codec64 for MzOffset {
311    fn codec_name() -> String {
312        "MzOffset".to_string()
313    }
314
315    fn encode(&self) -> [u8; 8] {
316        mz_persist_types::Codec64::encode(&self.offset)
317    }
318
319    fn decode(buf: [u8; 8]) -> Self {
320        Self {
321            offset: mz_persist_types::Codec64::decode(buf),
322        }
323    }
324}
325
326impl columnation::Columnation for MzOffset {
327    type InnerRegion = columnation::CopyRegion<MzOffset>;
328}
329
330impl MzOffset {
331    pub fn checked_sub(self, other: Self) -> Option<Self> {
332        self.offset
333            .checked_sub(other.offset)
334            .map(|offset| Self { offset })
335    }
336}
337
338/// Convert from MzOffset to Kafka::Offset as long as
339/// the offset is not negative
340impl From<u64> for MzOffset {
341    fn from(offset: u64) -> Self {
342        Self { offset }
343    }
344}
345
346impl std::fmt::Display for MzOffset {
347    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
348        write!(f, "{}", self.offset)
349    }
350}
351
352// Assume overflow does not occur for addition
353impl Add<u64> for MzOffset {
354    type Output = MzOffset;
355
356    fn add(self, x: u64) -> MzOffset {
357        MzOffset {
358            offset: self.offset + x,
359        }
360    }
361}
362impl Add<Self> for MzOffset {
363    type Output = Self;
364
365    fn add(self, x: Self) -> Self {
366        MzOffset {
367            offset: self.offset + x.offset,
368        }
369    }
370}
371impl AddAssign<u64> for MzOffset {
372    fn add_assign(&mut self, x: u64) {
373        self.offset += x;
374    }
375}
376impl AddAssign<Self> for MzOffset {
377    fn add_assign(&mut self, x: Self) {
378        self.offset += x.offset;
379    }
380}
381
382/// Convert from `PgLsn` to MzOffset
383impl From<tokio_postgres::types::PgLsn> for MzOffset {
384    fn from(lsn: tokio_postgres::types::PgLsn) -> Self {
385        MzOffset { offset: lsn.into() }
386    }
387}
388
389impl Timestamp for MzOffset {
390    type Summary = MzOffset;
391
392    fn minimum() -> Self {
393        MzOffset {
394            offset: Timestamp::minimum(),
395        }
396    }
397}
398
399impl PathSummary<MzOffset> for MzOffset {
400    fn results_in(&self, src: &MzOffset) -> Option<MzOffset> {
401        Some(MzOffset {
402            offset: self.offset.results_in(&src.offset)?,
403        })
404    }
405
406    fn followed_by(&self, other: &Self) -> Option<Self> {
407        Some(MzOffset {
408            offset: PathSummary::<u64>::followed_by(&self.offset, &other.offset)?,
409        })
410    }
411}
412
413impl Refines<()> for MzOffset {
414    fn to_inner(_: ()) -> Self {
415        MzOffset::minimum()
416    }
417    fn to_outer(self) {}
418    fn summarize(_: Self::Summary) {}
419}
420
421impl PartialOrder for MzOffset {
422    #[inline]
423    fn less_equal(&self, other: &Self) -> bool {
424        self.offset.less_equal(&other.offset)
425    }
426}
427
428impl TotalOrder for MzOffset {}
429
430/// The meaning of the timestamp number produced by data sources. This type
431/// is not concerned with the source of the timestamp (like if the data came
432/// from a Debezium consistency topic or a CDCv2 stream), instead only what the
433/// timestamp number means.
434///
435/// Some variants here have attached data used to differentiate incomparable
436/// instantiations. These attached data types should be expanded in the future
437/// if we need to tell apart more kinds of sources.
438#[derive(
439    Clone,
440    Debug,
441    Ord,
442    PartialOrd,
443    Eq,
444    PartialEq,
445    Serialize,
446    Deserialize,
447    Hash
448)]
449pub enum Timeline {
450    /// EpochMilliseconds means the timestamp is the number of milliseconds since
451    /// the Unix epoch.
452    EpochMilliseconds,
453    /// External means the timestamp comes from an external data source and we
454    /// don't know what the number means. The attached String is the source's name,
455    /// which will result in different sources being incomparable.
456    External(String),
457    /// User means the user has manually specified a timeline. The attached
458    /// String is specified by the user, allowing them to decide sources that are
459    /// joinable.
460    User(String),
461}
462
463impl Timeline {
464    const EPOCH_MILLISECOND_ID_CHAR: char = 'M';
465    const EXTERNAL_ID_CHAR: char = 'E';
466    const USER_ID_CHAR: char = 'U';
467
468    fn id_char(&self) -> char {
469        match self {
470            Self::EpochMilliseconds => Self::EPOCH_MILLISECOND_ID_CHAR,
471            Self::External(_) => Self::EXTERNAL_ID_CHAR,
472            Self::User(_) => Self::USER_ID_CHAR,
473        }
474    }
475}
476
477impl ToString for Timeline {
478    fn to_string(&self) -> String {
479        match self {
480            Self::EpochMilliseconds => format!("{}", self.id_char()),
481            Self::External(id) => format!("{}.{id}", self.id_char()),
482            Self::User(id) => format!("{}.{id}", self.id_char()),
483        }
484    }
485}
486
487impl FromStr for Timeline {
488    type Err = String;
489
490    fn from_str(s: &str) -> Result<Self, Self::Err> {
491        if s.is_empty() {
492            return Err("empty timeline".to_string());
493        }
494        let mut chars = s.chars();
495        match chars.next().expect("non-empty string") {
496            Self::EPOCH_MILLISECOND_ID_CHAR => match chars.next() {
497                None => Ok(Self::EpochMilliseconds),
498                Some(_) => Err(format!("unknown timeline: {s}")),
499            },
500            Self::EXTERNAL_ID_CHAR => match chars.next() {
501                Some('.') => Ok(Self::External(chars.as_str().to_string())),
502                _ => Err(format!("unknown timeline: {s}")),
503            },
504            Self::USER_ID_CHAR => match chars.next() {
505                Some('.') => Ok(Self::User(chars.as_str().to_string())),
506                _ => Err(format!("unknown timeline: {s}")),
507            },
508            _ => Err(format!("unknown timeline: {s}")),
509        }
510    }
511}
512
513/// A connection to an external system
514pub trait SourceConnection: Debug + Clone + PartialEq + AlterCompatible {
515    /// The name of the external system (e.g kafka, postgres, etc).
516    fn name(&self) -> &'static str;
517
518    /// The name of the resource in the external system (e.g kafka topic) if any
519    fn external_reference(&self) -> Option<&str>;
520
521    /// Defines the key schema to use by default for this source connection type.
522    /// This will be used for the primary export of the source and as the default
523    /// pre-encoding key schema for the source.
524    fn default_key_desc(&self) -> RelationDesc;
525
526    /// Defines the value schema to use by default for this source connection type.
527    /// This will be used for the primary export of the source and as the default
528    /// pre-encoding value schema for the source.
529    fn default_value_desc(&self) -> RelationDesc;
530
531    /// The schema of this connection's timestamp type. This will also be the schema of the
532    /// progress relation.
533    fn timestamp_desc(&self) -> RelationDesc;
534
535    /// The id of the connection object (i.e the one obtained from running `CREATE CONNECTION`) in
536    /// the catalog, if any.
537    fn connection_id(&self) -> Option<CatalogItemId>;
538
539    /// Whether the source type supports read only mode.
540    fn supports_read_only(&self) -> bool;
541
542    /// Whether the source type prefers to run on only one replica of a multi-replica cluster.
543    fn prefers_single_replica(&self) -> bool;
544}
545
546#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
547pub enum Compression {
548    Gzip,
549    None,
550}
551
552/// Defines the configuration for how to handle data that is exported for a given
553/// Source Export.
554#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq)]
555pub struct SourceExportDataConfig<C: ConnectionAccess = InlinedConnection> {
556    pub encoding: Option<encoding::SourceDataEncoding<C>>,
557    pub envelope: SourceEnvelope,
558}
559
560impl<R: ConnectionResolver> IntoInlineConnection<SourceExportDataConfig, R>
561    for SourceExportDataConfig<ReferencedConnection>
562{
563    fn into_inline_connection(self, r: R) -> SourceExportDataConfig {
564        let SourceExportDataConfig { encoding, envelope } = self;
565
566        SourceExportDataConfig {
567            encoding: encoding.map(|e| e.into_inline_connection(r)),
568            envelope,
569        }
570    }
571}
572
573impl<C: ConnectionAccess> AlterCompatible for SourceExportDataConfig<C> {
574    fn alter_compatible(&self, id: GlobalId, other: &Self) -> Result<(), AlterError> {
575        if self == other {
576            return Ok(());
577        }
578        let Self { encoding, envelope } = &self;
579
580        let compatibility_checks = [
581            (
582                match (encoding, &other.encoding) {
583                    (Some(s), Some(o)) => s.alter_compatible(id, o).is_ok(),
584                    (s, o) => s == o,
585                },
586                "encoding",
587            ),
588            (envelope == &other.envelope, "envelope"),
589        ];
590
591        for (compatible, field) in compatibility_checks {
592            if !compatible {
593                tracing::warn!(
594                    "SourceDesc incompatible {field}:\nself:\n{:#?}\n\nother\n{:#?}",
595                    self,
596                    other
597                );
598
599                return Err(AlterError { id });
600            }
601        }
602        Ok(())
603    }
604}
605
606impl<C: ConnectionAccess> SourceExportDataConfig<C> {
607    /// Returns `true` if this connection yields data that is
608    /// append-only/monotonic. Append-monly means the source
609    /// never produces retractions.
610    // TODO(guswynn): consider enforcing this more completely at the
611    // parsing/typechecking level, by not using an `envelope`
612    // for sources like pg
613    pub fn monotonic(&self, connection: &GenericSourceConnection<C>) -> bool {
614        match &self.envelope {
615            // Upsert and CdcV2 may produce retractions.
616            SourceEnvelope::Upsert(_) | SourceEnvelope::CdcV2 => false,
617            SourceEnvelope::None(_) => {
618                match connection {
619                    // Postgres can produce retractions (deletes).
620                    GenericSourceConnection::Postgres(_) => false,
621                    // MySQL can produce retractions (deletes).
622                    GenericSourceConnection::MySql(_) => false,
623                    // SQL Server can produce retractions (deletes).
624                    GenericSourceConnection::SqlServer(_) => false,
625                    // Whether or not a Loadgen source can produce retractions varies.
626                    GenericSourceConnection::LoadGenerator(g) => g.load_generator.is_monotonic(),
627                    // Kafka exports with `None` envelope are append-only.
628                    GenericSourceConnection::Kafka(_) => true,
629                }
630            }
631        }
632    }
633}
634
635/// An external source of updates for a relational collection.
636#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq)]
637pub struct SourceDesc<C: ConnectionAccess = InlinedConnection> {
638    pub connection: GenericSourceConnection<C>,
639    pub timestamp_interval: Duration,
640}
641
642impl<R: ConnectionResolver> IntoInlineConnection<SourceDesc, R>
643    for SourceDesc<ReferencedConnection>
644{
645    fn into_inline_connection(self, r: R) -> SourceDesc {
646        let SourceDesc {
647            connection,
648            timestamp_interval,
649        } = self;
650
651        SourceDesc {
652            connection: connection.into_inline_connection(&r),
653            timestamp_interval,
654        }
655    }
656}
657
658impl<C: ConnectionAccess> AlterCompatible for SourceDesc<C> {
659    /// Determines if `self` is compatible with another `SourceDesc`, in such a
660    /// way that it is possible to turn `self` into `other` through a valid
661    /// series of transformations (e.g. no transformation or `ALTER SOURCE`).
662    fn alter_compatible(&self, id: GlobalId, other: &Self) -> Result<(), AlterError> {
663        if self == other {
664            return Ok(());
665        }
666        let Self {
667            connection,
668            // timestamp_interval is allowed to change via ALTER SOURCE
669            timestamp_interval: _,
670        } = &self;
671
672        let compatibility_checks = [(
673            connection.alter_compatible(id, &other.connection).is_ok(),
674            "connection",
675        )];
676
677        for (compatible, field) in compatibility_checks {
678            if !compatible {
679                tracing::warn!(
680                    "SourceDesc incompatible {field}:\nself:\n{:#?}\n\nother\n{:#?}",
681                    self,
682                    other
683                );
684
685                return Err(AlterError { id });
686            }
687        }
688
689        Ok(())
690    }
691}
692
693#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
694pub enum GenericSourceConnection<C: ConnectionAccess = InlinedConnection> {
695    Kafka(KafkaSourceConnection<C>),
696    Postgres(PostgresSourceConnection<C>),
697    MySql(MySqlSourceConnection<C>),
698    SqlServer(SqlServerSourceConnection<C>),
699    LoadGenerator(LoadGeneratorSourceConnection),
700}
701
702impl<C: ConnectionAccess> From<KafkaSourceConnection<C>> for GenericSourceConnection<C> {
703    fn from(conn: KafkaSourceConnection<C>) -> Self {
704        Self::Kafka(conn)
705    }
706}
707
708impl<C: ConnectionAccess> From<PostgresSourceConnection<C>> for GenericSourceConnection<C> {
709    fn from(conn: PostgresSourceConnection<C>) -> Self {
710        Self::Postgres(conn)
711    }
712}
713
714impl<C: ConnectionAccess> From<MySqlSourceConnection<C>> for GenericSourceConnection<C> {
715    fn from(conn: MySqlSourceConnection<C>) -> Self {
716        Self::MySql(conn)
717    }
718}
719
720impl<C: ConnectionAccess> From<SqlServerSourceConnection<C>> for GenericSourceConnection<C> {
721    fn from(conn: SqlServerSourceConnection<C>) -> Self {
722        Self::SqlServer(conn)
723    }
724}
725
726impl<C: ConnectionAccess> From<LoadGeneratorSourceConnection> for GenericSourceConnection<C> {
727    fn from(conn: LoadGeneratorSourceConnection) -> Self {
728        Self::LoadGenerator(conn)
729    }
730}
731
732impl<R: ConnectionResolver> IntoInlineConnection<GenericSourceConnection, R>
733    for GenericSourceConnection<ReferencedConnection>
734{
735    fn into_inline_connection(self, r: R) -> GenericSourceConnection {
736        match self {
737            GenericSourceConnection::Kafka(kafka) => {
738                GenericSourceConnection::Kafka(kafka.into_inline_connection(r))
739            }
740            GenericSourceConnection::Postgres(pg) => {
741                GenericSourceConnection::Postgres(pg.into_inline_connection(r))
742            }
743            GenericSourceConnection::MySql(mysql) => {
744                GenericSourceConnection::MySql(mysql.into_inline_connection(r))
745            }
746            GenericSourceConnection::SqlServer(sql_server) => {
747                GenericSourceConnection::SqlServer(sql_server.into_inline_connection(r))
748            }
749            GenericSourceConnection::LoadGenerator(lg) => {
750                GenericSourceConnection::LoadGenerator(lg)
751            }
752        }
753    }
754}
755
756impl<C: ConnectionAccess> SourceConnection for GenericSourceConnection<C> {
757    fn name(&self) -> &'static str {
758        match self {
759            Self::Kafka(conn) => conn.name(),
760            Self::Postgres(conn) => conn.name(),
761            Self::MySql(conn) => conn.name(),
762            Self::SqlServer(conn) => conn.name(),
763            Self::LoadGenerator(conn) => conn.name(),
764        }
765    }
766
767    fn external_reference(&self) -> Option<&str> {
768        match self {
769            Self::Kafka(conn) => conn.external_reference(),
770            Self::Postgres(conn) => conn.external_reference(),
771            Self::MySql(conn) => conn.external_reference(),
772            Self::SqlServer(conn) => conn.external_reference(),
773            Self::LoadGenerator(conn) => conn.external_reference(),
774        }
775    }
776
777    fn default_key_desc(&self) -> RelationDesc {
778        match self {
779            Self::Kafka(conn) => conn.default_key_desc(),
780            Self::Postgres(conn) => conn.default_key_desc(),
781            Self::MySql(conn) => conn.default_key_desc(),
782            Self::SqlServer(conn) => conn.default_key_desc(),
783            Self::LoadGenerator(conn) => conn.default_key_desc(),
784        }
785    }
786
787    fn default_value_desc(&self) -> RelationDesc {
788        match self {
789            Self::Kafka(conn) => conn.default_value_desc(),
790            Self::Postgres(conn) => conn.default_value_desc(),
791            Self::MySql(conn) => conn.default_value_desc(),
792            Self::SqlServer(conn) => conn.default_value_desc(),
793            Self::LoadGenerator(conn) => conn.default_value_desc(),
794        }
795    }
796
797    fn timestamp_desc(&self) -> RelationDesc {
798        match self {
799            Self::Kafka(conn) => conn.timestamp_desc(),
800            Self::Postgres(conn) => conn.timestamp_desc(),
801            Self::MySql(conn) => conn.timestamp_desc(),
802            Self::SqlServer(conn) => conn.timestamp_desc(),
803            Self::LoadGenerator(conn) => conn.timestamp_desc(),
804        }
805    }
806
807    fn connection_id(&self) -> Option<CatalogItemId> {
808        match self {
809            Self::Kafka(conn) => conn.connection_id(),
810            Self::Postgres(conn) => conn.connection_id(),
811            Self::MySql(conn) => conn.connection_id(),
812            Self::SqlServer(conn) => conn.connection_id(),
813            Self::LoadGenerator(conn) => conn.connection_id(),
814        }
815    }
816
817    fn supports_read_only(&self) -> bool {
818        match self {
819            GenericSourceConnection::Kafka(conn) => conn.supports_read_only(),
820            GenericSourceConnection::Postgres(conn) => conn.supports_read_only(),
821            GenericSourceConnection::MySql(conn) => conn.supports_read_only(),
822            GenericSourceConnection::SqlServer(conn) => conn.supports_read_only(),
823            GenericSourceConnection::LoadGenerator(conn) => conn.supports_read_only(),
824        }
825    }
826
827    fn prefers_single_replica(&self) -> bool {
828        match self {
829            GenericSourceConnection::Kafka(conn) => conn.prefers_single_replica(),
830            GenericSourceConnection::Postgres(conn) => conn.prefers_single_replica(),
831            GenericSourceConnection::MySql(conn) => conn.prefers_single_replica(),
832            GenericSourceConnection::SqlServer(conn) => conn.prefers_single_replica(),
833            GenericSourceConnection::LoadGenerator(conn) => conn.prefers_single_replica(),
834        }
835    }
836}
837impl<C: ConnectionAccess> crate::AlterCompatible for GenericSourceConnection<C> {
838    fn alter_compatible(&self, id: GlobalId, other: &Self) -> Result<(), AlterError> {
839        if self == other {
840            return Ok(());
841        }
842        let r = match (self, other) {
843            (Self::Kafka(conn), Self::Kafka(other)) => conn.alter_compatible(id, other),
844            (Self::Postgres(conn), Self::Postgres(other)) => conn.alter_compatible(id, other),
845            (Self::MySql(conn), Self::MySql(other)) => conn.alter_compatible(id, other),
846            (Self::SqlServer(conn), Self::SqlServer(other)) => conn.alter_compatible(id, other),
847            (Self::LoadGenerator(conn), Self::LoadGenerator(other)) => {
848                conn.alter_compatible(id, other)
849            }
850            _ => Err(AlterError { id }),
851        };
852
853        if r.is_err() {
854            tracing::warn!(
855                "GenericSourceConnection incompatible:\nself:\n{:#?}\n\nother\n{:#?}",
856                self,
857                other
858            );
859        }
860
861        r
862    }
863}
864
865/// Details necessary for each source export to allow the source implementations
866/// to export data to the export's collection.
867#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
868pub enum SourceExportDetails {
869    /// Used when the primary collection of a source isn't an export to
870    /// output to.
871    None,
872    Kafka(KafkaSourceExportDetails),
873    Postgres(PostgresSourceExportDetails),
874    MySql(MySqlSourceExportDetails),
875    SqlServer(SqlServerSourceExportDetails),
876    LoadGenerator(LoadGeneratorSourceExportDetails),
877}
878
879impl crate::AlterCompatible for SourceExportDetails {
880    fn alter_compatible(&self, id: GlobalId, other: &Self) -> Result<(), AlterError> {
881        if self == other {
882            return Ok(());
883        }
884        let r = match (self, other) {
885            (Self::None, Self::None) => Ok(()),
886            (Self::Kafka(s), Self::Kafka(o)) => s.alter_compatible(id, o),
887            (Self::Postgres(s), Self::Postgres(o)) => s.alter_compatible(id, o),
888            (Self::MySql(s), Self::MySql(o)) => s.alter_compatible(id, o),
889            (Self::SqlServer(s), Self::SqlServer(o)) => s.alter_compatible(id, o),
890            (Self::LoadGenerator(s), Self::LoadGenerator(o)) => s.alter_compatible(id, o),
891            _ => Err(AlterError { id }),
892        };
893
894        if r.is_err() {
895            tracing::warn!(
896                "SourceExportDetails incompatible:\nself:\n{:#?}\n\nother\n{:#?}",
897                self,
898                other
899            );
900        }
901
902        r
903    }
904}
905
906/// Details necessary to store in the `Details` option of a source export
907/// statement (`CREATE SUBSOURCE` and `CREATE TABLE .. FROM SOURCE` statements),
908/// to generate the appropriate `SourceExportDetails` struct during planning.
909/// NOTE that this is serialized as proto to the catalog, so any changes here
910/// must be backwards compatible or will require a migration.
911#[derive(Debug, Eq, PartialEq)]
912pub enum SourceExportStatementDetails {
913    Postgres {
914        table: mz_postgres_util::desc::PostgresTableDesc,
915        /// Whether the text-to-oid cast for this export accepts the full `u32`
916        /// range. Exports created before the cast was widened decode as
917        /// `false` and must keep the legacy `i32`-range cast forever, because
918        /// replication re-casts old tuples on delete and the persisted rows
919        /// were ingested under the legacy semantics.
920        cast_oid_full_range: bool,
921        /// An upper bound on the upstream WAL position whose schema `table`
922        /// describes. `None` for exports created before this was recorded.
923        initial_lsn: Option<MzOffset>,
924    },
925    MySql {
926        table: mz_mysql_util::MySqlTableDesc,
927        initial_gtid_set: String,
928        binlog_full_metadata: bool,
929    },
930    SqlServer {
931        table: mz_sql_server_util::desc::SqlServerTableDesc,
932        capture_instance: Arc<str>,
933        initial_lsn: mz_sql_server_util::cdc::Lsn,
934    },
935    LoadGenerator {
936        output: LoadGeneratorOutput,
937    },
938    Kafka {},
939}
940
941impl RustType<ProtoSourceExportStatementDetails> for SourceExportStatementDetails {
942    fn into_proto(&self) -> ProtoSourceExportStatementDetails {
943        match self {
944            SourceExportStatementDetails::Postgres {
945                table,
946                cast_oid_full_range,
947                initial_lsn,
948            } => ProtoSourceExportStatementDetails {
949                kind: Some(proto_source_export_statement_details::Kind::Postgres(
950                    postgres::ProtoPostgresSourceExportStatementDetails {
951                        table: Some(table.into_proto()),
952                        cast_oid_full_range: *cast_oid_full_range,
953                        initial_lsn: initial_lsn.map(|lsn| lsn.offset),
954                    },
955                )),
956            },
957            SourceExportStatementDetails::MySql {
958                table,
959                initial_gtid_set,
960                binlog_full_metadata,
961            } => ProtoSourceExportStatementDetails {
962                kind: Some(proto_source_export_statement_details::Kind::Mysql(
963                    mysql::ProtoMySqlSourceExportStatementDetails {
964                        table: Some(table.into_proto()),
965                        initial_gtid_set: initial_gtid_set.clone(),
966                        binlog_full_metadata: *binlog_full_metadata,
967                    },
968                )),
969            },
970            SourceExportStatementDetails::SqlServer {
971                table,
972                capture_instance,
973                initial_lsn,
974            } => ProtoSourceExportStatementDetails {
975                kind: Some(proto_source_export_statement_details::Kind::SqlServer(
976                    sql_server::ProtoSqlServerSourceExportStatementDetails {
977                        table: Some(table.into_proto()),
978                        capture_instance: capture_instance.to_string(),
979                        initial_lsn: initial_lsn.as_bytes().to_vec(),
980                    },
981                )),
982            },
983            SourceExportStatementDetails::LoadGenerator { output } => {
984                ProtoSourceExportStatementDetails {
985                    kind: Some(proto_source_export_statement_details::Kind::Loadgen(
986                        load_generator::ProtoLoadGeneratorSourceExportStatementDetails {
987                            output: output.into_proto().into(),
988                        },
989                    )),
990                }
991            }
992            SourceExportStatementDetails::Kafka {} => ProtoSourceExportStatementDetails {
993                kind: Some(proto_source_export_statement_details::Kind::Kafka(
994                    kafka::ProtoKafkaSourceExportStatementDetails {},
995                )),
996            },
997        }
998    }
999
1000    fn from_proto(proto: ProtoSourceExportStatementDetails) -> Result<Self, TryFromProtoError> {
1001        use proto_source_export_statement_details::Kind;
1002        Ok(match proto.kind {
1003            Some(Kind::Postgres(details)) => SourceExportStatementDetails::Postgres {
1004                table: details
1005                    .table
1006                    .into_rust_if_some("ProtoPostgresSourceExportStatementDetails::table")?,
1007                cast_oid_full_range: details.cast_oid_full_range,
1008                initial_lsn: details.initial_lsn.map(MzOffset::from),
1009            },
1010            Some(Kind::Mysql(details)) => SourceExportStatementDetails::MySql {
1011                table: details
1012                    .table
1013                    .into_rust_if_some("ProtoMySqlSourceExportStatementDetails::table")?,
1014
1015                initial_gtid_set: details.initial_gtid_set,
1016                binlog_full_metadata: details.binlog_full_metadata,
1017            },
1018            Some(Kind::SqlServer(details)) => SourceExportStatementDetails::SqlServer {
1019                table: details
1020                    .table
1021                    .into_rust_if_some("ProtoSqlServerSourceExportStatementDetails::table")?,
1022                capture_instance: details.capture_instance.into(),
1023                initial_lsn: mz_sql_server_util::cdc::Lsn::try_from(details.initial_lsn.as_slice())
1024                    .map_err(|e| TryFromProtoError::InvalidFieldError(e.to_string()))?,
1025            },
1026            Some(Kind::Loadgen(details)) => SourceExportStatementDetails::LoadGenerator {
1027                output: details
1028                    .output
1029                    .into_rust_if_some("ProtoLoadGeneratorSourceExportStatementDetails::output")?,
1030            },
1031            Some(Kind::Kafka(_details)) => SourceExportStatementDetails::Kafka {},
1032            None => {
1033                return Err(TryFromProtoError::missing_field(
1034                    "ProtoSourceExportStatementDetails::kind",
1035                ));
1036            }
1037        })
1038    }
1039}
1040
1041#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
1042#[repr(transparent)]
1043pub struct SourceData(pub Result<Row, DataflowError>);
1044
1045impl Default for SourceData {
1046    fn default() -> Self {
1047        SourceData(Ok(Row::default()))
1048    }
1049}
1050
1051impl Deref for SourceData {
1052    type Target = Result<Row, DataflowError>;
1053
1054    fn deref(&self) -> &Self::Target {
1055        &self.0
1056    }
1057}
1058
1059impl DerefMut for SourceData {
1060    fn deref_mut(&mut self) -> &mut Self::Target {
1061        &mut self.0
1062    }
1063}
1064
1065impl RustType<ProtoSourceData> for SourceData {
1066    fn into_proto(&self) -> ProtoSourceData {
1067        use proto_source_data::Kind;
1068        ProtoSourceData {
1069            kind: Some(match &**self {
1070                Ok(row) => Kind::Ok(row.into_proto()),
1071                Err(err) => Kind::Err(err.into_proto()),
1072            }),
1073        }
1074    }
1075
1076    fn from_proto(proto: ProtoSourceData) -> Result<Self, TryFromProtoError> {
1077        use proto_source_data::Kind;
1078        match proto.kind {
1079            Some(kind) => match kind {
1080                Kind::Ok(row) => Ok(SourceData(Ok(row.into_rust()?))),
1081                Kind::Err(err) => Ok(SourceData(Err(err.into_rust()?))),
1082            },
1083            None => Result::Err(TryFromProtoError::missing_field("ProtoSourceData::kind")),
1084        }
1085    }
1086}
1087
1088impl Codec for SourceData {
1089    type Storage = ProtoRow;
1090    type Schema = RelationDesc;
1091
1092    fn codec_name() -> String {
1093        "protobuf[SourceData]".into()
1094    }
1095
1096    fn encode<B: BufMut>(&self, buf: &mut B) {
1097        self.into_proto()
1098            .encode(buf)
1099            .expect("no required fields means no initialization errors");
1100    }
1101
1102    fn decode(buf: &[u8], schema: &RelationDesc) -> Result<Self, String> {
1103        let mut val = SourceData::default();
1104        <Self as Codec>::decode_from(&mut val, buf, &mut None, schema)?;
1105        Ok(val)
1106    }
1107
1108    fn decode_from<'a>(
1109        &mut self,
1110        buf: &'a [u8],
1111        storage: &mut Option<ProtoRow>,
1112        schema: &RelationDesc,
1113    ) -> Result<(), String> {
1114        // Optimize for common case of `Ok` by leaving a (cleared) `ProtoRow` in
1115        // the `Ok` variant of `ProtoSourceData`. prost's `Message::merge` impl
1116        // is smart about reusing the `Vec<Datum>` when it can.
1117        let mut proto = storage.take().unwrap_or_default();
1118        proto.clear();
1119        let mut proto = ProtoSourceData {
1120            kind: Some(proto_source_data::Kind::Ok(proto)),
1121        };
1122        proto.merge(buf).map_err(|err| err.to_string())?;
1123        match (proto.kind, &mut self.0) {
1124            // Again, optimize for the common case...
1125            (Some(proto_source_data::Kind::Ok(proto)), Ok(row)) => {
1126                let ret = row.decode_from_proto(&proto, schema);
1127                storage.replace(proto);
1128                ret
1129            }
1130            // ...otherwise fall back to the obvious thing.
1131            (kind, _) => {
1132                let proto = ProtoSourceData { kind };
1133                *self = proto.into_rust().map_err(|err| err.to_string())?;
1134                // Nothing to put back in storage.
1135                Ok(())
1136            }
1137        }
1138    }
1139
1140    fn validate(val: &Self, desc: &Self::Schema) -> Result<(), String> {
1141        match &val.0 {
1142            Ok(row) => Row::validate(row, desc),
1143            Err(_) => Ok(()),
1144        }
1145    }
1146
1147    fn encode_schema(schema: &Self::Schema) -> Bytes {
1148        schema.into_proto().encode_to_vec().into()
1149    }
1150
1151    fn decode_schema(buf: &Bytes) -> Self::Schema {
1152        let proto = ProtoRelationDesc::decode(buf.as_ref()).expect("valid schema");
1153        proto.into_rust().expect("valid schema")
1154    }
1155}
1156
1157/// Given a [`RelationDesc`] returns an arbitrary [`SourceData`].
1158#[cfg(any(test, feature = "proptest"))]
1159pub fn arb_source_data_for_relation_desc(
1160    desc: &RelationDesc,
1161) -> impl Strategy<Value = SourceData> + use<> {
1162    let row_strat = arb_row_for_relation(desc).no_shrink();
1163
1164    proptest::strategy::Union::new_weighted(vec![
1165        (50, row_strat.prop_map(|row| SourceData(Ok(row))).boxed()),
1166        (
1167            1,
1168            any::<DataflowError>()
1169                .prop_map(|err| SourceData(Err(err)))
1170                .no_shrink()
1171                .boxed(),
1172        ),
1173    ])
1174}
1175
1176/// Describes how external references should be organized in a multi-level
1177/// hierarchy.
1178///
1179/// For both PostgreSQL and MySQL sources, these levels of reference are
1180/// intrinsic to the items which we're referencing. If there are other naming
1181/// schemas for other types of sources we discover, we might need to revisit
1182/// this.
1183pub trait ExternalCatalogReference {
1184    /// The "second" level of namespacing for the reference.
1185    fn schema_name(&self) -> &str;
1186    /// The lowest level of namespacing for the reference.
1187    fn item_name(&self) -> &str;
1188}
1189
1190impl ExternalCatalogReference for &mz_mysql_util::MySqlTableDesc {
1191    fn schema_name(&self) -> &str {
1192        &self.schema_name
1193    }
1194
1195    fn item_name(&self) -> &str {
1196        &self.name
1197    }
1198}
1199
1200impl ExternalCatalogReference for mz_postgres_util::desc::PostgresTableDesc {
1201    fn schema_name(&self) -> &str {
1202        &self.namespace
1203    }
1204
1205    fn item_name(&self) -> &str {
1206        &self.name
1207    }
1208}
1209
1210impl ExternalCatalogReference for &mz_sql_server_util::desc::SqlServerTableDesc {
1211    fn schema_name(&self) -> &str {
1212        &*self.schema_name
1213    }
1214
1215    fn item_name(&self) -> &str {
1216        &*self.name
1217    }
1218}
1219
1220// This implementation provides a means of converting arbitrary objects into a
1221// `SubsourceCatalogReference`, e.g. load generator view names.
1222impl<'a> ExternalCatalogReference for (&'a str, &'a str) {
1223    fn schema_name(&self) -> &str {
1224        self.0
1225    }
1226
1227    fn item_name(&self) -> &str {
1228        self.1
1229    }
1230}
1231
1232/// Stores and resolves references to a `&[T: ExternalCatalogReference]`.
1233///
1234/// This is meant to provide an API to quickly look up a source's subsources.
1235///
1236/// For sources that do not provide any subsources, use the `Default`
1237/// implementation, which is empty and will not be able to resolve any
1238/// references.
1239#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
1240pub struct SourceReferenceResolver {
1241    inner: BTreeMap<Ident, BTreeMap<Ident, BTreeMap<Ident, usize>>>,
1242}
1243
1244#[derive(Debug, Clone, thiserror::Error)]
1245pub enum ExternalReferenceResolutionError {
1246    #[error("reference to {name} not found in source")]
1247    DoesNotExist { name: String },
1248    #[error(
1249        "reference {name} is ambiguous, consider specifying an additional \
1250    layer of qualification"
1251    )]
1252    Ambiguous { name: String },
1253    #[error("invalid identifier: {0}")]
1254    Ident(#[from] IdentError),
1255}
1256
1257impl<'a> SourceReferenceResolver {
1258    /// Constructs a new `SourceReferenceResolver` from a slice of `T:
1259    /// SubsourceCatalogReference`.
1260    ///
1261    /// # Errors
1262    /// - If any `&str` provided cannot be taken to an [`Ident`].
1263    pub fn new<T: ExternalCatalogReference>(
1264        database: &str,
1265        referenceable_items: &'a [T],
1266    ) -> Result<SourceReferenceResolver, ExternalReferenceResolutionError> {
1267        // An index from table name -> schema name -> database name -> index in
1268        // `referenceable_items`.
1269        let mut inner = BTreeMap::new();
1270
1271        let database = Ident::new(database)?;
1272
1273        for (reference_idx, item) in referenceable_items.iter().enumerate() {
1274            let item_name = Ident::new(item.item_name())?;
1275            let schema_name = Ident::new(item.schema_name())?;
1276
1277            inner
1278                .entry(item_name)
1279                .or_insert_with(BTreeMap::new)
1280                .entry(schema_name)
1281                .or_insert_with(BTreeMap::new)
1282                .entry(database.clone())
1283                .or_insert(reference_idx);
1284        }
1285
1286        Ok(SourceReferenceResolver { inner })
1287    }
1288
1289    /// Returns the canonical reference and index from which it originated in
1290    /// the `referenceable_items` provided to [`Self::new`].
1291    ///
1292    /// # Args
1293    /// - `name` is `&[Ident]` to let users provide the inner element of
1294    ///   [`UnresolvedItemName`].
1295    /// - `canonicalize_to_width` limits the number of elements in the returned
1296    ///   [`UnresolvedItemName`];this is useful if the source type requires
1297    ///   contriving database and schema names that a subsource should not
1298    ///   persist as its reference.
1299    ///
1300    /// # Errors
1301    /// - If `name` does not resolve to an item in `self.inner`.
1302    ///
1303    /// # Panics
1304    /// - If `canonicalize_to_width`` is not in `1..=3`.
1305    pub fn resolve(
1306        &self,
1307        name: &[Ident],
1308        canonicalize_to_width: usize,
1309    ) -> Result<(UnresolvedItemName, usize), ExternalReferenceResolutionError> {
1310        let (db, schema, idx) = self.resolve_inner(name)?;
1311
1312        let item = name.last().expect("must have provided at least 1 element");
1313
1314        let canonical_name = match canonicalize_to_width {
1315            1 => vec![item.clone()],
1316            2 => vec![schema.clone(), item.clone()],
1317            3 => vec![db.clone(), schema.clone(), item.clone()],
1318            o => panic!("canonicalize_to_width values must be 1..=3, but got {}", o),
1319        };
1320
1321        Ok((UnresolvedItemName(canonical_name), idx))
1322    }
1323
1324    /// Returns the index from which it originated in the `referenceable_items`
1325    /// provided to [`Self::new`].
1326    ///
1327    /// # Args
1328    /// `name` is `&[Ident]` to let users provide the inner element of
1329    /// [`UnresolvedItemName`].
1330    ///
1331    /// # Errors
1332    /// - If `name` does not resolve to an item in `self.inner`.
1333    pub fn resolve_idx(&self, name: &[Ident]) -> Result<usize, ExternalReferenceResolutionError> {
1334        let (_db, _schema, idx) = self.resolve_inner(name)?;
1335        Ok(idx)
1336    }
1337
1338    /// Returns the index from which it originated in the `referenceable_items`
1339    /// provided to [`Self::new`].
1340    ///
1341    /// # Args
1342    /// `name` is `&[Ident]` to let users provide the inner element of
1343    /// [`UnresolvedItemName`].
1344    ///
1345    /// # Return
1346    /// Returns a tuple whose elements are:
1347    /// 1. The "database"- or top-level namespace of the reference.
1348    /// 2. The "schema"- or second-level namespace of the reference.
1349    /// 3. The index to find the item in `referenceable_items` argument provided
1350    ///    to `SourceReferenceResolver::new`.
1351    ///
1352    /// # Errors
1353    /// - If `name` does not resolve to an item in `self.inner`.
1354    fn resolve_inner<'name: 'a>(
1355        &'a self,
1356        name: &'name [Ident],
1357    ) -> Result<(&'a Ident, &'a Ident, usize), ExternalReferenceResolutionError> {
1358        let get_provided_name = || UnresolvedItemName(name.to_vec()).to_string();
1359
1360        // Names must be composed of 1..=3 elements.
1361        if !(1..=3).contains(&name.len()) {
1362            Err(ExternalReferenceResolutionError::DoesNotExist {
1363                name: get_provided_name(),
1364            })?;
1365        }
1366
1367        // Fill on the leading elements with `None` if they aren't present.
1368        let mut names = std::iter::repeat(None)
1369            .take(3 - name.len())
1370            .chain(name.iter().map(Some));
1371
1372        let database = names.next().flatten();
1373        let schema = names.next().flatten();
1374        let item = names
1375            .next()
1376            .flatten()
1377            .expect("must have provided the item name");
1378
1379        assert_none!(names.next(), "expected a 3-element iterator");
1380
1381        let schemas =
1382            self.inner
1383                .get(item)
1384                .ok_or_else(|| ExternalReferenceResolutionError::DoesNotExist {
1385                    name: get_provided_name(),
1386                })?;
1387
1388        let schema = match schema {
1389            Some(schema) => schema,
1390            None => schemas.keys().exactly_one().map_err(|_e| {
1391                ExternalReferenceResolutionError::Ambiguous {
1392                    name: get_provided_name(),
1393                }
1394            })?,
1395        };
1396
1397        let databases =
1398            schemas
1399                .get(schema)
1400                .ok_or_else(|| ExternalReferenceResolutionError::DoesNotExist {
1401                    name: get_provided_name(),
1402                })?;
1403
1404        let database = match database {
1405            Some(database) => database,
1406            None => databases.keys().exactly_one().map_err(|_e| {
1407                ExternalReferenceResolutionError::Ambiguous {
1408                    name: get_provided_name(),
1409                }
1410            })?,
1411        };
1412
1413        let reference_idx = databases.get(database).ok_or_else(|| {
1414            ExternalReferenceResolutionError::DoesNotExist {
1415                name: get_provided_name(),
1416            }
1417        })?;
1418
1419        Ok((database, schema, *reference_idx))
1420    }
1421}
1422
1423/// A decoder for [`Row`]s within [`SourceData`].
1424///
1425/// This type exists as a wrapper around [`RowColumnarDecoder`] to handle the
1426/// case where the [`RelationDesc`] we're encoding with has no columns. See
1427/// [`SourceDataRowColumnarEncoder`] for more details.
1428#[derive(Debug)]
1429pub enum SourceDataRowColumnarDecoder {
1430    Row(RowColumnarDecoder),
1431    EmptyRow,
1432}
1433
1434impl SourceDataRowColumnarDecoder {
1435    pub fn decode(&self, idx: usize, row: &mut Row) {
1436        match self {
1437            SourceDataRowColumnarDecoder::Row(decoder) => decoder.decode(idx, row),
1438            SourceDataRowColumnarDecoder::EmptyRow => {
1439                // Create a packer just to clear the Row.
1440                row.packer();
1441            }
1442        }
1443    }
1444
1445    pub fn goodbytes(&self) -> usize {
1446        match self {
1447            SourceDataRowColumnarDecoder::Row(decoder) => decoder.goodbytes(),
1448            SourceDataRowColumnarDecoder::EmptyRow => 0,
1449        }
1450    }
1451}
1452
1453#[derive(Debug)]
1454pub struct SourceDataColumnarDecoder {
1455    row_decoder: SourceDataRowColumnarDecoder,
1456    err_decoder: BinaryArray,
1457}
1458
1459impl SourceDataColumnarDecoder {
1460    pub fn new(col: StructArray, desc: &RelationDesc) -> Result<Self, anyhow::Error> {
1461        // TODO(parkmcar): We should validate the fields here.
1462        let (_fields, arrays, nullability) = col.into_parts();
1463
1464        if nullability.is_some() {
1465            anyhow::bail!("SourceData is not nullable, but found {nullability:?}");
1466        }
1467        if arrays.len() != 2 {
1468            anyhow::bail!("SourceData should only have two fields, found {arrays:?}");
1469        }
1470
1471        let errs = arrays[1]
1472            .as_any()
1473            .downcast_ref::<BinaryArray>()
1474            .ok_or_else(|| anyhow::anyhow!("expected BinaryArray, found {:?}", arrays[1]))?;
1475
1476        let row_decoder = match arrays[0].data_type() {
1477            arrow::datatypes::DataType::Struct(_) => {
1478                let rows = arrays[0]
1479                    .as_any()
1480                    .downcast_ref::<StructArray>()
1481                    .ok_or_else(|| {
1482                        anyhow::anyhow!("expected StructArray, found {:?}", arrays[0])
1483                    })?;
1484                let decoder = RowColumnarDecoder::new(rows.clone(), desc)?;
1485                SourceDataRowColumnarDecoder::Row(decoder)
1486            }
1487            arrow::datatypes::DataType::Null => SourceDataRowColumnarDecoder::EmptyRow,
1488            other => anyhow::bail!("expected Struct or Null Array, found {other:?}"),
1489        };
1490
1491        Ok(SourceDataColumnarDecoder {
1492            row_decoder,
1493            err_decoder: errs.clone(),
1494        })
1495    }
1496}
1497
1498impl ColumnDecoder<SourceData> for SourceDataColumnarDecoder {
1499    fn decode(&self, idx: usize, val: &mut SourceData) {
1500        let err_null = self.err_decoder.is_null(idx);
1501        let row_null = match &self.row_decoder {
1502            SourceDataRowColumnarDecoder::Row(decoder) => decoder.is_null(idx),
1503            SourceDataRowColumnarDecoder::EmptyRow => !err_null,
1504        };
1505
1506        match (row_null, err_null) {
1507            (true, false) => {
1508                let err = self.err_decoder.value(idx);
1509                let err = ProtoDataflowError::decode(err)
1510                    .expect("proto should be valid")
1511                    .into_rust()
1512                    .expect("error should be valid");
1513                val.0 = Err(err);
1514            }
1515            (false, true) => {
1516                let row = match val.0.as_mut() {
1517                    Ok(row) => row,
1518                    Err(_) => {
1519                        val.0 = Ok(Row::default());
1520                        val.0.as_mut().unwrap()
1521                    }
1522                };
1523                self.row_decoder.decode(idx, row);
1524            }
1525            (true, true) => panic!("should have one of 'ok' or 'err'"),
1526            (false, false) => panic!("cannot have both 'ok' and 'err'"),
1527        }
1528    }
1529
1530    fn is_null(&self, idx: usize) -> bool {
1531        let err_null = self.err_decoder.is_null(idx);
1532        let row_null = match &self.row_decoder {
1533            SourceDataRowColumnarDecoder::Row(decoder) => decoder.is_null(idx),
1534            SourceDataRowColumnarDecoder::EmptyRow => !err_null,
1535        };
1536        assert!(!err_null || !row_null, "SourceData should never be null!");
1537
1538        false
1539    }
1540
1541    fn goodbytes(&self) -> usize {
1542        self.row_decoder.goodbytes() + ArrayOrd::Binary(self.err_decoder.clone()).goodbytes()
1543    }
1544
1545    fn stats(&self) -> StructStats {
1546        let len = self.err_decoder.len();
1547        let err_stats = ColumnarStats {
1548            nulls: Some(ColumnNullStats {
1549                count: self.err_decoder.null_count(),
1550            }),
1551            values: PrimitiveStats::<Vec<u8>>::from_column(&self.err_decoder).into(),
1552        };
1553        // The top level struct is non-nullable and every entry is either an
1554        // `Ok(Row)` or an `Err(String)`. As a result, we can compute the number
1555        // of `Ok` entries by subtracting the number of `Err` entries from the
1556        // total count.
1557        let row_null_count = len - self.err_decoder.null_count();
1558        let row_stats = match &self.row_decoder {
1559            SourceDataRowColumnarDecoder::Row(encoder) => {
1560                // Sanity check that the number of row nulls/nones we calculated
1561                // using the error column matches what the row column thinks it
1562                // has.
1563                assert_eq!(encoder.null_count(), row_null_count);
1564                encoder.stats()
1565            }
1566            SourceDataRowColumnarDecoder::EmptyRow => StructStats {
1567                len,
1568                cols: BTreeMap::default(),
1569            },
1570        };
1571        let row_stats = ColumnarStats {
1572            nulls: Some(ColumnNullStats {
1573                count: row_null_count,
1574            }),
1575            values: ColumnStatKinds::Struct(row_stats),
1576        };
1577
1578        let stats = [
1579            (
1580                SourceDataColumnarEncoder::OK_COLUMN_NAME.to_string(),
1581                row_stats,
1582            ),
1583            (
1584                SourceDataColumnarEncoder::ERR_COLUMN_NAME.to_string(),
1585                err_stats,
1586            ),
1587        ];
1588        StructStats {
1589            len,
1590            cols: stats.into_iter().map(|(name, s)| (name, s)).collect(),
1591        }
1592    }
1593}
1594
1595/// An encoder for [`Row`]s within [`SourceData`].
1596///
1597/// This type exists as a wrapper around [`RowColumnarEncoder`] to support
1598/// encoding empty [`Row`]s. A [`RowColumnarEncoder`] finishes as a
1599/// [`StructArray`] which is required to have at least one column, and thus
1600/// cannot support empty [`Row`]s.
1601#[derive(Debug)]
1602pub enum SourceDataRowColumnarEncoder {
1603    Row(RowColumnarEncoder),
1604    EmptyRow,
1605}
1606
1607impl SourceDataRowColumnarEncoder {
1608    pub(crate) fn goodbytes(&self) -> usize {
1609        match self {
1610            SourceDataRowColumnarEncoder::Row(e) => e.goodbytes(),
1611            SourceDataRowColumnarEncoder::EmptyRow => 0,
1612        }
1613    }
1614
1615    pub fn append(&mut self, row: &Row) {
1616        match self {
1617            SourceDataRowColumnarEncoder::Row(encoder) => encoder.append(row),
1618            SourceDataRowColumnarEncoder::EmptyRow => {
1619                assert_eq!(row.iter().count(), 0)
1620            }
1621        }
1622    }
1623
1624    pub fn append_null(&mut self) {
1625        match self {
1626            SourceDataRowColumnarEncoder::Row(encoder) => encoder.append_null(),
1627            SourceDataRowColumnarEncoder::EmptyRow => (),
1628        }
1629    }
1630}
1631
1632#[derive(Debug)]
1633pub struct SourceDataColumnarEncoder {
1634    row_encoder: SourceDataRowColumnarEncoder,
1635    err_encoder: BinaryBuilder,
1636}
1637
1638impl SourceDataColumnarEncoder {
1639    const OK_COLUMN_NAME: &'static str = "ok";
1640    const ERR_COLUMN_NAME: &'static str = "err";
1641
1642    pub fn new(desc: &RelationDesc) -> Self {
1643        let row_encoder = match RowColumnarEncoder::new(desc) {
1644            Some(encoder) => SourceDataRowColumnarEncoder::Row(encoder),
1645            None => {
1646                assert!(desc.typ().columns().is_empty());
1647                SourceDataRowColumnarEncoder::EmptyRow
1648            }
1649        };
1650        let err_encoder = BinaryBuilder::new();
1651
1652        SourceDataColumnarEncoder {
1653            row_encoder,
1654            err_encoder,
1655        }
1656    }
1657}
1658
1659impl ColumnEncoder<SourceData> for SourceDataColumnarEncoder {
1660    type FinishedColumn = StructArray;
1661
1662    fn goodbytes(&self) -> usize {
1663        self.row_encoder.goodbytes() + self.err_encoder.values_slice().len()
1664    }
1665
1666    #[inline]
1667    fn append(&mut self, val: &SourceData) {
1668        match val.0.as_ref() {
1669            Ok(row) => {
1670                self.row_encoder.append(row);
1671                self.err_encoder.append_null();
1672            }
1673            Err(err) => {
1674                self.row_encoder.append_null();
1675                self.err_encoder
1676                    .append_value(err.into_proto().encode_to_vec());
1677            }
1678        }
1679    }
1680
1681    #[inline]
1682    fn append_null(&mut self) {
1683        panic!("appending a null into SourceDataColumnarEncoder is not supported");
1684    }
1685
1686    fn finish(self) -> Self::FinishedColumn {
1687        let SourceDataColumnarEncoder {
1688            row_encoder,
1689            mut err_encoder,
1690        } = self;
1691
1692        let err_column = BinaryBuilder::finish(&mut err_encoder);
1693        let row_column: ArrayRef = match row_encoder {
1694            SourceDataRowColumnarEncoder::Row(encoder) => {
1695                let column = encoder.finish();
1696                Arc::new(column)
1697            }
1698            SourceDataRowColumnarEncoder::EmptyRow => Arc::new(NullArray::new(err_column.len())),
1699        };
1700
1701        assert_eq!(row_column.len(), err_column.len());
1702
1703        let fields = vec![
1704            Field::new(Self::OK_COLUMN_NAME, row_column.data_type().clone(), true),
1705            Field::new(Self::ERR_COLUMN_NAME, err_column.data_type().clone(), true),
1706        ];
1707        let arrays: Vec<Arc<dyn Array>> = vec![row_column, Arc::new(err_column)];
1708        StructArray::new(Fields::from(fields), arrays, None)
1709    }
1710}
1711
1712impl Schema<SourceData> for RelationDesc {
1713    type ArrowColumn = StructArray;
1714    type Statistics = StructStats;
1715
1716    type Decoder = SourceDataColumnarDecoder;
1717    type Encoder = SourceDataColumnarEncoder;
1718
1719    fn decoder(&self, col: Self::ArrowColumn) -> Result<Self::Decoder, anyhow::Error> {
1720        SourceDataColumnarDecoder::new(col, self)
1721    }
1722
1723    fn encoder(&self) -> Result<Self::Encoder, anyhow::Error> {
1724        Ok(SourceDataColumnarEncoder::new(self))
1725    }
1726}
1727
1728#[cfg(test)]
1729mod tests {
1730    use arrow::array::{ArrayData, make_comparator};
1731    use base64::Engine;
1732    use bytes::Bytes;
1733    use mz_expr::EvalError;
1734    use mz_ore::assert_err;
1735    use mz_ore::metrics::MetricsRegistry;
1736    use mz_persist::indexed::columnar::arrow::{realloc_any, realloc_array};
1737    use mz_persist::metrics::ColumnarMetrics;
1738    use mz_persist_types::parquet::EncodingConfig;
1739    use mz_persist_types::schema::{Migration, backward_compatible};
1740    use mz_persist_types::stats::{PartStats, PartStatsMetrics};
1741    use mz_repr::{
1742        ColumnIndex, DatumVec, PropRelationDescDiff, ProtoRelationDesc, RelationDescBuilder,
1743        RowArena, SqlScalarType, arb_relation_desc_diff, arb_relation_desc_projection,
1744    };
1745    use proptest::prelude::*;
1746    use proptest::strategy::{Union, ValueTree};
1747
1748    use crate::stats::RelationPartStats;
1749
1750    use super::*;
1751
1752    #[mz_ore::test]
1753    fn test_timeline_parsing() {
1754        assert_eq!(Ok(Timeline::EpochMilliseconds), "M".parse());
1755        assert_eq!(Ok(Timeline::External("JOE".to_string())), "E.JOE".parse());
1756        assert_eq!(Ok(Timeline::User("MIKE".to_string())), "U.MIKE".parse());
1757
1758        assert_err!("Materialize".parse::<Timeline>());
1759        assert_err!("Ejoe".parse::<Timeline>());
1760        assert_err!("Umike".parse::<Timeline>());
1761        assert_err!("Dance".parse::<Timeline>());
1762        assert_err!("".parse::<Timeline>());
1763    }
1764
1765    #[track_caller]
1766    fn roundtrip_source_data(
1767        desc: &RelationDesc,
1768        datas: Vec<SourceData>,
1769        read_desc: &RelationDesc,
1770        config: &EncodingConfig,
1771    ) {
1772        let metrics = ColumnarMetrics::disconnected();
1773        let mut encoder = <RelationDesc as Schema<SourceData>>::encoder(desc).unwrap();
1774        for data in &datas {
1775            encoder.append(data);
1776        }
1777        let col = encoder.finish();
1778
1779        // The top-level StructArray for SourceData should always be non-nullable.
1780        assert!(!col.is_nullable());
1781
1782        // Reallocate our arrays with lgalloc.
1783        let col = realloc_array(&col, &metrics);
1784
1785        // Roundtrip through ProtoArray format.
1786        {
1787            let proto = col.to_data().into_proto();
1788            let bytes = proto.encode_to_vec();
1789            let proto = mz_persist_types::arrow::ProtoArrayData::decode(&bytes[..]).unwrap();
1790            let array_data: ArrayData = proto.into_rust().unwrap();
1791
1792            let col_rnd = StructArray::from(array_data.clone());
1793            assert_eq!(col, col_rnd);
1794
1795            let col_dyn = arrow::array::make_array(array_data);
1796            let col_dyn = col_dyn.as_any().downcast_ref::<StructArray>().unwrap();
1797            assert_eq!(&col, col_dyn);
1798        }
1799
1800        // Encode to Parquet.
1801        let mut buf = Vec::new();
1802        let fields = Fields::from(vec![Field::new("k", col.data_type().clone(), false)]);
1803        let arrays: Vec<Arc<dyn Array>> = vec![Arc::new(col.clone())];
1804        mz_persist_types::parquet::encode_arrays(&mut buf, fields, arrays, config).unwrap();
1805
1806        // Decode from Parquet.
1807        let buf = Bytes::from(buf);
1808        let mut reader = mz_persist_types::parquet::decode_arrays(buf).unwrap();
1809        let maybe_batch = reader.next();
1810
1811        // If we didn't encode any data then our record_batch will be empty.
1812        let Some(record_batch) = maybe_batch else {
1813            assert!(datas.is_empty());
1814            return;
1815        };
1816        let record_batch = record_batch.unwrap();
1817
1818        assert_eq!(record_batch.columns().len(), 1);
1819        let rnd_col = &record_batch.columns()[0];
1820        let rnd_col = realloc_any(Arc::clone(rnd_col), &metrics);
1821        let rnd_col = rnd_col
1822            .as_any()
1823            .downcast_ref::<StructArray>()
1824            .unwrap()
1825            .clone();
1826
1827        // Try generating stats for the data, just to make sure we don't panic.
1828        let stats = <RelationDesc as Schema<SourceData>>::decoder_any(desc, &rnd_col)
1829            .expect("valid decoder")
1830            .stats();
1831
1832        // Read back all of our data and assert it roundtrips.
1833        let mut rnd_data = SourceData(Ok(Row::default()));
1834        let decoder = <RelationDesc as Schema<SourceData>>::decoder(desc, rnd_col.clone()).unwrap();
1835        for (idx, og_data) in datas.iter().enumerate() {
1836            decoder.decode(idx, &mut rnd_data);
1837            assert_eq!(og_data, &rnd_data);
1838        }
1839
1840        // Read back all of our data a second time with a projection applied, and make sure the
1841        // stats are valid.
1842        let stats_metrics = PartStatsMetrics::new(&MetricsRegistry::new());
1843        let stats = RelationPartStats {
1844            name: "test",
1845            metrics: &stats_metrics,
1846            stats: &PartStats { key: stats },
1847            desc: read_desc,
1848        };
1849        let mut datum_vec = DatumVec::new();
1850        let arena = RowArena::default();
1851        let decoder = <RelationDesc as Schema<SourceData>>::decoder(read_desc, rnd_col).unwrap();
1852
1853        for (idx, og_data) in datas.iter().enumerate() {
1854            decoder.decode(idx, &mut rnd_data);
1855            match (&og_data.0, &rnd_data.0) {
1856                (Ok(og_row), Ok(rnd_row)) => {
1857                    // Filter down to just the Datums in the projection schema.
1858                    {
1859                        let datums = datum_vec.borrow_with(og_row);
1860                        let projected_datums =
1861                            datums.iter().enumerate().filter_map(|(idx, datum)| {
1862                                read_desc
1863                                    .contains_index(&ColumnIndex::from_raw(idx))
1864                                    .then_some(datum)
1865                            });
1866                        let og_projected_row = Row::pack(projected_datums);
1867                        assert_eq!(&og_projected_row, rnd_row);
1868                    }
1869
1870                    // Validate the stats for all of our projected columns.
1871                    {
1872                        let proj_datums = datum_vec.borrow_with(rnd_row);
1873                        for (pos, (idx, _, _)) in read_desc.iter_all().enumerate() {
1874                            let spec = stats.col_stats(idx, &arena);
1875                            assert!(spec.may_contain(proj_datums[pos]));
1876                        }
1877                    }
1878                }
1879                (Err(_), Err(_)) => assert_eq!(og_data, &rnd_data),
1880                (_, _) => panic!("decoded to a different type? {og_data:?} {rnd_data:?}"),
1881            }
1882        }
1883
1884        // Verify that the RelationDesc itself roundtrips through
1885        // {encode,decode}_schema.
1886        let encoded_schema = SourceData::encode_schema(desc);
1887        let roundtrip_desc = SourceData::decode_schema(&encoded_schema);
1888        assert_eq!(desc, &roundtrip_desc);
1889
1890        // Verify that the RelationDesc is backward compatible with itself (this
1891        // mostly checks for `unimplemented!` type panics).
1892        let migration =
1893            mz_persist_types::schema::backward_compatible(col.data_type(), col.data_type());
1894        let migration = migration.expect("should be backward compatible with self");
1895        // Also verify that the Fn doesn't do anything wonky.
1896        let migrated = migration.migrate(Arc::new(col.clone()));
1897        assert_eq!(col.data_type(), migrated.data_type());
1898    }
1899
1900    #[mz_ore::test]
1901    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
1902    fn all_source_data_roundtrips() {
1903        let mut weights = vec![(500, Just(0..8)), (50, Just(8..32))];
1904        if std::env::var("PROPTEST_LARGE_DATA").is_ok() {
1905            weights.extend([
1906                (10, Just(32..128)),
1907                (5, Just(128..512)),
1908                (3, Just(512..2048)),
1909                (1, Just(2048..8192)),
1910            ]);
1911        }
1912        let num_rows = Union::new_weighted(weights);
1913
1914        // TODO(parkmycar): There are so many clones going on here, and maybe we can avoid them?
1915        let strat = (any::<RelationDesc>(), num_rows)
1916            .prop_flat_map(|(desc, num_rows)| {
1917                arb_relation_desc_projection(desc.clone())
1918                    .prop_map(move |read_desc| (desc.clone(), read_desc, num_rows.clone()))
1919            })
1920            .prop_flat_map(|(desc, read_desc, num_rows)| {
1921                proptest::collection::vec(arb_source_data_for_relation_desc(&desc), num_rows)
1922                    .prop_map(move |datas| (desc.clone(), datas, read_desc.clone()))
1923            });
1924
1925        let combined_strat = (any::<EncodingConfig>(), strat);
1926        proptest!(|((config, (desc, source_datas, read_desc)) in combined_strat)| {
1927            roundtrip_source_data(&desc, source_datas, &read_desc, &config);
1928        });
1929    }
1930
1931    #[mz_ore::test]
1932    fn roundtrip_error_nulls() {
1933        let desc = RelationDescBuilder::default()
1934            .with_column(
1935                "ts",
1936                SqlScalarType::TimestampTz { precision: None }.nullable(false),
1937            )
1938            .finish();
1939        let source_datas = vec![SourceData(Err(DataflowError::EvalError(
1940            EvalError::DateOutOfRange.into(),
1941        )))];
1942        let config = EncodingConfig::default();
1943        roundtrip_source_data(&desc, source_datas, &desc, &config);
1944    }
1945
1946    fn is_sorted(array: &dyn Array) -> bool {
1947        let sort_options = arrow::compute::SortOptions::default();
1948        let Ok(cmp) = make_comparator(array, array, sort_options) else {
1949            // TODO: arrow v51.0.0 doesn't support comparing structs. When
1950            // we migrate to v52+, the `build_compare` function is
1951            // deprecated and replaced by `make_comparator`, which does
1952            // support structs. At which point, this will work (and we
1953            // should switch this early return to an expect, if possible).
1954            return false;
1955        };
1956        (0..array.len())
1957            .tuple_windows()
1958            .all(|(i, j)| cmp(i, j).is_le())
1959    }
1960
1961    fn get_data_type(schema: &impl Schema<SourceData>) -> arrow::datatypes::DataType {
1962        use mz_persist_types::columnar::ColumnEncoder;
1963        let array = Schema::encoder(schema).expect("valid schema").finish();
1964        Array::data_type(&array).clone()
1965    }
1966
1967    #[track_caller]
1968    fn backward_compatible_testcase(
1969        old: &RelationDesc,
1970        new: &RelationDesc,
1971        migration: Migration,
1972        datas: &[SourceData],
1973    ) {
1974        let mut encoder = Schema::<SourceData>::encoder(old).expect("valid schema");
1975        for data in datas {
1976            encoder.append(data);
1977        }
1978        let old = encoder.finish();
1979        let new = Schema::<SourceData>::encoder(new)
1980            .expect("valid schema")
1981            .finish();
1982        let old: Arc<dyn Array> = Arc::new(old);
1983        let new: Arc<dyn Array> = Arc::new(new);
1984        let migrated = migration.migrate(Arc::clone(&old));
1985        assert_eq!(migrated.data_type(), new.data_type());
1986
1987        // Check the sortedness preservation, if we can.
1988        if migration.preserves_order() && is_sorted(&old) {
1989            assert!(is_sorted(&new))
1990        }
1991    }
1992
1993    #[mz_ore::test]
1994    fn backward_compatible_empty_add_column() {
1995        let old = RelationDesc::empty();
1996        let new = RelationDesc::from_names_and_types([("a", SqlScalarType::Bool.nullable(true))]);
1997
1998        let old_data_type = get_data_type(&old);
1999        let new_data_type = get_data_type(&new);
2000
2001        let migration = backward_compatible(&old_data_type, &new_data_type);
2002        assert!(migration.is_some());
2003    }
2004
2005    #[mz_ore::test]
2006    fn backward_compatible_project_away_all() {
2007        let old = RelationDesc::from_names_and_types([("a", SqlScalarType::Bool.nullable(true))]);
2008        let new = RelationDesc::empty();
2009
2010        let old_data_type = get_data_type(&old);
2011        let new_data_type = get_data_type(&new);
2012
2013        let migration = backward_compatible(&old_data_type, &new_data_type);
2014        assert!(migration.is_some());
2015    }
2016
2017    #[mz_ore::test]
2018    #[cfg_attr(miri, ignore)]
2019    fn backward_compatible_migrate() {
2020        let strat = (any::<RelationDesc>(), any::<RelationDesc>()).prop_flat_map(|(old, new)| {
2021            proptest::collection::vec(arb_source_data_for_relation_desc(&old), 2)
2022                .prop_map(move |datas| (old.clone(), new.clone(), datas))
2023        });
2024
2025        proptest!(|((old, new, datas) in strat)| {
2026            let old_data_type = get_data_type(&old);
2027            let new_data_type = get_data_type(&new);
2028
2029            if let Some(migration) = backward_compatible(&old_data_type, &new_data_type) {
2030                backward_compatible_testcase(&old, &new, migration, &datas);
2031            };
2032        });
2033    }
2034
2035    #[mz_ore::test]
2036    #[cfg_attr(miri, ignore)]
2037    fn backward_compatible_migrate_from_common() {
2038        use mz_repr::SqlColumnType;
2039        fn test_case(old: RelationDesc, diffs: Vec<PropRelationDescDiff>, datas: Vec<SourceData>) {
2040            // TODO(parkmycar): As we iterate on schema migrations more things should become compatible.
2041            let should_be_compatible = diffs.iter().all(|diff| match diff {
2042                // We only support adding nullable columns.
2043                PropRelationDescDiff::AddColumn {
2044                    typ: SqlColumnType { nullable, .. },
2045                    ..
2046                } => *nullable,
2047                PropRelationDescDiff::DropColumn { .. } => true,
2048                _ => false,
2049            });
2050
2051            let mut new = old.clone();
2052            for diff in diffs.into_iter() {
2053                diff.apply(&mut new)
2054            }
2055
2056            let old_data_type = get_data_type(&old);
2057            let new_data_type = get_data_type(&new);
2058
2059            if let Some(migration) = backward_compatible(&old_data_type, &new_data_type) {
2060                backward_compatible_testcase(&old, &new, migration, &datas);
2061            } else if should_be_compatible {
2062                panic!("new DataType was not compatible when it should have been!");
2063            }
2064        }
2065
2066        let strat = any::<RelationDesc>()
2067            .prop_flat_map(|desc| {
2068                proptest::collection::vec(arb_source_data_for_relation_desc(&desc), 2)
2069                    .no_shrink()
2070                    .prop_map(move |datas| (desc.clone(), datas))
2071            })
2072            .prop_flat_map(|(desc, datas)| {
2073                arb_relation_desc_diff(&desc)
2074                    .prop_map(move |diffs| (desc.clone(), diffs, datas.clone()))
2075            });
2076
2077        proptest!(|((old, diffs, datas) in strat)| {
2078            test_case(old, diffs, datas);
2079        });
2080    }
2081
2082    #[mz_ore::test]
2083    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
2084    fn empty_relation_desc_roundtrips() {
2085        let empty = RelationDesc::empty();
2086        let rows = proptest::collection::vec(arb_source_data_for_relation_desc(&empty), 0..8)
2087            .prop_map(move |datas| (empty.clone(), datas));
2088
2089        // Note: This case should be covered by the `all_source_data_roundtrips` test above, but
2090        // it's a special case that we explicitly want to exercise.
2091        proptest!(|((config, (desc, source_datas)) in (any::<EncodingConfig>(), rows))| {
2092            roundtrip_source_data(&desc, source_datas, &desc, &config);
2093        });
2094    }
2095
2096    #[mz_ore::test]
2097    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
2098    fn arrow_datatype_consistent() {
2099        fn test_case(desc: RelationDesc, datas: Vec<SourceData>) {
2100            let half = datas.len() / 2;
2101
2102            let mut encoder_a = <RelationDesc as Schema<SourceData>>::encoder(&desc).unwrap();
2103            for data in &datas[..half] {
2104                encoder_a.append(data);
2105            }
2106            let col_a = encoder_a.finish();
2107
2108            let mut encoder_b = <RelationDesc as Schema<SourceData>>::encoder(&desc).unwrap();
2109            for data in &datas[half..] {
2110                encoder_b.append(data);
2111            }
2112            let col_b = encoder_b.finish();
2113
2114            // The DataType of the resulting column should not change based on what data was
2115            // encoded.
2116            assert_eq!(col_a.data_type(), col_b.data_type());
2117        }
2118
2119        let num_rows = 12;
2120        let strat = any::<RelationDesc>().prop_flat_map(|desc| {
2121            proptest::collection::vec(arb_source_data_for_relation_desc(&desc), num_rows)
2122                .prop_map(move |datas| (desc.clone(), datas))
2123        });
2124
2125        proptest!(|((desc, data) in strat)| {
2126            test_case(desc, data);
2127        });
2128    }
2129
2130    #[mz_ore::test]
2131    #[cfg_attr(miri, ignore)] // too slow
2132    fn source_proto_serialization_stability() {
2133        let min_protos = 10;
2134        let encoded = include_str!("snapshots/source-datas.txt");
2135
2136        // Decode the pre-generated source datas
2137        let mut decoded: Vec<(RelationDesc, SourceData)> = encoded
2138            .lines()
2139            .map(|s| {
2140                let (desc, data) = s.split_once(',').expect("comma separated data");
2141                let desc = base64::engine::general_purpose::STANDARD
2142                    .decode(desc)
2143                    .expect("valid base64");
2144                let data = base64::engine::general_purpose::STANDARD
2145                    .decode(data)
2146                    .expect("valid base64");
2147                (desc, data)
2148            })
2149            .map(|(desc, data)| {
2150                let desc = ProtoRelationDesc::decode(&desc[..]).expect("valid proto");
2151                let desc = desc.into_rust().expect("valid proto");
2152                let data = SourceData::decode(&data, &desc).expect("valid proto");
2153                (desc, data)
2154            })
2155            .collect();
2156
2157        // If there are fewer than the minimum examples, generate some new ones arbitrarily
2158        let mut runner = proptest::test_runner::TestRunner::deterministic();
2159        let strategy = RelationDesc::arbitrary().prop_flat_map(|desc| {
2160            arb_source_data_for_relation_desc(&desc).prop_map(move |data| (desc.clone(), data))
2161        });
2162        while decoded.len() < min_protos {
2163            let arbitrary_data = strategy
2164                .new_tree(&mut runner)
2165                .expect("source data")
2166                .current();
2167            decoded.push(arbitrary_data);
2168        }
2169
2170        // Reencode and compare the strings
2171        let mut reencoded = String::new();
2172        let mut buf = vec![];
2173        for (desc, data) in decoded {
2174            buf.clear();
2175            desc.into_proto().encode(&mut buf).expect("success");
2176            base64::engine::general_purpose::STANDARD.encode_string(buf.as_slice(), &mut reencoded);
2177            reencoded.push(',');
2178
2179            buf.clear();
2180            data.encode(&mut buf);
2181            base64::engine::general_purpose::STANDARD.encode_string(buf.as_slice(), &mut reencoded);
2182            reencoded.push('\n');
2183        }
2184
2185        // Optimizations in Persist, particularly consolidation on read,
2186        // depend on a stable serialization for the serialized data.
2187        // For example, reordering proto fields could cause us
2188        // to generate a different (equivalent) serialization for a record,
2189        // and the two versions would not consolidate out.
2190        // This can impact correctness!
2191        //
2192        // If you need to change how SourceDatas are encoded, that's still fine...
2193        // but we'll also need to increase
2194        // the MINIMUM_CONSOLIDATED_VERSION as part of the same release.
2195        assert_eq!(
2196            encoded,
2197            reencoded.as_str(),
2198            "SourceData serde should be stable"
2199        )
2200    }
2201}