Skip to main content

mz_sql/
pure.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! SQL purification.
11//!
12//! See the [crate-level documentation](crate) for details.
13
14use 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
129/// Generates a subsource name by prepending source schema name if present
130///
131/// For eg. if source is `a.b`, then `a` will be prepended to the subsource name
132/// so that it's generated in the same schema as source
133fn 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
142// TODO(database-issues#8620): Remove once subsources are removed
143/// Validates the requested subsources do not have name conflicts with each other
144/// and that the same upstream table is not referenced multiple times.
145fn validate_source_export_names<T>(
146    requested_source_exports: &[RequestedSourceExport<T>],
147) -> Result<(), PlanError> {
148    // This condition would get caught during the catalog transaction, but produces a
149    // vague, non-contextual error. Instead, error here so we can suggest to the user
150    // how to fix the problem.
151    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    // We disallow subsource statements from referencing the same upstream table, but we allow
178    // `CREATE TABLE .. FROM SOURCE` statements to do so. Since `CREATE TABLE .. FROM SOURCE`
179    // purification will only provide 1 `requested_source_export`, we can leave this here without
180    // needing to differentiate between the two types of statements.
181    // TODO(roshan): Remove this when auto-generated subsources are deprecated.
182    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        // The progress subsource, if we are offloading progress info to a separate relation
212        create_progress_subsource_stmt: Option<CreateSubsourceStatement<Aug>>,
213        create_source_stmt: CreateSourceStatement<Aug>,
214        // Map of subsource names to external details
215        subsources: BTreeMap<UnresolvedItemName, PurifiedSourceExport>,
216        /// All the available upstream references that can be added as tables
217        /// to this primary source.
218        available_source_references: SourceReferences,
219    },
220    PurifiedAlterSource {
221        alter_source_stmt: AlterSourceStatement<Aug>,
222    },
223    PurifiedAlterSourceAddSubsources {
224        // This just saves us an annoying catalog lookup
225        source_name: ResolvedItemName,
226        /// Options that we will need the values of to update the source's
227        /// definition.
228        options: Vec<AlterSourceAddSubsourceOption<Aug>>,
229        // Map of subsource names to external details
230        subsources: BTreeMap<UnresolvedItemName, PurifiedSourceExport>,
231    },
232    PurifiedAlterSourceRefreshReferences {
233        source_name: ResolvedItemName,
234        /// The updated available upstream references for the primary source.
235        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/// The existing source a statement drives, and how the statement uses it.
279///
280/// Resolved from the statement rather than from `ResolvedIds`: `ALTER SOURCE`
281/// carries its target as an `UnresolvedItemName` that name resolution never
282/// records.
283#[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
297/// Resolves the existing source that `stmt` drives, mirroring how the
298/// purification paths below resolve it themselves.
299///
300/// `None` means the statement drives no such source: it names its connection
301/// outright (`CREATE SOURCE`, `CREATE SINK`), or the name does not resolve to a
302/// source. Purification reports the latter itself, and bails on it before
303/// opening any connection. A [`Statement`] variant that instead reaches
304/// upstream through an existing catalog item needs an arm here, or callers see
305/// no source at all.
306pub 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
327/// Purifies a statement, removing any dependencies on external state.
328///
329/// See the section on [purification](crate#purification) in the crate
330/// documentation for details.
331///
332/// Note that this doesn't handle CREATE MATERIALIZED VIEW, which is
333/// handled by [purify_create_materialized_view_options] instead.
334/// This could be made more consistent by a refactoring discussed here:
335/// <https://github.com/MaterializeInc/materialize/pull/23870#discussion_r1435922709>
336pub 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
375/// Injects `DOC ON` comments into all Avro formats that are using a schema
376/// registry by finding all SQL `COMMENT`s that are attached to the sink's
377/// underlying materialized view (as well as any types referenced by that
378/// underlying materialized view).
379pub(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    // Collect all objects referenced by the sink.
385    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    // Collect all Avro formats that use a schema registry, as well as a set of
394    // all identifiers named in user-provided `DOC ON` options.
395    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 each Avro format in the sink, inject the appropriate `DOC ON` options
413    // for each item referenced by the sink.
414    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        // Adding existing comments if not already provided by user
424        for object_id in &object_ids {
425            // Always add comments to the latest version of the item.
426            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                // Attach comment for the item itself, if the user has not
439                // already provided an overriding `DOC ON` option for the item.
440                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                // Attach comment for each column in the item, if the user has
456                // not already provided an overriding `DOC ON` option for the
457                // column.
458                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
494/// Checks that the sink described in the statement can connect to its external
495/// resources.
496async fn purify_create_sink(
497    catalog: impl SessionCatalog,
498    mut create_sink_stmt: CreateSinkStatement<Aug>,
499    storage_configuration: &StorageConfiguration,
500) -> Result<PurifiedStatement, PlanError> {
501    // General purification
502    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    // The list of options that the user is allowed to specify.
515    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            // We must not leave any state behind in the Kafka broker, so just ensure that
538            // we can connect. This means we don't ensure that we can create the topic and
539            // introduces TOCTOU errors, but creating an inoperable sink is infinitely
540            // preferable to leaking state in users' environments.
541            let scx = StatementContext::new(None, &catalog);
542            let connection = {
543                let item = scx.get_item_by_resolved_name(connection)?;
544                // Get Kafka connection
545                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                    // anyhow doesn't support Clone, so not trivial to move into PlanError
572                    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                // Get Iceberg connection
603                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            // For S3 Tables connections in the Materialize Cloud product, verify the
615            // AWS region matches the environment's region. This check only applies when
616            // the enable_s3_tables_region_check dyncfg is set.
617            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                        // Later on we default to "us-east-1" if the region is not set on the S3 Tables
625                        // connection, so we need to do the same check here.
626                        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            // Validate the sink's (optional) AWS connection even though we never use it.
643            // TODO(kynan): If we do start using the sink's creds, check again that this validation
644            //   accurately reflects what we need.
645            //   Consider rolling the storage creds validation into the catalog connection's "connect" fn,
646            //   which already validates the catalog creds (currently also used for the storage layer).
647            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                    // Get AWS connection
652                    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            // Now that we've validated the sink's storage creds (if they exist)
674            // we _could_ use them to build a complete Iceberg client (both catalog and storage).
675            // TODO(kynan): Actually use those sink-specific creds here instead of ignoring them.
676            // Purification only proves the catalog is reachable, so it needs no table-scoped
677            // storage credentials.
678            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            // Get Kafka connection
710            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
734/// Runs a function on each format within the format specifier.
735///
736/// The function is provided with a `DocOnSchema` that indicates whether the
737/// format is for the key, value, or both.
738///
739// TODO(benesch): rename `DocOnSchema` to the more general `FormatRestriction`.
740fn 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/// Defines whether purification should enforce that at least one valid source
755/// reference is provided on the provided statement.
756#[derive(Debug, Copy, Clone, PartialEq, Eq)]
757pub(crate) enum SourceReferencePolicy {
758    /// Don't allow source references to be provided. This is used for
759    /// enforcing that `CREATE SOURCE` statements don't create subsources.
760    NotAllowed,
761    /// Allow empty references, such as when creating a source that
762    /// will have tables added afterwards.
763    Optional,
764    /// Require that at least one reference is resolved to an upstream
765    /// object.
766    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    // Depending on if the user must or can use the `CREATE TABLE .. FROM SOURCE` statement
823    // to add tables after this source is created we might need to enforce that
824    // auto-generated subsources are created or not created by this source statement.
825    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                // Get Kafka connection
854                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                    // anyhow doesn't support Clone, so not trivial to move into PlanError
881                    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                    // Validate that the topic at least exists.
893                    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                    // Validate the start offsets.
908                    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                    // Translate `START TIMESTAMP` to a start offset.
921                    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            // Record whether the upstream server is a physical replica, which changes how we can determine
994            // the latest LSN.
995            let is_physical_replica = mz_postgres_util::get_is_in_recovery(&client).await?;
996
997            // Read after the references above, so it bounds the LSN their schemas belong to.
998            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            // Record the active replication timeline_id to allow detection of a future upstream
1029            // point-in-time-recovery that will put the source into an error state.
1030            let timeline_id = mz_postgres_util::get_timeline_id(&client).await?;
1031
1032            // Remove any old detail references
1033            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            // Connect to the upstream SQL Server instance so we can validate
1055            // we're compatible with CDC.
1056            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            // Update our set of requested source exports.
1111            requested_subsource_map.extend(source_exports);
1112
1113            // Record the most recent restore_history_id, or none if the system has never been
1114            // restored.
1115
1116            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            // Update our 'TEXT' and 'EXCLUDE' column options with the purified and normalized set.
1131            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            // Retrieve the current @gtid_executed value of the server to mark as the effective
1174            // initial snapshot point such that we can ensure consistency if the initial source
1175            // snapshot is broken up over multiple points in time.
1176            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            // We don't have any fields in this details struct but keep this around for
1206            // conformity with postgres and in-case we end up needing it again in the future.
1207            let details = MySqlSourceDetails {};
1208            // Update options with the purified details
1209            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            // Filter to the references that need to be created as 'subsources', which
1241            // doesn't include the default output for single-output sources.
1242            // TODO(database-issues#8620): Remove once subsources are removed
1243            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    // Now that we know which subsources to create alongside this
1328    // statement, remove the references so it is not canonicalized as
1329    // part of the `CREATE SOURCE` statement in the catalog.
1330    *external_references = None;
1331
1332    // Generate progress subsource for old syntax
1333    let create_progress_subsource_stmt = if uses_old_syntax {
1334        // Take name from input or generate name
1335        let name = match progress_subsource {
1336            Some(name) => match name {
1337                DeferredItemName::Deferred(name) => name.clone(),
1338                // Already checked for this value above.
1339                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        // Create the subsource statement
1374        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            // Progress subsources do not refer to the source to which they belong.
1393            // Instead the primary source depends on it (the opposite is true of
1394            // ingestion exports, which depend on the primary source).
1395            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
1421/// On success, returns the details on new subsources and updated
1422/// 'options' that sequencing expects for handling `ALTER SOURCE` statements.
1423async 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    // Get name.
1436    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    // Ensure it's an ingestion-based and alterable source.
1451    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
1512// TODO(database-issues#8620): Remove once subsources are removed
1513/// Equivalent to `purify_create_source` but for `AlterSourceStatement`.
1514async 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    // Validate this is a source that can have subsources added.
1524    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            // Get PostgresConnection for generating subsources.
1549            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            // Read after the references above, so it bounds the LSN their schemas belong to.
1563            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            // Retrieve the current @gtid_executed value of the server to mark as the effective
1617            // initial snapshot point for these subsources.
1618            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            // Update options with the purified details
1650            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            // Open a connection to the upstream SQL Server instance.
1666            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            // Query the upstream SQL Server instance for available tables to replicate.
1677            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            // Add the new exports to our subsource map.
1708            requested_subsource_map.extend(source_exports);
1709
1710            // Update options on the CREATE SOURCE statement with the purified details.
1711            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        // The source kind was already established earlier in this function;
1726        // reaching this catch-all is an internal invariant violation.
1727        _ => 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            // Get PostgresConnection for generating subsources.
1745            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            // Open a connection to the upstream SQL Server instance.
1793            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            // Query the upstream SQL Server instance for available tables to replicate.
1804            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    // Columns and constraints cannot be specified by the user but will be populated below.
1851    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    // Get the source item
1861    let item = match scx.get_item_by_resolved_name(source_name) {
1862        Ok(item) => item,
1863        Err(e) => return Err(e),
1864    };
1865
1866    // Ensure it's an ingestion-based and alterable source.
1867    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    // Our text column values are unqualified (just column names), but the purification methods below
1898    // expect to match the fully-qualified names against the full set of tables in upstream, so we
1899    // need to qualify them using the external reference first.
1900    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    // Should be overriden below if a source-specific format is required.
1924    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    // Run purification work specific to each source type: resolve the external reference to
1945    // a fully qualified name and obtain the appropriate details for the source-export statement
1946    let purified_export = match desc.connection {
1947        GenericSourceConnection::Postgres(pg_source_connection) => {
1948            // Get PostgresConnection for generating subsources.
1949            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            // Read after the references above, so it bounds the LSN their schemas belong to.
1963            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                // TODO(database-issues#8620): Remove once subsources are removed
1976                // This `normalized_text_columns` is not relevant for us and is only returned for
1977                // `CREATE SOURCE` statements that automatically generate subsources
1978                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            // There should be exactly one source_export returned for this statement
1993            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            // Retrieve the current @gtid_executed value of the server to mark as the effective
2017            // initial snapshot point for this table.
2018            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                // TODO(database-issues#8620): Remove once subsources are removed
2030                // `normalized_text/exclude_columns` is not relevant for us and is only returned for
2031                // `CREATE SOURCE` statements that automatically generate subsources
2032                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            // There should be exactly one source_export returned for this statement
2047            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            // There should be exactly one source_export returned for this statement
2086            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            // There should be exactly one source_export returned
2098            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            // There should be exactly one source_export returned
2123            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    // Update the external reference in the statement to the resolved fully-qualified
2145    // external reference
2146    *external_reference = Some(purified_export.external_reference.clone());
2147
2148    // Update options in the statement using the purified export details
2149    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                    // Purification is what populates the `Defined` variant;
2191                    // reaching this is an internal invariant violation.
2192                    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                    // Purification is what populates the `Defined` variant;
2241                    // reaching this is an internal invariant violation.
2242                    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                    // Purification is what populates the `Defined` variant;
2298                    // reaching this is an internal invariant violation.
2299                    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            // We only determine the table description for multi-output load generator sources here,
2318            // whereas single-output load generators will have their relation description
2319            // determined during statment planning as envelope and format options may affect their
2320            // schema.
2321            if let Some(desc) = desc {
2322                let (gen_columns, gen_constraints) = scx.relation_desc_into_table_defs(&desc)?;
2323                match columns {
2324                    // Purification is what populates the `Defined` variant;
2325                    // reaching this is an internal invariant violation.
2326                    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            // NOTE: Kafka tables have their 'schemas' purified into the statement inside the
2348            // format field, so we don't specify any columns or constraints to be stored
2349            // on the statement here. The RelationDesc will be determined during planning.
2350            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    // TODO: We might as well use the retrieved available references to update the source
2361    // available references table in the catalog, so plumb this through.
2362    // available_source_references: retrieved_source_references.available_source_references(),
2363    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    // get the first subsource to determine the connection type
2465    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                // Create the subsource statement
2509                let subsource = CreateSubsourceStatement {
2510                    name: subsource_name,
2511                    columns,
2512                    of_source: Some(source_name.clone()),
2513                    // unlike sources that come from an external upstream, we
2514                    // have more leniency to introduce different constraints
2515                    // every time the load generator is run; i.e. we are not as
2516                    // worried about introducing junk data.
2517                    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            // TODO: as part of database-issues#8322, Kafka sources will begin
2541            // producing data––we'll need to understand the schema
2542            // of the output here.
2543            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    // Structural validation below runs unconditionally — even when a seed is
2673    // already present. The seed is not proof the rest of the statement is
2674    // well-formed: the grammar accepts a user-written `SEED VALUE SCHEMA`, so a
2675    // seeded statement may have arrived straight from a user rather than from
2676    // our own re-parsed `create_sql`. Only the Glue *fetch* is safe to skip
2677    // when seeded; the connection-type and `SCHEMA NAME` checks are not.
2678    let Connection::GlueSchemaRegistry(gsr_connection) = item.connection()? else {
2679        return Err(GluePurificationError::NotGlueConnection(full_name).into());
2680    };
2681
2682    // Pull `SCHEMA NAME` out of the option bag. Required. Use the shared
2683    // extractor (rather than a hand-rolled match) so non-string values get a
2684    // proper "invalid value" error instead of being silently dropped.
2685    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    // The per-side and compatibility options are sink-only. A source reads a
2694    // single writer schema, so reject them here rather than silently ignoring.
2695    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        // The schema is already pinned — either we purified this statement
2710        // before and re-parsed it from persisted create_sql, or a user supplied
2711        // the seed directly. Either way the reader schema is fixed, so skip the
2712        // Glue lookup. (The structural checks above still ran.)
2713        return Ok(());
2714    }
2715    let gsr_connection = gsr_connection.into_inline_connection(catalog);
2716
2717    // Build the SDK config the same way the storage decoder does at runtime
2718    // — same auth, region, and endpoint override resolution.
2719    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            // We are in a normal tokio context during purification.
2728            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    // The runtime decode path only handles Avro; reject other formats at
2744    // planning time so the failure mode is a clear SQL error, not a
2745    // permanent decode error on every record.
2746    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    /// Reference schemas for the key schema, in dependency order.
2777    pub key_reference_schemas: Vec<String>,
2778    /// Reference schemas for the value schema, in dependency order.
2779    pub value_reference_schemas: Vec<String>,
2780}
2781
2782/// Result of fetching a schema, including any referenced schemas.
2783struct SchemaWithReferences {
2784    /// The primary schema.
2785    schema: String,
2786    /// Reference schemas in dependency order.
2787    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            // Use get_subject_and_references to also fetch referenced schemas
2798            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        // It's possible that CSR was provided with an inline strategy, but without a subject
2812        // or schema id to look up, there isn't a clean way to lookup references for the schema.
2813        // This could be done with extra work, but at this point, it isn't clear this is needed.
2814        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
2856/// Collect protobuf message descriptor from CSR and compile the descriptor.
2857async 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    // Compile .proto files into a file descriptor set.
2870    let mut source_tree = VirtualSourceTree::new();
2871
2872    // Add well-known types (e.g., google/protobuf/timestamp.proto) to the source
2873    // tree. These are implicitly available to protoc but are typically not
2874    // registered in the schema registry.
2875    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    // Ensure there is exactly one message in the file.
2890    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    // Encode the file descriptor set into a SQL byte string.
2898    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
2913/// Purifies a CREATE MATERIALIZED VIEW statement. Additionally, it adjusts `resolved_ids` if
2914/// references to ids appear or disappear during the purification.
2915///
2916/// Note that in contrast with [`purify_statement`], this doesn't need to be async, because
2917/// this function is not making any network calls.
2918pub 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    // 0. Preparations:
2925    // Prepare an expression that calls `mz_now()`, which we can insert in various later steps.
2926    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    // Prepare the `mz_timestamp` type.
2955    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        // 1. Purify `REFRESH AT CREATION` to `REFRESH AT mz_now()`.
2974        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        // 2. If `REFRESH EVERY` doesn't have an `ALIGNED TO`, then add `ALIGNED TO mz_now()`.
2986        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        // 3. Substitute `mz_now()` with the timestamp chosen for the CREATE MATERIALIZED VIEW
2996        // statement. (This has to happen after the above steps, which might introduce `mz_now()`.)
2997        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    // 4. If the user didn't give any REFRESH option, then default to ON COMMIT.
3020    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    // 5. Attend to `resolved_ids`: The purification might have
3036    // - added references to `mz_timestamp`;
3037    // - removed references to `mz_now`.
3038    if introduced_mz_timestamp {
3039        resolved_ids.add_item(mz_timestamp_id);
3040    }
3041    // Even though we always remove `mz_now()` from the `with_options`, there might be `mz_now()`
3042    // remaining in the main query expression of the MV, so let's visit the entire statement to look
3043    // for `mz_now()` everywhere.
3044    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
3051/// Returns true if the [MaterializedViewOption] either already involves `mz_now()` or will involve
3052/// after purification.
3053pub 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            // For a `REFRESH EVERY` without an `ALIGNED TO`, purification will default the
3073            // `ALIGNED TO` to `mz_now()`.
3074            true
3075        }
3076        Some(WithOptionValue::Refresh(RefreshOptionValue::AtCreation)) => {
3077            // `REFRESH AT CREATION` will be purified to `REFRESH AT mz_now()`.
3078            true
3079        }
3080        _ => false,
3081    }
3082}
3083
3084/// Determines whether the AST involves `mz_now()`.
3085struct 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                // We substitute `mz_now()` with number + a cast to `mz_timestamp`. The cast is to
3138                // not alter the type of the expression.
3139                *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}