1use 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#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq)]
84pub struct IngestionDescription<S: 'static = (), C: ConnectionAccess = InlinedConnection> {
85 pub desc: SourceDesc<C>,
87 pub source_exports: BTreeMap<GlobalId, SourceExport<S>>,
103 pub instance_id: StorageInstanceId,
105 pub remap_collection_id: GlobalId,
107 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 pub fn collection_ids(&self) -> impl Iterator<Item = GlobalId> + '_ {
133 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 pub storage_metadata: S,
251 pub details: SourceExportDetails,
253 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#[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
338impl 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
352impl 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
382impl 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#[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,
453 External(String),
457 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
513pub trait SourceConnection: Debug + Clone + PartialEq + AlterCompatible {
515 fn name(&self) -> &'static str;
517
518 fn external_reference(&self) -> Option<&str>;
520
521 fn default_key_desc(&self) -> RelationDesc;
525
526 fn default_value_desc(&self) -> RelationDesc;
530
531 fn timestamp_desc(&self) -> RelationDesc;
534
535 fn connection_id(&self) -> Option<CatalogItemId>;
538
539 fn supports_read_only(&self) -> bool;
541
542 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#[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 pub fn monotonic(&self, connection: &GenericSourceConnection<C>) -> bool {
614 match &self.envelope {
615 SourceEnvelope::Upsert(_) | SourceEnvelope::CdcV2 => false,
617 SourceEnvelope::None(_) => {
618 match connection {
619 GenericSourceConnection::Postgres(_) => false,
621 GenericSourceConnection::MySql(_) => false,
623 GenericSourceConnection::SqlServer(_) => false,
625 GenericSourceConnection::LoadGenerator(g) => g.load_generator.is_monotonic(),
627 GenericSourceConnection::Kafka(_) => true,
629 }
630 }
631 }
632 }
633}
634
635#[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 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: _,
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#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
868pub enum SourceExportDetails {
869 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#[derive(Debug, Eq, PartialEq)]
912pub enum SourceExportStatementDetails {
913 Postgres {
914 table: mz_postgres_util::desc::PostgresTableDesc,
915 cast_oid_full_range: bool,
921 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 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 (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 (kind, _) => {
1132 let proto = ProtoSourceData { kind };
1133 *self = proto.into_rust().map_err(|err| err.to_string())?;
1134 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#[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
1176pub trait ExternalCatalogReference {
1184 fn schema_name(&self) -> &str;
1186 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
1220impl<'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#[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 pub fn new<T: ExternalCatalogReference>(
1264 database: &str,
1265 referenceable_items: &'a [T],
1266 ) -> Result<SourceReferenceResolver, ExternalReferenceResolutionError> {
1267 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 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 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 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 if !(1..=3).contains(&name.len()) {
1362 Err(ExternalReferenceResolutionError::DoesNotExist {
1363 name: get_provided_name(),
1364 })?;
1365 }
1366
1367 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#[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 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 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 let row_null_count = len - self.err_decoder.null_count();
1558 let row_stats = match &self.row_decoder {
1559 SourceDataRowColumnarDecoder::Row(encoder) => {
1560 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#[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 assert!(!col.is_nullable());
1781
1782 let col = realloc_array(&col, &metrics);
1784
1785 {
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 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 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 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 let stats = <RelationDesc as Schema<SourceData>>::decoder_any(desc, &rnd_col)
1829 .expect("valid decoder")
1830 .stats();
1831
1832 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 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 {
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 {
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 let encoded_schema = SourceData::encode_schema(desc);
1887 let roundtrip_desc = SourceData::decode_schema(&encoded_schema);
1888 assert_eq!(desc, &roundtrip_desc);
1889
1890 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 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)] 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 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 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 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 let should_be_compatible = diffs.iter().all(|diff| match diff {
2042 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)] 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 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)] 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 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)] fn source_proto_serialization_stability() {
2133 let min_protos = 10;
2134 let encoded = include_str!("snapshots/source-datas.txt");
2135
2136 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 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 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 assert_eq!(
2196 encoded,
2197 reencoded.as_str(),
2198 "SourceData serde should be stable"
2199 )
2200 }
2201}