1use std::collections::{BTreeMap, BTreeSet};
15use std::fmt;
16use std::iter;
17use std::path::Path;
18use std::sync::Arc;
19
20use anyhow::anyhow;
21use itertools::Itertools;
22use mz_adapter_types::dyncfgs::ENABLE_S3_TABLES_REGION_CHECK;
23use mz_ccsr::{Client, GetBySubjectError};
24use mz_cloud_provider::CloudProvider;
25use mz_controller_types::ClusterId;
26use mz_kafka_util::client::MzClientContext;
27use mz_mysql_util::MySqlTableDesc;
28use mz_ore::collections::CollectionExt;
29use mz_ore::error::ErrorExt;
30use mz_ore::future::InTask;
31use mz_ore::iter::IteratorExt;
32use mz_ore::str::StrExt;
33use mz_postgres_util::desc::PostgresTableDesc;
34use mz_proto::RustType;
35use mz_repr::{CatalogItemId, RelationDesc, RelationVersionSelector, Timestamp, strconv};
36use mz_sql_parser::ast::display::AstDisplay;
37use mz_sql_parser::ast::visit::{Visit, visit_function};
38use mz_sql_parser::ast::visit_mut::{VisitMut, visit_expr_mut};
39use mz_sql_parser::ast::{
40 AlterSourceAction, AlterSourceAddSubsourceOptionName, AlterSourceStatement, AvroDocOn,
41 ColumnName, CreateMaterializedViewStatement, CreateSinkConnection, CreateSinkOptionName,
42 CreateSinkStatement, CreateSourceOptionName, CreateSubsourceOption, CreateSubsourceOptionName,
43 CreateTableFromSourceStatement, CsrConfigOption, CsrConfigOptionName, CsrConnection,
44 CsrSeedAvro, CsrSeedProtobuf, CsrSeedProtobufSchema, DeferredItemName, DocOnIdentifier,
45 DocOnSchema, Expr, Function, FunctionArgs, GlueAvroOption, GlueAvroSeed, Ident,
46 KafkaSourceConfigOption, KafkaSourceConfigOptionName, LoadGenerator, LoadGeneratorOption,
47 LoadGeneratorOptionName, MaterializedViewOption, MaterializedViewOptionName, MySqlConfigOption,
48 MySqlConfigOptionName, PgConfigOption, PgConfigOptionName, RawItemName,
49 ReaderSchemaSelectionStrategy, RefreshAtOptionValue, RefreshEveryOptionValue,
50 RefreshOptionValue, SourceEnvelope, SqlServerConfigOption, SqlServerConfigOptionName,
51 Statement, TableFromSourceColumns, TableFromSourceOption, TableFromSourceOptionName,
52 UnresolvedItemName,
53};
54use mz_sql_server_util::desc::SqlServerTableDesc;
55use mz_storage_types::configuration::StorageConfiguration;
56use mz_storage_types::connections::Connection;
57use mz_storage_types::connections::inline::IntoInlineConnection;
58use mz_storage_types::errors::ContextCreationError;
59use mz_storage_types::sources::load_generator::LoadGeneratorOutput;
60use mz_storage_types::sources::mysql::MySqlSourceDetails;
61use mz_storage_types::sources::postgres::PostgresSourcePublicationDetails;
62use mz_storage_types::sources::{
63 GenericSourceConnection, MzOffset, SourceConnection, SourceDesc, SourceExportStatementDetails,
64 SqlServerSourceExtras,
65};
66use prost::Message;
67use protobuf_native::MessageLite;
68use protobuf_native::compiler::{SourceTreeDescriptorDatabase, VirtualSourceTree};
69use rdkafka::admin::AdminClient;
70use references::{RetrievedSourceReferences, SourceReferenceClient};
71use uuid::Uuid;
72
73use crate::ast::{
74 AlterSourceAddSubsourceOption, AvroSchema, CreateSourceConnection, CreateSourceStatement,
75 CreateSubsourceStatement, CsrConnectionAvro, CsrConnectionProtobuf, ExternalReferenceExport,
76 ExternalReferences, Format, FormatSpecifier, ProtobufSchema, Value, WithOptionValue,
77};
78use crate::catalog::{CatalogItemType, SessionCatalog};
79use crate::kafka_util::{KafkaSinkConfigOptionExtracted, KafkaSourceConfigOptionExtracted};
80use crate::names::{
81 Aug, FullItemName, PartialItemName, ResolvedColumnReference, ResolvedDataType, ResolvedIds,
82 ResolvedItemName,
83};
84use crate::plan::error::PlanError;
85use crate::plan::statement::ddl::load_generator_ast_to_generator;
86use crate::plan::{SourceReferences, StatementContext};
87use crate::pure::error::{IcebergSinkPurificationError, SqlServerSourcePurificationError};
88use crate::pure::mysql::{ensure_binlog_full_metadata, is_binlog_full_metadata};
89use crate::{kafka_util, normalize};
90
91use self::error::{
92 CsrPurificationError, KafkaSinkPurificationError, KafkaSourcePurificationError,
93 LoadGeneratorSourcePurificationError, MySqlSourcePurificationError, PgSourcePurificationError,
94};
95
96pub(crate) mod error;
97mod references;
98
99pub mod mysql;
100pub mod postgres;
101pub mod sql_server;
102
103pub(crate) struct RequestedSourceExport<T> {
104 external_reference: UnresolvedItemName,
105 name: UnresolvedItemName,
106 meta: T,
107}
108
109impl<T: fmt::Debug> fmt::Debug for RequestedSourceExport<T> {
110 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
111 f.debug_struct("RequestedSourceExport")
112 .field("external_reference", &self.external_reference)
113 .field("name", &self.name)
114 .field("meta", &self.meta)
115 .finish()
116 }
117}
118
119impl<T> RequestedSourceExport<T> {
120 fn change_meta<F>(self, new_meta: F) -> RequestedSourceExport<F> {
121 RequestedSourceExport {
122 external_reference: self.external_reference,
123 name: self.name,
124 meta: new_meta,
125 }
126 }
127}
128
129fn source_export_name_gen(
134 source_name: &UnresolvedItemName,
135 subsource_name: &str,
136) -> Result<UnresolvedItemName, PlanError> {
137 let mut partial = normalize::unresolved_item_name(source_name.clone())?;
138 partial.item = subsource_name.to_string();
139 Ok(UnresolvedItemName::from(partial))
140}
141
142fn validate_source_export_names<T>(
146 requested_source_exports: &[RequestedSourceExport<T>],
147) -> Result<(), PlanError> {
148 if let Some(name) = requested_source_exports
152 .iter()
153 .map(|subsource| &subsource.name)
154 .duplicates()
155 .next()
156 .cloned()
157 {
158 let mut upstream_references: Vec<_> = requested_source_exports
159 .into_iter()
160 .filter_map(|subsource| {
161 if &subsource.name == &name {
162 Some(subsource.external_reference.clone())
163 } else {
164 None
165 }
166 })
167 .collect();
168
169 upstream_references.sort();
170
171 Err(PlanError::SubsourceNameConflict {
172 name,
173 upstream_references,
174 })?;
175 }
176
177 if let Some(name) = requested_source_exports
183 .iter()
184 .map(|export| &export.external_reference)
185 .duplicates()
186 .next()
187 .cloned()
188 {
189 let mut target_names: Vec<_> = requested_source_exports
190 .into_iter()
191 .filter_map(|export| {
192 if &export.external_reference == &name {
193 Some(export.name.clone())
194 } else {
195 None
196 }
197 })
198 .collect();
199
200 target_names.sort();
201
202 Err(PlanError::SubsourceDuplicateReference { name, target_names })?;
203 }
204
205 Ok(())
206}
207
208#[derive(Debug, Clone, PartialEq, Eq)]
209pub enum PurifiedStatement {
210 PurifiedCreateSource {
211 create_progress_subsource_stmt: Option<CreateSubsourceStatement<Aug>>,
213 create_source_stmt: CreateSourceStatement<Aug>,
214 subsources: BTreeMap<UnresolvedItemName, PurifiedSourceExport>,
216 available_source_references: SourceReferences,
219 },
220 PurifiedAlterSource {
221 alter_source_stmt: AlterSourceStatement<Aug>,
222 },
223 PurifiedAlterSourceAddSubsources {
224 source_name: ResolvedItemName,
226 options: Vec<AlterSourceAddSubsourceOption<Aug>>,
229 subsources: BTreeMap<UnresolvedItemName, PurifiedSourceExport>,
231 },
232 PurifiedAlterSourceRefreshReferences {
233 source_name: ResolvedItemName,
234 available_source_references: SourceReferences,
236 },
237 PurifiedCreateSink(CreateSinkStatement<Aug>),
238 PurifiedCreateTableFromSource {
239 stmt: CreateTableFromSourceStatement<Aug>,
240 },
241}
242
243#[derive(Debug, Clone, PartialEq, Eq)]
244pub struct PurifiedSourceExport {
245 pub external_reference: UnresolvedItemName,
246 pub details: PurifiedExportDetails,
247}
248
249#[derive(Debug, Clone, PartialEq, Eq)]
250pub enum PurifiedExportDetails {
251 MySql {
252 table: MySqlTableDesc,
253 text_columns: Option<Vec<Ident>>,
254 exclude_columns: Option<Vec<Ident>>,
255 initial_gtid_set: String,
256 binlog_full_metadata: bool,
257 },
258 Postgres {
259 table: PostgresTableDesc,
260 text_columns: Option<Vec<Ident>>,
261 exclude_columns: Option<Vec<Ident>>,
262 initial_lsn: MzOffset,
263 },
264 SqlServer {
265 table: SqlServerTableDesc,
266 text_columns: Option<Vec<Ident>>,
267 excl_columns: Option<Vec<Ident>>,
268 capture_instance: Arc<str>,
269 initial_lsn: mz_sql_server_util::cdc::Lsn,
270 },
271 Kafka {},
272 LoadGenerator {
273 table: Option<RelationDesc>,
274 output: LoadGeneratorOutput,
275 },
276}
277
278#[derive(Debug, Clone, Copy)]
284pub enum StatementSource {
285 Altered(CatalogItemId),
286 Read(CatalogItemId),
287}
288
289impl StatementSource {
290 pub fn id(&self) -> CatalogItemId {
291 match self {
292 StatementSource::Altered(id) | StatementSource::Read(id) => *id,
293 }
294 }
295}
296
297pub fn statement_source(
307 catalog: &impl SessionCatalog,
308 stmt: &Statement<Aug>,
309) -> Option<StatementSource> {
310 let scx = StatementContext::new(None, catalog);
311 match stmt {
312 Statement::AlterSource(stmt) => {
313 let item = scx
314 .resolve_item(RawItemName::Name(stmt.source_name.clone()))
315 .ok()?;
316 (item.item_type() == CatalogItemType::Source)
317 .then(|| StatementSource::Altered(item.id()))
318 }
319 Statement::CreateTableFromSource(stmt) => {
320 let item = scx.get_item_by_resolved_name(&stmt.source).ok()?;
321 (item.item_type() == CatalogItemType::Source).then(|| StatementSource::Read(item.id()))
322 }
323 _ => None,
324 }
325}
326
327pub async fn purify_statement(
337 catalog: impl SessionCatalog,
338 now: u64,
339 stmt: Statement<Aug>,
340 storage_configuration: &StorageConfiguration,
341) -> (Result<PurifiedStatement, PlanError>, Option<ClusterId>) {
342 match stmt {
343 Statement::CreateSource(stmt) => {
344 let cluster_id = stmt.in_cluster.as_ref().map(|cluster| cluster.id.clone());
345 (
346 purify_create_source(catalog, now, stmt, storage_configuration).await,
347 cluster_id,
348 )
349 }
350 Statement::AlterSource(stmt) => (
351 purify_alter_source(catalog, stmt, storage_configuration).await,
352 None,
353 ),
354 Statement::CreateSink(stmt) => {
355 let cluster_id = stmt.in_cluster.as_ref().map(|cluster| cluster.id.clone());
356 (
357 purify_create_sink(catalog, stmt, storage_configuration).await,
358 cluster_id,
359 )
360 }
361 Statement::CreateTableFromSource(stmt) => (
362 purify_create_table_from_source(catalog, stmt, storage_configuration).await,
363 None,
364 ),
365 o => (
366 Err(internal_err!(
367 "unexpected statement type in purification: {:?}",
368 o
369 )),
370 None,
371 ),
372 }
373}
374
375pub(crate) fn purify_create_sink_avro_doc_on_options(
380 catalog: &dyn SessionCatalog,
381 from_id: CatalogItemId,
382 format: &mut Option<FormatSpecifier<Aug>>,
383) -> Result<(), PlanError> {
384 let from = catalog.get_item(&from_id);
386 let object_ids = from
387 .references()
388 .items()
389 .copied()
390 .chain_one(from.id())
391 .collect::<Vec<_>>();
392
393 let mut avro_format_options = vec![];
396 for_each_format(format, |doc_on_schema, fmt| match fmt {
397 Format::Avro(AvroSchema::InlineSchema { .. })
398 | Format::Avro(AvroSchema::Glue { .. })
399 | Format::Bytes
400 | Format::Csv { .. }
401 | Format::Json { .. }
402 | Format::Protobuf(..)
403 | Format::Regex(..)
404 | Format::Text => (),
405 Format::Avro(AvroSchema::Csr {
406 csr_connection: CsrConnectionAvro { connection, .. },
407 }) => {
408 avro_format_options.push((doc_on_schema, &mut connection.options));
409 }
410 });
411
412 for (for_schema, options) in avro_format_options {
415 let user_provided_comments = options
416 .iter()
417 .filter_map(|CsrConfigOption { name, .. }| match name {
418 CsrConfigOptionName::AvroDocOn(doc_on) => Some(doc_on.clone()),
419 _ => None,
420 })
421 .collect::<BTreeSet<_>>();
422
423 for object_id in &object_ids {
425 let item = catalog
427 .get_item(object_id)
428 .at_version(RelationVersionSelector::Latest);
429 let full_resolved_name = ResolvedItemName::Item {
430 id: *object_id,
431 qualifiers: item.name().qualifiers.clone(),
432 full_name: catalog.resolve_full_name(item.name()),
433 print_id: !matches!(item.item_type(), CatalogItemType::Func),
434 version: RelationVersionSelector::Latest,
435 };
436
437 if let Some(comments_map) = catalog.get_item_comments(object_id) {
438 let doc_on_item_key = AvroDocOn {
441 identifier: DocOnIdentifier::Type(full_resolved_name.clone()),
442 for_schema,
443 };
444 if !user_provided_comments.contains(&doc_on_item_key) {
445 if let Some(root_comment) = comments_map.get(&None) {
446 options.push(CsrConfigOption {
447 name: CsrConfigOptionName::AvroDocOn(doc_on_item_key),
448 value: Some(mz_sql_parser::ast::WithOptionValue::Value(Value::String(
449 root_comment.clone(),
450 ))),
451 });
452 }
453 }
454
455 let column_descs = match item.type_details() {
459 Some(details) => details.typ.desc(catalog).unwrap_or_default(),
460 None => item.relation_desc().map(|d| d.into_owned()),
461 };
462
463 if let Some(desc) = column_descs {
464 for (pos, column_name) in desc.iter_names().enumerate() {
465 if let Some(comment_str) = comments_map.get(&Some(pos + 1)) {
466 let doc_on_column_key = AvroDocOn {
467 identifier: DocOnIdentifier::Column(ColumnName {
468 relation: full_resolved_name.clone(),
469 column: ResolvedColumnReference::Column {
470 name: column_name.to_owned(),
471 index: pos,
472 },
473 }),
474 for_schema,
475 };
476 if !user_provided_comments.contains(&doc_on_column_key) {
477 options.push(CsrConfigOption {
478 name: CsrConfigOptionName::AvroDocOn(doc_on_column_key),
479 value: Some(mz_sql_parser::ast::WithOptionValue::Value(
480 Value::String(comment_str.clone()),
481 )),
482 });
483 }
484 }
485 }
486 }
487 }
488 }
489 }
490
491 Ok(())
492}
493
494async fn purify_create_sink(
497 catalog: impl SessionCatalog,
498 mut create_sink_stmt: CreateSinkStatement<Aug>,
499 storage_configuration: &StorageConfiguration,
500) -> Result<PurifiedStatement, PlanError> {
501 let CreateSinkStatement {
503 connection,
504 format,
505 with_options,
506 name: _,
507 in_cluster: _,
508 if_not_exists: _,
509 from,
510 envelope: _,
511 mode: _,
512 } = &mut create_sink_stmt;
513
514 const USER_ALLOWED_WITH_OPTIONS: &[CreateSinkOptionName] = &[
516 CreateSinkOptionName::Snapshot,
517 CreateSinkOptionName::CommitInterval,
518 ];
519
520 if let Some(op) = with_options
521 .iter()
522 .find(|op| !USER_ALLOWED_WITH_OPTIONS.contains(&op.name))
523 {
524 sql_bail!(
525 "CREATE SINK...WITH ({}..) is not allowed",
526 op.name.to_ast_string_simple(),
527 )
528 }
529
530 match &connection {
531 CreateSinkConnection::Kafka {
532 connection,
533 options,
534 key: _,
535 headers: _,
536 } => {
537 let scx = StatementContext::new(None, &catalog);
542 let connection = {
543 let item = scx.get_item_by_resolved_name(connection)?;
544 match item.connection()? {
546 Connection::Kafka(connection) => {
547 connection.clone().into_inline_connection(scx.catalog)
548 }
549 _ => sql_bail!(
550 "{} is not a kafka connection",
551 scx.catalog.resolve_full_name(item.name())
552 ),
553 }
554 };
555
556 let extracted_options: KafkaSinkConfigOptionExtracted = options.clone().try_into()?;
557
558 if extracted_options.legacy_ids == Some(true) {
559 sql_bail!("LEGACY IDs option is not supported");
560 }
561
562 let client: AdminClient<_> = connection
563 .create_with_context(
564 storage_configuration,
565 MzClientContext::default(),
566 &BTreeMap::new(),
567 InTask::No,
568 )
569 .await
570 .map_err(|e| {
571 KafkaSinkPurificationError::AdminClientError(Arc::new(e))
573 })?;
574
575 let metadata = client
576 .inner()
577 .fetch_metadata(
578 None,
579 storage_configuration
580 .parameters
581 .kafka_timeout_config
582 .fetch_metadata_timeout,
583 )
584 .map_err(|e| {
585 KafkaSinkPurificationError::AdminClientError(Arc::new(
586 ContextCreationError::KafkaError(e),
587 ))
588 })?;
589
590 if metadata.brokers().len() == 0 {
591 Err(KafkaSinkPurificationError::ZeroBrokers)?;
592 }
593 }
594 CreateSinkConnection::Iceberg {
595 catalog_connection,
596 aws_connection,
597 ..
598 } => {
599 let scx = StatementContext::new(None, &catalog);
600 let connection = {
601 let item = scx.get_item_by_resolved_name(catalog_connection)?;
602 match item.connection()? {
604 Connection::IcebergCatalog(connection) => {
605 connection.clone().into_inline_connection(scx.catalog)
606 }
607 _ => sql_bail!(
608 "{} is not an iceberg connection",
609 scx.catalog.resolve_full_name(item.name())
610 ),
611 }
612 };
613
614 if let Some(s3tables) = connection.s3tables_catalog() {
618 let enable_region_check =
619 ENABLE_S3_TABLES_REGION_CHECK.get(scx.catalog.system_vars().dyncfgs());
620 if enable_region_check {
621 let env_id = &catalog.config().environment_id;
622 if matches!(env_id.cloud_provider(), CloudProvider::Aws) {
623 let env_region = env_id.cloud_provider_region();
624 let s3_tables_region = s3tables
627 .aws_connection
628 .connection
629 .region
630 .clone()
631 .unwrap_or_else(|| "us-east-1".to_string());
632 if s3_tables_region != env_region {
633 Err(IcebergSinkPurificationError::S3TablesRegionMismatch {
634 s3_tables_region,
635 environment_region: env_region.to_string(),
636 })?;
637 }
638 }
639 }
640 }
641
642 if let Some(aws_connection) = aws_connection {
648 let aws_conn_id = aws_connection.item_id();
649 let aws_connection = {
650 let item = scx.get_item_by_resolved_name(aws_connection)?;
651 match item.connection()? {
653 Connection::Aws(aws_connection) => aws_connection.clone(),
654 _ => sql_bail!(
655 "{} is not an aws connection",
656 scx.catalog.resolve_full_name(item.name())
657 ),
658 }
659 };
660
661 let _sdk_config = aws_connection
662 .load_sdk_config(
663 &storage_configuration.connection_context,
664 aws_conn_id.clone(),
665 InTask::No,
666 mz_storage_types::dyncfgs::ENFORCE_EXTERNAL_ADDRESSES
667 .get(storage_configuration.config_set()),
668 )
669 .await
670 .map_err(|e| IcebergSinkPurificationError::AwsSdkContextError(Arc::new(e)))?;
671 }
672
673 let _catalog = connection
679 .connect(storage_configuration, InTask::No, None)
680 .await
681 .map_err(|e| IcebergSinkPurificationError::CatalogError(Arc::new(e)))?;
682 }
683 }
684
685 let mut csr_connection_ids = BTreeSet::new();
686 for_each_format(format, |_, fmt| match fmt {
687 Format::Avro(AvroSchema::InlineSchema { .. })
688 | Format::Avro(AvroSchema::Glue { .. })
689 | Format::Bytes
690 | Format::Csv { .. }
691 | Format::Json { .. }
692 | Format::Protobuf(ProtobufSchema::InlineSchema { .. })
693 | Format::Regex(..)
694 | Format::Text => (),
695 Format::Avro(AvroSchema::Csr {
696 csr_connection: CsrConnectionAvro { connection, .. },
697 })
698 | Format::Protobuf(ProtobufSchema::Csr {
699 csr_connection: CsrConnectionProtobuf { connection, .. },
700 }) => {
701 csr_connection_ids.insert(*connection.connection.item_id());
702 }
703 });
704
705 let scx = StatementContext::new(None, &catalog);
706 for csr_connection_id in csr_connection_ids {
707 let connection = {
708 let item = scx.get_item(&csr_connection_id);
709 match item.connection()? {
711 Connection::Csr(connection) => connection.clone().into_inline_connection(&catalog),
712 _ => Err(CsrPurificationError::NotCsrConnection(
713 scx.catalog.resolve_full_name(item.name()),
714 ))?,
715 }
716 };
717
718 let client = connection
719 .connect(storage_configuration, InTask::No)
720 .await
721 .map_err(|e| CsrPurificationError::ClientError(Arc::new(e)))?;
722
723 client
724 .list_subjects()
725 .await
726 .map_err(|e| CsrPurificationError::ListSubjectsError(Arc::new(e)))?;
727 }
728
729 purify_create_sink_avro_doc_on_options(&catalog, *from.item_id(), format)?;
730
731 Ok(PurifiedStatement::PurifiedCreateSink(create_sink_stmt))
732}
733
734fn for_each_format<'a, F>(format: &'a mut Option<FormatSpecifier<Aug>>, mut f: F)
741where
742 F: FnMut(DocOnSchema, &'a mut Format<Aug>),
743{
744 match format {
745 None => (),
746 Some(FormatSpecifier::Bare(fmt)) => f(DocOnSchema::All, fmt),
747 Some(FormatSpecifier::KeyValue { key, value }) => {
748 f(DocOnSchema::KeyOnly, key);
749 f(DocOnSchema::ValueOnly, value);
750 }
751 }
752}
753
754#[derive(Debug, Copy, Clone, PartialEq, Eq)]
757pub(crate) enum SourceReferencePolicy {
758 NotAllowed,
761 Optional,
764 Required,
767}
768
769async fn purify_create_source(
770 catalog: impl SessionCatalog,
771 now: u64,
772 mut create_source_stmt: CreateSourceStatement<Aug>,
773 storage_configuration: &StorageConfiguration,
774) -> Result<PurifiedStatement, PlanError> {
775 let CreateSourceStatement {
776 name: source_name,
777 col_names,
778 key_constraint,
779 connection: source_connection,
780 format,
781 envelope,
782 include_metadata,
783 external_references,
784 progress_subsource,
785 with_options,
786 ..
787 } = &mut create_source_stmt;
788
789 let uses_old_syntax = !col_names.is_empty()
790 || key_constraint.is_some()
791 || format.is_some()
792 || envelope.is_some()
793 || !include_metadata.is_empty()
794 || external_references.is_some()
795 || progress_subsource.is_some();
796
797 if let Some(DeferredItemName::Named(_)) = progress_subsource {
798 sql_bail!("Cannot manually ID qualify progress subsource")
799 }
800
801 let mut requested_subsource_map = BTreeMap::new();
802
803 let progress_desc = match &source_connection {
804 CreateSourceConnection::Kafka { .. } => {
805 &mz_storage_types::sources::kafka::KAFKA_PROGRESS_DESC
806 }
807 CreateSourceConnection::Postgres { .. } => {
808 &mz_storage_types::sources::postgres::PG_PROGRESS_DESC
809 }
810 CreateSourceConnection::SqlServer { .. } => {
811 &mz_storage_types::sources::sql_server::SQL_SERVER_PROGRESS_DESC
812 }
813 CreateSourceConnection::MySql { .. } => {
814 &mz_storage_types::sources::mysql::MYSQL_PROGRESS_DESC
815 }
816 CreateSourceConnection::LoadGenerator { .. } => {
817 &mz_storage_types::sources::load_generator::LOAD_GEN_PROGRESS_DESC
818 }
819 };
820 let scx = StatementContext::new(None, &catalog);
821
822 let reference_policy = if scx.catalog.system_vars().enable_create_table_from_source()
826 && scx.catalog.system_vars().force_source_table_syntax()
827 {
828 SourceReferencePolicy::NotAllowed
829 } else if scx.catalog.system_vars().enable_create_table_from_source() {
830 SourceReferencePolicy::Optional
831 } else {
832 SourceReferencePolicy::Required
833 };
834
835 let mut format_options = SourceFormatOptions::Default;
836
837 let retrieved_source_references: RetrievedSourceReferences;
838
839 match source_connection {
840 CreateSourceConnection::Kafka {
841 connection,
842 options: base_with_options,
843 ..
844 } => {
845 if let Some(external_references) = external_references {
846 Err(KafkaSourcePurificationError::ReferencedSubsources(
847 external_references.clone(),
848 ))?;
849 }
850
851 let connection = {
852 let item = scx.get_item_by_resolved_name(connection)?;
853 match item.connection()? {
855 Connection::Kafka(connection) => {
856 connection.clone().into_inline_connection(&catalog)
857 }
858 _ => Err(KafkaSourcePurificationError::NotKafkaConnection(
859 scx.catalog.resolve_full_name(item.name()),
860 ))?,
861 }
862 };
863
864 let extracted_options: KafkaSourceConfigOptionExtracted =
865 base_with_options.clone().try_into()?;
866
867 let topic = extracted_options
868 .topic
869 .ok_or(KafkaSourcePurificationError::ConnectionMissingTopic)?;
870
871 let consumer = connection
872 .create_with_context(
873 storage_configuration,
874 MzClientContext::default(),
875 &BTreeMap::new(),
876 InTask::No,
877 )
878 .await
879 .map_err(|e| {
880 KafkaSourcePurificationError::KafkaConsumerError(
882 e.display_with_causes().to_string(),
883 )
884 })?;
885 let consumer = Arc::new(consumer);
886
887 match (
888 extracted_options.start_offset,
889 extracted_options.start_timestamp,
890 ) {
891 (None, None) => {
892 kafka_util::ensure_topic_exists(
894 Arc::clone(&consumer),
895 &topic,
896 storage_configuration
897 .parameters
898 .kafka_timeout_config
899 .fetch_metadata_timeout,
900 )
901 .await?;
902 }
903 (Some(_), Some(_)) => {
904 sql_bail!("cannot specify START TIMESTAMP and START OFFSET at same time")
905 }
906 (Some(start_offsets), None) => {
907 kafka_util::validate_start_offsets(
909 Arc::clone(&consumer),
910 &topic,
911 start_offsets,
912 storage_configuration
913 .parameters
914 .kafka_timeout_config
915 .fetch_metadata_timeout,
916 )
917 .await?;
918 }
919 (None, Some(time_offset)) => {
920 let start_offsets = kafka_util::lookup_start_offsets(
922 Arc::clone(&consumer),
923 &topic,
924 time_offset,
925 now,
926 storage_configuration
927 .parameters
928 .kafka_timeout_config
929 .fetch_metadata_timeout,
930 )
931 .await?;
932
933 base_with_options.retain(|val| {
934 !matches!(val.name, KafkaSourceConfigOptionName::StartTimestamp)
935 });
936 base_with_options.push(KafkaSourceConfigOption {
937 name: KafkaSourceConfigOptionName::StartOffset,
938 value: Some(WithOptionValue::Sequence(
939 start_offsets
940 .iter()
941 .map(|offset| {
942 WithOptionValue::Value(Value::Number(offset.to_string()))
943 })
944 .collect(),
945 )),
946 });
947 }
948 }
949
950 let reference_client = SourceReferenceClient::Kafka { topic: &topic };
951 retrieved_source_references = reference_client.get_source_references().await?;
952
953 format_options = SourceFormatOptions::Kafka { topic };
954 }
955 CreateSourceConnection::Postgres {
956 connection,
957 options,
958 } => {
959 let connection_item = scx.get_item_by_resolved_name(connection)?;
960 let connection = match connection_item.connection().map_err(PlanError::from)? {
961 Connection::Postgres(connection) => {
962 connection.clone().into_inline_connection(&catalog)
963 }
964 _ => Err(PgSourcePurificationError::NotPgConnection(
965 scx.catalog.resolve_full_name(connection_item.name()),
966 ))?,
967 };
968 let crate::plan::statement::PgConfigOptionExtracted {
969 publication,
970 text_columns,
971 exclude_columns,
972 details,
973 ..
974 } = options.clone().try_into()?;
975 let publication =
976 publication.ok_or(PgSourcePurificationError::ConnectionMissingPublication)?;
977
978 if details.is_some() {
979 Err(PgSourcePurificationError::UserSpecifiedDetails)?;
980 }
981
982 let client = connection
983 .validate(connection_item.id(), storage_configuration)
984 .await?;
985
986 let reference_client = SourceReferenceClient::Postgres {
987 client: &client,
988 publication: &publication,
989 database: &connection.database,
990 };
991 retrieved_source_references = reference_client.get_source_references().await?;
992
993 let is_physical_replica = mz_postgres_util::get_is_in_recovery(&client).await?;
996
997 let initial_lsn = MzOffset::from(
999 mz_postgres_util::fetch_max_lsn(&client, is_physical_replica).await?,
1000 );
1001
1002 let postgres::PurifiedSourceExports {
1003 source_exports: subsources,
1004 normalized_text_columns,
1005 } = postgres::purify_source_exports(
1006 &client,
1007 &retrieved_source_references,
1008 external_references,
1009 text_columns,
1010 exclude_columns,
1011 &BTreeSet::new(),
1012 false,
1013 source_name,
1014 &reference_policy,
1015 initial_lsn,
1016 )
1017 .await?;
1018
1019 if let Some(text_cols_option) = options
1020 .iter_mut()
1021 .find(|option| option.name == PgConfigOptionName::TextColumns)
1022 {
1023 text_cols_option.value = Some(WithOptionValue::Sequence(normalized_text_columns));
1024 }
1025
1026 requested_subsource_map.extend(subsources);
1027
1028 let timeline_id = mz_postgres_util::get_timeline_id(&client).await?;
1031
1032 options.retain(|PgConfigOption { name, .. }| name != &PgConfigOptionName::Details);
1034 let details = PostgresSourcePublicationDetails {
1035 slot: format!(
1036 "materialize_{}",
1037 Uuid::new_v4().to_string().replace('-', "")
1038 ),
1039 timeline_id: Some(timeline_id),
1040 database: connection.database,
1041 is_physical_replica: Some(is_physical_replica),
1042 };
1043 options.push(PgConfigOption {
1044 name: PgConfigOptionName::Details,
1045 value: Some(WithOptionValue::Value(Value::String(hex::encode(
1046 details.into_proto().encode_to_vec(),
1047 )))),
1048 })
1049 }
1050 CreateSourceConnection::SqlServer {
1051 connection,
1052 options,
1053 } => {
1054 let connection_item = scx.get_item_by_resolved_name(connection)?;
1057 let connection = match connection_item.connection()? {
1058 Connection::SqlServer(connection) => {
1059 connection.clone().into_inline_connection(&catalog)
1060 }
1061 _ => Err(SqlServerSourcePurificationError::NotSqlServerConnection(
1062 scx.catalog.resolve_full_name(connection_item.name()),
1063 ))?,
1064 };
1065 let crate::plan::statement::ddl::SqlServerConfigOptionExtracted {
1066 details,
1067 text_columns,
1068 exclude_columns,
1069 seen: _,
1070 } = options.clone().try_into()?;
1071
1072 if details.is_some() {
1073 Err(SqlServerSourcePurificationError::UserSpecifiedDetails)?;
1074 }
1075
1076 let mut client = connection
1077 .validate(connection_item.id(), storage_configuration)
1078 .await?;
1079
1080 let database: Arc<str> = connection.database.into();
1081 let reference_client = SourceReferenceClient::SqlServer {
1082 client: &mut client,
1083 database: Arc::clone(&database),
1084 };
1085 retrieved_source_references = reference_client.get_source_references().await?;
1086 tracing::debug!(?retrieved_source_references, "got source references");
1087
1088 let timeout = mz_storage_types::sources::sql_server::MAX_LSN_WAIT
1089 .get(storage_configuration.config_set());
1090
1091 let purified_source_exports = sql_server::purify_source_exports(
1092 &*database,
1093 &mut client,
1094 &retrieved_source_references,
1095 external_references,
1096 &text_columns,
1097 &exclude_columns,
1098 source_name,
1099 timeout,
1100 &reference_policy,
1101 )
1102 .await?;
1103
1104 let sql_server::PurifiedSourceExports {
1105 source_exports,
1106 normalized_text_columns,
1107 normalized_excl_columns,
1108 } = purified_source_exports;
1109
1110 requested_subsource_map.extend(source_exports);
1112
1113 let restore_history_id =
1117 mz_sql_server_util::inspect::get_latest_restore_history_id(&mut client).await?;
1118 let details = SqlServerSourceExtras { restore_history_id };
1119
1120 options.retain(|SqlServerConfigOption { name, .. }| {
1121 name != &SqlServerConfigOptionName::Details
1122 });
1123 options.push(SqlServerConfigOption {
1124 name: SqlServerConfigOptionName::Details,
1125 value: Some(WithOptionValue::Value(Value::String(hex::encode(
1126 details.into_proto().encode_to_vec(),
1127 )))),
1128 });
1129
1130 if let Some(text_cols_option) = options
1132 .iter_mut()
1133 .find(|option| option.name == SqlServerConfigOptionName::TextColumns)
1134 {
1135 text_cols_option.value = Some(WithOptionValue::Sequence(normalized_text_columns));
1136 }
1137 if let Some(excl_cols_option) = options
1138 .iter_mut()
1139 .find(|option| option.name == SqlServerConfigOptionName::ExcludeColumns)
1140 {
1141 excl_cols_option.value = Some(WithOptionValue::Sequence(normalized_excl_columns));
1142 }
1143 }
1144 CreateSourceConnection::MySql {
1145 connection,
1146 options,
1147 } => {
1148 let connection_item = scx.get_item_by_resolved_name(connection)?;
1149 let connection = match connection_item.connection()? {
1150 Connection::MySql(connection) => {
1151 connection.clone().into_inline_connection(&catalog)
1152 }
1153 _ => Err(MySqlSourcePurificationError::NotMySqlConnection(
1154 scx.catalog.resolve_full_name(connection_item.name()),
1155 ))?,
1156 };
1157 let crate::plan::statement::ddl::MySqlConfigOptionExtracted {
1158 details,
1159 text_columns,
1160 exclude_columns,
1161 seen: _,
1162 } = options.clone().try_into()?;
1163
1164 if details.is_some() {
1165 Err(MySqlSourcePurificationError::UserSpecifiedDetails)?;
1166 }
1167
1168 let mut conn = connection
1169 .validate(connection_item.id(), storage_configuration)
1170 .await
1171 .map_err(MySqlSourcePurificationError::InvalidConnection)?;
1172
1173 let initial_gtid_set =
1177 mz_mysql_util::query_sys_var(&mut conn, "global.gtid_executed").await?;
1178
1179 let binlog_full_metadata = is_binlog_full_metadata(&mut conn).await?;
1180
1181 let reference_client = SourceReferenceClient::MySql {
1182 conn: &mut conn,
1183 include_system_schemas: mysql::references_system_schemas(external_references),
1184 };
1185 retrieved_source_references = reference_client.get_source_references().await?;
1186
1187 let mysql::PurifiedSourceExports {
1188 source_exports: subsources,
1189 normalized_text_columns,
1190 normalized_exclude_columns,
1191 } = mysql::purify_source_exports(
1192 &mut conn,
1193 &retrieved_source_references,
1194 external_references,
1195 text_columns,
1196 exclude_columns,
1197 source_name,
1198 initial_gtid_set.clone(),
1199 &reference_policy,
1200 binlog_full_metadata,
1201 )
1202 .await?;
1203 requested_subsource_map.extend(subsources);
1204
1205 let details = MySqlSourceDetails {};
1208 options
1210 .retain(|MySqlConfigOption { name, .. }| name != &MySqlConfigOptionName::Details);
1211 options.push(MySqlConfigOption {
1212 name: MySqlConfigOptionName::Details,
1213 value: Some(WithOptionValue::Value(Value::String(hex::encode(
1214 details.into_proto().encode_to_vec(),
1215 )))),
1216 });
1217
1218 if let Some(text_cols_option) = options
1219 .iter_mut()
1220 .find(|option| option.name == MySqlConfigOptionName::TextColumns)
1221 {
1222 text_cols_option.value = Some(WithOptionValue::Sequence(normalized_text_columns));
1223 }
1224 if let Some(exclude_cols_option) = options
1225 .iter_mut()
1226 .find(|option| option.name == MySqlConfigOptionName::ExcludeColumns)
1227 {
1228 exclude_cols_option.value =
1229 Some(WithOptionValue::Sequence(normalized_exclude_columns));
1230 }
1231 }
1232 CreateSourceConnection::LoadGenerator { generator, options } => {
1233 let load_generator =
1234 load_generator_ast_to_generator(&scx, generator, options, include_metadata)?;
1235
1236 let reference_client = SourceReferenceClient::LoadGenerator {
1237 generator: &load_generator,
1238 };
1239 retrieved_source_references = reference_client.get_source_references().await?;
1240 let multi_output_sources =
1244 retrieved_source_references
1245 .all_references()
1246 .iter()
1247 .any(|r| {
1248 matches!(
1249 r.load_generator_output(),
1250 Some(output) if output != &LoadGeneratorOutput::Default
1251 )
1252 });
1253
1254 match external_references {
1255 Some(requested)
1256 if matches!(reference_policy, SourceReferencePolicy::NotAllowed) =>
1257 {
1258 Err(PlanError::UseTablesForSources(requested.to_string()))?
1259 }
1260 Some(requested) if !multi_output_sources => match requested {
1261 ExternalReferences::SubsetTables(_) => {
1262 Err(LoadGeneratorSourcePurificationError::ForTables)?
1263 }
1264 ExternalReferences::SubsetSchemas(_) => {
1265 Err(LoadGeneratorSourcePurificationError::ForSchemas)?
1266 }
1267 ExternalReferences::All => {
1268 Err(LoadGeneratorSourcePurificationError::ForAllTables)?
1269 }
1270 },
1271 Some(requested) => {
1272 let requested_exports = retrieved_source_references
1273 .requested_source_exports(Some(requested), source_name)?;
1274 for export in requested_exports {
1275 requested_subsource_map.insert(
1276 export.name,
1277 PurifiedSourceExport {
1278 external_reference: export.external_reference,
1279 details: PurifiedExportDetails::LoadGenerator {
1280 table: export
1281 .meta
1282 .load_generator_desc()
1283 .ok_or_else(|| {
1284 internal_err!(
1285 "expected load generator source reference"
1286 )
1287 })?
1288 .clone(),
1289 output: export
1290 .meta
1291 .load_generator_output()
1292 .ok_or_else(|| {
1293 internal_err!(
1294 "expected load generator source reference"
1295 )
1296 })?
1297 .clone(),
1298 },
1299 },
1300 );
1301 }
1302 }
1303 None => {
1304 if multi_output_sources
1305 && matches!(reference_policy, SourceReferencePolicy::Required)
1306 {
1307 Err(LoadGeneratorSourcePurificationError::MultiOutputRequiresForAllTables)?
1308 }
1309 }
1310 }
1311
1312 if let LoadGenerator::Clock = generator {
1313 if !options
1314 .iter()
1315 .any(|p| p.name == LoadGeneratorOptionName::AsOf)
1316 {
1317 let now = catalog.now();
1318 options.push(LoadGeneratorOption {
1319 name: LoadGeneratorOptionName::AsOf,
1320 value: Some(WithOptionValue::Value(Value::Number(now.to_string()))),
1321 });
1322 }
1323 }
1324 }
1325 }
1326
1327 *external_references = None;
1331
1332 let create_progress_subsource_stmt = if uses_old_syntax {
1334 let name = match progress_subsource {
1336 Some(name) => match name {
1337 DeferredItemName::Deferred(name) => name.clone(),
1338 DeferredItemName::Named(_) => {
1340 sql_bail!("progress subsource name cannot be a resolved name")
1341 }
1342 },
1343 None => {
1344 let (item, prefix) = source_name
1345 .0
1346 .split_last()
1347 .ok_or_else(|| sql_err!("source name must have at least one component"))?;
1348 let item_name =
1349 Ident::try_generate_name(item.to_string(), "_progress", |candidate| {
1350 let mut suggested_name = prefix.to_vec();
1351 suggested_name.push(candidate.clone());
1352
1353 let partial =
1354 normalize::unresolved_item_name(UnresolvedItemName(suggested_name))?;
1355 let qualified = scx.allocate_qualified_name(partial)?;
1356 let item_exists = scx.catalog.get_item_by_name(&qualified).is_some();
1357 let type_exists = scx.catalog.get_type_by_name(&qualified).is_some();
1358 Ok::<_, PlanError>(!item_exists && !type_exists)
1359 })?;
1360
1361 let mut full_name = prefix.to_vec();
1362 full_name.push(item_name);
1363 let full_name = normalize::unresolved_item_name(UnresolvedItemName(full_name))?;
1364 let qualified_name = scx.allocate_qualified_name(full_name)?;
1365 let full_name = scx.catalog.resolve_full_name(&qualified_name);
1366
1367 UnresolvedItemName::from(full_name.clone())
1368 }
1369 };
1370
1371 let (columns, constraints) = scx.relation_desc_into_table_defs(progress_desc)?;
1372
1373 let mut progress_with_options: Vec<_> = with_options
1375 .iter()
1376 .filter_map(|opt| match opt.name {
1377 CreateSourceOptionName::TimestampInterval => None,
1378 CreateSourceOptionName::RetainHistory => Some(CreateSubsourceOption {
1379 name: CreateSubsourceOptionName::RetainHistory,
1380 value: opt.value.clone(),
1381 }),
1382 })
1383 .collect();
1384 progress_with_options.push(CreateSubsourceOption {
1385 name: CreateSubsourceOptionName::Progress,
1386 value: Some(WithOptionValue::Value(Value::Boolean(true))),
1387 });
1388
1389 Some(CreateSubsourceStatement {
1390 name,
1391 columns,
1392 of_source: None,
1396 constraints,
1397 if_not_exists: false,
1398 with_options: progress_with_options,
1399 })
1400 } else {
1401 None
1402 };
1403
1404 purify_source_format(
1405 &catalog,
1406 format,
1407 &format_options,
1408 envelope,
1409 storage_configuration,
1410 )
1411 .await?;
1412
1413 Ok(PurifiedStatement::PurifiedCreateSource {
1414 create_progress_subsource_stmt,
1415 create_source_stmt,
1416 subsources: requested_subsource_map,
1417 available_source_references: retrieved_source_references.available_source_references(),
1418 })
1419}
1420
1421async fn purify_alter_source(
1424 catalog: impl SessionCatalog,
1425 stmt: AlterSourceStatement<Aug>,
1426 storage_configuration: &StorageConfiguration,
1427) -> Result<PurifiedStatement, PlanError> {
1428 let scx = StatementContext::new(None, &catalog);
1429 let AlterSourceStatement {
1430 source_name: unresolved_source_name,
1431 action,
1432 if_exists,
1433 } = stmt;
1434
1435 let item = match scx.resolve_item(RawItemName::Name(unresolved_source_name.clone())) {
1437 Ok(item) => item,
1438 Err(_) if if_exists => {
1439 return Ok(PurifiedStatement::PurifiedAlterSource {
1440 alter_source_stmt: AlterSourceStatement {
1441 source_name: unresolved_source_name,
1442 action,
1443 if_exists,
1444 },
1445 });
1446 }
1447 Err(e) => return Err(e),
1448 };
1449
1450 let desc = match item.source_desc()? {
1452 Some(desc) => desc.clone().into_inline_connection(scx.catalog),
1453 None => {
1454 sql_bail!("cannot ALTER this type of source")
1455 }
1456 };
1457
1458 let source_name = item.name();
1459
1460 let resolved_source_name = ResolvedItemName::Item {
1461 id: item.id(),
1462 qualifiers: item.name().qualifiers.clone(),
1463 full_name: scx.catalog.resolve_full_name(source_name),
1464 print_id: true,
1465 version: RelationVersionSelector::Latest,
1466 };
1467
1468 let partial_name = scx.catalog.minimal_qualification(source_name);
1469
1470 match action {
1471 AlterSourceAction::AddSubsources {
1472 external_references,
1473 options,
1474 } => {
1475 if scx.catalog.system_vars().enable_create_table_from_source()
1476 && scx.catalog.system_vars().force_source_table_syntax()
1477 {
1478 Err(PlanError::UseTablesForSources(
1479 "ALTER SOURCE .. ADD SUBSOURCES ..".to_string(),
1480 ))?;
1481 }
1482
1483 purify_alter_source_add_subsources(
1484 external_references,
1485 options,
1486 desc,
1487 partial_name,
1488 unresolved_source_name,
1489 resolved_source_name,
1490 storage_configuration,
1491 )
1492 .await
1493 }
1494 AlterSourceAction::RefreshReferences => {
1495 purify_alter_source_refresh_references(
1496 desc,
1497 resolved_source_name,
1498 storage_configuration,
1499 )
1500 .await
1501 }
1502 _ => Ok(PurifiedStatement::PurifiedAlterSource {
1503 alter_source_stmt: AlterSourceStatement {
1504 source_name: unresolved_source_name,
1505 action,
1506 if_exists,
1507 },
1508 }),
1509 }
1510}
1511
1512async fn purify_alter_source_add_subsources(
1515 external_references: Vec<ExternalReferenceExport>,
1516 mut options: Vec<AlterSourceAddSubsourceOption<Aug>>,
1517 desc: SourceDesc,
1518 partial_source_name: PartialItemName,
1519 unresolved_source_name: UnresolvedItemName,
1520 resolved_source_name: ResolvedItemName,
1521 storage_configuration: &StorageConfiguration,
1522) -> Result<PurifiedStatement, PlanError> {
1523 let connection_id = match &desc.connection {
1525 GenericSourceConnection::Postgres(c) => c.connection_id,
1526 GenericSourceConnection::MySql(c) => c.connection_id,
1527 GenericSourceConnection::SqlServer(c) => c.connection_id,
1528 _ => sql_bail!(
1529 "source {} does not support ALTER SOURCE.",
1530 partial_source_name
1531 ),
1532 };
1533
1534 let crate::plan::statement::ddl::AlterSourceAddSubsourceOptionExtracted {
1535 text_columns,
1536 exclude_columns,
1537 details,
1538 seen: _,
1539 } = options.clone().try_into()?;
1540 if details.is_some() {
1541 sql_bail!("DETAILS option cannot be explicitly set");
1542 }
1543
1544 let mut requested_subsource_map = BTreeMap::new();
1545
1546 match desc.connection {
1547 GenericSourceConnection::Postgres(pg_source_connection) => {
1548 let pg_connection = &pg_source_connection.connection;
1550
1551 let client = pg_connection
1552 .validate(connection_id, storage_configuration)
1553 .await?;
1554
1555 let reference_client = SourceReferenceClient::Postgres {
1556 client: &client,
1557 publication: &pg_source_connection.publication,
1558 database: &pg_connection.database,
1559 };
1560 let retrieved_source_references = reference_client.get_source_references().await?;
1561
1562 let initial_lsn = MzOffset::from(
1564 mz_postgres_util::fetch_max_lsn(
1565 &client,
1566 pg_source_connection
1567 .publication_details
1568 .get_is_physical_replica(),
1569 )
1570 .await?,
1571 );
1572
1573 let postgres::PurifiedSourceExports {
1574 source_exports: subsources,
1575 normalized_text_columns,
1576 } = postgres::purify_source_exports(
1577 &client,
1578 &retrieved_source_references,
1579 &Some(ExternalReferences::SubsetTables(external_references)),
1580 text_columns,
1581 exclude_columns,
1582 &BTreeSet::new(),
1583 false,
1584 &unresolved_source_name,
1585 &SourceReferencePolicy::Required,
1586 initial_lsn,
1587 )
1588 .await?;
1589
1590 if let Some(text_cols_option) = options
1591 .iter_mut()
1592 .find(|option| option.name == AlterSourceAddSubsourceOptionName::TextColumns)
1593 {
1594 text_cols_option.value = Some(WithOptionValue::Sequence(normalized_text_columns));
1595 }
1596
1597 requested_subsource_map.extend(subsources);
1598 }
1599 GenericSourceConnection::MySql(mysql_source_connection) => {
1600 let mysql_connection = &mysql_source_connection.connection;
1601 let config = mysql_connection
1602 .config(
1603 &storage_configuration.connection_context.secrets_reader,
1604 storage_configuration,
1605 InTask::No,
1606 )
1607 .await?;
1608
1609 let mut conn = config
1610 .connect(
1611 "mysql purification",
1612 &storage_configuration.connection_context.ssh_tunnel_manager,
1613 )
1614 .await?;
1615
1616 let initial_gtid_set =
1619 mz_mysql_util::query_sys_var(&mut conn, "global.gtid_executed").await?;
1620
1621 let binlog_full_metadata = is_binlog_full_metadata(&mut conn).await?;
1622
1623 let requested_references = Some(ExternalReferences::SubsetTables(external_references));
1624
1625 let reference_client = SourceReferenceClient::MySql {
1626 conn: &mut conn,
1627 include_system_schemas: mysql::references_system_schemas(&requested_references),
1628 };
1629 let retrieved_source_references = reference_client.get_source_references().await?;
1630
1631 let mysql::PurifiedSourceExports {
1632 source_exports: subsources,
1633 normalized_text_columns,
1634 normalized_exclude_columns,
1635 } = mysql::purify_source_exports(
1636 &mut conn,
1637 &retrieved_source_references,
1638 &requested_references,
1639 text_columns,
1640 exclude_columns,
1641 &unresolved_source_name,
1642 initial_gtid_set,
1643 &SourceReferencePolicy::Required,
1644 binlog_full_metadata,
1645 )
1646 .await?;
1647 requested_subsource_map.extend(subsources);
1648
1649 if let Some(text_cols_option) = options
1651 .iter_mut()
1652 .find(|option| option.name == AlterSourceAddSubsourceOptionName::TextColumns)
1653 {
1654 text_cols_option.value = Some(WithOptionValue::Sequence(normalized_text_columns));
1655 }
1656 if let Some(exclude_cols_option) = options
1657 .iter_mut()
1658 .find(|option| option.name == AlterSourceAddSubsourceOptionName::ExcludeColumns)
1659 {
1660 exclude_cols_option.value =
1661 Some(WithOptionValue::Sequence(normalized_exclude_columns));
1662 }
1663 }
1664 GenericSourceConnection::SqlServer(sql_server_source) => {
1665 let sql_server_connection = &sql_server_source.connection;
1667 let config = sql_server_connection
1668 .resolve_config(
1669 &storage_configuration.connection_context.secrets_reader,
1670 storage_configuration,
1671 InTask::No,
1672 )
1673 .await?;
1674 let mut client = mz_sql_server_util::Client::connect(config).await?;
1675
1676 let database = sql_server_connection.database.clone().into();
1678 let source_references = SourceReferenceClient::SqlServer {
1679 client: &mut client,
1680 database: Arc::clone(&database),
1681 }
1682 .get_source_references()
1683 .await?;
1684 let requested_references = Some(ExternalReferences::SubsetTables(external_references));
1685
1686 let timeout = mz_storage_types::sources::sql_server::MAX_LSN_WAIT
1687 .get(storage_configuration.config_set());
1688
1689 let result = sql_server::purify_source_exports(
1690 &*database,
1691 &mut client,
1692 &source_references,
1693 &requested_references,
1694 &text_columns,
1695 &exclude_columns,
1696 &unresolved_source_name,
1697 timeout,
1698 &SourceReferencePolicy::Required,
1699 )
1700 .await;
1701 let sql_server::PurifiedSourceExports {
1702 source_exports,
1703 normalized_text_columns,
1704 normalized_excl_columns,
1705 } = result?;
1706
1707 requested_subsource_map.extend(source_exports);
1709
1710 if let Some(text_cols_option) = options
1712 .iter_mut()
1713 .find(|option| option.name == AlterSourceAddSubsourceOptionName::TextColumns)
1714 {
1715 text_cols_option.value = Some(WithOptionValue::Sequence(normalized_text_columns));
1716 }
1717 if let Some(exclude_cols_option) = options
1718 .iter_mut()
1719 .find(|option| option.name == AlterSourceAddSubsourceOptionName::ExcludeColumns)
1720 {
1721 exclude_cols_option.value =
1722 Some(WithOptionValue::Sequence(normalized_excl_columns));
1723 }
1724 }
1725 _ => bail_internal!("source does not support ALTER SOURCE...ADD SUBSOURCE"),
1728 };
1729
1730 Ok(PurifiedStatement::PurifiedAlterSourceAddSubsources {
1731 source_name: resolved_source_name,
1732 options,
1733 subsources: requested_subsource_map,
1734 })
1735}
1736
1737async fn purify_alter_source_refresh_references(
1738 desc: SourceDesc,
1739 resolved_source_name: ResolvedItemName,
1740 storage_configuration: &StorageConfiguration,
1741) -> Result<PurifiedStatement, PlanError> {
1742 let retrieved_source_references = match desc.connection {
1743 GenericSourceConnection::Postgres(pg_source_connection) => {
1744 let pg_connection = &pg_source_connection.connection;
1746
1747 let config = pg_connection
1748 .config(
1749 &storage_configuration.connection_context.secrets_reader,
1750 storage_configuration,
1751 InTask::No,
1752 )
1753 .await?;
1754
1755 let client = config
1756 .connect(
1757 "postgres_purification",
1758 &storage_configuration.connection_context.ssh_tunnel_manager,
1759 )
1760 .await?;
1761 let reference_client = SourceReferenceClient::Postgres {
1762 client: &client,
1763 publication: &pg_source_connection.publication,
1764 database: &pg_connection.database,
1765 };
1766 reference_client.get_source_references().await?
1767 }
1768 GenericSourceConnection::MySql(mysql_source_connection) => {
1769 let mysql_connection = &mysql_source_connection.connection;
1770 let config = mysql_connection
1771 .config(
1772 &storage_configuration.connection_context.secrets_reader,
1773 storage_configuration,
1774 InTask::No,
1775 )
1776 .await?;
1777
1778 let mut conn = config
1779 .connect(
1780 "mysql purification",
1781 &storage_configuration.connection_context.ssh_tunnel_manager,
1782 )
1783 .await?;
1784
1785 let reference_client = SourceReferenceClient::MySql {
1786 conn: &mut conn,
1787 include_system_schemas: false,
1788 };
1789 reference_client.get_source_references().await?
1790 }
1791 GenericSourceConnection::SqlServer(sql_server_source) => {
1792 let sql_server_connection = &sql_server_source.connection;
1794 let config = sql_server_connection
1795 .resolve_config(
1796 &storage_configuration.connection_context.secrets_reader,
1797 storage_configuration,
1798 InTask::No,
1799 )
1800 .await?;
1801 let mut client = mz_sql_server_util::Client::connect(config).await?;
1802
1803 let source_references = SourceReferenceClient::SqlServer {
1805 client: &mut client,
1806 database: sql_server_connection.database.clone().into(),
1807 }
1808 .get_source_references()
1809 .await?;
1810 source_references
1811 }
1812 GenericSourceConnection::LoadGenerator(load_gen_connection) => {
1813 let reference_client = SourceReferenceClient::LoadGenerator {
1814 generator: &load_gen_connection.load_generator,
1815 };
1816 reference_client.get_source_references().await?
1817 }
1818 GenericSourceConnection::Kafka(kafka_conn) => {
1819 let reference_client = SourceReferenceClient::Kafka {
1820 topic: &kafka_conn.topic,
1821 };
1822 reference_client.get_source_references().await?
1823 }
1824 };
1825 Ok(PurifiedStatement::PurifiedAlterSourceRefreshReferences {
1826 source_name: resolved_source_name,
1827 available_source_references: retrieved_source_references.available_source_references(),
1828 })
1829}
1830
1831async fn purify_create_table_from_source(
1832 catalog: impl SessionCatalog,
1833 mut stmt: CreateTableFromSourceStatement<Aug>,
1834 storage_configuration: &StorageConfiguration,
1835) -> Result<PurifiedStatement, PlanError> {
1836 let scx = StatementContext::new(None, &catalog);
1837 let CreateTableFromSourceStatement {
1838 name: _,
1839 columns,
1840 constraints,
1841 source: source_name,
1842 if_not_exists: _,
1843 external_reference,
1844 format,
1845 envelope,
1846 include_metadata: _,
1847 with_options,
1848 } = &mut stmt;
1849
1850 if matches!(columns, TableFromSourceColumns::Defined(_)) {
1852 sql_bail!("CREATE TABLE .. FROM SOURCE column definitions cannot be specified directly");
1853 }
1854 if !constraints.is_empty() {
1855 sql_bail!(
1856 "CREATE TABLE .. FROM SOURCE constraint definitions cannot be specified directly"
1857 );
1858 }
1859
1860 let item = match scx.get_item_by_resolved_name(source_name) {
1862 Ok(item) => item,
1863 Err(e) => return Err(e),
1864 };
1865
1866 let desc = match item.source_desc()? {
1868 Some(desc) => desc.clone().into_inline_connection(scx.catalog),
1869 None => {
1870 sql_bail!("cannot ALTER this type of source")
1871 }
1872 };
1873 let unresolved_source_name: UnresolvedItemName = source_name.full_item_name().clone().into();
1874
1875 let crate::plan::statement::ddl::TableFromSourceOptionExtracted {
1876 text_columns,
1877 exclude_columns,
1878 exclude_constraints,
1879 exclude_all_constraints,
1880 retain_history: _,
1881 details,
1882 partition_by: _,
1883 seen: _,
1884 } = with_options.clone().try_into()?;
1885 if details.is_some() {
1886 sql_bail!("DETAILS option cannot be explicitly set");
1887 }
1888
1889 if !exclude_constraints.is_empty() || exclude_all_constraints {
1890 scx.require_feature_flag(&crate::session::vars::ENABLE_EXCLUDE_CONSTRAINTS_OPTION)?;
1891 }
1892 if !exclude_constraints.is_empty() && exclude_all_constraints {
1893 sql_bail!("EXCLUDE ALL CONSTRAINTS cannot be combined with EXCLUDE CONSTRAINTS");
1894 }
1895 let exclude_constraints: BTreeSet<String> = exclude_constraints.into_iter().collect();
1896
1897 let qualified_text_columns = text_columns
1901 .iter()
1902 .map(|col| {
1903 UnresolvedItemName(
1904 external_reference
1905 .as_ref()
1906 .map(|er| er.0.iter().chain_one(col).map(|i| i.clone()).collect())
1907 .unwrap_or_else(|| vec![col.clone()]),
1908 )
1909 })
1910 .collect_vec();
1911 let qualified_exclude_columns = exclude_columns
1912 .iter()
1913 .map(|col| {
1914 UnresolvedItemName(
1915 external_reference
1916 .as_ref()
1917 .map(|er| er.0.iter().chain_one(col).map(|i| i.clone()).collect())
1918 .unwrap_or_else(|| vec![col.clone()]),
1919 )
1920 })
1921 .collect_vec();
1922
1923 let mut format_options = SourceFormatOptions::Default;
1925
1926 let retrieved_source_references: RetrievedSourceReferences;
1927
1928 let requested_references = external_reference.as_ref().map(|ref_name| {
1929 ExternalReferences::SubsetTables(vec![ExternalReferenceExport {
1930 reference: ref_name.clone(),
1931 alias: None,
1932 }])
1933 });
1934
1935 if (!exclude_constraints.is_empty() || exclude_all_constraints)
1936 && !matches!(desc.connection, GenericSourceConnection::Postgres(_))
1937 {
1938 sql_bail!(
1939 "EXCLUDE CONSTRAINTS is not supported for {} sources",
1940 desc.connection.name()
1941 );
1942 }
1943
1944 let purified_export = match desc.connection {
1947 GenericSourceConnection::Postgres(pg_source_connection) => {
1948 let pg_connection = &pg_source_connection.connection;
1950
1951 let client = pg_connection
1952 .validate(pg_source_connection.connection_id, storage_configuration)
1953 .await?;
1954
1955 let reference_client = SourceReferenceClient::Postgres {
1956 client: &client,
1957 publication: &pg_source_connection.publication,
1958 database: &pg_connection.database,
1959 };
1960 retrieved_source_references = reference_client.get_source_references().await?;
1961
1962 let initial_lsn = MzOffset::from(
1964 mz_postgres_util::fetch_max_lsn(
1965 &client,
1966 pg_source_connection
1967 .publication_details
1968 .get_is_physical_replica(),
1969 )
1970 .await?,
1971 );
1972
1973 let postgres::PurifiedSourceExports {
1974 source_exports,
1975 normalized_text_columns: _,
1979 } = postgres::purify_source_exports(
1980 &client,
1981 &retrieved_source_references,
1982 &requested_references,
1983 qualified_text_columns,
1984 qualified_exclude_columns,
1985 &exclude_constraints,
1986 exclude_all_constraints,
1987 &unresolved_source_name,
1988 &SourceReferencePolicy::Required,
1989 initial_lsn,
1990 )
1991 .await?;
1992 let (_, purified_export) = source_exports.into_element();
1994 purified_export
1995 }
1996 GenericSourceConnection::MySql(mysql_source_connection) => {
1997 let mysql_connection = &mysql_source_connection.connection;
1998 let config = mysql_connection
1999 .config(
2000 &storage_configuration.connection_context.secrets_reader,
2001 storage_configuration,
2002 InTask::No,
2003 )
2004 .await?;
2005
2006 let mut conn = config
2007 .connect(
2008 "mysql purification",
2009 &storage_configuration.connection_context.ssh_tunnel_manager,
2010 )
2011 .await?;
2012
2013 ensure_binlog_full_metadata(&mut conn).await?;
2014 let binlog_full_metadata = true;
2015
2016 let initial_gtid_set =
2019 mz_mysql_util::query_sys_var(&mut conn, "global.gtid_executed").await?;
2020
2021 let reference_client = SourceReferenceClient::MySql {
2022 conn: &mut conn,
2023 include_system_schemas: mysql::references_system_schemas(&requested_references),
2024 };
2025 retrieved_source_references = reference_client.get_source_references().await?;
2026
2027 let mysql::PurifiedSourceExports {
2028 source_exports,
2029 normalized_text_columns: _,
2033 normalized_exclude_columns: _,
2034 } = mysql::purify_source_exports(
2035 &mut conn,
2036 &retrieved_source_references,
2037 &requested_references,
2038 qualified_text_columns,
2039 qualified_exclude_columns,
2040 &unresolved_source_name,
2041 initial_gtid_set,
2042 &SourceReferencePolicy::Required,
2043 binlog_full_metadata,
2044 )
2045 .await?;
2046 let (_, purified_export) = source_exports.into_element();
2048 purified_export
2049 }
2050 GenericSourceConnection::SqlServer(sql_server_source) => {
2051 let connection = sql_server_source.connection;
2052 let config = connection
2053 .resolve_config(
2054 &storage_configuration.connection_context.secrets_reader,
2055 storage_configuration,
2056 InTask::No,
2057 )
2058 .await?;
2059 let mut client = mz_sql_server_util::Client::connect(config).await?;
2060
2061 let database: Arc<str> = connection.database.into();
2062 let reference_client = SourceReferenceClient::SqlServer {
2063 client: &mut client,
2064 database: Arc::clone(&database),
2065 };
2066 retrieved_source_references = reference_client.get_source_references().await?;
2067 tracing::debug!(?retrieved_source_references, "got source references");
2068
2069 let timeout = mz_storage_types::sources::sql_server::MAX_LSN_WAIT
2070 .get(storage_configuration.config_set());
2071
2072 let purified_source_exports = sql_server::purify_source_exports(
2073 &*database,
2074 &mut client,
2075 &retrieved_source_references,
2076 &requested_references,
2077 &qualified_text_columns,
2078 &qualified_exclude_columns,
2079 &unresolved_source_name,
2080 timeout,
2081 &SourceReferencePolicy::Required,
2082 )
2083 .await?;
2084
2085 let (_, purified_export) = purified_source_exports.source_exports.into_element();
2087 purified_export
2088 }
2089 GenericSourceConnection::LoadGenerator(load_gen_connection) => {
2090 let reference_client = SourceReferenceClient::LoadGenerator {
2091 generator: &load_gen_connection.load_generator,
2092 };
2093 retrieved_source_references = reference_client.get_source_references().await?;
2094
2095 let requested_exports = retrieved_source_references
2096 .requested_source_exports(requested_references.as_ref(), &unresolved_source_name)?;
2097 let export = requested_exports.into_element();
2099 PurifiedSourceExport {
2100 external_reference: export.external_reference,
2101 details: PurifiedExportDetails::LoadGenerator {
2102 table: export
2103 .meta
2104 .load_generator_desc()
2105 .ok_or_else(|| internal_err!("expected load generator source reference"))?
2106 .clone(),
2107 output: export
2108 .meta
2109 .load_generator_output()
2110 .ok_or_else(|| internal_err!("expected load generator source reference"))?
2111 .clone(),
2112 },
2113 }
2114 }
2115 GenericSourceConnection::Kafka(kafka_conn) => {
2116 let reference_client = SourceReferenceClient::Kafka {
2117 topic: &kafka_conn.topic,
2118 };
2119 retrieved_source_references = reference_client.get_source_references().await?;
2120 let requested_exports = retrieved_source_references
2121 .requested_source_exports(requested_references.as_ref(), &unresolved_source_name)?;
2122 let export = requested_exports.into_element();
2124
2125 format_options = SourceFormatOptions::Kafka {
2126 topic: kafka_conn.topic.clone(),
2127 };
2128 PurifiedSourceExport {
2129 external_reference: export.external_reference,
2130 details: PurifiedExportDetails::Kafka {},
2131 }
2132 }
2133 };
2134
2135 purify_source_format(
2136 &catalog,
2137 format,
2138 &format_options,
2139 envelope,
2140 storage_configuration,
2141 )
2142 .await?;
2143
2144 *external_reference = Some(purified_export.external_reference.clone());
2147
2148 match &purified_export.details {
2150 PurifiedExportDetails::Postgres { .. } => {
2151 let mut unsupported_cols = vec![];
2152 let postgres::PostgresExportStatementValues {
2153 columns: gen_columns,
2154 constraints: gen_constraints,
2155 text_columns: gen_text_columns,
2156 exclude_columns: gen_exclude_columns,
2157 details: gen_details,
2158 external_reference: _,
2159 } = postgres::generate_source_export_statement_values(
2160 &scx,
2161 purified_export,
2162 &mut unsupported_cols,
2163 )?;
2164 if !unsupported_cols.is_empty() {
2165 unsupported_cols.sort();
2166 Err(PgSourcePurificationError::UnrecognizedTypes {
2167 cols: unsupported_cols,
2168 })?;
2169 }
2170
2171 if let Some(text_cols_option) = with_options
2172 .iter_mut()
2173 .find(|option| option.name == TableFromSourceOptionName::TextColumns)
2174 {
2175 if let Some(gen_text_columns) = gen_text_columns {
2176 text_cols_option.value = Some(WithOptionValue::Sequence(gen_text_columns));
2177 }
2178 }
2179 if let Some(exclude_cols_option) = with_options
2180 .iter_mut()
2181 .find(|option| option.name == TableFromSourceOptionName::ExcludeColumns)
2182 {
2183 if let Some(gen_exclude_columns) = gen_exclude_columns {
2184 exclude_cols_option.value =
2185 Some(WithOptionValue::Sequence(gen_exclude_columns));
2186 }
2187 }
2188 match columns {
2189 TableFromSourceColumns::Defined(_) => {
2190 bail_internal!(
2193 "column definitions cannot be explicitly set for this source type"
2194 )
2195 }
2196 TableFromSourceColumns::NotSpecified => {
2197 *columns = TableFromSourceColumns::Defined(gen_columns);
2198 *constraints = gen_constraints;
2199 }
2200 TableFromSourceColumns::Named(_) => {
2201 sql_bail!("columns cannot be named for Postgres sources")
2202 }
2203 }
2204 with_options.push(TableFromSourceOption {
2205 name: TableFromSourceOptionName::Details,
2206 value: Some(WithOptionValue::Value(Value::String(hex::encode(
2207 gen_details.into_proto().encode_to_vec(),
2208 )))),
2209 })
2210 }
2211 PurifiedExportDetails::MySql { .. } => {
2212 let mysql::MySqlExportStatementValues {
2213 columns: gen_columns,
2214 constraints: gen_constraints,
2215 text_columns: gen_text_columns,
2216 exclude_columns: gen_exclude_columns,
2217 details: gen_details,
2218 external_reference: _,
2219 } = mysql::generate_source_export_statement_values(&scx, purified_export)?;
2220
2221 if let Some(text_cols_option) = with_options
2222 .iter_mut()
2223 .find(|option| option.name == TableFromSourceOptionName::TextColumns)
2224 {
2225 if let Some(gen_text_columns) = gen_text_columns {
2226 text_cols_option.value = Some(WithOptionValue::Sequence(gen_text_columns));
2227 }
2228 }
2229 if let Some(exclude_cols_option) = with_options
2230 .iter_mut()
2231 .find(|option| option.name == TableFromSourceOptionName::ExcludeColumns)
2232 {
2233 if let Some(gen_exclude_columns) = gen_exclude_columns {
2234 exclude_cols_option.value =
2235 Some(WithOptionValue::Sequence(gen_exclude_columns));
2236 }
2237 }
2238 match columns {
2239 TableFromSourceColumns::Defined(_) => {
2240 bail_internal!(
2243 "column definitions cannot be explicitly set for this source type"
2244 )
2245 }
2246 TableFromSourceColumns::NotSpecified => {
2247 *columns = TableFromSourceColumns::Defined(gen_columns);
2248 *constraints = gen_constraints;
2249 }
2250 TableFromSourceColumns::Named(_) => {
2251 sql_bail!("columns cannot be named for MySQL sources")
2252 }
2253 }
2254 with_options.push(TableFromSourceOption {
2255 name: TableFromSourceOptionName::Details,
2256 value: Some(WithOptionValue::Value(Value::String(hex::encode(
2257 gen_details.into_proto().encode_to_vec(),
2258 )))),
2259 })
2260 }
2261 PurifiedExportDetails::SqlServer { .. } => {
2262 let sql_server::SqlServerExportStatementValues {
2263 columns: gen_columns,
2264 constraints: gen_constraints,
2265 text_columns: gen_text_columns,
2266 excl_columns: gen_excl_columns,
2267 details: gen_details,
2268 external_reference: _,
2269 } = sql_server::generate_source_export_statement_values(&scx, purified_export)?;
2270
2271 if let Some(text_cols_option) = with_options
2272 .iter_mut()
2273 .find(|opt| opt.name == TableFromSourceOptionName::TextColumns)
2274 {
2275 if let Some(gen_text_columns) = gen_text_columns {
2276 text_cols_option.value = Some(WithOptionValue::Sequence(gen_text_columns));
2277 }
2278 }
2279 if let Some(exclude_cols_option) = with_options
2280 .iter_mut()
2281 .find(|opt| opt.name == TableFromSourceOptionName::ExcludeColumns)
2282 {
2283 if let Some(gen_excl_columns) = gen_excl_columns {
2284 exclude_cols_option.value = Some(WithOptionValue::Sequence(gen_excl_columns));
2285 }
2286 }
2287
2288 match columns {
2289 TableFromSourceColumns::NotSpecified => {
2290 *columns = TableFromSourceColumns::Defined(gen_columns);
2291 *constraints = gen_constraints;
2292 }
2293 TableFromSourceColumns::Named(_) => {
2294 sql_bail!("columns cannot be named for SQL Server sources")
2295 }
2296 TableFromSourceColumns::Defined(_) => {
2297 bail_internal!(
2300 "column definitions cannot be explicitly set for this source type"
2301 )
2302 }
2303 }
2304
2305 with_options.push(TableFromSourceOption {
2306 name: TableFromSourceOptionName::Details,
2307 value: Some(WithOptionValue::Value(Value::String(hex::encode(
2308 gen_details.into_proto().encode_to_vec(),
2309 )))),
2310 })
2311 }
2312 PurifiedExportDetails::LoadGenerator { .. } => {
2313 let (desc, output) = match purified_export.details {
2314 PurifiedExportDetails::LoadGenerator { table, output } => (table, output),
2315 _ => bail_internal!("purified export details must be load generator"),
2316 };
2317 if let Some(desc) = desc {
2322 let (gen_columns, gen_constraints) = scx.relation_desc_into_table_defs(&desc)?;
2323 match columns {
2324 TableFromSourceColumns::Defined(_) => bail_internal!(
2327 "column definitions cannot be explicitly set for this source type"
2328 ),
2329 TableFromSourceColumns::NotSpecified => {
2330 *columns = TableFromSourceColumns::Defined(gen_columns);
2331 *constraints = gen_constraints;
2332 }
2333 TableFromSourceColumns::Named(_) => {
2334 sql_bail!("columns cannot be named for multi-output load generator sources")
2335 }
2336 }
2337 }
2338 let details = SourceExportStatementDetails::LoadGenerator { output };
2339 with_options.push(TableFromSourceOption {
2340 name: TableFromSourceOptionName::Details,
2341 value: Some(WithOptionValue::Value(Value::String(hex::encode(
2342 details.into_proto().encode_to_vec(),
2343 )))),
2344 })
2345 }
2346 PurifiedExportDetails::Kafka {} => {
2347 let details = SourceExportStatementDetails::Kafka {};
2351 with_options.push(TableFromSourceOption {
2352 name: TableFromSourceOptionName::Details,
2353 value: Some(WithOptionValue::Value(Value::String(hex::encode(
2354 details.into_proto().encode_to_vec(),
2355 )))),
2356 })
2357 }
2358 };
2359
2360 Ok(PurifiedStatement::PurifiedCreateTableFromSource { stmt })
2364}
2365
2366enum SourceFormatOptions {
2367 Default,
2368 Kafka { topic: String },
2369}
2370
2371async fn purify_source_format(
2372 catalog: &dyn SessionCatalog,
2373 format: &mut Option<FormatSpecifier<Aug>>,
2374 options: &SourceFormatOptions,
2375 envelope: &Option<SourceEnvelope>,
2376 storage_configuration: &StorageConfiguration,
2377) -> Result<(), PlanError> {
2378 if matches!(format, Some(FormatSpecifier::KeyValue { .. }))
2379 && !matches!(options, SourceFormatOptions::Kafka { .. })
2380 {
2381 sql_bail!("Kafka sources are the only source type that can provide KEY/VALUE formats")
2382 }
2383
2384 match format.as_mut() {
2385 None => {}
2386 Some(FormatSpecifier::Bare(format)) => {
2387 purify_source_format_single(catalog, format, options, envelope, storage_configuration)
2388 .await?;
2389 }
2390
2391 Some(FormatSpecifier::KeyValue { key, value: val }) => {
2392 purify_source_format_single(catalog, key, options, envelope, storage_configuration)
2393 .await?;
2394 purify_source_format_single(catalog, val, options, envelope, storage_configuration)
2395 .await?;
2396 }
2397 }
2398 Ok(())
2399}
2400
2401async fn purify_source_format_single(
2402 catalog: &dyn SessionCatalog,
2403 format: &mut Format<Aug>,
2404 options: &SourceFormatOptions,
2405 envelope: &Option<SourceEnvelope>,
2406 storage_configuration: &StorageConfiguration,
2407) -> Result<(), PlanError> {
2408 match format {
2409 Format::Avro(schema) => match schema {
2410 AvroSchema::Csr { csr_connection } => {
2411 purify_csr_connection_avro(
2412 catalog,
2413 options,
2414 csr_connection,
2415 envelope,
2416 storage_configuration,
2417 )
2418 .await?
2419 }
2420 AvroSchema::InlineSchema { .. } => {}
2421 AvroSchema::Glue {
2422 connection,
2423 with_options,
2424 seed,
2425 } => {
2426 purify_glue_connection_avro(
2427 catalog,
2428 options,
2429 connection,
2430 with_options,
2431 seed,
2432 storage_configuration,
2433 )
2434 .await?
2435 }
2436 },
2437 Format::Protobuf(schema) => match schema {
2438 ProtobufSchema::Csr { csr_connection } => {
2439 purify_csr_connection_proto(
2440 catalog,
2441 options,
2442 csr_connection,
2443 envelope,
2444 storage_configuration,
2445 )
2446 .await?;
2447 }
2448 ProtobufSchema::InlineSchema { .. } => {}
2449 },
2450 Format::Bytes
2451 | Format::Regex(_)
2452 | Format::Json { .. }
2453 | Format::Text
2454 | Format::Csv { .. } => (),
2455 }
2456 Ok(())
2457}
2458
2459pub fn generate_subsource_statements(
2460 scx: &StatementContext,
2461 source_name: ResolvedItemName,
2462 subsources: BTreeMap<UnresolvedItemName, PurifiedSourceExport>,
2463) -> Result<Vec<CreateSubsourceStatement<Aug>>, PlanError> {
2464 if subsources.is_empty() {
2466 return Ok(vec![]);
2467 }
2468 let (_, purified_export) = subsources
2469 .iter()
2470 .next()
2471 .ok_or_else(|| internal_err!("expected at least one subsource"))?;
2472
2473 let statements = match &purified_export.details {
2474 PurifiedExportDetails::Postgres { .. } => {
2475 crate::pure::postgres::generate_create_subsource_statements(
2476 scx,
2477 source_name,
2478 subsources,
2479 )?
2480 }
2481 PurifiedExportDetails::MySql { .. } => {
2482 crate::pure::mysql::generate_create_subsource_statements(scx, source_name, subsources)?
2483 }
2484 PurifiedExportDetails::SqlServer { .. } => {
2485 crate::pure::sql_server::generate_create_subsource_statements(
2486 scx,
2487 source_name,
2488 subsources,
2489 )?
2490 }
2491 PurifiedExportDetails::LoadGenerator { .. } => {
2492 let mut subsource_stmts = Vec::with_capacity(subsources.len());
2493 for (subsource_name, purified_export) in subsources {
2494 let (desc, output) = match purified_export.details {
2495 PurifiedExportDetails::LoadGenerator { table, output } => (table, output),
2496 _ => {
2497 bail_internal!("purified export details must be load generator")
2498 }
2499 };
2500 let desc = desc.ok_or_else(|| {
2501 internal_err!(
2502 "subsources cannot be generated for single-output load generators"
2503 )
2504 })?;
2505
2506 let (columns, table_constraints) = scx.relation_desc_into_table_defs(&desc)?;
2507 let details = SourceExportStatementDetails::LoadGenerator { output };
2508 let subsource = CreateSubsourceStatement {
2510 name: subsource_name,
2511 columns,
2512 of_source: Some(source_name.clone()),
2513 constraints: table_constraints,
2518 if_not_exists: false,
2519 with_options: vec![
2520 CreateSubsourceOption {
2521 name: CreateSubsourceOptionName::ExternalReference,
2522 value: Some(WithOptionValue::UnresolvedItemName(
2523 purified_export.external_reference,
2524 )),
2525 },
2526 CreateSubsourceOption {
2527 name: CreateSubsourceOptionName::Details,
2528 value: Some(WithOptionValue::Value(Value::String(hex::encode(
2529 details.into_proto().encode_to_vec(),
2530 )))),
2531 },
2532 ],
2533 };
2534 subsource_stmts.push(subsource);
2535 }
2536
2537 subsource_stmts
2538 }
2539 PurifiedExportDetails::Kafka { .. } => {
2540 if !subsources.is_empty() {
2544 bail_internal!("Kafka sources do not produce data-bearing subsources");
2545 }
2546 vec![]
2547 }
2548 };
2549 Ok(statements)
2550}
2551
2552async fn purify_csr_connection_proto(
2553 catalog: &dyn SessionCatalog,
2554 options: &SourceFormatOptions,
2555 csr_connection: &mut CsrConnectionProtobuf<Aug>,
2556 envelope: &Option<SourceEnvelope>,
2557 storage_configuration: &StorageConfiguration,
2558) -> Result<(), PlanError> {
2559 let SourceFormatOptions::Kafka { topic } = options else {
2560 sql_bail!("Confluent Schema Registry is only supported with Kafka sources")
2561 };
2562
2563 let CsrConnectionProtobuf {
2564 seed,
2565 connection: CsrConnection {
2566 connection,
2567 options: _,
2568 },
2569 } = csr_connection;
2570 match seed {
2571 None => {
2572 let scx = StatementContext::new(None, &*catalog);
2573
2574 let ccsr_connection = match scx.get_item_by_resolved_name(connection)?.connection()? {
2575 Connection::Csr(connection) => connection.clone().into_inline_connection(catalog),
2576 _ => sql_bail!("{} is not a schema registry connection", connection),
2577 };
2578
2579 let ccsr_client = ccsr_connection
2580 .connect(storage_configuration, InTask::No)
2581 .await
2582 .map_err(|e| CsrPurificationError::ClientError(Arc::new(e)))?;
2583
2584 let value = compile_proto(&format!("{}-value", topic), &ccsr_client).await?;
2585 let key = compile_proto(&format!("{}-key", topic), &ccsr_client)
2586 .await
2587 .ok();
2588
2589 if matches!(envelope, Some(SourceEnvelope::Debezium)) && key.is_none() {
2590 sql_bail!("Key schema is required for ENVELOPE DEBEZIUM");
2591 }
2592
2593 *seed = Some(CsrSeedProtobuf { value, key });
2594 }
2595 Some(_) => (),
2596 }
2597
2598 Ok(())
2599}
2600
2601async fn purify_csr_connection_avro(
2602 catalog: &dyn SessionCatalog,
2603 options: &SourceFormatOptions,
2604 csr_connection: &mut CsrConnectionAvro<Aug>,
2605 envelope: &Option<SourceEnvelope>,
2606 storage_configuration: &StorageConfiguration,
2607) -> Result<(), PlanError> {
2608 let SourceFormatOptions::Kafka { topic } = options else {
2609 sql_bail!("Confluent Schema Registry is only supported with Kafka sources")
2610 };
2611
2612 let CsrConnectionAvro {
2613 connection: CsrConnection { connection, .. },
2614 seed,
2615 key_strategy,
2616 value_strategy,
2617 } = csr_connection;
2618 if seed.is_none() {
2619 let scx = StatementContext::new(None, &*catalog);
2620 let csr_connection = match scx.get_item_by_resolved_name(connection)?.connection()? {
2621 Connection::Csr(connection) => connection.clone().into_inline_connection(catalog),
2622 _ => sql_bail!("{} is not a schema registry connection", connection),
2623 };
2624 let ccsr_client = csr_connection
2625 .connect(storage_configuration, InTask::No)
2626 .await
2627 .map_err(|e| CsrPurificationError::ClientError(Arc::new(e)))?;
2628
2629 let Schema {
2630 key_schema,
2631 value_schema,
2632 key_reference_schemas,
2633 value_reference_schemas,
2634 } = get_remote_csr_schema(
2635 &ccsr_client,
2636 key_strategy.clone().unwrap_or_default(),
2637 value_strategy.clone().unwrap_or_default(),
2638 topic,
2639 )
2640 .await?;
2641 if matches!(envelope, Some(SourceEnvelope::Debezium)) && key_schema.is_none() {
2642 sql_bail!("Key schema is required for ENVELOPE DEBEZIUM");
2643 }
2644
2645 *seed = Some(CsrSeedAvro {
2646 key_schema,
2647 value_schema,
2648 key_reference_schemas,
2649 value_reference_schemas,
2650 })
2651 }
2652
2653 Ok(())
2654}
2655
2656async fn purify_glue_connection_avro(
2657 catalog: &dyn SessionCatalog,
2658 options: &SourceFormatOptions,
2659 connection: &ResolvedItemName,
2660 with_options: &[GlueAvroOption<Aug>],
2661 seed: &mut Option<GlueAvroSeed>,
2662 storage_configuration: &StorageConfiguration,
2663) -> Result<(), PlanError> {
2664 use crate::pure::error::GluePurificationError;
2665 let SourceFormatOptions::Kafka { .. } = options else {
2666 sql_bail!("AWS Glue Schema Registry is only supported with Kafka sources")
2667 };
2668
2669 let scx = StatementContext::new(None, &*catalog);
2670 let item = scx.get_item_by_resolved_name(connection)?;
2671 let full_name = scx.catalog.resolve_full_name(item.name());
2672 let Connection::GlueSchemaRegistry(gsr_connection) = item.connection()? else {
2679 return Err(GluePurificationError::NotGlueConnection(full_name).into());
2680 };
2681
2682 let crate::plan::statement::ddl::GlueAvroOptionExtracted {
2686 schema_name,
2687 key_schema_name,
2688 value_schema_name,
2689 key_compatibility_level,
2690 value_compatibility_level,
2691 seen: _,
2692 } = with_options.to_vec().try_into()?;
2693 if key_schema_name.is_some()
2696 || value_schema_name.is_some()
2697 || key_compatibility_level.is_some()
2698 || value_compatibility_level.is_some()
2699 {
2700 sql_bail!(
2701 "KEY SCHEMA NAME, VALUE SCHEMA NAME, KEY COMPATIBILITY LEVEL, and VALUE \
2702 COMPATIBILITY LEVEL are not supported for AWS Glue Schema Registry sources, \
2703 use SCHEMA NAME instead"
2704 );
2705 }
2706 let schema_name = schema_name.ok_or(GluePurificationError::MissingSchemaName)?;
2707
2708 if seed.is_some() {
2709 return Ok(());
2714 }
2715 let gsr_connection = gsr_connection.into_inline_connection(catalog);
2716
2717 let enforce_external_addresses = mz_storage_types::dyncfgs::ENFORCE_EXTERNAL_ADDRESSES
2720 .get(storage_configuration.config_set());
2721 let sdk_config = gsr_connection
2722 .aws_connection
2723 .connection
2724 .load_sdk_config(
2725 &storage_configuration.connection_context,
2726 gsr_connection.aws_connection.connection_id,
2727 InTask::No,
2729 enforce_external_addresses,
2730 )
2731 .await
2732 .map_err(|e| GluePurificationError::LoadSdkConfigError(Arc::new(e)))?;
2733 let glue_client = mz_aws_glue_schema_registry::ClientConfig::new(sdk_config).build();
2734
2735 let version = glue_client
2736 .get_schema_version_latest_by_name(&gsr_connection.registry_name, &schema_name)
2737 .await
2738 .map_err(|e| GluePurificationError::SchemaLookupError {
2739 registry: gsr_connection.registry_name.clone(),
2740 schema: schema_name.clone(),
2741 cause: Arc::new(e),
2742 })?;
2743 match &version.data_format {
2747 Some(mz_aws_glue_schema_registry::DataFormat::Avro) => {}
2748 other => {
2749 return Err(GluePurificationError::UnsupportedDataFormat {
2750 registry: gsr_connection.registry_name.clone(),
2751 schema: schema_name.clone(),
2752 format: other
2753 .as_ref()
2754 .map(|f| f.as_str().to_string())
2755 .unwrap_or_else(|| "<unspecified>".to_string()),
2756 }
2757 .into());
2758 }
2759 }
2760 let value_schema =
2761 version
2762 .definition
2763 .ok_or_else(|| GluePurificationError::EmptyDefinition {
2764 registry: gsr_connection.registry_name.clone(),
2765 schema: schema_name.clone(),
2766 })?;
2767
2768 *seed = Some(GlueAvroSeed { value_schema });
2769 Ok(())
2770}
2771
2772#[derive(Debug)]
2773pub struct Schema {
2774 pub key_schema: Option<String>,
2775 pub value_schema: String,
2776 pub key_reference_schemas: Vec<String>,
2778 pub value_reference_schemas: Vec<String>,
2780}
2781
2782struct SchemaWithReferences {
2784 schema: String,
2786 references: Vec<String>,
2788}
2789
2790async fn get_schema_with_strategy(
2791 client: &Client,
2792 strategy: ReaderSchemaSelectionStrategy,
2793 subject: &str,
2794) -> Result<Option<SchemaWithReferences>, PlanError> {
2795 match strategy {
2796 ReaderSchemaSelectionStrategy::Latest => {
2797 match client.get_subject_and_references(subject).await {
2799 Ok((primary, dependencies)) => Ok(Some(SchemaWithReferences {
2800 schema: primary.schema.raw,
2801 references: dependencies.into_iter().map(|s| s.schema.raw).collect(),
2802 })),
2803 Err(GetBySubjectError::SubjectNotFound)
2804 | Err(GetBySubjectError::VersionNotFound(_)) => Ok(None),
2805 Err(e) => Err(PlanError::FetchingCsrSchemaFailed {
2806 schema_lookup: format!("subject {}", subject.quoted()),
2807 cause: Arc::new(e),
2808 }),
2809 }
2810 }
2811 ReaderSchemaSelectionStrategy::Inline(raw) => Ok(Some(SchemaWithReferences {
2815 schema: raw,
2816 references: vec![],
2817 })),
2818 ReaderSchemaSelectionStrategy::ById(id) => {
2819 match client.get_subject_and_references_by_id(id).await {
2820 Ok((primary, dependencies)) => Ok(Some(SchemaWithReferences {
2821 schema: primary.schema.raw,
2822 references: dependencies.into_iter().map(|s| s.schema.raw).collect(),
2823 })),
2824 Err(GetBySubjectError::SubjectNotFound)
2825 | Err(GetBySubjectError::VersionNotFound(_)) => Ok(None),
2826 Err(e) => Err(PlanError::FetchingCsrSchemaFailed {
2827 schema_lookup: format!("subject {}", subject.quoted()),
2828 cause: Arc::new(e),
2829 }),
2830 }
2831 }
2832 }
2833}
2834
2835async fn get_remote_csr_schema(
2836 ccsr_client: &mz_ccsr::Client,
2837 key_strategy: ReaderSchemaSelectionStrategy,
2838 value_strategy: ReaderSchemaSelectionStrategy,
2839 topic: &str,
2840) -> Result<Schema, PlanError> {
2841 let value_schema_name = format!("{}-value", topic);
2842 let value_result =
2843 get_schema_with_strategy(ccsr_client, value_strategy, &value_schema_name).await?;
2844 let value_result = value_result.ok_or_else(|| anyhow!("No value schema found"))?;
2845
2846 let key_subject = format!("{}-key", topic);
2847 let key_result = get_schema_with_strategy(ccsr_client, key_strategy, &key_subject).await?;
2848 Ok(Schema {
2849 key_schema: key_result.as_ref().map(|r| r.schema.clone()),
2850 value_schema: value_result.schema,
2851 key_reference_schemas: key_result.map(|r| r.references).unwrap_or_default(),
2852 value_reference_schemas: value_result.references,
2853 })
2854}
2855
2856async fn compile_proto(
2858 subject_name: &String,
2859 ccsr_client: &Client,
2860) -> Result<CsrSeedProtobufSchema, PlanError> {
2861 let (primary_subject, dependency_subjects) = ccsr_client
2862 .get_subject_and_references(subject_name)
2863 .await
2864 .map_err(|e| PlanError::FetchingCsrSchemaFailed {
2865 schema_lookup: format!("subject {}", subject_name.quoted()),
2866 cause: Arc::new(e),
2867 })?;
2868
2869 let mut source_tree = VirtualSourceTree::new();
2871
2872 source_tree.as_mut().map_well_known_types();
2876
2877 for subject in iter::once(&primary_subject).chain(dependency_subjects.iter()) {
2878 source_tree.as_mut().add_file(
2879 Path::new(&subject.name),
2880 subject.schema.raw.as_bytes().to_vec(),
2881 );
2882 }
2883 let mut db = SourceTreeDescriptorDatabase::new(source_tree.as_mut());
2884 let fds = db
2885 .as_mut()
2886 .build_file_descriptor_set(&[Path::new(&primary_subject.name)])
2887 .map_err(|cause| PlanError::InvalidProtobufSchema { cause })?;
2888
2889 let primary_fd = fds.file(0);
2891 let message_name = match primary_fd.message_type_size() {
2892 1 => String::from_utf8_lossy(primary_fd.message_type(0).name()).into_owned(),
2893 0 => bail_unsupported!(29603, "Protobuf schemas with no messages"),
2894 _ => bail_unsupported!(29603, "Protobuf schemas with multiple messages"),
2895 };
2896
2897 let bytes = &fds
2899 .serialize()
2900 .map_err(|cause| PlanError::InvalidProtobufSchema { cause })?;
2901 let mut schema = String::new();
2902 strconv::format_bytes(&mut schema, bytes);
2903
2904 Ok(CsrSeedProtobufSchema {
2905 schema,
2906 message_name,
2907 })
2908}
2909
2910const MZ_NOW_NAME: &str = "mz_now";
2911const MZ_NOW_SCHEMA: &str = "mz_catalog";
2912
2913pub fn purify_create_materialized_view_options(
2919 catalog: impl SessionCatalog,
2920 mz_now: Option<Timestamp>,
2921 cmvs: &mut CreateMaterializedViewStatement<Aug>,
2922 resolved_ids: &mut ResolvedIds,
2923) {
2924 let (mz_now_id, mz_now_expr) = {
2927 let item = catalog
2928 .resolve_function(&PartialItemName {
2929 database: None,
2930 schema: Some(MZ_NOW_SCHEMA.to_string()),
2931 item: MZ_NOW_NAME.to_string(),
2932 })
2933 .expect("we should be able to resolve mz_now");
2934 (
2935 item.id(),
2936 Expr::Function(Function {
2937 name: ResolvedItemName::Item {
2938 id: item.id(),
2939 qualifiers: item.name().qualifiers.clone(),
2940 full_name: catalog.resolve_full_name(item.name()),
2941 print_id: false,
2942 version: RelationVersionSelector::Latest,
2943 },
2944 args: FunctionArgs::Args {
2945 args: Vec::new(),
2946 order_by: Vec::new(),
2947 },
2948 filter: None,
2949 over: None,
2950 distinct: false,
2951 }),
2952 )
2953 };
2954 let (mz_timestamp_id, mz_timestamp_type) = {
2956 let item = catalog.get_system_type("mz_timestamp");
2957 let full_name = catalog.resolve_full_name(item.name());
2958 (
2959 item.id(),
2960 ResolvedDataType::Named {
2961 id: item.id(),
2962 qualifiers: item.name().qualifiers.clone(),
2963 full_name,
2964 modifiers: vec![],
2965 print_id: true,
2966 },
2967 )
2968 };
2969
2970 let mut introduced_mz_timestamp = false;
2971
2972 for option in cmvs.with_options.iter_mut() {
2973 if matches!(
2975 option.value,
2976 Some(WithOptionValue::Refresh(RefreshOptionValue::AtCreation))
2977 ) {
2978 option.value = Some(WithOptionValue::Refresh(RefreshOptionValue::At(
2979 RefreshAtOptionValue {
2980 time: mz_now_expr.clone(),
2981 },
2982 )));
2983 }
2984
2985 if let Some(WithOptionValue::Refresh(RefreshOptionValue::Every(
2987 RefreshEveryOptionValue { aligned_to, .. },
2988 ))) = &mut option.value
2989 {
2990 if aligned_to.is_none() {
2991 *aligned_to = Some(mz_now_expr.clone());
2992 }
2993 }
2994
2995 match &mut option.value {
2998 Some(WithOptionValue::Refresh(RefreshOptionValue::At(RefreshAtOptionValue {
2999 time,
3000 }))) => {
3001 let mut visitor = MzNowPurifierVisitor::new(mz_now, mz_timestamp_type.clone());
3002 visitor.visit_expr_mut(time);
3003 introduced_mz_timestamp |= visitor.introduced_mz_timestamp;
3004 }
3005 Some(WithOptionValue::Refresh(RefreshOptionValue::Every(
3006 RefreshEveryOptionValue {
3007 interval: _,
3008 aligned_to: Some(aligned_to),
3009 },
3010 ))) => {
3011 let mut visitor = MzNowPurifierVisitor::new(mz_now, mz_timestamp_type.clone());
3012 visitor.visit_expr_mut(aligned_to);
3013 introduced_mz_timestamp |= visitor.introduced_mz_timestamp;
3014 }
3015 _ => {}
3016 }
3017 }
3018
3019 if !cmvs.with_options.iter().any(|o| {
3021 matches!(
3022 o,
3023 MaterializedViewOption {
3024 value: Some(WithOptionValue::Refresh(..)),
3025 ..
3026 }
3027 )
3028 }) {
3029 cmvs.with_options.push(MaterializedViewOption {
3030 name: MaterializedViewOptionName::Refresh,
3031 value: Some(WithOptionValue::Refresh(RefreshOptionValue::OnCommit)),
3032 })
3033 }
3034
3035 if introduced_mz_timestamp {
3039 resolved_ids.add_item(mz_timestamp_id);
3040 }
3041 let mut visitor = ExprContainsTemporalVisitor::new();
3045 visitor.visit_create_materialized_view_statement(cmvs);
3046 if !visitor.contains_temporal {
3047 resolved_ids.remove_item(&mz_now_id);
3048 }
3049}
3050
3051pub fn materialized_view_option_contains_temporal(mvo: &MaterializedViewOption<Aug>) -> bool {
3054 match &mvo.value {
3055 Some(WithOptionValue::Refresh(RefreshOptionValue::At(RefreshAtOptionValue { time }))) => {
3056 let mut visitor = ExprContainsTemporalVisitor::new();
3057 visitor.visit_expr(time);
3058 visitor.contains_temporal
3059 }
3060 Some(WithOptionValue::Refresh(RefreshOptionValue::Every(RefreshEveryOptionValue {
3061 interval: _,
3062 aligned_to: Some(aligned_to),
3063 }))) => {
3064 let mut visitor = ExprContainsTemporalVisitor::new();
3065 visitor.visit_expr(aligned_to);
3066 visitor.contains_temporal
3067 }
3068 Some(WithOptionValue::Refresh(RefreshOptionValue::Every(RefreshEveryOptionValue {
3069 interval: _,
3070 aligned_to: None,
3071 }))) => {
3072 true
3075 }
3076 Some(WithOptionValue::Refresh(RefreshOptionValue::AtCreation)) => {
3077 true
3079 }
3080 _ => false,
3081 }
3082}
3083
3084struct ExprContainsTemporalVisitor {
3086 pub contains_temporal: bool,
3087}
3088
3089impl ExprContainsTemporalVisitor {
3090 pub fn new() -> ExprContainsTemporalVisitor {
3091 ExprContainsTemporalVisitor {
3092 contains_temporal: false,
3093 }
3094 }
3095}
3096
3097impl Visit<'_, Aug> for ExprContainsTemporalVisitor {
3098 fn visit_function(&mut self, func: &Function<Aug>) {
3099 self.contains_temporal |= func.name.full_item_name().item == MZ_NOW_NAME;
3100 visit_function(self, func);
3101 }
3102}
3103
3104struct MzNowPurifierVisitor {
3105 pub mz_now: Option<Timestamp>,
3106 pub mz_timestamp_type: ResolvedDataType,
3107 pub introduced_mz_timestamp: bool,
3108}
3109
3110impl MzNowPurifierVisitor {
3111 pub fn new(
3112 mz_now: Option<Timestamp>,
3113 mz_timestamp_type: ResolvedDataType,
3114 ) -> MzNowPurifierVisitor {
3115 MzNowPurifierVisitor {
3116 mz_now,
3117 mz_timestamp_type,
3118 introduced_mz_timestamp: false,
3119 }
3120 }
3121}
3122
3123impl VisitMut<'_, Aug> for MzNowPurifierVisitor {
3124 fn visit_expr_mut(&mut self, expr: &'_ mut Expr<Aug>) {
3125 match expr {
3126 Expr::Function(Function {
3127 name:
3128 ResolvedItemName::Item {
3129 full_name: FullItemName { item, .. },
3130 ..
3131 },
3132 ..
3133 }) if item == &MZ_NOW_NAME.to_string() => {
3134 let mz_now = self.mz_now.expect(
3135 "we should have chosen a timestamp if the expression contains mz_now()",
3136 );
3137 *expr = Expr::Cast {
3140 expr: Box::new(Expr::Value(Value::Number(mz_now.to_string()))),
3141 data_type: self.mz_timestamp_type.clone(),
3142 };
3143 self.introduced_mz_timestamp = true;
3144 }
3145 _ => visit_expr_mut(self, expr),
3146 }
3147 }
3148}