1use std::collections::{BTreeMap, BTreeSet};
16use std::fmt::Write;
17use std::iter;
18use std::num::NonZeroU32;
19use std::time::Duration;
20
21use chrono::DateTime;
22use itertools::Itertools;
23use mz_adapter_types::compaction::{CompactionWindow, DEFAULT_LOGICAL_COMPACTION_WINDOW_DURATION};
24use mz_arrow_util::builder::ArrowBuilder;
25use mz_auth::password::Password;
26use mz_controller_types::{ClusterId, DEFAULT_REPLICA_LOGGING_INTERVAL, ReplicaId};
27use mz_expr::{CollectionPlan, UnmaterializableFunc};
28use mz_interchange::avro::{AvroSchemaGenerator, DocTarget};
29use mz_ore::cast::{CastFrom, TryCastFrom};
30use mz_ore::collections::{CollectionExt, HashSet};
31use mz_ore::num::NonNeg;
32use mz_ore::str::StrExt;
33use mz_ore::{soft_assert_or_log, soft_panic_or_log};
34use mz_proto::RustType;
35use mz_repr::adt::interval::Interval;
36use mz_repr::adt::mz_acl_item::{MzAclItem, PrivilegeMap};
37use mz_repr::adt::timestamp::CheckedTimestamp;
38use mz_repr::network_policy_id::NetworkPolicyId;
39use mz_repr::optimize::OptimizerFeatureOverrides;
40use mz_repr::refresh_schedule::{RefreshEvery, RefreshSchedule};
41use mz_repr::role_id::RoleId;
42use mz_repr::{
43 CatalogItemId, ColumnName, RelationDesc, RelationVersion, RelationVersionSelector,
44 SqlColumnType, SqlRelationType, SqlScalarType, Timestamp, VersionedRelationDesc,
45 preserves_order, strconv,
46};
47use mz_sql_parser::ast::{
48 self, AlterClusterAction, AlterClusterStatement, AlterConnectionAction, AlterConnectionOption,
49 AlterConnectionOptionName, AlterConnectionStatement, AlterIndexAction, AlterIndexStatement,
50 AlterMaterializedViewApplyReplacementStatement, AlterNetworkPolicyStatement,
51 AlterObjectRenameStatement, AlterObjectSwapStatement, AlterRetainHistoryStatement,
52 AlterRoleOption, AlterRoleStatement, AlterSecretStatement, AlterSetClusterStatement,
53 AlterSinkAction, AlterSinkStatement, AlterSourceAction, AlterSourceAddSubsourceOption,
54 AlterSourceAddSubsourceOptionName, AlterSourceStatement, AlterSystemResetAllStatement,
55 AlterSystemResetStatement, AlterSystemSetStatement, AlterTableAddColumnStatement, AvroSchema,
56 AvroSchemaOption, AvroSchemaOptionName, ClusterAlterOption, ClusterAlterOptionName,
57 ClusterAlterOptionValue, ClusterAlterUntilReadyOption, ClusterAlterUntilReadyOptionName,
58 ClusterAutoScalingStrategyOptionValue, ClusterFeature, ClusterFeatureName, ClusterOption,
59 ClusterOptionName, ClusterScheduleOptionValue, ColumnDef, ColumnOption, CommentObjectType,
60 CommentStatement, ConnectionOption, ConnectionOptionName, CreateClusterReplicaStatement,
61 CreateClusterStatement, CreateConnectionOption, CreateConnectionOptionName,
62 CreateConnectionStatement, CreateConnectionType, CreateDatabaseStatement, CreateIndexStatement,
63 CreateMaterializedViewStatement, CreateMetricSinkOption, CreateMetricSinkOptionName,
64 CreateMetricSinkStatement, CreateNetworkPolicyStatement, CreateRoleStatement,
65 CreateSchemaStatement, CreateSecretStatement, CreateSinkConnection, CreateSinkOption,
66 CreateSinkOptionName, CreateSinkStatement, CreateSourceConnection, CreateSourceOption,
67 CreateSourceOptionName, CreateSourceStatement, CreateSubsourceOption,
68 CreateSubsourceOptionName, CreateSubsourceStatement, CreateTableFromSourceStatement,
69 CreateTableStatement, CreateTypeAs, CreateTypeListOption, CreateTypeListOptionName,
70 CreateTypeMapOption, CreateTypeMapOptionName, CreateTypeStatement, CreateViewStatement,
71 CreateWebhookSourceStatement, CsrConfigOption, CsrConfigOptionName, CsrConnection,
72 CsrConnectionAvro, CsrConnectionProtobuf, CsrSeedProtobuf, CsvColumns, DeferredItemName,
73 DocOnIdentifier, DocOnSchema, DropObjectsStatement, DropOwnedStatement, Expr, Format,
74 FormatSpecifier, GlueAvroOption, GlueAvroOptionName, IcebergSinkConfigOption, Ident,
75 IfExistsBehavior, IndexOption, IndexOptionName, KafkaSinkConfigOption, KeyConstraint,
76 LoadGeneratorOption, LoadGeneratorOptionName, MaterializedViewOption,
77 MaterializedViewOptionName, MySqlConfigOption, MySqlConfigOptionName, NetworkPolicyOption,
78 NetworkPolicyOptionName, NetworkPolicyRuleDefinition, NetworkPolicyRuleOption,
79 NetworkPolicyRuleOptionName, OnHydrationOptionValue, PgConfigOption, PgConfigOptionName,
80 ProtobufSchema, QualifiedReplica, RefreshAtOptionValue, RefreshEveryOptionValue,
81 RefreshOptionValue, ReplicaDefinition, ReplicaOption, ReplicaOptionName, RoleAttribute,
82 SetRoleVar, SourceErrorPolicy, SourceIncludeMetadata, SqlServerConfigOption,
83 SqlServerConfigOptionName, Statement, TableConstraint, TableFromSourceColumns,
84 TableFromSourceOption, TableFromSourceOptionName, TableOption, TableOptionName,
85 UnresolvedDatabaseName, UnresolvedItemName, UnresolvedObjectName, UnresolvedSchemaName, Value,
86 ViewDefinition, WithOptionValue,
87};
88use mz_sql_parser::ident;
89use mz_sql_parser::parser::StatementParseResult;
90use mz_storage_types::connections::inline::ReferencedConnection;
91use mz_storage_types::connections::{Connection, KafkaTopicOptions};
92use mz_storage_types::sinks::{
93 IcebergSinkConnection, KafkaIdStyle, KafkaSinkConnection, KafkaSinkFormat, KafkaSinkFormatType,
94 SinkEnvelope, StorageSinkConnection, iceberg_type_overrides,
95};
96use mz_storage_types::sources::encoding::{
97 AvroEncoding, ColumnSpec, CsvEncoding, DataEncoding, ProtobufEncoding, RegexEncoding,
98 SourceDataEncoding, included_column_desc,
99};
100use mz_storage_types::sources::envelope::{
101 KeyEnvelope, NoneEnvelope, SourceEnvelope, UnplannedSourceEnvelope, UpsertStyle,
102};
103use mz_storage_types::sources::kafka::{
104 KafkaMetadataKind, KafkaSourceConnection, KafkaSourceExportDetails, kafka_metadata_columns_desc,
105};
106use mz_storage_types::sources::load_generator::{
107 KeyValueLoadGenerator, LOAD_GENERATOR_KEY_VALUE_OFFSET_DEFAULT, LoadGenerator,
108 LoadGeneratorOutput, LoadGeneratorSourceConnection, LoadGeneratorSourceExportDetails,
109};
110use mz_storage_types::sources::mysql::{
111 MySqlSourceConnection, MySqlSourceDetails, ProtoMySqlSourceDetails,
112};
113use mz_storage_types::sources::postgres::{
114 PostgresSourceConnection, PostgresSourcePublicationDetails,
115 ProtoPostgresSourcePublicationDetails,
116};
117use mz_storage_types::sources::sql_server::{
118 ProtoSqlServerSourceExtras, SqlServerSourceExportDetails,
119};
120use mz_storage_types::sources::{
121 GenericSourceConnection, MySqlSourceExportDetails, PostgresSourceExportDetails,
122 ProtoSourceExportStatementDetails, SourceConnection, SourceDesc, SourceExportDataConfig,
123 SourceExportDetails, SourceExportStatementDetails, SqlServerSourceConnection,
124 SqlServerSourceExtras, Timeline,
125};
126use mz_storage_types::wire_format::WireFormat;
127use prost::Message;
128
129use crate::ast::display::AstDisplay;
130use crate::catalog::{
131 CatalogCluster, CatalogDatabase, CatalogError, CatalogItem, CatalogItemType,
132 CatalogRecordField, CatalogType, CatalogTypeDetails, ObjectType, SystemObjectType,
133};
134use crate::iceberg::IcebergSinkConfigOptionExtracted;
135use crate::kafka_util::{KafkaSinkConfigOptionExtracted, KafkaSourceConfigOptionExtracted};
136use crate::names::{
137 Aug, CommentObjectId, DatabaseId, DependencyIds, ObjectId, PartialItemName, QualifiedItemName,
138 ResolvedClusterName, ResolvedColumnReference, ResolvedDataType, ResolvedDatabaseSpecifier,
139 ResolvedItemName, ResolvedNetworkPolicyName, SchemaSpecifier, SystemObjectId,
140};
141use crate::normalize::{self, ident};
142use crate::plan::error::PlanError;
143use crate::plan::query::{
144 ExprContext, QueryLifetime, TypeResolutionBudget, plan_expr, scalar_type_from_sql,
145};
146use crate::plan::scope::Scope;
147use crate::plan::statement::ddl::connection::{INALTERABLE_OPTIONS, MUTUALLY_EXCLUSIVE_SETS};
148use crate::plan::statement::{StatementContext, StatementDesc, scl};
149use crate::plan::typeconv::CastContext;
150use crate::plan::with_options::{OptionalDuration, OptionalString, TryFromValue};
151use crate::plan::{
152 AlterClusterPlan, AlterClusterPlanStrategy, AlterClusterRenamePlan,
153 AlterClusterReplicaRenamePlan, AlterClusterSwapPlan, AlterConnectionPlan, AlterItemRenamePlan,
154 AlterMaterializedViewApplyReplacementPlan, AlterNetworkPolicyPlan, AlterNoopPlan,
155 AlterOptionParameter, AlterRetainHistoryPlan, AlterRolePlan, AlterSchemaRenamePlan,
156 AlterSchemaSwapPlan, AlterSecretPlan, AlterSetClusterPlan, AlterSinkPlan,
157 AlterSourceTimestampIntervalPlan, AlterSystemResetAllPlan, AlterSystemResetPlan,
158 AlterSystemSetPlan, AlterTablePlan, AutoScalingStrategy, ClusterSchedule, CommentPlan,
159 ComputeReplicaConfig, ComputeReplicaIntrospectionConfig, ConnectionDetails,
160 CreateClusterManagedPlan, CreateClusterPlan, CreateClusterReplicaPlan,
161 CreateClusterUnmanagedPlan, CreateClusterVariant, CreateConnectionPlan, CreateDatabasePlan,
162 CreateIndexPlan, CreateMaterializedViewPlan, CreateMetricSinkPlan, CreateNetworkPolicyPlan,
163 CreateRolePlan, CreateSchemaPlan, CreateSecretPlan, CreateSinkPlan, CreateSourcePlan,
164 CreateTablePlan, CreateTypePlan, CreateViewPlan, DataSourceDesc, DropObjectsPlan,
165 DropOwnedPlan, HirRelationExpr, Index, MaterializedView, MetricSink, NetworkPolicyRule,
166 NetworkPolicyRuleAction, NetworkPolicyRuleDirection, OnHydration, Plan, PlanClusterOption,
167 PlanNotice, PolicyAddress, QueryContext, ReplicaConfig, Secret, Sink, Source, Table,
168 TableDataSource, Type, VariableValue, View, WebhookBodyFormat, WebhookHeaderFilters,
169 WebhookHeaders, WebhookValidation, literal, plan_utils, query, transform_ast,
170};
171use crate::session::vars::{
172 self, ENABLE_AUTO_SCALING_STRATEGY, ENABLE_CLUSTER_SCHEDULE_REFRESH,
173 ENABLE_COLLECTION_PARTITION_BY, ENABLE_CREATE_TABLE_FROM_SOURCE, ENABLE_KAFKA_SINK_HEADERS,
174 ENABLE_METRIC_SINK, ENABLE_REFRESH_EVERY_MVS, ENABLE_REPLICA_TARGETED_MATERIALIZED_VIEWS,
175 VarInput,
176};
177use crate::{names, parse};
178
179mod connection;
180
181const MAX_NUM_COLUMNS: usize = 256;
186
187const MAX_KAFKA_TOPIC_METADATA_REFRESH_INTERVAL: Duration = Duration::from_secs(60 * 60);
188const MIN_KAFKA_TOPIC_METADATA_REFRESH_INTERVAL: Duration = Duration::from_secs(1);
189
190static MANAGED_REPLICA_PATTERN: std::sync::LazyLock<regex::Regex> =
191 std::sync::LazyLock::new(|| regex::Regex::new(r"^r(\d)+$").unwrap());
192
193fn check_partition_by(desc: &RelationDesc, mut partition_by: Vec<Ident>) -> Result<(), PlanError> {
197 if partition_by.len() > desc.len() {
198 tracing::error!(
199 "PARTITION BY contains more columns than the relation. (expected at most {}, got {})",
200 desc.len(),
201 partition_by.len()
202 );
203 partition_by.truncate(desc.len());
204 }
205
206 let desc_prefix = desc.iter().take(partition_by.len());
207 for (idx, ((desc_name, desc_type), partition_name)) in
208 desc_prefix.zip_eq(partition_by).enumerate()
209 {
210 let partition_name = normalize::column_name(partition_name);
211 if *desc_name != partition_name {
212 sql_bail!(
213 "PARTITION BY columns should be a prefix of the relation's columns (expected {desc_name} at index {idx}, got {partition_name})"
214 );
215 }
216 if !preserves_order(&desc_type.scalar_type) {
217 sql_bail!("PARTITION BY column {partition_name} has unsupported type");
218 }
219 }
220 Ok(())
221}
222
223pub fn describe_create_database(
224 _: &StatementContext,
225 _: CreateDatabaseStatement,
226) -> Result<StatementDesc, PlanError> {
227 Ok(StatementDesc::new(None))
228}
229
230pub fn plan_create_database(
231 _: &StatementContext,
232 CreateDatabaseStatement {
233 name,
234 if_not_exists,
235 }: CreateDatabaseStatement,
236) -> Result<Plan, PlanError> {
237 Ok(Plan::CreateDatabase(CreateDatabasePlan {
238 name: normalize::ident(name.0),
239 if_not_exists,
240 }))
241}
242
243pub fn describe_create_schema(
244 _: &StatementContext,
245 _: CreateSchemaStatement,
246) -> Result<StatementDesc, PlanError> {
247 Ok(StatementDesc::new(None))
248}
249
250pub fn plan_create_schema(
251 scx: &StatementContext,
252 CreateSchemaStatement {
253 mut name,
254 if_not_exists,
255 }: CreateSchemaStatement,
256) -> Result<Plan, PlanError> {
257 if name.0.len() > 2 {
258 sql_bail!("schema name {} has more than two components", name);
259 }
260 let schema_name = normalize::ident(
261 name.0
262 .pop()
263 .expect("names always have at least one component"),
264 );
265 let database_spec = match name.0.pop() {
266 None => match scx.catalog.active_database() {
267 Some(id) => ResolvedDatabaseSpecifier::Id(id.clone()),
268 None => sql_bail!("no database specified and no active database"),
269 },
270 Some(n) => match scx.resolve_database(&UnresolvedDatabaseName(n.clone())) {
271 Ok(database) => ResolvedDatabaseSpecifier::Id(database.id()),
272 Err(_) => sql_bail!("invalid database {}", n.as_str()),
273 },
274 };
275 Ok(Plan::CreateSchema(CreateSchemaPlan {
276 database_spec,
277 schema_name,
278 if_not_exists,
279 }))
280}
281
282pub fn describe_create_table(
283 _: &StatementContext,
284 _: CreateTableStatement<Aug>,
285) -> Result<StatementDesc, PlanError> {
286 Ok(StatementDesc::new(None))
287}
288
289pub fn plan_create_table(
290 scx: &StatementContext,
291 stmt: CreateTableStatement<Aug>,
292) -> Result<Plan, PlanError> {
293 let CreateTableStatement {
294 name,
295 columns,
296 constraints,
297 if_not_exists,
298 temporary,
299 with_options,
300 } = &stmt;
301
302 let names: Vec<_> = columns
303 .iter()
304 .filter(|c| {
305 let is_versioned = c
309 .options
310 .iter()
311 .any(|o| matches!(o.option, ColumnOption::Versioned { .. }));
312 !is_versioned
313 })
314 .map(|c| normalize::column_name(c.name.clone()))
315 .collect();
316
317 if let Some(dup) = names.iter().duplicates().next() {
318 sql_bail!("column {} specified more than once", dup.quoted());
319 }
320
321 let mut column_types = Vec::with_capacity(columns.len());
324 let mut defaults = Vec::with_capacity(columns.len());
325 let mut changes = BTreeMap::new();
326 let mut keys = Vec::new();
327
328 for (i, c) in columns.into_iter().enumerate() {
329 let aug_data_type = &c.data_type;
330 let ty = query::scalar_type_from_sql(scx, aug_data_type)?;
331 let mut nullable = true;
332 let mut default = Expr::null();
333 let mut versioned = false;
334 for option in &c.options {
335 match &option.option {
336 ColumnOption::NotNull => nullable = false,
337 ColumnOption::Default(expr) => {
338 let mut expr = expr.clone();
341 transform_ast::transform(scx, &mut expr)?;
342 let _ = query::plan_default_expr(scx, &expr, &ty)?;
343 default = expr.clone();
344 }
345 ColumnOption::Unique { is_primary } => {
346 keys.push(vec![i]);
347 if *is_primary {
348 nullable = false;
349 }
350 }
351 ColumnOption::Versioned { action, version } => {
352 let version = RelationVersion::from(*version);
353 versioned = true;
354
355 let name = normalize::column_name(c.name.clone());
356 let typ = ty.clone().nullable(nullable);
357
358 changes.insert(version, (action.clone(), name, typ));
359 }
360 other => {
361 bail_unsupported!(format!("CREATE TABLE with column constraint: {}", other))
362 }
363 }
364 }
365 if !versioned {
368 column_types.push(ty.nullable(nullable));
369 }
370 defaults.push(default);
371 }
372
373 let mut seen_primary = false;
374 'c: for constraint in constraints {
375 match constraint {
376 TableConstraint::Unique {
377 name: _,
378 columns,
379 is_primary,
380 nulls_not_distinct,
381 } => {
382 if seen_primary && *is_primary {
383 sql_bail!(
384 "multiple primary keys for table {} are not allowed",
385 name.to_ast_string_stable()
386 );
387 }
388 seen_primary = *is_primary || seen_primary;
389
390 let mut key = vec![];
391 for column in columns {
392 let column = normalize::column_name(column.clone());
393 match names.iter().position(|name| *name == column) {
394 None => sql_bail!("unknown column in constraint: {}", column),
395 Some(i) => {
396 let nullable = &mut column_types[i].nullable;
397 if *is_primary {
398 if *nulls_not_distinct {
399 sql_bail!(
400 "[internal error] PRIMARY KEY does not support NULLS NOT DISTINCT"
401 );
402 }
403
404 *nullable = false;
405 } else if !(*nulls_not_distinct || !*nullable) {
406 break 'c;
409 }
410
411 key.push(i);
412 }
413 }
414 }
415
416 if *is_primary {
417 keys.insert(0, key);
418 } else {
419 keys.push(key);
420 }
421 }
422 TableConstraint::ForeignKey { .. } => {
423 scx.require_feature_flag(&vars::UNSAFE_ENABLE_TABLE_FOREIGN_KEY)?
426 }
427 TableConstraint::Check { .. } => {
428 scx.require_feature_flag(&vars::UNSAFE_ENABLE_TABLE_CHECK_CONSTRAINT)?
431 }
432 }
433 }
434
435 if !keys.is_empty() {
436 scx.require_feature_flag(&vars::UNSAFE_ENABLE_TABLE_KEYS)?
439 }
440
441 let typ = SqlRelationType::new(column_types).with_keys(keys);
442
443 let temporary = *temporary;
444 let name = if temporary {
445 scx.allocate_temporary_qualified_name(normalize::unresolved_item_name(name.to_owned())?)?
446 } else {
447 scx.allocate_qualified_name(normalize::unresolved_item_name(name.to_owned())?)?
448 };
449
450 let full_name = scx.catalog.resolve_full_name(&name);
452 let partial_name = PartialItemName::from(full_name.clone());
453 if let (false, Ok(item)) = (
456 if_not_exists,
457 scx.catalog.resolve_item_or_type(&partial_name),
458 ) {
459 return Err(PlanError::ItemAlreadyExists {
460 name: full_name.to_string(),
461 item_type: item.item_type(),
462 });
463 }
464
465 let desc = RelationDesc::new(typ, names);
466 let mut desc = VersionedRelationDesc::new(desc);
467 for (version, (_action, name, typ)) in changes.into_iter() {
468 let new_version = desc.add_column(name, typ);
469 if version != new_version {
470 return Err(PlanError::InvalidTable {
471 name: full_name.item,
472 });
473 }
474 }
475
476 let create_sql = normalize::create_statement(scx, Statement::CreateTable(stmt.clone()))?;
477
478 let original_desc = desc.at_version(RelationVersionSelector::Specific(RelationVersion::root()));
484 let options = plan_table_options(scx, &original_desc, with_options.clone())?;
485
486 let compaction_window = options.iter().find_map(|o| {
487 #[allow(irrefutable_let_patterns)]
488 if let crate::plan::TableOption::RetainHistory(lcw) = o {
489 Some(lcw.clone())
490 } else {
491 None
492 }
493 });
494
495 let table = Table {
496 create_sql,
497 desc,
498 temporary,
499 compaction_window,
500 data_source: TableDataSource::TableWrites { defaults },
501 };
502 Ok(Plan::CreateTable(CreateTablePlan {
503 name,
504 table,
505 if_not_exists: *if_not_exists,
506 }))
507}
508
509pub fn describe_create_table_from_source(
510 _: &StatementContext,
511 _: CreateTableFromSourceStatement<Aug>,
512) -> Result<StatementDesc, PlanError> {
513 Ok(StatementDesc::new(None))
514}
515
516pub fn describe_create_webhook_source(
517 _: &StatementContext,
518 _: CreateWebhookSourceStatement<Aug>,
519) -> Result<StatementDesc, PlanError> {
520 Ok(StatementDesc::new(None))
521}
522
523pub fn describe_create_source(
524 _: &StatementContext,
525 _: CreateSourceStatement<Aug>,
526) -> Result<StatementDesc, PlanError> {
527 Ok(StatementDesc::new(None))
528}
529
530pub fn describe_create_subsource(
531 _: &StatementContext,
532 _: CreateSubsourceStatement<Aug>,
533) -> Result<StatementDesc, PlanError> {
534 Ok(StatementDesc::new(None))
535}
536
537generate_extracted_config!(
538 CreateSourceOption,
539 (TimestampInterval, Duration),
540 (RetainHistory, OptionalDuration)
541);
542
543generate_extracted_config!(
544 PgConfigOption,
545 (Details, String),
546 (Publication, String),
547 (TextColumns, Vec::<UnresolvedItemName>, Default(vec![])),
548 (ExcludeColumns, Vec::<UnresolvedItemName>, Default(vec![]))
549);
550
551generate_extracted_config!(
552 MySqlConfigOption,
553 (Details, String),
554 (TextColumns, Vec::<UnresolvedItemName>, Default(vec![])),
555 (ExcludeColumns, Vec::<UnresolvedItemName>, Default(vec![]))
556);
557
558generate_extracted_config!(
559 SqlServerConfigOption,
560 (Details, String),
561 (TextColumns, Vec::<UnresolvedItemName>, Default(vec![])),
562 (ExcludeColumns, Vec::<UnresolvedItemName>, Default(vec![]))
563);
564
565pub fn plan_create_webhook_source(
566 scx: &StatementContext,
567 mut stmt: CreateWebhookSourceStatement<Aug>,
568) -> Result<Plan, PlanError> {
569 if stmt.is_table {
570 scx.require_feature_flag(&ENABLE_CREATE_TABLE_FROM_SOURCE)?;
571 }
572
573 let in_cluster = source_sink_cluster_config(scx, &mut stmt.in_cluster)?;
576 let create_sql =
577 normalize::create_statement(scx, Statement::CreateWebhookSource(stmt.clone()))?;
578
579 let CreateWebhookSourceStatement {
580 name,
581 if_not_exists,
582 body_format,
583 include_headers,
584 validate_using,
585 is_table,
586 in_cluster: _,
588 } = stmt;
589
590 let validate_using = validate_using
591 .map(|stmt| query::plan_webhook_validate_using(scx, stmt))
592 .transpose()?;
593 if let Some(WebhookValidation { expression, .. }) = &validate_using {
594 if !expression.contains_column() {
597 return Err(PlanError::WebhookValidationDoesNotUseColumns);
598 }
599 if expression.contains_unmaterializable_except(&[UnmaterializableFunc::CurrentTimestamp]) {
603 return Err(PlanError::WebhookValidationNonDeterministic);
604 }
605 }
606
607 let body_format = match body_format {
608 Format::Bytes => WebhookBodyFormat::Bytes,
609 Format::Json { array } => WebhookBodyFormat::Json { array },
610 Format::Text => WebhookBodyFormat::Text,
611 ty => {
613 return Err(PlanError::Unsupported {
614 feature: format!("{ty} is not a valid BODY FORMAT for a WEBHOOK source"),
615 discussion_no: None,
616 });
617 }
618 };
619
620 let mut column_ty = vec![
621 SqlColumnType {
623 scalar_type: SqlScalarType::from(body_format),
624 nullable: false,
625 },
626 ];
627 let mut column_names = vec!["body".to_string()];
628
629 let mut headers = WebhookHeaders::default();
630
631 if let Some(filters) = include_headers.column {
633 column_ty.push(SqlColumnType {
634 scalar_type: SqlScalarType::Map {
635 value_type: Box::new(SqlScalarType::String),
636 custom_id: None,
637 },
638 nullable: false,
639 });
640 column_names.push("headers".to_string());
641
642 let (allow, block): (BTreeSet<_>, BTreeSet<_>) =
643 filters.into_iter().partition_map(|filter| {
644 if filter.block {
645 itertools::Either::Right(filter.header_name)
646 } else {
647 itertools::Either::Left(filter.header_name)
648 }
649 });
650 headers.header_column = Some(WebhookHeaderFilters { allow, block });
651 }
652
653 for header in include_headers.mappings {
655 let scalar_type = header
656 .use_bytes
657 .then_some(SqlScalarType::Bytes)
658 .unwrap_or(SqlScalarType::String);
659 column_ty.push(SqlColumnType {
660 scalar_type,
661 nullable: true,
662 });
663 column_names.push(header.column_name.into_string());
664
665 let column_idx = column_ty.len() - 1;
666 assert_eq!(
668 column_idx,
669 column_names.len() - 1,
670 "header column names and types don't match"
671 );
672 headers
673 .mapped_headers
674 .insert(column_idx, (header.header_name, header.use_bytes));
675 }
676
677 let mut unique_check = HashSet::with_capacity(column_names.len());
679 for name in &column_names {
680 if !unique_check.insert(name) {
681 return Err(PlanError::AmbiguousColumn(name.clone().into()));
682 }
683 }
684 if column_names.len() > MAX_NUM_COLUMNS {
685 return Err(PlanError::TooManyColumns {
686 max_num_columns: MAX_NUM_COLUMNS,
687 req_num_columns: column_names.len(),
688 });
689 }
690
691 let typ = SqlRelationType::new(column_ty);
692 let desc = RelationDesc::new(typ, column_names);
693
694 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name)?)?;
696 let full_name = scx.catalog.resolve_full_name(&name);
697 let partial_name = PartialItemName::from(full_name.clone());
698 if let (false, Ok(item)) = (if_not_exists, scx.catalog.resolve_item(&partial_name)) {
699 return Err(PlanError::ItemAlreadyExists {
700 name: full_name.to_string(),
701 item_type: item.item_type(),
702 });
703 }
704
705 let timeline = Timeline::EpochMilliseconds;
708
709 let plan = if is_table {
710 let data_source = DataSourceDesc::Webhook {
711 validate_using,
712 body_format,
713 headers,
714 cluster_id: Some(in_cluster.id()),
715 };
716 let data_source = TableDataSource::DataSource {
717 desc: data_source,
718 timeline,
719 };
720 Plan::CreateTable(CreateTablePlan {
721 name,
722 if_not_exists,
723 table: Table {
724 create_sql,
725 desc: VersionedRelationDesc::new(desc),
726 temporary: false,
727 compaction_window: None,
728 data_source,
729 },
730 })
731 } else {
732 let data_source = DataSourceDesc::Webhook {
733 validate_using,
734 body_format,
735 headers,
736 cluster_id: None,
738 };
739 Plan::CreateSource(CreateSourcePlan {
740 name,
741 source: Source {
742 create_sql,
743 data_source,
744 desc,
745 compaction_window: None,
746 },
747 if_not_exists,
748 timeline,
749 in_cluster: Some(in_cluster.id()),
750 })
751 };
752
753 Ok(plan)
754}
755
756pub fn plan_create_source(
757 scx: &StatementContext,
758 mut stmt: CreateSourceStatement<Aug>,
759) -> Result<Plan, PlanError> {
760 let CreateSourceStatement {
761 name,
762 in_cluster: _,
763 col_names,
764 connection: source_connection,
765 envelope,
766 if_not_exists,
767 format,
768 key_constraint,
769 include_metadata,
770 with_options,
771 external_references: referenced_subsources,
772 progress_subsource,
773 } = &stmt;
774
775 mz_ore::soft_assert_or_log!(
776 referenced_subsources.is_none(),
777 "referenced subsources must be cleared in purification"
778 );
779
780 let force_source_table_syntax = scx.catalog.system_vars().enable_create_table_from_source()
781 && scx.catalog.system_vars().force_source_table_syntax();
782
783 if force_source_table_syntax {
786 if envelope.is_some() || format.is_some() || !include_metadata.is_empty() {
787 Err(PlanError::UseTablesForSources(
788 "CREATE SOURCE (ENVELOPE|FORMAT|INCLUDE)".to_string(),
789 ))?;
790 }
791 }
792
793 let envelope = envelope.clone().unwrap_or(ast::SourceEnvelope::None);
794
795 if !matches!(source_connection, CreateSourceConnection::Kafka { .. })
796 && include_metadata
797 .iter()
798 .any(|sic| matches!(sic, SourceIncludeMetadata::Headers { .. }))
799 {
800 sql_bail!("INCLUDE HEADERS with non-Kafka sources not supported");
802 }
803 if !matches!(
804 source_connection,
805 CreateSourceConnection::Kafka { .. } | CreateSourceConnection::LoadGenerator { .. }
806 ) && !include_metadata.is_empty()
807 {
808 bail_unsupported!("INCLUDE metadata with non-Kafka sources");
809 }
810
811 if !include_metadata.is_empty()
812 && !matches!(
813 envelope,
814 ast::SourceEnvelope::Upsert { .. }
815 | ast::SourceEnvelope::None
816 | ast::SourceEnvelope::Debezium
817 )
818 {
819 sql_bail!("INCLUDE <metadata> requires ENVELOPE (NONE|UPSERT|DEBEZIUM)");
820 }
821
822 let external_connection =
823 plan_generic_source_connection(scx, source_connection, include_metadata)?;
824
825 let CreateSourceOptionExtracted {
826 timestamp_interval,
827 retain_history,
828 seen: _,
829 } = CreateSourceOptionExtracted::try_from(with_options.clone())?;
830
831 let metadata_columns_desc = match external_connection {
832 GenericSourceConnection::Kafka(KafkaSourceConnection {
833 ref metadata_columns,
834 ..
835 }) => kafka_metadata_columns_desc(metadata_columns),
836 _ => vec![],
837 };
838
839 let (mut desc, envelope, encoding) = apply_source_envelope_encoding(
841 scx,
842 &envelope,
843 format,
844 Some(external_connection.default_key_desc()),
845 external_connection.default_value_desc(),
846 include_metadata,
847 metadata_columns_desc,
848 &external_connection,
849 )?;
850 plan_utils::maybe_rename_columns(format!("source {}", name), &mut desc, col_names)?;
851
852 let names: Vec<_> = desc.iter_names().cloned().collect();
853 if let Some(dup) = names.iter().duplicates().next() {
854 sql_bail!("column {} specified more than once", dup.quoted());
855 }
856
857 if let Some(KeyConstraint::PrimaryKeyNotEnforced { columns }) = key_constraint.clone() {
859 scx.require_feature_flag(&vars::ENABLE_PRIMARY_KEY_NOT_ENFORCED)?;
862
863 let key_columns = columns
864 .into_iter()
865 .map(normalize::column_name)
866 .collect::<Vec<_>>();
867
868 let mut uniq = BTreeSet::new();
869 for col in key_columns.iter() {
870 if !uniq.insert(col) {
871 sql_bail!("Repeated column name in source key constraint: {}", col);
872 }
873 }
874
875 let key_indices = key_columns
876 .iter()
877 .map(|col| {
878 let name_idx = desc
879 .get_by_name(col)
880 .map(|(idx, _type)| idx)
881 .ok_or_else(|| sql_err!("No such column in source key constraint: {}", col))?;
882 if desc.get_unambiguous_name(name_idx).is_none() {
883 sql_bail!("Ambiguous column in source key constraint: {}", col);
884 }
885 Ok(name_idx)
886 })
887 .collect::<Result<Vec<_>, _>>()?;
888
889 if !desc.typ().keys.is_empty() {
890 return Err(key_constraint_err(&desc, &key_columns));
891 } else {
892 desc = desc.with_key(key_indices);
893 }
894 }
895
896 let timestamp_interval = match timestamp_interval {
897 Some(duration) => {
898 if scx.pcx.is_some() {
902 let min = scx.catalog.system_vars().min_timestamp_interval();
903 let max = scx.catalog.system_vars().max_timestamp_interval();
904 if duration < min || duration > max {
905 return Err(PlanError::InvalidTimestampInterval {
906 min,
907 max,
908 requested: duration,
909 });
910 }
911 }
912 duration
913 }
914 None => scx.catalog.system_vars().default_timestamp_interval(),
915 };
916
917 let (desc, data_source) = match progress_subsource {
918 Some(name) => {
919 let DeferredItemName::Named(name) = name else {
920 sql_bail!("[internal error] progress subsource must be named during purification");
921 };
922 let ResolvedItemName::Item { id, .. } = name else {
923 sql_bail!("[internal error] invalid target id");
924 };
925
926 let details = match external_connection {
927 GenericSourceConnection::Kafka(ref c) => {
928 SourceExportDetails::Kafka(KafkaSourceExportDetails {
929 metadata_columns: c.metadata_columns.clone(),
930 })
931 }
932 GenericSourceConnection::LoadGenerator(ref c) => match c.load_generator {
933 LoadGenerator::Auction
934 | LoadGenerator::Marketing
935 | LoadGenerator::Tpch { .. } => SourceExportDetails::None,
936 LoadGenerator::Counter { .. }
937 | LoadGenerator::Clock
938 | LoadGenerator::Datums
939 | LoadGenerator::KeyValue(_) => {
940 SourceExportDetails::LoadGenerator(LoadGeneratorSourceExportDetails {
941 output: LoadGeneratorOutput::Default,
942 })
943 }
944 },
945 GenericSourceConnection::Postgres(_)
946 | GenericSourceConnection::MySql(_)
947 | GenericSourceConnection::SqlServer(_) => SourceExportDetails::None,
948 };
949
950 let data_source = DataSourceDesc::OldSyntaxIngestion {
951 desc: SourceDesc {
952 connection: external_connection,
953 timestamp_interval,
954 },
955 progress_subsource: *id,
956 data_config: SourceExportDataConfig {
957 encoding,
958 envelope: envelope.clone(),
959 },
960 details,
961 };
962 (desc, data_source)
963 }
964 None => {
965 let desc = external_connection.timestamp_desc();
966 let data_source = DataSourceDesc::Ingestion(SourceDesc {
967 connection: external_connection,
968 timestamp_interval,
969 });
970 (desc, data_source)
971 }
972 };
973
974 let if_not_exists = *if_not_exists;
975 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name.clone())?)?;
976
977 let full_name = scx.catalog.resolve_full_name(&name);
979 let partial_name = PartialItemName::from(full_name.clone());
980 if let (false, Ok(item)) = (
983 if_not_exists,
984 scx.catalog.resolve_item_or_type(&partial_name),
985 ) {
986 return Err(PlanError::ItemAlreadyExists {
987 name: full_name.to_string(),
988 item_type: item.item_type(),
989 });
990 }
991
992 let in_cluster = source_sink_cluster_config(scx, &mut stmt.in_cluster)?;
996
997 let create_sql = normalize::create_statement(scx, Statement::CreateSource(stmt))?;
998
999 let timeline = match envelope {
1001 SourceEnvelope::CdcV2 => {
1002 Timeline::External(scx.catalog.resolve_full_name(&name).to_string())
1003 }
1004 _ => Timeline::EpochMilliseconds,
1005 };
1006
1007 let compaction_window = plan_retain_history_option(scx, retain_history)?;
1008
1009 let source = Source {
1010 create_sql,
1011 data_source,
1012 desc,
1013 compaction_window,
1014 };
1015
1016 Ok(Plan::CreateSource(CreateSourcePlan {
1017 name,
1018 source,
1019 if_not_exists,
1020 timeline,
1021 in_cluster: Some(in_cluster.id()),
1022 }))
1023}
1024
1025pub fn plan_generic_source_connection(
1026 scx: &StatementContext<'_>,
1027 source_connection: &CreateSourceConnection<Aug>,
1028 include_metadata: &Vec<SourceIncludeMetadata>,
1029) -> Result<GenericSourceConnection<ReferencedConnection>, PlanError> {
1030 Ok(match source_connection {
1031 CreateSourceConnection::Kafka {
1032 connection,
1033 options,
1034 } => GenericSourceConnection::Kafka(plan_kafka_source_connection(
1035 scx,
1036 connection,
1037 options,
1038 include_metadata,
1039 )?),
1040 CreateSourceConnection::Postgres {
1041 connection,
1042 options,
1043 } => GenericSourceConnection::Postgres(plan_postgres_source_connection(
1044 scx, connection, options,
1045 )?),
1046 CreateSourceConnection::SqlServer {
1047 connection,
1048 options,
1049 } => GenericSourceConnection::SqlServer(plan_sqlserver_source_connection(
1050 scx, connection, options,
1051 )?),
1052 CreateSourceConnection::MySql {
1053 connection,
1054 options,
1055 } => {
1056 GenericSourceConnection::MySql(plan_mysql_source_connection(scx, connection, options)?)
1057 }
1058 CreateSourceConnection::LoadGenerator { generator, options } => {
1059 GenericSourceConnection::LoadGenerator(plan_load_generator_source_connection(
1060 scx,
1061 generator,
1062 options,
1063 include_metadata,
1064 )?)
1065 }
1066 })
1067}
1068
1069fn plan_load_generator_source_connection(
1070 scx: &StatementContext<'_>,
1071 generator: &ast::LoadGenerator,
1072 options: &Vec<LoadGeneratorOption<Aug>>,
1073 include_metadata: &Vec<SourceIncludeMetadata>,
1074) -> Result<LoadGeneratorSourceConnection, PlanError> {
1075 let load_generator =
1076 load_generator_ast_to_generator(scx, generator, options, include_metadata)?;
1077 let LoadGeneratorOptionExtracted {
1078 tick_interval,
1079 as_of,
1080 up_to,
1081 ..
1082 } = options.clone().try_into()?;
1083 let tick_micros = match tick_interval {
1084 Some(interval) => Some(interval.as_micros().try_into()?),
1085 None => None,
1086 };
1087 if up_to < as_of {
1088 sql_bail!("UP TO cannot be less than AS OF");
1089 }
1090 Ok(LoadGeneratorSourceConnection {
1091 load_generator,
1092 tick_micros,
1093 as_of,
1094 up_to,
1095 })
1096}
1097
1098fn plan_mysql_source_connection(
1099 scx: &StatementContext<'_>,
1100 connection: &ResolvedItemName,
1101 options: &Vec<MySqlConfigOption<Aug>>,
1102) -> Result<MySqlSourceConnection<ReferencedConnection>, PlanError> {
1103 let connection_item = scx.get_item_by_resolved_name(connection)?;
1104 match connection_item.connection()? {
1105 Connection::MySql(connection) => connection,
1106 _ => sql_bail!(
1107 "{} is not a MySQL connection",
1108 scx.catalog.resolve_full_name(connection_item.name())
1109 ),
1110 };
1111 let MySqlConfigOptionExtracted {
1112 details,
1113 text_columns: _,
1117 exclude_columns: _,
1118 seen: _,
1119 } = options.clone().try_into()?;
1120 let details = details
1121 .as_ref()
1122 .ok_or_else(|| internal_err!("MySQL source missing details"))?;
1123 let details = hex::decode(details).map_err(|e| sql_err!("{}", e))?;
1124 let details = ProtoMySqlSourceDetails::decode(&*details).map_err(|e| sql_err!("{}", e))?;
1125 let details = MySqlSourceDetails::from_proto(details).map_err(|e| sql_err!("{}", e))?;
1126 Ok(MySqlSourceConnection {
1127 connection: connection_item.id(),
1128 connection_id: connection_item.id(),
1129 details,
1130 })
1131}
1132
1133fn plan_sqlserver_source_connection(
1134 scx: &StatementContext<'_>,
1135 connection: &ResolvedItemName,
1136 options: &Vec<SqlServerConfigOption<Aug>>,
1137) -> Result<SqlServerSourceConnection<ReferencedConnection>, PlanError> {
1138 let connection_item = scx.get_item_by_resolved_name(connection)?;
1139 match connection_item.connection()? {
1140 Connection::SqlServer(connection) => connection,
1141 _ => sql_bail!(
1142 "{} is not a SQL Server connection",
1143 scx.catalog.resolve_full_name(connection_item.name())
1144 ),
1145 };
1146 let SqlServerConfigOptionExtracted { details, .. } = options.clone().try_into()?;
1147 let details = details
1148 .as_ref()
1149 .ok_or_else(|| internal_err!("SQL Server source missing details"))?;
1150 let extras = hex::decode(details)
1151 .map_err(|e| sql_err!("{e}"))
1152 .and_then(|raw| ProtoSqlServerSourceExtras::decode(&*raw).map_err(|e| sql_err!("{e}")))
1153 .and_then(|proto| SqlServerSourceExtras::from_proto(proto).map_err(|e| sql_err!("{e}")))?;
1154 Ok(SqlServerSourceConnection {
1155 connection_id: connection_item.id(),
1156 connection: connection_item.id(),
1157 extras,
1158 })
1159}
1160
1161fn plan_postgres_source_connection(
1162 scx: &StatementContext<'_>,
1163 connection: &ResolvedItemName,
1164 options: &Vec<PgConfigOption<Aug>>,
1165) -> Result<PostgresSourceConnection<ReferencedConnection>, PlanError> {
1166 let connection_item = scx.get_item_by_resolved_name(connection)?;
1167 let PgConfigOptionExtracted {
1168 details,
1169 publication,
1170 text_columns: _,
1174 exclude_columns: _,
1178 seen: _,
1179 } = options.clone().try_into()?;
1180 let details = details
1181 .as_ref()
1182 .ok_or_else(|| internal_err!("Postgres source missing details"))?;
1183 let details = hex::decode(details).map_err(|e| sql_err!("{}", e))?;
1184 let details =
1185 ProtoPostgresSourcePublicationDetails::decode(&*details).map_err(|e| sql_err!("{}", e))?;
1186 let publication_details =
1187 PostgresSourcePublicationDetails::from_proto(details).map_err(|e| sql_err!("{}", e))?;
1188 Ok(PostgresSourceConnection {
1189 connection: connection_item.id(),
1190 connection_id: connection_item.id(),
1191 publication: publication.ok_or_else(|| internal_err!("PUBLICATION option is required"))?,
1193 publication_details,
1194 })
1195}
1196
1197fn plan_kafka_source_connection(
1198 scx: &StatementContext<'_>,
1199 connection_name: &ResolvedItemName,
1200 options: &Vec<ast::KafkaSourceConfigOption<Aug>>,
1201 include_metadata: &Vec<SourceIncludeMetadata>,
1202) -> Result<KafkaSourceConnection<ReferencedConnection>, PlanError> {
1203 let connection_item = scx.get_item_by_resolved_name(connection_name)?;
1204 if !matches!(connection_item.connection()?, Connection::Kafka(_)) {
1205 sql_bail!(
1206 "{} is not a kafka connection",
1207 scx.catalog.resolve_full_name(connection_item.name())
1208 )
1209 }
1210 let KafkaSourceConfigOptionExtracted {
1211 group_id_prefix,
1212 topic,
1213 topic_metadata_refresh_interval,
1214 start_timestamp: _, start_offset,
1216 seen: _,
1217 }: KafkaSourceConfigOptionExtracted = options.clone().try_into()?;
1218 let topic = topic.ok_or_else(|| internal_err!("TOPIC option is required"))?;
1220 let mut start_offsets = BTreeMap::new();
1221 if let Some(offsets) = start_offset {
1222 for (part, offset) in offsets.iter().enumerate() {
1223 if *offset < 0 {
1224 sql_bail!("START OFFSET must be a nonnegative integer");
1225 }
1226 start_offsets.insert(i32::try_from(part)?, *offset);
1227 }
1228 }
1229 if topic_metadata_refresh_interval > Duration::from_secs(60 * 60) {
1230 sql_bail!("TOPIC METADATA REFRESH INTERVAL cannot be greater than 1 hour");
1233 }
1234 if topic_metadata_refresh_interval < Duration::from_secs(1) {
1235 sql_bail!("TOPIC METADATA REFRESH INTERVAL must be at least 1 second");
1238 }
1239 let metadata_columns = include_metadata
1240 .into_iter()
1241 .flat_map(|item| match item {
1242 SourceIncludeMetadata::Timestamp { alias } => {
1243 let name = match alias {
1244 Some(name) => name.to_string(),
1245 None => "timestamp".to_owned(),
1246 };
1247 Some((name, KafkaMetadataKind::Timestamp))
1248 }
1249 SourceIncludeMetadata::Partition { alias } => {
1250 let name = match alias {
1251 Some(name) => name.to_string(),
1252 None => "partition".to_owned(),
1253 };
1254 Some((name, KafkaMetadataKind::Partition))
1255 }
1256 SourceIncludeMetadata::Offset { alias } => {
1257 let name = match alias {
1258 Some(name) => name.to_string(),
1259 None => "offset".to_owned(),
1260 };
1261 Some((name, KafkaMetadataKind::Offset))
1262 }
1263 SourceIncludeMetadata::Headers { alias } => {
1264 let name = match alias {
1265 Some(name) => name.to_string(),
1266 None => "headers".to_owned(),
1267 };
1268 Some((name, KafkaMetadataKind::Headers))
1269 }
1270 SourceIncludeMetadata::Header {
1271 alias,
1272 key,
1273 use_bytes,
1274 } => Some((
1275 alias.to_string(),
1276 KafkaMetadataKind::Header {
1277 key: key.clone(),
1278 use_bytes: *use_bytes,
1279 },
1280 )),
1281 SourceIncludeMetadata::Key { .. } => {
1282 None
1284 }
1285 })
1286 .collect();
1287 Ok(KafkaSourceConnection {
1288 connection: connection_item.id(),
1289 connection_id: connection_item.id(),
1290 topic,
1291 start_offsets,
1292 group_id_prefix,
1293 topic_metadata_refresh_interval,
1294 metadata_columns,
1295 })
1296}
1297
1298fn apply_source_envelope_encoding(
1299 scx: &StatementContext,
1300 envelope: &ast::SourceEnvelope,
1301 format: &Option<FormatSpecifier<Aug>>,
1302 key_desc: Option<RelationDesc>,
1303 value_desc: RelationDesc,
1304 include_metadata: &[SourceIncludeMetadata],
1305 metadata_columns_desc: Vec<(&str, SqlColumnType)>,
1306 source_connection: &GenericSourceConnection<ReferencedConnection>,
1307) -> Result<
1308 (
1309 RelationDesc,
1310 SourceEnvelope,
1311 Option<SourceDataEncoding<ReferencedConnection>>,
1312 ),
1313 PlanError,
1314> {
1315 let encoding = match format {
1316 Some(format) => Some(get_encoding(scx, format, envelope)?),
1317 None => None,
1318 };
1319
1320 let (key_desc, value_desc) = match &encoding {
1321 Some(encoding) => {
1322 match value_desc.typ().columns() {
1325 [typ] => match typ.scalar_type {
1326 SqlScalarType::Bytes => {}
1327 _ => sql_bail!(
1328 "The schema produced by the source is incompatible with format decoding"
1329 ),
1330 },
1331 _ => sql_bail!(
1332 "The schema produced by the source is incompatible with format decoding"
1333 ),
1334 }
1335
1336 let (key_desc, value_desc) = encoding.desc()?;
1337
1338 let key_desc = key_desc.map(|desc| {
1345 let is_kafka = matches!(source_connection, GenericSourceConnection::Kafka(_));
1346 let is_envelope_none = matches!(envelope, ast::SourceEnvelope::None);
1347 if is_kafka && is_envelope_none {
1348 RelationDesc::from_names_and_types(
1349 desc.into_iter()
1350 .map(|(name, typ)| (name, typ.nullable(true))),
1351 )
1352 } else {
1353 desc
1354 }
1355 });
1356 (key_desc, value_desc)
1357 }
1358 None => (key_desc, value_desc),
1359 };
1360
1361 let key_envelope_no_encoding = matches!(
1374 source_connection,
1375 GenericSourceConnection::LoadGenerator(LoadGeneratorSourceConnection {
1376 load_generator: LoadGenerator::KeyValue(_),
1377 ..
1378 })
1379 );
1380 let mut key_envelope = get_key_envelope(
1381 include_metadata,
1382 encoding.as_ref(),
1383 key_envelope_no_encoding,
1384 )?;
1385
1386 match (&envelope, &key_envelope) {
1387 (ast::SourceEnvelope::Debezium, KeyEnvelope::None) => {}
1388 (ast::SourceEnvelope::Debezium, _) => sql_bail!(
1389 "Cannot use INCLUDE KEY with ENVELOPE DEBEZIUM: Debezium values include all keys."
1390 ),
1391 _ => {}
1392 };
1393
1394 let envelope = match &envelope {
1403 ast::SourceEnvelope::None => UnplannedSourceEnvelope::None(key_envelope),
1405 ast::SourceEnvelope::Debezium => {
1406 let after_idx = match typecheck_debezium(&value_desc) {
1408 Ok((_before_idx, after_idx)) => Ok(after_idx),
1409 Err(type_err) => match encoding.as_ref().map(|e| &e.value) {
1410 Some(DataEncoding::Avro(_)) => Err(type_err),
1411 _ => Err(sql_err!(
1412 "ENVELOPE DEBEZIUM requires that VALUE FORMAT is set to AVRO"
1413 )),
1414 },
1415 }?;
1416
1417 UnplannedSourceEnvelope::Upsert {
1418 style: UpsertStyle::Debezium { after_idx },
1419 }
1420 }
1421 ast::SourceEnvelope::Upsert {
1422 value_decode_err_policy,
1423 } => {
1424 let key_encoding = match encoding.as_ref().and_then(|e| e.key.as_ref()) {
1425 None => {
1426 if !key_envelope_no_encoding {
1427 bail_unsupported!(format!(
1428 "UPSERT requires a key/value format: {:?}",
1429 format
1430 ))
1431 }
1432 None
1433 }
1434 Some(key_encoding) => Some(key_encoding),
1435 };
1436 if key_envelope == KeyEnvelope::None {
1439 key_envelope = get_unnamed_key_envelope(key_encoding)?;
1440 }
1441 let style = match value_decode_err_policy.as_slice() {
1443 [] => UpsertStyle::Default(key_envelope),
1444 [SourceErrorPolicy::Inline { alias }] => {
1445 scx.require_feature_flag(&vars::ENABLE_ENVELOPE_UPSERT_INLINE_ERRORS)?;
1446 UpsertStyle::ValueErrInline {
1447 key_envelope,
1448 error_column: alias
1449 .as_ref()
1450 .map_or_else(|| "error".to_string(), |a| a.to_string()),
1451 }
1452 }
1453 _ => {
1454 bail_unsupported!("ENVELOPE UPSERT with unsupported value decode error policy")
1455 }
1456 };
1457
1458 UnplannedSourceEnvelope::Upsert { style }
1459 }
1460 ast::SourceEnvelope::CdcV2 => {
1461 scx.require_feature_flag(&vars::ENABLE_ENVELOPE_MATERIALIZE)?;
1462 match format {
1464 Some(FormatSpecifier::Bare(Format::Avro(_))) => {}
1465 _ => bail_unsupported!("non-Avro-encoded ENVELOPE MATERIALIZE"),
1466 }
1467 UnplannedSourceEnvelope::CdcV2
1468 }
1469 };
1470
1471 let metadata_desc = included_column_desc(metadata_columns_desc);
1472 let (envelope, desc) = envelope.desc(key_desc, value_desc, metadata_desc)?;
1473
1474 Ok((desc, envelope, encoding))
1475}
1476
1477fn plan_source_export_desc(
1480 scx: &StatementContext,
1481 name: &UnresolvedItemName,
1482 columns: &Vec<ColumnDef<Aug>>,
1483 constraints: &Vec<TableConstraint<Aug>>,
1484) -> Result<RelationDesc, PlanError> {
1485 let names: Vec<_> = columns
1486 .iter()
1487 .map(|c| normalize::column_name(c.name.clone()))
1488 .collect();
1489
1490 if let Some(dup) = names.iter().duplicates().next() {
1491 sql_bail!("column {} specified more than once", dup.quoted());
1492 }
1493
1494 let mut column_types = Vec::with_capacity(columns.len());
1497 let mut keys = Vec::new();
1498
1499 for (i, c) in columns.into_iter().enumerate() {
1500 let aug_data_type = &c.data_type;
1501 let ty = query::scalar_type_from_sql(scx, aug_data_type)?;
1502 let mut nullable = true;
1503 for option in &c.options {
1504 match &option.option {
1505 ColumnOption::NotNull => nullable = false,
1506 ColumnOption::Default(_) => {
1507 bail_unsupported!("Source export with default value")
1508 }
1509 ColumnOption::Unique { is_primary } => {
1510 keys.push(vec![i]);
1511 if *is_primary {
1512 nullable = false;
1513 }
1514 }
1515 other => {
1516 bail_unsupported!(format!("Source export with column constraint: {}", other))
1517 }
1518 }
1519 }
1520 column_types.push(ty.nullable(nullable));
1521 }
1522
1523 let mut seen_primary = false;
1524 'c: for constraint in constraints {
1525 match constraint {
1526 TableConstraint::Unique {
1527 name: _,
1528 columns,
1529 is_primary,
1530 nulls_not_distinct,
1531 } => {
1532 if seen_primary && *is_primary {
1533 sql_bail!(
1534 "multiple primary keys for source export {} are not allowed",
1535 name.to_ast_string_stable()
1536 );
1537 }
1538 seen_primary = *is_primary || seen_primary;
1539
1540 let mut key = vec![];
1541 for column in columns {
1542 let column = normalize::column_name(column.clone());
1543 match names.iter().position(|name| *name == column) {
1544 None => sql_bail!("unknown column in constraint: {}", column),
1545 Some(i) => {
1546 let nullable = &mut column_types[i].nullable;
1547 if *is_primary {
1548 if *nulls_not_distinct {
1549 sql_bail!(
1550 "[internal error] PRIMARY KEY does not support NULLS NOT DISTINCT"
1551 );
1552 }
1553 *nullable = false;
1554 } else if !(*nulls_not_distinct || !*nullable) {
1555 break 'c;
1558 }
1559
1560 key.push(i);
1561 }
1562 }
1563 }
1564
1565 if *is_primary {
1566 keys.insert(0, key);
1567 } else {
1568 keys.push(key);
1569 }
1570 }
1571 TableConstraint::ForeignKey { .. } => {
1572 bail_unsupported!("Source export with a foreign key")
1573 }
1574 TableConstraint::Check { .. } => {
1575 bail_unsupported!("Source export with a check constraint")
1576 }
1577 }
1578 }
1579
1580 let typ = SqlRelationType::new(column_types).with_keys(keys);
1581 let desc = RelationDesc::new(typ, names);
1582 Ok(desc)
1583}
1584
1585generate_extracted_config!(
1586 CreateSubsourceOption,
1587 (Progress, bool, Default(false)),
1588 (ExternalReference, UnresolvedItemName),
1589 (RetainHistory, OptionalDuration),
1590 (TextColumns, Vec::<Ident>, Default(vec![])),
1591 (ExcludeColumns, Vec::<Ident>, Default(vec![])),
1592 (Details, String)
1593);
1594
1595pub fn plan_create_subsource(
1596 scx: &StatementContext,
1597 stmt: CreateSubsourceStatement<Aug>,
1598) -> Result<Plan, PlanError> {
1599 let CreateSubsourceStatement {
1600 name,
1601 columns,
1602 of_source,
1603 constraints,
1604 if_not_exists,
1605 with_options,
1606 } = &stmt;
1607
1608 let CreateSubsourceOptionExtracted {
1609 progress,
1610 retain_history,
1611 external_reference,
1612 text_columns,
1613 exclude_columns,
1614 details,
1615 seen: _,
1616 } = with_options.clone().try_into()?;
1617
1618 if !(progress ^ (external_reference.is_some() && of_source.is_some())) {
1623 bail_internal!(
1624 "CREATE SUBSOURCE statement must specify either PROGRESS or REFERENCES option"
1625 );
1626 }
1627
1628 let desc = plan_source_export_desc(scx, name, columns, constraints)?;
1629
1630 let data_source = if let Some(source_reference) = of_source {
1631 if scx.catalog.system_vars().enable_create_table_from_source()
1634 && scx.catalog.system_vars().force_source_table_syntax()
1635 {
1636 Err(PlanError::UseTablesForSources(
1637 "CREATE SUBSOURCE".to_string(),
1638 ))?;
1639 }
1640
1641 let ingestion_id = *source_reference.item_id();
1644 let external_reference = external_reference.ok_or_else(|| {
1645 sql_err!("CREATE SUBSOURCE with REFERENCES requires EXTERNAL REFERENCE option")
1646 })?;
1647
1648 let details = details
1651 .as_ref()
1652 .ok_or_else(|| internal_err!("source-export subsource missing details"))?;
1653 let details = hex::decode(details).map_err(|e| sql_err!("{}", e))?;
1654 let details =
1655 ProtoSourceExportStatementDetails::decode(&*details).map_err(|e| sql_err!("{}", e))?;
1656 let details =
1657 SourceExportStatementDetails::from_proto(details).map_err(|e| sql_err!("{}", e))?;
1658 let details = match details {
1659 SourceExportStatementDetails::Postgres {
1660 table,
1661 cast_oid_full_range,
1662 initial_lsn,
1663 } => SourceExportDetails::Postgres(PostgresSourceExportDetails {
1664 column_casts: crate::pure::postgres::generate_column_casts(
1665 scx,
1666 &table,
1667 &text_columns,
1668 cast_oid_full_range,
1669 )?,
1670 table,
1671 initial_lsn,
1672 }),
1673 SourceExportStatementDetails::MySql {
1674 table,
1675 initial_gtid_set,
1676 binlog_full_metadata,
1677 } => SourceExportDetails::MySql(MySqlSourceExportDetails {
1678 table,
1679 initial_gtid_set,
1680 text_columns: text_columns.into_iter().map(|c| c.into_string()).collect(),
1681 exclude_columns: exclude_columns
1682 .into_iter()
1683 .map(|c| c.into_string())
1684 .collect(),
1685 binlog_full_metadata,
1686 }),
1687 SourceExportStatementDetails::SqlServer {
1688 table,
1689 capture_instance,
1690 initial_lsn,
1691 } => SourceExportDetails::SqlServer(SqlServerSourceExportDetails {
1692 capture_instance,
1693 table,
1694 initial_lsn,
1695 text_columns: text_columns.into_iter().map(|c| c.into_string()).collect(),
1696 exclude_columns: exclude_columns
1697 .into_iter()
1698 .map(|c| c.into_string())
1699 .collect(),
1700 }),
1701 SourceExportStatementDetails::LoadGenerator { output } => {
1702 SourceExportDetails::LoadGenerator(LoadGeneratorSourceExportDetails { output })
1703 }
1704 SourceExportStatementDetails::Kafka {} => {
1705 bail_unsupported!("subsources cannot reference Kafka sources")
1706 }
1707 };
1708 DataSourceDesc::IngestionExport {
1709 ingestion_id,
1710 external_reference,
1711 details,
1712 data_config: SourceExportDataConfig {
1714 envelope: SourceEnvelope::None(NoneEnvelope {
1715 key_envelope: KeyEnvelope::None,
1716 key_arity: 0,
1717 }),
1718 encoding: None,
1719 },
1720 }
1721 } else if progress {
1722 DataSourceDesc::Progress
1723 } else {
1724 sql_bail!("CREATE SUBSOURCE must specify one of PROGRESS or REFERENCES option")
1725 };
1726
1727 let if_not_exists = *if_not_exists;
1728 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name.clone())?)?;
1729
1730 let create_sql = normalize::create_statement(scx, Statement::CreateSubsource(stmt))?;
1731
1732 let compaction_window = plan_retain_history_option(scx, retain_history)?;
1733 let source = Source {
1734 create_sql,
1735 data_source,
1736 desc,
1737 compaction_window,
1738 };
1739
1740 Ok(Plan::CreateSource(CreateSourcePlan {
1741 name,
1742 source,
1743 if_not_exists,
1744 timeline: Timeline::EpochMilliseconds,
1745 in_cluster: None,
1746 }))
1747}
1748
1749generate_extracted_config!(
1750 TableFromSourceOption,
1751 (TextColumns, Vec::<Ident>, Default(vec![])),
1752 (ExcludeColumns, Vec::<Ident>, Default(vec![])),
1753 (ExcludeConstraints, Vec::<String>, Default(vec![])),
1754 (ExcludeAllConstraints, bool, Default(false)),
1755 (PartitionBy, Vec<Ident>),
1756 (RetainHistory, OptionalDuration),
1757 (Details, String)
1758);
1759
1760pub fn plan_create_table_from_source(
1761 scx: &StatementContext,
1762 stmt: CreateTableFromSourceStatement<Aug>,
1763) -> Result<Plan, PlanError> {
1764 if !scx.catalog.system_vars().enable_create_table_from_source() {
1765 sql_bail!("CREATE TABLE ... FROM SOURCE is not supported");
1766 }
1767
1768 let CreateTableFromSourceStatement {
1769 name,
1770 columns,
1771 constraints,
1772 if_not_exists,
1773 source,
1774 external_reference,
1775 envelope,
1776 format,
1777 include_metadata,
1778 with_options,
1779 } = &stmt;
1780
1781 let envelope = envelope.clone().unwrap_or(ast::SourceEnvelope::None);
1782
1783 let TableFromSourceOptionExtracted {
1784 text_columns,
1785 exclude_columns,
1786 exclude_constraints: _,
1787 exclude_all_constraints: _,
1788 retain_history,
1789 partition_by,
1790 details,
1791 seen: _,
1792 } = with_options.clone().try_into()?;
1793
1794 let source_item = scx.get_item_by_resolved_name(source)?;
1795 let ingestion_id = source_item.id();
1796
1797 let details = details
1800 .as_ref()
1801 .ok_or_else(|| internal_err!("source-export missing details"))?;
1802 let details = hex::decode(details).map_err(|e| sql_err!("{}", e))?;
1803 let details =
1804 ProtoSourceExportStatementDetails::decode(&*details).map_err(|e| sql_err!("{}", e))?;
1805 let details =
1806 SourceExportStatementDetails::from_proto(details).map_err(|e| sql_err!("{}", e))?;
1807
1808 if !matches!(details, SourceExportStatementDetails::Kafka { .. })
1809 && include_metadata
1810 .iter()
1811 .any(|sic| matches!(sic, SourceIncludeMetadata::Headers { .. }))
1812 {
1813 sql_bail!("INCLUDE HEADERS with non-Kafka source table not supported");
1815 }
1816 if !matches!(
1817 details,
1818 SourceExportStatementDetails::Kafka { .. }
1819 | SourceExportStatementDetails::LoadGenerator { .. }
1820 ) && !include_metadata.is_empty()
1821 {
1822 bail_unsupported!("INCLUDE metadata with non-Kafka source table");
1823 }
1824
1825 let details = match details {
1826 SourceExportStatementDetails::Postgres {
1827 table,
1828 cast_oid_full_range,
1829 initial_lsn,
1830 } => SourceExportDetails::Postgres(PostgresSourceExportDetails {
1831 column_casts: crate::pure::postgres::generate_column_casts(
1832 scx,
1833 &table,
1834 &text_columns,
1835 cast_oid_full_range,
1836 )?,
1837 table,
1838 initial_lsn,
1839 }),
1840 SourceExportStatementDetails::MySql {
1841 table,
1842 initial_gtid_set,
1843 binlog_full_metadata,
1844 } => SourceExportDetails::MySql(MySqlSourceExportDetails {
1845 table,
1846 initial_gtid_set,
1847 text_columns: text_columns.into_iter().map(|c| c.into_string()).collect(),
1848 exclude_columns: exclude_columns
1849 .into_iter()
1850 .map(|c| c.into_string())
1851 .collect(),
1852 binlog_full_metadata,
1853 }),
1854 SourceExportStatementDetails::SqlServer {
1855 table,
1856 capture_instance,
1857 initial_lsn,
1858 } => SourceExportDetails::SqlServer(SqlServerSourceExportDetails {
1859 table,
1860 capture_instance,
1861 initial_lsn,
1862 text_columns: text_columns.into_iter().map(|c| c.into_string()).collect(),
1863 exclude_columns: exclude_columns
1864 .into_iter()
1865 .map(|c| c.into_string())
1866 .collect(),
1867 }),
1868 SourceExportStatementDetails::LoadGenerator { output } => {
1869 SourceExportDetails::LoadGenerator(LoadGeneratorSourceExportDetails { output })
1870 }
1871 SourceExportStatementDetails::Kafka {} => {
1872 if !include_metadata.is_empty()
1873 && !matches!(
1874 envelope,
1875 ast::SourceEnvelope::Upsert { .. }
1876 | ast::SourceEnvelope::None
1877 | ast::SourceEnvelope::Debezium
1878 )
1879 {
1880 sql_bail!("INCLUDE <metadata> requires ENVELOPE (NONE|UPSERT|DEBEZIUM)");
1882 }
1883
1884 let metadata_columns = include_metadata
1885 .into_iter()
1886 .flat_map(|item| match item {
1887 SourceIncludeMetadata::Timestamp { alias } => {
1888 let name = match alias {
1889 Some(name) => name.to_string(),
1890 None => "timestamp".to_owned(),
1891 };
1892 Some((name, KafkaMetadataKind::Timestamp))
1893 }
1894 SourceIncludeMetadata::Partition { alias } => {
1895 let name = match alias {
1896 Some(name) => name.to_string(),
1897 None => "partition".to_owned(),
1898 };
1899 Some((name, KafkaMetadataKind::Partition))
1900 }
1901 SourceIncludeMetadata::Offset { alias } => {
1902 let name = match alias {
1903 Some(name) => name.to_string(),
1904 None => "offset".to_owned(),
1905 };
1906 Some((name, KafkaMetadataKind::Offset))
1907 }
1908 SourceIncludeMetadata::Headers { alias } => {
1909 let name = match alias {
1910 Some(name) => name.to_string(),
1911 None => "headers".to_owned(),
1912 };
1913 Some((name, KafkaMetadataKind::Headers))
1914 }
1915 SourceIncludeMetadata::Header {
1916 alias,
1917 key,
1918 use_bytes,
1919 } => Some((
1920 alias.to_string(),
1921 KafkaMetadataKind::Header {
1922 key: key.clone(),
1923 use_bytes: *use_bytes,
1924 },
1925 )),
1926 SourceIncludeMetadata::Key { .. } => {
1927 None
1929 }
1930 })
1931 .collect();
1932
1933 SourceExportDetails::Kafka(KafkaSourceExportDetails { metadata_columns })
1934 }
1935 };
1936
1937 let source_connection = &source_item
1938 .source_desc()?
1939 .ok_or_else(|| sql_err!("item is not a source"))?
1940 .connection;
1941
1942 let (key_desc, value_desc) =
1947 if matches!(columns, TableFromSourceColumns::Defined(_)) || !constraints.is_empty() {
1948 let columns = match columns {
1949 TableFromSourceColumns::Defined(columns) => columns,
1950 _ => bail_internal!("expected column definitions to be present"),
1951 };
1952 let desc = plan_source_export_desc(scx, name, columns, constraints)?;
1953 (None, desc)
1954 } else {
1955 let key_desc = source_connection.default_key_desc();
1956 let value_desc = source_connection.default_value_desc();
1957 (Some(key_desc), value_desc)
1958 };
1959
1960 let metadata_columns_desc = match &details {
1961 SourceExportDetails::Kafka(KafkaSourceExportDetails {
1962 metadata_columns, ..
1963 }) => kafka_metadata_columns_desc(metadata_columns),
1964 _ => vec![],
1965 };
1966
1967 let (mut desc, envelope, encoding) = apply_source_envelope_encoding(
1968 scx,
1969 &envelope,
1970 format,
1971 key_desc,
1972 value_desc,
1973 include_metadata,
1974 metadata_columns_desc,
1975 source_connection,
1976 )?;
1977 if let TableFromSourceColumns::Named(col_names) = columns {
1978 plan_utils::maybe_rename_columns(format!("source table {}", name), &mut desc, col_names)?;
1979 }
1980
1981 let names: Vec<_> = desc.iter_names().cloned().collect();
1982 if let Some(dup) = names.iter().duplicates().next() {
1983 sql_bail!("column {} specified more than once", dup.quoted());
1984 }
1985
1986 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name.clone())?)?;
1987
1988 let timeline = match envelope {
1991 SourceEnvelope::CdcV2 => {
1992 Timeline::External(scx.catalog.resolve_full_name(&name).to_string())
1993 }
1994 _ => Timeline::EpochMilliseconds,
1995 };
1996
1997 if let Some(partition_by) = partition_by {
1998 scx.require_feature_flag(&ENABLE_COLLECTION_PARTITION_BY)?;
1999 check_partition_by(&desc, partition_by)?;
2000 }
2001
2002 let data_source = DataSourceDesc::IngestionExport {
2003 ingestion_id,
2004 external_reference: external_reference
2006 .as_ref()
2007 .ok_or_else(|| sql_err!("EXTERNAL REFERENCE is required"))?
2008 .clone(),
2009 details,
2010 data_config: SourceExportDataConfig { envelope, encoding },
2011 };
2012
2013 let if_not_exists = *if_not_exists;
2014
2015 let create_sql = normalize::create_statement(scx, Statement::CreateTableFromSource(stmt))?;
2016
2017 let compaction_window = plan_retain_history_option(scx, retain_history)?;
2018 let table = Table {
2019 create_sql,
2020 desc: VersionedRelationDesc::new(desc),
2021 temporary: false,
2022 compaction_window,
2023 data_source: TableDataSource::DataSource {
2024 desc: data_source,
2025 timeline,
2026 },
2027 };
2028
2029 Ok(Plan::CreateTable(CreateTablePlan {
2030 name,
2031 table,
2032 if_not_exists,
2033 }))
2034}
2035
2036generate_extracted_config!(
2037 LoadGeneratorOption,
2038 (TickInterval, Duration),
2039 (AsOf, u64, Default(0_u64)),
2040 (UpTo, u64, Default(u64::MAX)),
2041 (ScaleFactor, f64),
2042 (MaxCardinality, u64),
2043 (Keys, u64),
2044 (SnapshotRounds, u64),
2045 (TransactionalSnapshot, bool),
2046 (ValueSize, u64),
2047 (Seed, u64),
2048 (Partitions, u64),
2049 (BatchSize, u64)
2050);
2051
2052impl LoadGeneratorOptionExtracted {
2053 pub(super) fn ensure_only_valid_options(
2054 &self,
2055 loadgen: &ast::LoadGenerator,
2056 ) -> Result<(), PlanError> {
2057 use mz_sql_parser::ast::LoadGeneratorOptionName::*;
2058
2059 let mut options = self.seen.clone();
2060
2061 let permitted_options: &[_] = match loadgen {
2062 ast::LoadGenerator::Auction => &[TickInterval, AsOf, UpTo],
2063 ast::LoadGenerator::Clock => &[TickInterval, AsOf, UpTo],
2064 ast::LoadGenerator::Counter => &[TickInterval, AsOf, UpTo, MaxCardinality],
2065 ast::LoadGenerator::Marketing => &[TickInterval, AsOf, UpTo],
2066 ast::LoadGenerator::Datums => &[TickInterval, AsOf, UpTo],
2067 ast::LoadGenerator::Tpch => &[TickInterval, AsOf, UpTo, ScaleFactor],
2068 ast::LoadGenerator::KeyValue => &[
2069 TickInterval,
2070 Keys,
2071 SnapshotRounds,
2072 TransactionalSnapshot,
2073 ValueSize,
2074 Seed,
2075 Partitions,
2076 BatchSize,
2077 ],
2078 };
2079
2080 for o in permitted_options {
2081 options.remove(o);
2082 }
2083
2084 if !options.is_empty() {
2085 sql_bail!(
2086 "{} load generators do not support {} values",
2087 loadgen,
2088 options.iter().join(", ")
2089 )
2090 }
2091
2092 Ok(())
2093 }
2094}
2095
2096pub(crate) fn load_generator_ast_to_generator(
2097 scx: &StatementContext,
2098 loadgen: &ast::LoadGenerator,
2099 options: &[LoadGeneratorOption<Aug>],
2100 include_metadata: &[SourceIncludeMetadata],
2101) -> Result<LoadGenerator, PlanError> {
2102 let extracted: LoadGeneratorOptionExtracted = options.to_vec().try_into()?;
2103 extracted.ensure_only_valid_options(loadgen)?;
2104
2105 if loadgen != &ast::LoadGenerator::KeyValue && !include_metadata.is_empty() {
2106 sql_bail!("INCLUDE metadata only supported with `KEY VALUE` load generators");
2107 }
2108
2109 let load_generator = match loadgen {
2110 ast::LoadGenerator::Auction => LoadGenerator::Auction,
2111 ast::LoadGenerator::Clock => {
2112 scx.require_feature_flag(&vars::ENABLE_LOAD_GENERATOR_CLOCK)?;
2113 LoadGenerator::Clock
2114 }
2115 ast::LoadGenerator::Counter => {
2116 scx.require_feature_flag(&vars::ENABLE_LOAD_GENERATOR_COUNTER)?;
2117 let LoadGeneratorOptionExtracted {
2118 max_cardinality, ..
2119 } = extracted;
2120 LoadGenerator::Counter { max_cardinality }
2121 }
2122 ast::LoadGenerator::Marketing => LoadGenerator::Marketing,
2123 ast::LoadGenerator::Datums => {
2124 scx.require_feature_flag(&vars::ENABLE_LOAD_GENERATOR_DATUMS)?;
2125 LoadGenerator::Datums
2126 }
2127 ast::LoadGenerator::Tpch => {
2128 let LoadGeneratorOptionExtracted { scale_factor, .. } = extracted;
2129
2130 let sf: f64 = scale_factor.unwrap_or(0.01);
2132 if !sf.is_finite() || sf < 0.0 {
2133 sql_bail!("unsupported scale factor {sf}");
2134 }
2135
2136 let f_to_i = |multiplier: f64| -> Result<i64, PlanError> {
2137 let total = (sf * multiplier).floor();
2138 let mut i = i64::try_cast_from(total)
2139 .ok_or_else(|| sql_err!("unsupported scale factor {sf}"))?;
2140 if i < 1 {
2141 i = 1;
2142 }
2143 Ok(i)
2144 };
2145
2146 let count_supplier = f_to_i(10_000f64)?;
2149 let count_part = f_to_i(200_000f64)?;
2150 let count_customer = f_to_i(150_000f64)?;
2151 let count_orders = f_to_i(150_000f64 * 10f64)?;
2152 let count_clerk = f_to_i(1_000f64)?;
2153
2154 LoadGenerator::Tpch {
2155 count_supplier,
2156 count_part,
2157 count_customer,
2158 count_orders,
2159 count_clerk,
2160 }
2161 }
2162 mz_sql_parser::ast::LoadGenerator::KeyValue => {
2163 scx.require_feature_flag(&vars::ENABLE_LOAD_GENERATOR_KEY_VALUE)?;
2164 let LoadGeneratorOptionExtracted {
2165 keys,
2166 snapshot_rounds,
2167 transactional_snapshot,
2168 value_size,
2169 tick_interval,
2170 seed,
2171 partitions,
2172 batch_size,
2173 ..
2174 } = extracted;
2175
2176 let mut include_offset = None;
2177 for im in include_metadata {
2178 match im {
2179 SourceIncludeMetadata::Offset { alias } => {
2180 include_offset = match alias {
2181 Some(alias) => Some(alias.to_string()),
2182 None => Some(LOAD_GENERATOR_KEY_VALUE_OFFSET_DEFAULT.to_string()),
2183 }
2184 }
2185 SourceIncludeMetadata::Key { .. } => continue,
2186
2187 _ => {
2188 sql_bail!("only `INCLUDE OFFSET` and `INCLUDE KEY` is supported");
2189 }
2190 };
2191 }
2192
2193 let lgkv = KeyValueLoadGenerator {
2194 keys: keys.ok_or_else(|| sql_err!("LOAD GENERATOR KEY VALUE requires KEYS"))?,
2195 snapshot_rounds: snapshot_rounds
2196 .ok_or_else(|| sql_err!("LOAD GENERATOR KEY VALUE requires SNAPSHOT ROUNDS"))?,
2197 transactional_snapshot: transactional_snapshot.unwrap_or(true),
2199 value_size: value_size
2200 .ok_or_else(|| sql_err!("LOAD GENERATOR KEY VALUE requires VALUE SIZE"))?,
2201 partitions: partitions
2202 .ok_or_else(|| sql_err!("LOAD GENERATOR KEY VALUE requires PARTITIONS"))?,
2203 tick_interval,
2204 batch_size: batch_size
2205 .ok_or_else(|| sql_err!("LOAD GENERATOR KEY VALUE requires BATCH SIZE"))?,
2206 seed: seed.ok_or_else(|| sql_err!("LOAD GENERATOR KEY VALUE requires SEED"))?,
2207 include_offset,
2208 };
2209
2210 if lgkv.keys == 0
2211 || lgkv.partitions == 0
2212 || lgkv.value_size == 0
2213 || lgkv.batch_size == 0
2214 {
2215 sql_bail!("LOAD GENERATOR KEY VALUE options must be non-zero")
2216 }
2217
2218 if lgkv.keys % lgkv.partitions != 0 {
2219 sql_bail!("KEYS must be a multiple of PARTITIONS")
2220 }
2221
2222 if lgkv.batch_size > lgkv.keys {
2223 sql_bail!("KEYS must be larger than BATCH SIZE")
2224 }
2225
2226 if (lgkv.keys / lgkv.partitions) % lgkv.batch_size != 0 {
2229 sql_bail!("PARTITIONS * BATCH SIZE must be a divisor of KEYS")
2230 }
2231
2232 if lgkv.snapshot_rounds == 0 {
2233 sql_bail!("SNAPSHOT ROUNDS must be larger than 0")
2234 }
2235
2236 LoadGenerator::KeyValue(lgkv)
2237 }
2238 };
2239
2240 Ok(load_generator)
2241}
2242
2243fn typecheck_debezium(value_desc: &RelationDesc) -> Result<(Option<usize>, usize), PlanError> {
2244 let before = value_desc.get_by_name(&"before".into());
2245 let (after_idx, after_ty) = value_desc
2246 .get_by_name(&"after".into())
2247 .ok_or_else(|| sql_err!("'after' column missing from debezium input"))?;
2248 let before_idx = if let Some((before_idx, before_ty)) = before {
2249 if !matches!(before_ty.scalar_type, SqlScalarType::Record { .. }) {
2250 sql_bail!("'before' column must be of type record");
2251 }
2252 if before_ty != after_ty {
2253 sql_bail!("'before' type differs from 'after' column");
2254 }
2255 Some(before_idx)
2256 } else {
2257 None
2258 };
2259 Ok((before_idx, after_idx))
2260}
2261
2262fn get_encoding(
2263 scx: &StatementContext,
2264 format: &FormatSpecifier<Aug>,
2265 envelope: &ast::SourceEnvelope,
2266) -> Result<SourceDataEncoding<ReferencedConnection>, PlanError> {
2267 let encoding = match format {
2268 FormatSpecifier::Bare(format) => get_encoding_inner(scx, format)?,
2269 FormatSpecifier::KeyValue { key, value } => {
2270 let key = {
2271 let encoding = get_encoding_inner(scx, key)?;
2272 Some(encoding.key.unwrap_or(encoding.value))
2273 };
2274 let value = get_encoding_inner(scx, value)?.value;
2275 SourceDataEncoding { key, value }
2276 }
2277 };
2278
2279 let requires_keyvalue = matches!(
2280 envelope,
2281 ast::SourceEnvelope::Debezium | ast::SourceEnvelope::Upsert { .. }
2282 );
2283 let is_keyvalue = encoding.key.is_some();
2284 if requires_keyvalue && !is_keyvalue {
2285 sql_bail!("ENVELOPE [DEBEZIUM] UPSERT requires that KEY FORMAT be specified");
2286 };
2287
2288 Ok(encoding)
2289}
2290
2291fn source_sink_cluster_config<'a, 'ctx>(
2297 scx: &'a StatementContext<'ctx>,
2298 in_cluster: &mut Option<ResolvedClusterName>,
2299) -> Result<&'a dyn CatalogCluster<'ctx>, PlanError> {
2300 let cluster = match in_cluster {
2301 None => {
2302 let cluster = scx.catalog.resolve_cluster(None)?;
2303 *in_cluster = Some(ResolvedClusterName {
2304 id: cluster.id(),
2305 print_name: None,
2306 });
2307 cluster
2308 }
2309 Some(in_cluster) => scx.catalog.get_cluster(in_cluster.id),
2310 };
2311
2312 Ok(cluster)
2313}
2314
2315generate_extracted_config!(AvroSchemaOption, (ConfluentWireFormat, bool, Default(true)));
2316
2317generate_extracted_config!(
2321 GlueAvroOption,
2322 (SchemaName, String),
2323 (KeySchemaName, String),
2324 (ValueSchemaName, String),
2325 (KeyCompatibilityLevel, String),
2326 (ValueCompatibilityLevel, String)
2327);
2328
2329#[derive(Debug)]
2330pub struct Schema {
2331 pub key_schema: Option<String>,
2332 pub value_schema: String,
2333 pub key_reference_schemas: Vec<String>,
2335 pub value_reference_schemas: Vec<String>,
2337 pub wire_format: WireFormat<ReferencedConnection>,
2341}
2342
2343fn get_encoding_inner(
2344 scx: &StatementContext,
2345 format: &Format<Aug>,
2346) -> Result<SourceDataEncoding<ReferencedConnection>, PlanError> {
2347 let value = match format {
2348 Format::Bytes => DataEncoding::Bytes,
2349 Format::Avro(schema) => {
2350 let Schema {
2351 key_schema,
2352 value_schema,
2353 key_reference_schemas,
2354 value_reference_schemas,
2355 wire_format,
2356 } = match schema {
2357 AvroSchema::InlineSchema {
2360 schema: ast::Schema { schema },
2361 with_options,
2362 } => {
2363 let AvroSchemaOptionExtracted {
2364 confluent_wire_format,
2365 ..
2366 } = with_options.clone().try_into()?;
2367 let wire_format = if confluent_wire_format {
2368 WireFormat::Confluent { registry: None }
2369 } else {
2370 WireFormat::None
2371 };
2372 Schema {
2373 key_schema: None,
2374 value_schema: schema.clone(),
2375 key_reference_schemas: vec![],
2376 value_reference_schemas: vec![],
2377 wire_format,
2378 }
2379 }
2380 AvroSchema::Csr {
2381 csr_connection:
2382 CsrConnectionAvro {
2383 connection,
2384 seed,
2385 key_strategy: _,
2386 value_strategy: _,
2387 },
2388 } => {
2389 let item = scx.get_item_by_resolved_name(&connection.connection)?;
2390 let csr_connection = match item.connection()? {
2391 Connection::Csr(_) => item.id(),
2392 _ => {
2393 sql_bail!(
2394 "{} is not a Confluent Schema Registry connection",
2395 scx.catalog
2396 .resolve_full_name(item.name())
2397 .to_string()
2398 .quoted()
2399 )
2400 }
2401 };
2402
2403 if let Some(seed) = seed {
2404 Schema {
2405 key_schema: seed.key_schema.clone(),
2406 value_schema: seed.value_schema.clone(),
2407 key_reference_schemas: seed.key_reference_schemas.clone(),
2408 value_reference_schemas: seed.value_reference_schemas.clone(),
2409 wire_format: WireFormat::Confluent {
2410 registry: Some(csr_connection),
2411 },
2412 }
2413 } else {
2414 sql_bail!("Avro CSR seed resolution has not been performed")
2415 }
2416 }
2417 AvroSchema::Glue {
2418 connection,
2419 with_options: _,
2420 seed,
2421 } => {
2422 let item = scx.get_item_by_resolved_name(connection)?;
2423 let glue_connection = match item.connection()? {
2424 Connection::GlueSchemaRegistry(_) => item.id(),
2425 _ => {
2426 sql_bail!(
2427 "{} is not an AWS Glue Schema Registry connection",
2428 scx.catalog
2429 .resolve_full_name(item.name())
2430 .to_string()
2431 .quoted()
2432 )
2433 }
2434 };
2435
2436 let Some(seed) = seed else {
2440 sql_bail!("Avro Glue seed resolution has not been performed");
2441 };
2442
2443 Schema {
2444 key_schema: None,
2445 value_schema: seed.value_schema.clone(),
2446 key_reference_schemas: vec![],
2447 value_reference_schemas: vec![],
2448 wire_format: WireFormat::Glue {
2449 registry: Some(glue_connection),
2450 },
2451 }
2452 }
2453 };
2454
2455 if let Some(key_schema) = key_schema {
2456 return Ok(SourceDataEncoding {
2457 key: Some(DataEncoding::Avro(AvroEncoding {
2458 schema: key_schema,
2459 reference_schemas: key_reference_schemas,
2460 wire_format: wire_format.clone(),
2461 })),
2462 value: DataEncoding::Avro(AvroEncoding {
2463 schema: value_schema,
2464 reference_schemas: value_reference_schemas,
2465 wire_format,
2466 }),
2467 });
2468 } else {
2469 DataEncoding::Avro(AvroEncoding {
2470 schema: value_schema,
2471 reference_schemas: value_reference_schemas,
2472 wire_format,
2473 })
2474 }
2475 }
2476 Format::Protobuf(schema) => match schema {
2477 ProtobufSchema::Csr {
2478 csr_connection:
2479 CsrConnectionProtobuf {
2480 connection:
2481 CsrConnection {
2482 connection,
2483 options,
2484 },
2485 seed,
2486 },
2487 } => {
2488 if let Some(CsrSeedProtobuf { key, value }) = seed {
2489 let item = scx.get_item_by_resolved_name(connection)?;
2490 let _ = match item.connection()? {
2491 Connection::Csr(connection) => connection,
2492 _ => {
2493 sql_bail!(
2494 "{} is not a schema registry connection",
2495 scx.catalog
2496 .resolve_full_name(item.name())
2497 .to_string()
2498 .quoted()
2499 )
2500 }
2501 };
2502
2503 if !options.is_empty() {
2504 sql_bail!("Protobuf CSR connections do not support any options");
2505 }
2506
2507 let value = DataEncoding::Protobuf(ProtobufEncoding {
2508 descriptors: strconv::parse_bytes(&value.schema)?,
2509 message_name: value.message_name.clone(),
2510 confluent_wire_format: true,
2511 });
2512 if let Some(key) = key {
2513 return Ok(SourceDataEncoding {
2514 key: Some(DataEncoding::Protobuf(ProtobufEncoding {
2515 descriptors: strconv::parse_bytes(&key.schema)?,
2516 message_name: key.message_name.clone(),
2517 confluent_wire_format: true,
2518 })),
2519 value,
2520 });
2521 }
2522 value
2523 } else {
2524 sql_bail!("Protobuf CSR seed resolution has not been performed")
2525 }
2526 }
2527 ProtobufSchema::InlineSchema {
2528 message_name,
2529 schema: ast::Schema { schema },
2530 } => {
2531 let descriptors = strconv::parse_bytes(schema)?;
2532
2533 DataEncoding::Protobuf(ProtobufEncoding {
2534 descriptors,
2535 message_name: message_name.to_owned(),
2536 confluent_wire_format: false,
2537 })
2538 }
2539 },
2540 Format::Regex(regex) => DataEncoding::Regex(RegexEncoding {
2541 regex: mz_repr::adt::regex::Regex::new(regex, false)
2542 .map_err(|e| sql_err!("parsing regex: {e}"))?,
2543 }),
2544 Format::Csv { columns, delimiter } => {
2545 let columns = match columns {
2546 CsvColumns::Header { names } => {
2547 if names.is_empty() {
2548 sql_bail!("[internal error] column spec should get names in purify")
2549 }
2550 ColumnSpec::Header {
2551 names: names.iter().cloned().map(|n| n.into_string()).collect(),
2552 }
2553 }
2554 CsvColumns::Count(n) => ColumnSpec::Count(usize::cast_from(*n)),
2555 };
2556 DataEncoding::Csv(CsvEncoding {
2557 columns,
2558 delimiter: u8::try_from(*delimiter)
2559 .map_err(|_| sql_err!("CSV delimiter must be an ASCII character"))?,
2560 })
2561 }
2562 Format::Json { array: false } => DataEncoding::Json,
2563 Format::Json { array: true } => bail_unsupported!("JSON ARRAY format in sources"),
2564 Format::Text => DataEncoding::Text,
2565 };
2566 Ok(SourceDataEncoding { key: None, value })
2567}
2568
2569fn get_key_envelope(
2571 included_items: &[SourceIncludeMetadata],
2572 encoding: Option<&SourceDataEncoding<ReferencedConnection>>,
2573 key_envelope_no_encoding: bool,
2574) -> Result<KeyEnvelope, PlanError> {
2575 let key_definition = included_items
2576 .iter()
2577 .find(|i| matches!(i, SourceIncludeMetadata::Key { .. }));
2578 if let Some(SourceIncludeMetadata::Key { alias }) = key_definition {
2579 match (alias, encoding.and_then(|e| e.key.as_ref())) {
2580 (Some(name), Some(_)) => Ok(KeyEnvelope::Named(name.as_str().to_string())),
2581 (None, Some(key)) => get_unnamed_key_envelope(Some(key)),
2582 (Some(name), _) if key_envelope_no_encoding => {
2583 Ok(KeyEnvelope::Named(name.as_str().to_string()))
2584 }
2585 (None, _) if key_envelope_no_encoding => get_unnamed_key_envelope(None),
2586 (_, None) => {
2587 sql_bail!(
2591 "INCLUDE KEY requires specifying KEY FORMAT .. VALUE FORMAT, \
2592 got bare FORMAT"
2593 );
2594 }
2595 }
2596 } else {
2597 Ok(KeyEnvelope::None)
2598 }
2599}
2600
2601fn get_unnamed_key_envelope(
2604 key: Option<&DataEncoding<ReferencedConnection>>,
2605) -> Result<KeyEnvelope, PlanError> {
2606 let is_composite = match key {
2610 Some(DataEncoding::Bytes | DataEncoding::Json | DataEncoding::Text) => false,
2611 Some(
2612 DataEncoding::Avro(_)
2613 | DataEncoding::Csv(_)
2614 | DataEncoding::Protobuf(_)
2615 | DataEncoding::Regex { .. },
2616 ) => true,
2617 None => false,
2618 };
2619
2620 if is_composite {
2621 Ok(KeyEnvelope::Flattened)
2622 } else {
2623 Ok(KeyEnvelope::Named("key".to_string()))
2624 }
2625}
2626
2627pub fn describe_create_view(
2628 _: &StatementContext,
2629 _: CreateViewStatement<Aug>,
2630) -> Result<StatementDesc, PlanError> {
2631 Ok(StatementDesc::new(None))
2632}
2633
2634pub fn plan_view(
2635 scx: &StatementContext,
2636 def: &mut ViewDefinition<Aug>,
2637 temporary: bool,
2638) -> Result<(QualifiedItemName, View), PlanError> {
2639 let create_sql = normalize::create_statement(
2640 scx,
2641 Statement::CreateView(CreateViewStatement {
2642 if_exists: IfExistsBehavior::Error,
2643 temporary,
2644 definition: def.clone(),
2645 }),
2646 )?;
2647
2648 let ViewDefinition {
2649 name,
2650 columns,
2651 query,
2652 } = def;
2653
2654 let query::PlannedRootQuery {
2655 expr,
2656 mut desc,
2657 finishing,
2658 scope: _,
2659 } = query::plan_root_query(scx, query.clone(), QueryLifetime::View)?;
2660 assert!(HirRelationExpr::is_trivial_row_set_finishing_hir(
2666 &finishing,
2667 expr.arity()
2668 ));
2669 if expr.contains_parameters()? {
2670 return Err(PlanError::ParameterNotAllowed("views".to_string()));
2671 }
2672
2673 let dependencies = expr
2674 .depends_on()
2675 .into_iter()
2676 .map(|gid| scx.catalog.resolve_item_id(&gid))
2677 .collect();
2678
2679 let name = if temporary {
2680 scx.allocate_temporary_qualified_name(normalize::unresolved_item_name(name.to_owned())?)?
2681 } else {
2682 scx.allocate_qualified_name(normalize::unresolved_item_name(name.to_owned())?)?
2683 };
2684
2685 plan_utils::maybe_rename_columns_exact(
2686 scx.catalog,
2687 format!("view {}", scx.catalog.resolve_full_name(&name)),
2688 &mut desc,
2689 columns,
2690 )?;
2691 let names: Vec<ColumnName> = desc.iter_names().cloned().collect();
2692
2693 if let Some(dup) = names.iter().duplicates().next() {
2694 sql_bail!("column {} specified more than once", dup.quoted());
2695 }
2696
2697 let view = View {
2698 create_sql,
2699 expr,
2700 dependencies,
2701 column_names: names,
2702 temporary,
2703 };
2704
2705 Ok((name, view))
2706}
2707
2708pub fn plan_create_view(
2709 scx: &StatementContext,
2710 mut stmt: CreateViewStatement<Aug>,
2711) -> Result<Plan, PlanError> {
2712 let CreateViewStatement {
2713 temporary,
2714 if_exists,
2715 definition,
2716 } = &mut stmt;
2717 let (name, view) = plan_view(scx, definition, *temporary)?;
2718
2719 let ignore_if_exists_errors = scx.pcx().map_or(false, |pcx| pcx.ignore_if_exists_errors);
2722
2723 let replace = if *if_exists == IfExistsBehavior::Replace && !ignore_if_exists_errors {
2724 let if_exists = true;
2725 let cascade = false;
2726 let maybe_item_to_drop = plan_drop_item(
2727 scx,
2728 ObjectType::View,
2729 if_exists,
2730 definition.name.clone(),
2731 cascade,
2732 )?;
2733
2734 if let Some(id) = maybe_item_to_drop {
2736 let dependencies = view.expr.depends_on();
2737 let invalid_drop = scx
2738 .get_item(&id)
2739 .global_ids()
2740 .any(|gid| dependencies.contains(&gid));
2741 if invalid_drop {
2742 let item = scx.catalog.get_item(&id);
2743 sql_bail!(
2744 "cannot replace view {0}: depended upon by new {0} definition",
2745 scx.catalog.resolve_full_name(item.name())
2746 );
2747 }
2748
2749 Some(id)
2750 } else {
2751 None
2752 }
2753 } else {
2754 None
2755 };
2756 let drop_ids = replace
2757 .map(|id| {
2758 scx.catalog
2759 .item_dependents(id)
2760 .into_iter()
2761 .map(|id| id.unwrap_item_id())
2762 .collect()
2763 })
2764 .unwrap_or_default();
2765
2766 validate_view_dependencies(scx, &view.dependencies.0)?;
2767
2768 let full_name = scx.catalog.resolve_full_name(&name);
2770 let partial_name = PartialItemName::from(full_name.clone());
2771 if let (Ok(item), IfExistsBehavior::Error, false) = (
2774 scx.catalog.resolve_item_or_type(&partial_name),
2775 *if_exists,
2776 ignore_if_exists_errors,
2777 ) {
2778 return Err(PlanError::ItemAlreadyExists {
2779 name: full_name.to_string(),
2780 item_type: item.item_type(),
2781 });
2782 }
2783
2784 Ok(Plan::CreateView(CreateViewPlan {
2785 name,
2786 view,
2787 replace,
2788 drop_ids,
2789 if_not_exists: *if_exists == IfExistsBehavior::Skip,
2790 ambiguous_columns: *scx.ambiguous_columns.borrow(),
2791 }))
2792}
2793
2794fn validate_view_dependencies(
2796 scx: &StatementContext,
2797 dependencies: &BTreeSet<CatalogItemId>,
2798) -> Result<(), PlanError> {
2799 for id in dependencies {
2800 let item = scx.catalog.get_item(id);
2801 if item.replacement_target().is_some() {
2802 let name = scx.catalog.minimal_qualification(item.name());
2803 return Err(PlanError::InvalidDependency {
2804 name: name.to_string(),
2805 item_type: format!("replacement {}", item.item_type()),
2806 });
2807 }
2808 }
2809
2810 Ok(())
2811}
2812
2813pub fn describe_create_materialized_view(
2814 _: &StatementContext,
2815 _: CreateMaterializedViewStatement<Aug>,
2816) -> Result<StatementDesc, PlanError> {
2817 Ok(StatementDesc::new(None))
2818}
2819
2820pub fn describe_create_network_policy(
2821 _: &StatementContext,
2822 _: CreateNetworkPolicyStatement<Aug>,
2823) -> Result<StatementDesc, PlanError> {
2824 Ok(StatementDesc::new(None))
2825}
2826
2827pub fn describe_alter_network_policy(
2828 _: &StatementContext,
2829 _: AlterNetworkPolicyStatement<Aug>,
2830) -> Result<StatementDesc, PlanError> {
2831 Ok(StatementDesc::new(None))
2832}
2833
2834fn check_refresh_time(option: &str, ts: Timestamp) -> Result<(), PlanError> {
2837 let renderable = i64::try_from(ts)
2838 .ok()
2839 .and_then(DateTime::from_timestamp_millis)
2840 .is_some_and(|dt| CheckedTimestamp::try_from(dt).is_ok());
2841 if !renderable {
2842 sql_bail!("{option} time too large: {ts}");
2843 }
2844 Ok(())
2845}
2846
2847pub fn plan_create_materialized_view(
2848 scx: &StatementContext,
2849 mut stmt: CreateMaterializedViewStatement<Aug>,
2850) -> Result<Plan, PlanError> {
2851 let cluster_id =
2852 crate::plan::statement::resolve_cluster_for_materialized_view(scx.catalog, &stmt)?;
2853 stmt.in_cluster = Some(ResolvedClusterName {
2854 id: cluster_id,
2855 print_name: None,
2856 });
2857
2858 let target_replica = match &stmt.in_cluster_replica {
2859 Some(replica_name) => {
2860 scx.require_feature_flag(&ENABLE_REPLICA_TARGETED_MATERIALIZED_VIEWS)?;
2861
2862 let cluster = scx.catalog.get_cluster(cluster_id);
2863 let replica_id = cluster
2864 .replica_ids()
2865 .get(replica_name.as_str())
2866 .copied()
2867 .ok_or_else(|| {
2868 CatalogError::UnknownClusterReplica(replica_name.as_str().to_string())
2869 })?;
2870 Some(replica_id)
2871 }
2872 None => None,
2873 };
2874
2875 let create_sql =
2876 normalize::create_statement(scx, Statement::CreateMaterializedView(stmt.clone()))?;
2877
2878 let partial_name = normalize::unresolved_item_name(stmt.name)?;
2879 let name = scx.allocate_qualified_name(partial_name.clone())?;
2880
2881 let query::PlannedRootQuery {
2882 expr,
2883 mut desc,
2884 finishing,
2885 scope: _,
2886 } = query::plan_root_query(scx, stmt.query, QueryLifetime::MaterializedView)?;
2887 assert!(HirRelationExpr::is_trivial_row_set_finishing_hir(
2889 &finishing,
2890 expr.arity()
2891 ));
2892 if expr.contains_parameters()? {
2893 return Err(PlanError::ParameterNotAllowed(
2894 "materialized views".to_string(),
2895 ));
2896 }
2897
2898 plan_utils::maybe_rename_columns_exact(
2899 scx.catalog,
2900 format!("materialized view {}", scx.catalog.resolve_full_name(&name)),
2901 &mut desc,
2902 &stmt.columns,
2903 )?;
2904 let column_names: Vec<ColumnName> = desc.iter_names().cloned().collect();
2905
2906 let MaterializedViewOptionExtracted {
2907 assert_not_null,
2908 partition_by,
2909 retain_history,
2910 refresh,
2911 seen: _,
2912 }: MaterializedViewOptionExtracted = stmt.with_options.try_into()?;
2913
2914 if let Some(partition_by) = partition_by {
2915 scx.require_feature_flag(&ENABLE_COLLECTION_PARTITION_BY)?;
2916 check_partition_by(&desc, partition_by)?;
2917 }
2918
2919 let refresh_schedule = {
2920 let mut refresh_schedule = RefreshSchedule::default();
2921 let mut on_commits_seen = 0;
2922 for refresh_option_value in refresh {
2923 if !matches!(refresh_option_value, RefreshOptionValue::OnCommit) {
2924 scx.require_feature_flag(&ENABLE_REFRESH_EVERY_MVS)?;
2925 }
2926 match refresh_option_value {
2927 RefreshOptionValue::OnCommit => {
2928 on_commits_seen += 1;
2929 }
2930 RefreshOptionValue::AtCreation => {
2931 soft_panic_or_log!("REFRESH AT CREATION should have been purified away");
2932 bail_internal!("REFRESH AT CREATION should have been purified away")
2933 }
2934 RefreshOptionValue::At(RefreshAtOptionValue { mut time }) => {
2935 transform_ast::transform(scx, &mut time)?; let ecx = &ExprContext {
2937 qcx: &QueryContext::root(scx, QueryLifetime::OneShot),
2938 name: "REFRESH AT",
2939 scope: &Scope::empty(),
2940 relation_type: &SqlRelationType::empty(),
2941 allow_aggregates: false,
2942 allow_subqueries: false,
2943 allow_parameters: false,
2944 allow_windows: false,
2945 };
2946 let hir = plan_expr(ecx, &time)?.cast_to(
2947 ecx,
2948 CastContext::Assignment,
2949 &SqlScalarType::MzTimestamp,
2950 )?;
2951 let timestamp = hir
2953 .into_literal_mz_timestamp()
2954 .ok_or_else(|| PlanError::InvalidRefreshAt)?;
2955 check_refresh_time("REFRESH AT", timestamp)?;
2956 refresh_schedule.ats.push(timestamp);
2957 }
2958 RefreshOptionValue::Every(RefreshEveryOptionValue {
2959 interval,
2960 aligned_to,
2961 }) => {
2962 let interval = Interval::try_from_value(Value::Interval(interval))?;
2963 if interval.as_microseconds() <= 0 {
2964 sql_bail!("REFRESH interval must be positive; got: {}", interval);
2965 }
2966 if interval.months != 0 {
2967 sql_bail!("REFRESH interval must not involve units larger than days");
2972 }
2973 let interval = interval.duration()?;
2974 if u64::try_from(interval.as_millis()).is_err()
2979 || Interval::from_duration(&interval).is_err()
2980 {
2981 sql_bail!("REFRESH interval too large");
2982 }
2983 if interval.as_micros() < 1000 {
2984 sql_bail!("REFRESH interval must be at least 1 ms")
2985 }
2986
2987 let mut aligned_to = match aligned_to {
2988 Some(aligned_to) => aligned_to,
2989 None => {
2990 soft_panic_or_log!(
2991 "ALIGNED TO should have been filled in by purification"
2992 );
2993 sql_bail!(
2994 "INTERNAL ERROR: ALIGNED TO should have been filled in by purification"
2995 )
2996 }
2997 };
2998
2999 transform_ast::transform(scx, &mut aligned_to)?;
3001
3002 let ecx = &ExprContext {
3003 qcx: &QueryContext::root(scx, QueryLifetime::OneShot),
3004 name: "REFRESH EVERY ... ALIGNED TO",
3005 scope: &Scope::empty(),
3006 relation_type: &SqlRelationType::empty(),
3007 allow_aggregates: false,
3008 allow_subqueries: false,
3009 allow_parameters: false,
3010 allow_windows: false,
3011 };
3012 let aligned_to_hir = plan_expr(ecx, &aligned_to)?.cast_to(
3013 ecx,
3014 CastContext::Assignment,
3015 &SqlScalarType::MzTimestamp,
3016 )?;
3017 let aligned_to_const = aligned_to_hir
3019 .into_literal_mz_timestamp()
3020 .ok_or_else(|| PlanError::InvalidRefreshEveryAlignedTo)?;
3021 check_refresh_time("REFRESH EVERY ... ALIGNED TO", aligned_to_const)?;
3022
3023 refresh_schedule.everies.push(RefreshEvery {
3024 interval,
3025 aligned_to: aligned_to_const,
3026 });
3027 }
3028 }
3029 }
3030
3031 if on_commits_seen > 1 {
3032 sql_bail!("REFRESH ON COMMIT cannot be specified multiple times");
3033 }
3034 if on_commits_seen > 0 && refresh_schedule != RefreshSchedule::default() {
3035 sql_bail!("REFRESH ON COMMIT is not compatible with any of the other REFRESH options");
3036 }
3037
3038 if refresh_schedule == RefreshSchedule::default() {
3039 None
3040 } else {
3041 Some(refresh_schedule)
3042 }
3043 };
3044
3045 let as_of = stmt.as_of.map(Timestamp::from);
3046 let compaction_window = plan_retain_history_option(scx, retain_history)?;
3047 let mut non_null_assertions = assert_not_null
3048 .into_iter()
3049 .map(normalize::column_name)
3050 .map(|assertion_name| {
3051 column_names
3052 .iter()
3053 .position(|col| col == &assertion_name)
3054 .ok_or_else(|| {
3055 sql_err!(
3056 "column {} in ASSERT NOT NULL option not found",
3057 assertion_name.quoted()
3058 )
3059 })
3060 })
3061 .collect::<Result<Vec<_>, _>>()?;
3062 non_null_assertions.sort();
3063 if let Some(dup) = non_null_assertions.iter().duplicates().next() {
3064 let dup = &column_names[*dup];
3065 sql_bail!("duplicate column {} in non-null assertions", dup.quoted());
3066 }
3067
3068 if let Some(dup) = column_names.iter().duplicates().next() {
3069 sql_bail!("column {} specified more than once", dup.quoted());
3070 }
3071
3072 let if_exists = match scx.pcx().map(|pcx| pcx.ignore_if_exists_errors) {
3075 Ok(true) => IfExistsBehavior::Skip,
3076 _ => stmt.if_exists,
3077 };
3078
3079 let mut replace = None;
3080 let mut if_not_exists = false;
3081 match if_exists {
3082 IfExistsBehavior::Replace => {
3083 let if_exists = true;
3084 let cascade = false;
3085 let replace_id = plan_drop_item(
3086 scx,
3087 ObjectType::MaterializedView,
3088 if_exists,
3089 partial_name.clone().into(),
3090 cascade,
3091 )?;
3092
3093 if let Some(id) = replace_id {
3095 let dependencies = expr.depends_on();
3096 let invalid_drop = scx
3097 .get_item(&id)
3098 .global_ids()
3099 .any(|gid| dependencies.contains(&gid));
3100 if invalid_drop {
3101 let item = scx.catalog.get_item(&id);
3102 sql_bail!(
3103 "cannot replace materialized view {0}: depended upon by new {0} definition",
3104 scx.catalog.resolve_full_name(item.name())
3105 );
3106 }
3107 replace = Some(id);
3108 }
3109 }
3110 IfExistsBehavior::Skip => if_not_exists = true,
3111 IfExistsBehavior::Error => (),
3112 }
3113 let drop_ids = replace
3114 .map(|id| {
3115 scx.catalog
3116 .item_dependents(id)
3117 .into_iter()
3118 .map(|id| id.unwrap_item_id())
3119 .collect()
3120 })
3121 .unwrap_or_default();
3122 let mut dependencies: BTreeSet<_> = expr
3123 .depends_on()
3124 .into_iter()
3125 .map(|gid| scx.catalog.resolve_item_id(&gid))
3126 .collect();
3127
3128 let mut replacement_target = None;
3130 if let Some(target_name) = &stmt.replacement_for {
3131 scx.require_feature_flag(&vars::ENABLE_REPLACEMENT_MATERIALIZED_VIEWS)?;
3132
3133 let target = scx.get_item_by_resolved_name(target_name)?;
3134 if target.item_type() != CatalogItemType::MaterializedView {
3135 return Err(PlanError::InvalidReplacement {
3136 item_type: target.item_type(),
3137 item_name: scx.catalog.minimal_qualification(target.name()),
3138 replacement_type: CatalogItemType::MaterializedView,
3139 replacement_name: partial_name,
3140 });
3141 }
3142 if target.id().is_system() {
3143 sql_bail!(
3144 "cannot replace {} because it is required by the database system",
3145 scx.catalog.minimal_qualification(target.name()),
3146 );
3147 }
3148
3149 for dependent in scx.catalog.item_dependents(target.id()) {
3151 if let ObjectId::Item(id) = dependent
3152 && dependencies.contains(&id)
3153 {
3154 sql_bail!(
3155 "replacement would cause {} to depend on itself",
3156 scx.catalog.minimal_qualification(target.name()),
3157 );
3158 }
3159 }
3160
3161 dependencies.insert(target.id());
3162
3163 for use_id in target.used_by() {
3164 let use_item = scx.get_item(use_id);
3165 if use_item.replacement_target() == Some(target.id()) {
3166 sql_bail!(
3167 "cannot replace {} because it already has a replacement: {}",
3168 scx.catalog.minimal_qualification(target.name()),
3169 scx.catalog.minimal_qualification(use_item.name()),
3170 );
3171 }
3172 }
3173
3174 replacement_target = Some(target.id());
3175 }
3176
3177 validate_view_dependencies(scx, &dependencies)?;
3178
3179 let full_name = scx.catalog.resolve_full_name(&name);
3181 let partial_name = PartialItemName::from(full_name.clone());
3182 if let (IfExistsBehavior::Error, Ok(item)) =
3185 (if_exists, scx.catalog.resolve_item_or_type(&partial_name))
3186 {
3187 return Err(PlanError::ItemAlreadyExists {
3188 name: full_name.to_string(),
3189 item_type: item.item_type(),
3190 });
3191 }
3192
3193 Ok(Plan::CreateMaterializedView(CreateMaterializedViewPlan {
3194 name,
3195 materialized_view: MaterializedView {
3196 create_sql,
3197 expr,
3198 dependencies: DependencyIds(dependencies),
3199 column_names,
3200 replacement_target,
3201 cluster_id,
3202 target_replica,
3203 non_null_assertions,
3204 compaction_window,
3205 refresh_schedule,
3206 as_of,
3207 },
3208 replace,
3209 drop_ids,
3210 if_not_exists,
3211 ambiguous_columns: *scx.ambiguous_columns.borrow(),
3212 }))
3213}
3214
3215generate_extracted_config!(
3216 MaterializedViewOption,
3217 (AssertNotNull, Ident, AllowMultiple),
3218 (PartitionBy, Vec<Ident>),
3219 (RetainHistory, OptionalDuration),
3220 (Refresh, RefreshOptionValue<Aug>, AllowMultiple)
3221);
3222
3223pub fn describe_create_sink(
3224 _: &StatementContext,
3225 _: CreateSinkStatement<Aug>,
3226) -> Result<StatementDesc, PlanError> {
3227 Ok(StatementDesc::new(None))
3228}
3229
3230generate_extracted_config!(
3231 CreateSinkOption,
3232 (Snapshot, bool),
3233 (PartitionStrategy, String),
3234 (Version, u64),
3235 (CommitInterval, Duration)
3236);
3237
3238pub fn plan_create_sink(
3239 scx: &StatementContext,
3240 stmt: CreateSinkStatement<Aug>,
3241) -> Result<Plan, PlanError> {
3242 let Some(name) = stmt.name.clone() else {
3244 return Err(PlanError::MissingName(CatalogItemType::Sink));
3245 };
3246 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name)?)?;
3247 let full_name = scx.catalog.resolve_full_name(&name);
3248 let partial_name = PartialItemName::from(full_name.clone());
3249 if let (false, Ok(item)) = (stmt.if_not_exists, scx.catalog.resolve_item(&partial_name)) {
3250 return Err(PlanError::ItemAlreadyExists {
3251 name: full_name.to_string(),
3252 item_type: item.item_type(),
3253 });
3254 }
3255
3256 plan_sink(scx, stmt)
3257}
3258
3259fn plan_sink(
3264 scx: &StatementContext,
3265 mut stmt: CreateSinkStatement<Aug>,
3266) -> Result<Plan, PlanError> {
3267 let CreateSinkStatement {
3268 name,
3269 in_cluster: _,
3270 from,
3271 connection,
3272 format,
3273 envelope,
3274 mode,
3275 if_not_exists,
3276 with_options,
3277 } = stmt.clone();
3278
3279 let Some(name) = name else {
3280 return Err(PlanError::MissingName(CatalogItemType::Sink));
3281 };
3282 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name)?)?;
3283
3284 let envelope = match (&connection, envelope, mode) {
3285 (CreateSinkConnection::Kafka { .. }, Some(ast::SinkEnvelope::Upsert), None) => {
3287 SinkEnvelope::Upsert
3288 }
3289 (CreateSinkConnection::Kafka { .. }, Some(ast::SinkEnvelope::Debezium), None) => {
3290 SinkEnvelope::Debezium
3291 }
3292 (CreateSinkConnection::Kafka { .. }, None, None) => {
3293 sql_bail!("ENVELOPE clause is required")
3294 }
3295 (CreateSinkConnection::Kafka { .. }, _, Some(_)) => {
3296 sql_bail!("MODE is not supported for Kafka sinks, use ENVELOPE instead")
3297 }
3298 (CreateSinkConnection::Iceberg { .. }, None, Some(ast::IcebergSinkMode::Upsert)) => {
3300 SinkEnvelope::Upsert
3301 }
3302 (CreateSinkConnection::Iceberg { .. }, None, Some(ast::IcebergSinkMode::Append)) => {
3303 SinkEnvelope::Append
3304 }
3305 (CreateSinkConnection::Iceberg { .. }, None, None) => {
3306 sql_bail!("MODE clause is required")
3307 }
3308 (CreateSinkConnection::Iceberg { .. }, Some(_), _) => {
3309 sql_bail!("ENVELOPE is not supported for Iceberg sinks, use MODE instead")
3310 }
3311 };
3312
3313 let from_name = &from;
3314 let from = scx.get_item_by_resolved_name(&from)?;
3315
3316 {
3317 use CatalogItemType::*;
3318 match from.item_type() {
3319 Table | Source | MaterializedView => {
3320 if from.replacement_target().is_some() {
3321 let name = scx.catalog.minimal_qualification(from.name());
3322 return Err(PlanError::InvalidSinkFrom {
3323 name: name.to_string(),
3324 item_type: format!("replacement {}", from.item_type()),
3325 });
3326 }
3327 }
3328 Sink | MetricSink | View | Index | Type | Func | Secret | Connection => {
3329 let name = scx.catalog.minimal_qualification(from.name());
3330 return Err(PlanError::InvalidSinkFrom {
3331 name: name.to_string(),
3332 item_type: from.item_type().to_string(),
3333 });
3334 }
3335 }
3336 }
3337
3338 if from.id().is_system() {
3339 bail_unsupported!("creating a sink directly on a catalog object");
3340 }
3341
3342 let desc = from
3343 .relation_desc()
3344 .ok_or_else(|| sql_err!("item does not have a relation description"))?;
3345 let key_indices = match &connection {
3346 CreateSinkConnection::Kafka { key: Some(key), .. }
3347 | CreateSinkConnection::Iceberg { key: Some(key), .. } => {
3348 let key_columns = key
3349 .key_columns
3350 .clone()
3351 .into_iter()
3352 .map(normalize::column_name)
3353 .collect::<Vec<_>>();
3354 let mut uniq = BTreeSet::new();
3355 for col in key_columns.iter() {
3356 if !uniq.insert(col) {
3357 sql_bail!("duplicate column referenced in KEY: {}", col);
3358 }
3359 }
3360 let indices = key_columns
3361 .iter()
3362 .map(|col| {
3363 let name_idx =
3364 desc.get_by_name(col)
3365 .map(|(idx, _type)| idx)
3366 .ok_or_else(|| {
3367 sql_err!("column referenced in KEY does not exist: {}", col)
3368 })?;
3369 if desc.get_unambiguous_name(name_idx).is_none() {
3370 sql_bail!("column referenced in KEY is ambiguous: {}", col);
3371 }
3372 Ok(name_idx)
3373 })
3374 .collect::<Result<Vec<_>, _>>()?;
3375
3376 if matches!(&connection, CreateSinkConnection::Iceberg { .. }) {
3379 let cols: Vec<_> = desc.iter().collect();
3380 for &idx in &indices {
3381 let (col_name, col_type) = cols[idx];
3382 let scalar = &col_type.scalar_type;
3383 let is_valid = matches!(
3384 scalar,
3385 SqlScalarType::Bool
3387 | SqlScalarType::Int16
3388 | SqlScalarType::Int32
3389 | SqlScalarType::Int64
3390 | SqlScalarType::UInt16
3391 | SqlScalarType::UInt32
3392 | SqlScalarType::UInt64
3393 | SqlScalarType::Numeric { .. }
3395 | SqlScalarType::Date
3397 | SqlScalarType::Time
3398 | SqlScalarType::Timestamp { .. }
3399 | SqlScalarType::TimestampTz { .. }
3400 | SqlScalarType::Interval
3401 | SqlScalarType::MzTimestamp
3402 | SqlScalarType::String
3404 | SqlScalarType::Char { .. }
3405 | SqlScalarType::VarChar { .. }
3406 | SqlScalarType::PgLegacyChar
3407 | SqlScalarType::PgLegacyName
3408 | SqlScalarType::Bytes
3409 | SqlScalarType::Jsonb
3410 | SqlScalarType::Uuid
3412 | SqlScalarType::Oid
3413 | SqlScalarType::RegProc
3414 | SqlScalarType::RegType
3415 | SqlScalarType::RegClass
3416 | SqlScalarType::MzAclItem
3417 | SqlScalarType::AclItem
3418 | SqlScalarType::Int2Vector
3419 );
3420 if !is_valid {
3421 return Err(PlanError::IcebergSinkUnsupportedKeyType {
3422 column: col_name.to_string(),
3423 column_type: format!("{:?}", scalar),
3424 });
3425 }
3426 }
3427 }
3428
3429 let is_valid_key = desc
3430 .typ()
3431 .keys
3432 .iter()
3433 .any(|key_columns| key_columns.iter().all(|column| indices.contains(column)));
3434
3435 if !is_valid_key && envelope == SinkEnvelope::Upsert {
3436 if key.not_enforced {
3437 scx.catalog
3438 .add_notice(PlanNotice::UpsertSinkKeyNotEnforced {
3439 key: key_columns.clone(),
3440 name: name.item.clone(),
3441 })
3442 } else {
3443 return Err(PlanError::UpsertSinkWithInvalidKey {
3444 name: from_name.full_name_str(),
3445 desired_key: key_columns.iter().map(|c| c.to_string()).collect(),
3446 valid_keys: desc
3447 .typ()
3448 .keys
3449 .iter()
3450 .map(|key| {
3451 key.iter()
3452 .map(|col| desc.get_name(*col).as_str().into())
3453 .collect()
3454 })
3455 .collect(),
3456 });
3457 }
3458 }
3459 Some(indices)
3460 }
3461 CreateSinkConnection::Kafka { key: None, .. }
3462 | CreateSinkConnection::Iceberg { key: None, .. } => None,
3463 };
3464
3465 if key_indices.is_some() && envelope == SinkEnvelope::Append {
3466 sql_bail!("KEY is not supported for MODE APPEND Iceberg sinks");
3467 }
3468
3469 if envelope == SinkEnvelope::Append {
3471 if let CreateSinkConnection::Iceberg { .. } = &connection {
3472 use mz_storage_types::sinks::{
3473 ICEBERG_APPEND_DIFF_COLUMN, ICEBERG_APPEND_TIMESTAMP_COLUMN,
3474 };
3475 for (col_name, _) in desc.iter() {
3476 if col_name.as_str() == ICEBERG_APPEND_DIFF_COLUMN
3477 || col_name.as_str() == ICEBERG_APPEND_TIMESTAMP_COLUMN
3478 {
3479 sql_bail!(
3480 "column {} conflicts with the system column that MODE APPEND \
3481 adds to the Iceberg table",
3482 col_name.quoted()
3483 );
3484 }
3485 }
3486 }
3487 }
3488
3489 let headers_index = match &connection {
3490 CreateSinkConnection::Kafka {
3491 headers: Some(headers),
3492 ..
3493 } => {
3494 scx.require_feature_flag(&ENABLE_KAFKA_SINK_HEADERS)?;
3495
3496 match envelope {
3497 SinkEnvelope::Upsert | SinkEnvelope::Append => (),
3498 SinkEnvelope::Debezium => {
3499 sql_bail!("HEADERS option is not supported with ENVELOPE DEBEZIUM")
3500 }
3501 };
3502
3503 let headers = normalize::column_name(headers.clone());
3504 let (idx, ty) = desc
3505 .get_by_name(&headers)
3506 .ok_or_else(|| sql_err!("HEADERS column ({}) is unknown", headers))?;
3507
3508 if desc.get_unambiguous_name(idx).is_none() {
3509 sql_bail!("HEADERS column ({}) is ambiguous", headers);
3510 }
3511
3512 match &ty.scalar_type {
3513 SqlScalarType::Map { value_type, .. }
3514 if matches!(&**value_type, SqlScalarType::String | SqlScalarType::Bytes) => {}
3515 _ => sql_bail!(
3516 "HEADERS column must have type map[text => text] or map[text => bytea]"
3517 ),
3518 }
3519
3520 Some(idx)
3521 }
3522 _ => None,
3523 };
3524
3525 let relation_key_indices = desc.typ().keys.get(0).cloned();
3527
3528 let key_desc_and_indices = key_indices.map(|key_indices| {
3529 let cols = desc
3530 .iter()
3531 .map(|(name, ty)| (name.clone(), ty.clone()))
3532 .collect::<Vec<_>>();
3533 let (names, types): (Vec<_>, Vec<_>) =
3534 key_indices.iter().map(|&idx| cols[idx].clone()).unzip();
3535 let typ = SqlRelationType::new(types);
3536 (RelationDesc::new(typ, names), key_indices)
3537 });
3538
3539 if key_desc_and_indices.is_none() && envelope == SinkEnvelope::Upsert {
3540 return Err(PlanError::UpsertSinkWithoutKey);
3541 }
3542
3543 let CreateSinkOptionExtracted {
3544 snapshot,
3545 version,
3546 partition_strategy: _,
3547 seen: _,
3548 commit_interval,
3549 } = with_options.try_into()?;
3550
3551 let connection_builder = match connection {
3552 CreateSinkConnection::Kafka {
3553 connection,
3554 options,
3555 ..
3556 } => kafka_sink_builder(
3557 scx,
3558 connection,
3559 options,
3560 format,
3561 relation_key_indices,
3562 key_desc_and_indices,
3563 headers_index,
3564 desc.into_owned(),
3565 envelope,
3566 from.id(),
3567 commit_interval,
3568 )?,
3569 CreateSinkConnection::Iceberg {
3570 catalog_connection,
3571 aws_connection,
3572 options,
3573 ..
3574 } => iceberg_sink_builder(
3575 scx,
3576 catalog_connection,
3577 aws_connection,
3578 options,
3579 relation_key_indices,
3580 key_desc_and_indices,
3581 commit_interval,
3582 &desc,
3583 )?,
3584 };
3585
3586 let with_snapshot = snapshot.unwrap_or(true);
3588 let version = version.unwrap_or(0);
3590
3591 let in_cluster = source_sink_cluster_config(scx, &mut stmt.in_cluster)?;
3595 let create_sql = normalize::create_statement(scx, Statement::CreateSink(stmt))?;
3596
3597 Ok(Plan::CreateSink(CreateSinkPlan {
3598 name,
3599 sink: Sink {
3600 create_sql,
3601 from: from.global_id(),
3602 connection: connection_builder,
3603 envelope,
3604 version,
3605 commit_interval,
3606 },
3607 with_snapshot,
3608 if_not_exists,
3609 in_cluster: in_cluster.id(),
3610 }))
3611}
3612
3613fn key_constraint_err(desc: &RelationDesc, user_keys: &[ColumnName]) -> PlanError {
3614 let user_keys = user_keys.iter().map(|column| column.as_str()).join(", ");
3615
3616 let existing_keys = desc
3617 .typ()
3618 .keys
3619 .iter()
3620 .map(|key_columns| {
3621 key_columns
3622 .iter()
3623 .map(|col| desc.get_name(*col).as_str())
3624 .join(", ")
3625 })
3626 .join(", ");
3627
3628 sql_err!(
3629 "Key constraint ({}) conflicts with existing key ({})",
3630 user_keys,
3631 existing_keys
3632 )
3633}
3634
3635#[derive(Debug, Default, PartialEq, Clone)]
3638pub struct CsrConfigOptionExtracted {
3639 seen: ::std::collections::BTreeSet<CsrConfigOptionName<Aug>>,
3640 pub(crate) avro_key_fullname: Option<String>,
3641 pub(crate) avro_value_fullname: Option<String>,
3642 pub(crate) null_defaults: bool,
3643 pub(crate) value_doc_options: BTreeMap<DocTarget, String>,
3644 pub(crate) key_doc_options: BTreeMap<DocTarget, String>,
3645 pub(crate) key_compatibility_level: Option<mz_ccsr::CompatibilityLevel>,
3646 pub(crate) value_compatibility_level: Option<mz_ccsr::CompatibilityLevel>,
3647}
3648
3649impl std::convert::TryFrom<Vec<CsrConfigOption<Aug>>> for CsrConfigOptionExtracted {
3650 type Error = crate::plan::PlanError;
3651 fn try_from(v: Vec<CsrConfigOption<Aug>>) -> Result<CsrConfigOptionExtracted, Self::Error> {
3652 let mut extracted = CsrConfigOptionExtracted::default();
3653 let mut common_doc_comments = BTreeMap::new();
3654 for option in v {
3655 if !extracted.seen.insert(option.name.clone()) {
3656 return Err(PlanError::Unstructured({
3657 format!("{} specified more than once", option.name)
3658 }));
3659 }
3660 let option_name = option.name.clone();
3661 let option_name_str = option_name.to_ast_string_simple();
3662 let better_error = |e: PlanError| PlanError::InvalidOptionValue {
3663 option_name: option_name.to_ast_string_simple(),
3664 err: e.into(),
3665 };
3666 let to_compatibility_level = |val: Option<WithOptionValue<Aug>>| {
3667 val.map(|s| match s {
3668 WithOptionValue::Value(Value::String(s)) => {
3669 mz_ccsr::CompatibilityLevel::try_from(s.to_uppercase().as_str())
3670 }
3671 _ => Err("must be a string".to_string()),
3672 })
3673 .transpose()
3674 .map_err(PlanError::Unstructured)
3675 .map_err(better_error)
3676 };
3677 match option.name {
3678 CsrConfigOptionName::AvroKeyFullname => {
3679 extracted.avro_key_fullname =
3680 <Option<String>>::try_from_value(option.value).map_err(better_error)?;
3681 }
3682 CsrConfigOptionName::AvroValueFullname => {
3683 extracted.avro_value_fullname =
3684 <Option<String>>::try_from_value(option.value).map_err(better_error)?;
3685 }
3686 CsrConfigOptionName::NullDefaults => {
3687 extracted.null_defaults =
3688 <bool>::try_from_value(option.value).map_err(better_error)?;
3689 }
3690 CsrConfigOptionName::AvroDocOn(doc_on) => {
3691 let value = String::try_from_value(option.value.ok_or_else(|| {
3692 PlanError::InvalidOptionValue {
3693 option_name: option_name_str,
3694 err: Box::new(PlanError::Unstructured("cannot be empty".to_string())),
3695 }
3696 })?)
3697 .map_err(better_error)?;
3698 let key = match doc_on.identifier {
3699 DocOnIdentifier::Column(ast::ColumnName {
3700 relation: ResolvedItemName::Item { id, .. },
3701 column: ResolvedColumnReference::Column { name, index: _ },
3702 }) => DocTarget::Field {
3703 object_id: id,
3704 column_name: name,
3705 },
3706 DocOnIdentifier::Type(ResolvedItemName::Item { id, .. }) => {
3707 DocTarget::Type(id)
3708 }
3709 _ => sql_bail!("invalid DOC ON identifier"),
3710 };
3711
3712 match doc_on.for_schema {
3713 DocOnSchema::KeyOnly => {
3714 extracted.key_doc_options.insert(key, value);
3715 }
3716 DocOnSchema::ValueOnly => {
3717 extracted.value_doc_options.insert(key, value);
3718 }
3719 DocOnSchema::All => {
3720 common_doc_comments.insert(key, value);
3721 }
3722 }
3723 }
3724 CsrConfigOptionName::KeyCompatibilityLevel => {
3725 extracted.key_compatibility_level = to_compatibility_level(option.value)?;
3726 }
3727 CsrConfigOptionName::ValueCompatibilityLevel => {
3728 extracted.value_compatibility_level = to_compatibility_level(option.value)?;
3729 }
3730 }
3731 }
3732
3733 for (key, value) in common_doc_comments {
3734 if !extracted.key_doc_options.contains_key(&key) {
3735 extracted.key_doc_options.insert(key.clone(), value.clone());
3736 }
3737 if !extracted.value_doc_options.contains_key(&key) {
3738 extracted.value_doc_options.insert(key, value);
3739 }
3740 }
3741 Ok(extracted)
3742 }
3743}
3744
3745fn iceberg_sink_builder(
3746 scx: &StatementContext,
3747 catalog_connection: ResolvedItemName,
3748 storage_connection: Option<ResolvedItemName>,
3749 options: Vec<IcebergSinkConfigOption<Aug>>,
3750 relation_key_indices: Option<Vec<usize>>,
3751 key_desc_and_indices: Option<(RelationDesc, Vec<usize>)>,
3752 commit_interval: Option<Duration>,
3753 desc: &RelationDesc,
3754) -> Result<StorageSinkConnection<ReferencedConnection>, PlanError> {
3755 ArrowBuilder::validate_desc_for_parquet(desc, iceberg_type_overrides)
3759 .map_err(|e| sql_err!("{}", e))?;
3760
3761 let catalog_connection_item = scx.get_item_by_resolved_name(&catalog_connection)?;
3762 let catalog_connection_id = catalog_connection_item.id();
3763 if !matches!(
3764 catalog_connection_item.connection()?,
3765 Connection::IcebergCatalog(_)
3766 ) {
3767 sql_bail!(
3768 "{} is not an iceberg catalog connection",
3769 scx.catalog
3770 .resolve_full_name(catalog_connection_item.name())
3771 .to_string()
3772 .quoted()
3773 );
3774 };
3775
3776 let storage_connection_item = storage_connection
3777 .map(|c| scx.get_item_by_resolved_name(&c))
3778 .transpose()?;
3779 let storage_connection_id = storage_connection_item.as_ref().map(|c| c.id());
3780 if let Some(c) = &storage_connection_item
3781 && !matches!(c.connection()?, Connection::Aws(_))
3782 {
3783 sql_bail!(
3784 "{} is not an AWS connection",
3785 scx.catalog.resolve_full_name(c.name()).to_string().quoted()
3786 );
3787 }
3788
3789 let IcebergSinkConfigOptionExtracted {
3790 table,
3791 namespace,
3792 seen: _,
3793 }: IcebergSinkConfigOptionExtracted = options.try_into()?;
3794
3795 let Some(table) = table else {
3796 sql_bail!("Iceberg sink must specify TABLE");
3797 };
3798 let Some(namespace) = namespace else {
3799 sql_bail!("Iceberg sink must specify NAMESPACE");
3800 };
3801 match commit_interval {
3802 None => sql_bail!("Iceberg sink must specify COMMIT INTERVAL"),
3803 Some(interval) if interval < Duration::from_secs(1) => {
3809 sql_bail!("COMMIT INTERVAL must be at least 1 second")
3810 }
3811 Some(_) => {}
3812 }
3813
3814 Ok(StorageSinkConnection::Iceberg(IcebergSinkConnection {
3815 catalog_connection_id,
3816 catalog_connection: catalog_connection_id,
3817 storage_connection_id,
3818 storage_connection: storage_connection_id,
3819 table,
3820 namespace,
3821 relation_key_indices,
3822 key_desc_and_indices,
3823 }))
3824}
3825
3826fn kafka_sink_builder(
3827 scx: &StatementContext,
3828 connection: ResolvedItemName,
3829 options: Vec<KafkaSinkConfigOption<Aug>>,
3830 format: Option<FormatSpecifier<Aug>>,
3831 relation_key_indices: Option<Vec<usize>>,
3832 key_desc_and_indices: Option<(RelationDesc, Vec<usize>)>,
3833 headers_index: Option<usize>,
3834 value_desc: RelationDesc,
3835 envelope: SinkEnvelope,
3836 sink_from: CatalogItemId,
3837 commit_interval: Option<Duration>,
3838) -> Result<StorageSinkConnection<ReferencedConnection>, PlanError> {
3839 let connection_item = scx.get_item_by_resolved_name(&connection)?;
3841 let connection_id = connection_item.id();
3842 match connection_item.connection()? {
3843 Connection::Kafka(_) => (),
3844 _ => sql_bail!(
3845 "{} is not a kafka connection",
3846 scx.catalog.resolve_full_name(connection_item.name())
3847 ),
3848 };
3849
3850 if commit_interval.is_some() {
3851 sql_bail!("COMMIT INTERVAL option is not supported with KAFKA sinks");
3852 }
3853
3854 let KafkaSinkConfigOptionExtracted {
3855 topic,
3856 compression_type,
3857 partition_by,
3858 progress_group_id_prefix,
3859 transactional_id_prefix,
3860 legacy_ids,
3861 topic_config,
3862 topic_metadata_refresh_interval,
3863 topic_partition_count,
3864 topic_replication_factor,
3865 seen: _,
3866 }: KafkaSinkConfigOptionExtracted = options.try_into()?;
3867
3868 let transactional_id = match (transactional_id_prefix, legacy_ids) {
3869 (Some(_), Some(true)) => {
3870 sql_bail!("LEGACY IDS cannot be used at the same time as TRANSACTIONAL ID PREFIX")
3871 }
3872 (None, Some(true)) => KafkaIdStyle::Legacy,
3873 (prefix, _) => KafkaIdStyle::Prefix(prefix),
3874 };
3875
3876 let progress_group_id = match (progress_group_id_prefix, legacy_ids) {
3877 (Some(_), Some(true)) => {
3878 sql_bail!("LEGACY IDS cannot be used at the same time as PROGRESS GROUP ID PREFIX")
3879 }
3880 (None, Some(true)) => KafkaIdStyle::Legacy,
3881 (prefix, _) => KafkaIdStyle::Prefix(prefix),
3882 };
3883
3884 let topic_name = topic.ok_or_else(|| sql_err!("KAFKA CONNECTION must specify TOPIC"))?;
3885
3886 if topic_metadata_refresh_interval > MAX_KAFKA_TOPIC_METADATA_REFRESH_INTERVAL {
3887 sql_bail!("TOPIC METADATA REFRESH INTERVAL cannot be greater than 1 hour");
3890 } else if topic_metadata_refresh_interval < MIN_KAFKA_TOPIC_METADATA_REFRESH_INTERVAL {
3891 sql_bail!("TOPIC METADATA REFRESH INTERVAL must be at least 1 second");
3894 }
3895
3896 let assert_positive = |val: Option<i32>, name: &str| {
3897 if let Some(val) = val {
3898 if val <= 0 {
3899 sql_bail!("{} must be a positive integer", name);
3900 }
3901 }
3902 val.map(NonNeg::try_from)
3903 .transpose()
3904 .map_err(|_| PlanError::Unstructured(format!("{} must be a positive integer", name)))
3905 };
3906 let topic_partition_count = assert_positive(topic_partition_count, "TOPIC PARTITION COUNT")?;
3907 let topic_replication_factor =
3908 assert_positive(topic_replication_factor, "TOPIC REPLICATION FACTOR")?;
3909
3910 let gen_avro_schema_options = |conn| {
3913 let CsrConnectionAvro {
3914 connection:
3915 CsrConnection {
3916 connection,
3917 options,
3918 },
3919 seed,
3920 key_strategy,
3921 value_strategy,
3922 } = conn;
3923 if seed.is_some() {
3924 sql_bail!("SEED option does not make sense with sinks");
3925 }
3926 if key_strategy.is_some() {
3927 sql_bail!("KEY STRATEGY option does not make sense with sinks");
3928 }
3929 if value_strategy.is_some() {
3930 sql_bail!("VALUE STRATEGY option does not make sense with sinks");
3931 }
3932
3933 let item = scx.get_item_by_resolved_name(&connection)?;
3934 let csr_connection = match item.connection()? {
3935 Connection::Csr(_) => item.id(),
3936 _ => {
3937 sql_bail!(
3938 "{} is not a schema registry connection",
3939 scx.catalog
3940 .resolve_full_name(item.name())
3941 .to_string()
3942 .quoted()
3943 )
3944 }
3945 };
3946 let extracted_options: CsrConfigOptionExtracted = options.try_into()?;
3947
3948 if key_desc_and_indices.is_none() && extracted_options.avro_key_fullname.is_some() {
3949 sql_bail!("Cannot specify AVRO KEY FULLNAME without a corresponding KEY field");
3950 }
3951
3952 if key_desc_and_indices.is_some()
3953 && (extracted_options.avro_key_fullname.is_some()
3954 ^ extracted_options.avro_value_fullname.is_some())
3955 {
3956 sql_bail!(
3957 "Must specify both AVRO KEY FULLNAME and AVRO VALUE FULLNAME when specifying generated schema names"
3958 );
3959 }
3960
3961 Ok((csr_connection, extracted_options))
3962 };
3963
3964 let map_format = |format: Format<Aug>, desc: &RelationDesc, is_key: bool| match format {
3965 Format::Json { array: false } => Ok::<_, PlanError>(KafkaSinkFormatType::Json),
3966 Format::Bytes if desc.arity() == 1 => {
3967 let col_type = &desc.typ().column_types[0].scalar_type;
3968 if !mz_pgrepr::Value::can_encode_binary(col_type) {
3969 bail_unsupported!(format!(
3970 "BYTES format with non-encodable type: {:?}",
3971 col_type
3972 ));
3973 }
3974
3975 Ok(KafkaSinkFormatType::Bytes)
3976 }
3977 Format::Text if desc.arity() == 1 => Ok(KafkaSinkFormatType::Text),
3978 Format::Bytes | Format::Text => {
3979 bail_unsupported!("BYTES or TEXT format with multiple columns")
3980 }
3981 Format::Json { array: true } => bail_unsupported!("JSON ARRAY format in sinks"),
3982 Format::Avro(AvroSchema::Csr { csr_connection }) => {
3983 let (csr_connection, options) = gen_avro_schema_options(csr_connection)?;
3984 let schema = if is_key {
3985 AvroSchemaGenerator::new(
3986 desc.clone(),
3987 false,
3988 options.key_doc_options,
3989 options.avro_key_fullname.as_deref().unwrap_or("row"),
3990 options.null_defaults,
3991 Some(sink_from),
3992 false,
3993 )?
3994 .schema()
3995 .to_string()
3996 } else {
3997 AvroSchemaGenerator::new(
3998 desc.clone(),
3999 matches!(envelope, SinkEnvelope::Debezium),
4000 options.value_doc_options,
4001 options.avro_value_fullname.as_deref().unwrap_or("envelope"),
4002 options.null_defaults,
4003 Some(sink_from),
4004 true,
4005 )?
4006 .schema()
4007 .to_string()
4008 };
4009 Ok(KafkaSinkFormatType::Avro {
4010 schema,
4011 compatibility_level: if is_key {
4012 options.key_compatibility_level
4013 } else {
4014 options.value_compatibility_level
4015 },
4016 schema_name: None,
4018 wire_format: WireFormat::Confluent {
4019 registry: Some(csr_connection),
4020 },
4021 })
4022 }
4023 Format::Avro(AvroSchema::Glue {
4024 connection,
4025 with_options,
4026 seed,
4027 }) => {
4028 if seed.is_some() {
4029 sql_bail!("SEED option does not make sense with sinks");
4030 }
4031
4032 let extracted: GlueAvroOptionExtracted = with_options.try_into()?;
4033 let GlueAvroOptionExtracted {
4034 schema_name,
4035 key_schema_name,
4036 value_schema_name,
4037 key_compatibility_level,
4038 value_compatibility_level,
4039 seen: _,
4040 } = extracted;
4041
4042 if schema_name.is_some() {
4047 sql_bail!(
4048 "SCHEMA NAME is not supported for AWS Glue Schema Registry sinks, \
4049 use KEY SCHEMA NAME and VALUE SCHEMA NAME instead"
4050 );
4051 }
4052
4053 if key_desc_and_indices.is_none()
4059 && (key_schema_name.is_some() || key_compatibility_level.is_some())
4060 {
4061 sql_bail!(
4062 "KEY SCHEMA NAME and KEY COMPATIBILITY LEVEL require a corresponding KEY field"
4063 );
4064 }
4065
4066 let item = scx.get_item_by_resolved_name(&connection)?;
4067 let glue_connection = match item.connection()? {
4068 Connection::GlueSchemaRegistry(_) => item.id(),
4069 _ => {
4070 sql_bail!(
4071 "{} is not an AWS Glue Schema Registry connection",
4072 scx.catalog
4073 .resolve_full_name(item.name())
4074 .to_string()
4075 .quoted()
4076 )
4077 }
4078 };
4079
4080 let compatibility_level = {
4084 let raw = if is_key {
4085 key_compatibility_level
4086 } else {
4087 value_compatibility_level
4088 };
4089 raw.map(|s| {
4090 mz_ccsr::CompatibilityLevel::try_from(s.to_uppercase().as_str())
4091 .map_err(PlanError::Unstructured)
4092 })
4093 .transpose()?
4094 };
4095
4096 let schema_name = if is_key {
4100 key_schema_name
4101 } else {
4102 value_schema_name
4103 };
4104
4105 let schema = if is_key {
4108 AvroSchemaGenerator::new(
4109 desc.clone(),
4110 false,
4111 Default::default(),
4112 "row",
4113 false,
4114 Some(sink_from),
4115 false,
4116 )?
4117 .schema()
4118 .to_string()
4119 } else {
4120 AvroSchemaGenerator::new(
4121 desc.clone(),
4122 matches!(envelope, SinkEnvelope::Debezium),
4123 Default::default(),
4124 "envelope",
4125 false,
4126 Some(sink_from),
4127 true,
4128 )?
4129 .schema()
4130 .to_string()
4131 };
4132 Ok(KafkaSinkFormatType::Avro {
4133 schema,
4134 compatibility_level,
4135 schema_name,
4136 wire_format: WireFormat::Glue {
4137 registry: Some(glue_connection),
4138 },
4139 })
4140 }
4141 format => bail_unsupported!(format!("sink format {:?}", format)),
4142 };
4143
4144 let partition_by = match &partition_by {
4145 Some(partition_by) => {
4146 let mut scope = Scope::from_source(None, value_desc.iter_names());
4147
4148 match envelope {
4149 SinkEnvelope::Upsert | SinkEnvelope::Append => (),
4150 SinkEnvelope::Debezium => {
4151 let key_indices: HashSet<_> = key_desc_and_indices
4152 .as_ref()
4153 .map(|(_desc, indices)| indices.as_slice())
4154 .unwrap_or_default()
4155 .into_iter()
4156 .collect();
4157 for (i, item) in scope.items.iter_mut().enumerate() {
4158 if !key_indices.contains(&i) {
4159 item.error_if_referenced = Some(|_table, column| {
4160 PlanError::InvalidPartitionByEnvelopeDebezium {
4161 column_name: column.to_string(),
4162 }
4163 });
4164 }
4165 }
4166 }
4167 };
4168
4169 let ecx = &ExprContext {
4170 qcx: &QueryContext::root(scx, QueryLifetime::OneShot),
4171 name: "PARTITION BY",
4172 scope: &scope,
4173 relation_type: value_desc.typ(),
4174 allow_aggregates: false,
4175 allow_subqueries: false,
4176 allow_parameters: false,
4177 allow_windows: false,
4178 };
4179 let expr = plan_expr(ecx, partition_by)?.cast_to(
4180 ecx,
4181 CastContext::Assignment,
4182 &SqlScalarType::UInt64,
4183 )?;
4184 let expr = expr.lower_uncorrelated(scx.catalog.system_vars())?;
4185
4186 Some(expr)
4187 }
4188 _ => None,
4189 };
4190
4191 let format = match format {
4193 Some(FormatSpecifier::KeyValue { key, value }) => {
4194 let key_format = match key_desc_and_indices.as_ref() {
4195 Some((desc, _indices)) => Some(map_format(key, desc, true)?),
4196 None => None,
4197 };
4198 KafkaSinkFormat {
4199 value_format: map_format(value, &value_desc, false)?,
4200 key_format,
4201 }
4202 }
4203 Some(FormatSpecifier::Bare(format)) => {
4204 let key_format = match key_desc_and_indices.as_ref() {
4205 Some((desc, _indices)) => Some(map_format(format.clone(), desc, true)?),
4206 None => None,
4207 };
4208 KafkaSinkFormat {
4209 value_format: map_format(format, &value_desc, false)?,
4210 key_format,
4211 }
4212 }
4213 None => bail_unsupported!("sink without format"),
4214 };
4215
4216 Ok(StorageSinkConnection::Kafka(KafkaSinkConnection {
4217 connection_id,
4218 connection: connection_id,
4219 format,
4220 topic: topic_name,
4221 relation_key_indices,
4222 key_desc_and_indices,
4223 headers_index,
4224 value_desc,
4225 partition_by,
4226 compression_type,
4227 progress_group_id,
4228 transactional_id,
4229 topic_options: KafkaTopicOptions {
4230 partition_count: topic_partition_count,
4231 replication_factor: topic_replication_factor,
4232 topic_config: topic_config.unwrap_or_default(),
4233 },
4234 topic_metadata_refresh_interval,
4235 }))
4236}
4237
4238pub fn describe_create_metric_sink(
4239 _: &StatementContext,
4240 _: CreateMetricSinkStatement<Aug>,
4241) -> Result<StatementDesc, PlanError> {
4242 Ok(StatementDesc::new(None))
4243}
4244
4245const METRIC_SINK_SOURCE_COLUMNS: &[(&str, fn(&SqlScalarType) -> bool)] = &[
4249 ("metric_name", |t| matches!(t, SqlScalarType::String)),
4250 ("metric_type", |t| matches!(t, SqlScalarType::String)),
4251 (
4252 "labels",
4253 |t| matches!(t, SqlScalarType::Map { value_type, .. } if matches!(**value_type, SqlScalarType::String)),
4254 ),
4255 ("value", |t| matches!(t, SqlScalarType::Float64)),
4256 ("help", |t| matches!(t, SqlScalarType::String)),
4257];
4258
4259generate_extracted_config!(CreateMetricSinkOption, (Prefix, String));
4260
4261const METRIC_SINK_PREFIX_MARKER: &str = "mz_metric_sink_";
4268
4269pub const METRIC_SINK_CURATED_PREFIX_MARKER: &str = "mz_metric_sink_curated_";
4272
4273pub fn validate_metric_sink_prefix(prefix: &str) -> Result<(), PlanError> {
4278 if prefix.is_empty() {
4279 return Err(sql_err!("metric sink prefix must not be empty"));
4280 }
4281 let mut chars = prefix.chars();
4282 let valid = match chars.next() {
4283 Some(c) if c.is_ascii_alphabetic() || c == '_' || c == ':' => {
4284 chars.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == ':')
4285 }
4286 _ => false,
4287 };
4288 if !valid {
4289 return Err(sql_err!(
4290 "metric sink prefix {:?} is not a valid start of a Prometheus metric name",
4291 prefix
4292 ));
4293 }
4294 if !prefix.starts_with(METRIC_SINK_PREFIX_MARKER) {
4295 return Err(sql_err!(
4296 "metric sink prefix {:?} must start with {:?}",
4297 prefix,
4298 METRIC_SINK_PREFIX_MARKER
4299 ));
4300 }
4301 Ok(())
4302}
4303
4304pub fn validate_user_metric_sink_prefix(prefix: &str) -> Result<(), PlanError> {
4313 validate_metric_sink_prefix(prefix)?;
4314 if prefix.starts_with(METRIC_SINK_CURATED_PREFIX_MARKER)
4315 || METRIC_SINK_CURATED_PREFIX_MARKER.starts_with(prefix)
4316 {
4317 return Err(sql_err!(
4318 "metric sink prefix {:?} overlaps {:?}, which is reserved for the built-in metric sinks",
4319 prefix,
4320 METRIC_SINK_CURATED_PREFIX_MARKER
4321 ));
4322 }
4323 Ok(())
4324}
4325
4326pub fn validate_metric_sink_desc(desc: &RelationDesc) -> Result<(), PlanError> {
4332 for (name, type_ok) in METRIC_SINK_SOURCE_COLUMNS {
4333 let col = ColumnName::from(*name);
4334 let (_, column_type) = desc
4335 .get_by_name(&col)
4336 .ok_or_else(|| sql_err!("metric sink source must expose column {:?}", name))?;
4337 if !type_ok(&column_type.scalar_type) {
4338 return Err(sql_err!(
4339 "metric sink source column {:?} is not of the required type",
4340 name
4341 ));
4342 }
4343 }
4344 Ok(())
4345}
4346
4347pub fn plan_create_metric_sink(
4355 scx: &StatementContext,
4356 mut stmt: CreateMetricSinkStatement<Aug>,
4357) -> Result<Plan, PlanError> {
4358 scx.require_feature_flag(&ENABLE_METRIC_SINK)?;
4359
4360 let CreateMetricSinkStatement {
4361 name,
4362 in_cluster,
4363 if_not_exists,
4364 from,
4365 with_options,
4366 } = &mut stmt;
4367
4368 let if_not_exists = *if_not_exists;
4369 let CreateMetricSinkOptionExtracted { prefix, seen: _ } = with_options.clone().try_into()?;
4370 let Some(prefix) = prefix else {
4373 sql_bail!("CREATE METRIC SINK requires a PREFIX option");
4374 };
4375 validate_user_metric_sink_prefix(&prefix)?;
4376 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name.clone())?)?;
4377 let full_name = scx.catalog.resolve_full_name(&name);
4378 let partial_name = PartialItemName::from(full_name.clone());
4379 if let (false, Ok(item)) = (if_not_exists, scx.catalog.resolve_item(&partial_name)) {
4380 return Err(PlanError::ItemAlreadyExists {
4381 name: full_name.to_string(),
4382 item_type: item.item_type(),
4383 });
4384 }
4385
4386 let from_item = scx.get_item_by_resolved_name(from)?;
4387 let desc = from_item.relation_desc().ok_or_else(|| {
4391 sql_err!(
4392 "cannot create metric sink from {} because it is a {}",
4393 scx.catalog.minimal_qualification(from_item.name()),
4394 from_item.item_type(),
4395 )
4396 })?;
4397 validate_metric_sink_desc(&desc)?;
4398
4399 let cluster_id = match in_cluster {
4400 None => scx.resolve_cluster(None)?.id(),
4401 Some(in_cluster) => in_cluster.id,
4402 };
4403 *in_cluster = Some(ResolvedClusterName {
4404 id: cluster_id,
4405 print_name: None,
4406 });
4407 let from_global_id = from_item.global_id();
4408
4409 let create_sql = normalize::create_statement(scx, Statement::CreateMetricSink(stmt))?;
4410
4411 Ok(Plan::CreateMetricSink(CreateMetricSinkPlan {
4412 name,
4413 metric_sink: MetricSink {
4414 create_sql,
4415 from: from_global_id,
4416 cluster_id,
4417 prefix,
4418 },
4419 if_not_exists,
4420 }))
4421}
4422
4423pub fn describe_create_index(
4424 _: &StatementContext,
4425 _: CreateIndexStatement<Aug>,
4426) -> Result<StatementDesc, PlanError> {
4427 Ok(StatementDesc::new(None))
4428}
4429
4430pub fn plan_create_index(
4431 scx: &StatementContext,
4432 mut stmt: CreateIndexStatement<Aug>,
4433) -> Result<Plan, PlanError> {
4434 let CreateIndexStatement {
4435 name,
4436 on_name,
4437 in_cluster,
4438 key_parts,
4439 with_options,
4440 if_not_exists,
4441 } = &mut stmt;
4442 let on = scx.get_item_by_resolved_name(on_name)?;
4443
4444 {
4445 use CatalogItemType::*;
4446 match on.item_type() {
4447 Table | Source | View | MaterializedView => {
4448 if on.replacement_target().is_some() {
4449 sql_bail!(
4450 "index cannot be created on {} because it is a replacement {}",
4451 on_name.full_name_str(),
4452 on.item_type(),
4453 );
4454 }
4455 }
4456 Sink | MetricSink | Index | Type | Func | Secret | Connection => {
4457 sql_bail!(
4458 "index cannot be created on {} because it is a {}",
4459 on_name.full_name_str(),
4460 on.item_type(),
4461 );
4462 }
4463 }
4464 }
4465
4466 let on_desc = on
4467 .relation_desc()
4468 .ok_or_else(|| sql_err!("item does not have a relation description"))?;
4469
4470 let filled_key_parts = match key_parts {
4471 Some(kp) => kp.to_vec(),
4472 None => {
4473 let mut name_counts = BTreeMap::new();
4478 for name in on_desc.iter_names() {
4479 *name_counts.entry(name).or_insert(0usize) += 1;
4480 }
4481 let key = on_desc.typ().default_key();
4482 key.iter()
4483 .map(|i| {
4484 let name = on_desc.get_name(*i);
4485 if name_counts.get(name).copied() == Some(1) {
4486 Expr::Identifier(vec![name.clone().into()])
4487 } else {
4488 Expr::Value(Value::Number((i + 1).to_string()))
4489 }
4490 })
4491 .collect()
4492 }
4493 };
4494 let keys = query::plan_index_exprs(scx, &on_desc, filled_key_parts.clone())?;
4495
4496 let index_name = if let Some(name) = name {
4497 QualifiedItemName {
4498 qualifiers: on.name().qualifiers.clone(),
4499 item: normalize::ident(name.clone()),
4500 }
4501 } else {
4502 let mut idx_name = QualifiedItemName {
4503 qualifiers: on.name().qualifiers.clone(),
4504 item: on.name().item.clone(),
4505 };
4506 if key_parts.is_none() {
4507 idx_name.item += "_primary_idx";
4509 } else {
4510 let index_name_col_suffix = keys
4513 .iter()
4514 .map(|k| match k {
4515 mz_expr::MirScalarExpr::Column(i, name) => {
4516 match (on_desc.get_unambiguous_name(*i), &name.0) {
4517 (Some(col_name), _) => col_name.to_string(),
4518 (None, Some(name)) => name.to_string(),
4519 (None, None) => format!("{}", i + 1),
4520 }
4521 }
4522 _ => "expr".to_string(),
4523 })
4524 .join("_");
4525 write!(idx_name.item, "_{index_name_col_suffix}_idx")
4526 .expect("write on strings cannot fail");
4527 idx_name.item = normalize::ident(Ident::new(&idx_name.item)?)
4528 }
4529
4530 if !*if_not_exists {
4531 scx.catalog.find_available_name(idx_name)
4532 } else {
4533 idx_name
4534 }
4535 };
4536
4537 let full_name = scx.catalog.resolve_full_name(&index_name);
4539 let partial_name = PartialItemName::from(full_name.clone());
4540 if let (Ok(item), false, false) = (
4548 scx.catalog.resolve_item_or_type(&partial_name),
4549 *if_not_exists,
4550 scx.pcx().map_or(false, |pcx| pcx.ignore_if_exists_errors),
4551 ) {
4552 return Err(PlanError::ItemAlreadyExists {
4553 name: full_name.to_string(),
4554 item_type: item.item_type(),
4555 });
4556 }
4557
4558 let options = plan_index_options(scx, with_options.clone())?;
4559 let cluster_id = match in_cluster {
4560 None => scx.resolve_cluster(None)?.id(),
4561 Some(in_cluster) => in_cluster.id,
4562 };
4563
4564 *in_cluster = Some(ResolvedClusterName {
4565 id: cluster_id,
4566 print_name: None,
4567 });
4568
4569 *name = Some(Ident::new(index_name.item.clone())?);
4571 *key_parts = Some(filled_key_parts);
4572 let if_not_exists = *if_not_exists;
4573
4574 let create_sql = normalize::create_statement(scx, Statement::CreateIndex(stmt))?;
4575 let compaction_window = options.iter().find_map(|o| {
4576 #[allow(irrefutable_let_patterns)]
4577 if let crate::plan::IndexOption::RetainHistory(lcw) = o {
4578 Some(lcw.clone())
4579 } else {
4580 None
4581 }
4582 });
4583
4584 Ok(Plan::CreateIndex(CreateIndexPlan {
4585 name: index_name,
4586 index: Index {
4587 create_sql,
4588 on: on.global_id(),
4589 keys,
4590 cluster_id,
4591 compaction_window,
4592 },
4593 if_not_exists,
4594 }))
4595}
4596
4597pub fn describe_create_type(
4598 _: &StatementContext,
4599 _: CreateTypeStatement<Aug>,
4600) -> Result<StatementDesc, PlanError> {
4601 Ok(StatementDesc::new(None))
4602}
4603
4604pub fn plan_create_type(
4605 scx: &StatementContext,
4606 stmt: CreateTypeStatement<Aug>,
4607) -> Result<Plan, PlanError> {
4608 let create_sql = normalize::create_statement(scx, Statement::CreateType(stmt.clone()))?;
4609 let CreateTypeStatement { name, as_type, .. } = stmt;
4610
4611 fn validate_data_type(
4619 scx: &StatementContext,
4620 data_type: ResolvedDataType,
4621 as_type: &str,
4622 key: &str,
4623 budget: &mut TypeResolutionBudget,
4624 ) -> Result<(CatalogItemId, Vec<i64>), PlanError> {
4625 let (id, modifiers) = match data_type {
4626 ResolvedDataType::Named { id, modifiers, .. } => (id, modifiers),
4627 _ => sql_bail!(
4628 "CREATE TYPE ... AS {}option {} can only use named data types, but \
4629 found unnamed data type {}. Use CREATE TYPE to create a named type first",
4630 as_type,
4631 key,
4632 data_type.human_readable_name(),
4633 ),
4634 };
4635
4636 let item = scx.catalog.get_item(&id);
4637 match item.type_details() {
4638 None => sql_bail!(
4639 "{} must be of class type, but received {} which is of class {}",
4640 key,
4641 scx.catalog.resolve_full_name(item.name()),
4642 item.item_type()
4643 ),
4644 Some(CatalogTypeDetails {
4645 typ: CatalogType::Char,
4646 ..
4647 }) => {
4648 bail_unsupported!("embedding char type in a list or map")
4649 }
4650 _ => {
4651 budget.resolve_child(scx.catalog, id, &modifiers)?;
4654
4655 Ok((id, modifiers))
4656 }
4657 }
4658 }
4659
4660 let mut budget = TypeResolutionBudget::for_root(scx.catalog);
4661 let inner = match as_type {
4662 CreateTypeAs::List { options } => {
4663 let CreateTypeListOptionExtracted {
4664 element_type,
4665 seen: _,
4666 } = CreateTypeListOptionExtracted::try_from(options)?;
4667 let element_type =
4668 element_type.ok_or_else(|| sql_err!("ELEMENT TYPE option is required"))?;
4669 let (id, modifiers) =
4670 validate_data_type(scx, element_type, "LIST ", "ELEMENT TYPE", &mut budget)?;
4671 CatalogType::List {
4672 element_reference: id,
4673 element_modifiers: modifiers,
4674 }
4675 }
4676 CreateTypeAs::Map { options } => {
4677 let CreateTypeMapOptionExtracted {
4678 key_type,
4679 value_type,
4680 seen: _,
4681 } = CreateTypeMapOptionExtracted::try_from(options)?;
4682 let key_type = key_type.ok_or_else(|| sql_err!("KEY TYPE option is required"))?;
4683 let value_type = value_type.ok_or_else(|| sql_err!("VALUE TYPE option is required"))?;
4684 let (key_id, key_modifiers) = validate_data_type(
4688 scx,
4689 key_type,
4690 "MAP ",
4691 "KEY TYPE",
4692 &mut TypeResolutionBudget::for_root(scx.catalog),
4693 )?;
4694 let (value_id, value_modifiers) =
4695 validate_data_type(scx, value_type, "MAP ", "VALUE TYPE", &mut budget)?;
4696 CatalogType::Map {
4697 key_reference: key_id,
4698 key_modifiers,
4699 value_reference: value_id,
4700 value_modifiers,
4701 }
4702 }
4703 CreateTypeAs::Record { column_defs } => {
4704 let mut fields = vec![];
4705 for column_def in column_defs {
4706 let data_type = column_def.data_type;
4707 let key = ident(column_def.name.clone());
4708 let (id, modifiers) = validate_data_type(scx, data_type, "", &key, &mut budget)?;
4709 fields.push(CatalogRecordField {
4710 name: ColumnName::from(key.clone()),
4711 type_reference: id,
4712 type_modifiers: modifiers,
4713 });
4714 }
4715 CatalogType::Record { fields }
4716 }
4717 };
4718
4719 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name)?)?;
4720
4721 let full_name = scx.catalog.resolve_full_name(&name);
4723 let partial_name = PartialItemName::from(full_name.clone());
4724 if let Ok(item) = scx.catalog.resolve_item_or_type(&partial_name) {
4727 if item.item_type().conflicts_with_type() {
4728 return Err(PlanError::ItemAlreadyExists {
4729 name: full_name.to_string(),
4730 item_type: item.item_type(),
4731 });
4732 }
4733 }
4734
4735 Ok(Plan::CreateType(CreateTypePlan {
4736 name,
4737 typ: Type { create_sql, inner },
4738 }))
4739}
4740
4741generate_extracted_config!(CreateTypeListOption, (ElementType, ResolvedDataType));
4742
4743generate_extracted_config!(
4744 CreateTypeMapOption,
4745 (KeyType, ResolvedDataType),
4746 (ValueType, ResolvedDataType)
4747);
4748
4749#[derive(Debug)]
4750pub enum PlannedAlterRoleOption {
4751 Attributes(PlannedRoleAttributes),
4752 Variable(PlannedRoleVariable),
4753}
4754
4755#[derive(Debug, Clone)]
4756pub struct PlannedRoleAttributes {
4757 pub inherit: Option<bool>,
4758 pub password: Option<Password>,
4759 pub scram_iterations: Option<NonZeroU32>,
4760 pub nopassword: Option<bool>,
4764 pub superuser: Option<bool>,
4765 pub login: Option<bool>,
4766}
4767
4768fn plan_role_attributes(
4769 options: Vec<RoleAttribute>,
4770 scx: &StatementContext,
4771) -> Result<PlannedRoleAttributes, PlanError> {
4772 let mut planned_attributes = PlannedRoleAttributes {
4773 inherit: None,
4774 password: None,
4775 scram_iterations: None,
4776 superuser: None,
4777 login: None,
4778 nopassword: None,
4779 };
4780
4781 for option in options {
4782 match option {
4783 RoleAttribute::Inherit | RoleAttribute::NoInherit
4784 if planned_attributes.inherit.is_some() =>
4785 {
4786 sql_bail!("conflicting or redundant options");
4787 }
4788 RoleAttribute::CreateCluster | RoleAttribute::NoCreateCluster => {
4789 bail_never_supported!(
4790 "CREATECLUSTER attribute",
4791 "sql/create-role/#details",
4792 "Use system privileges instead."
4793 );
4794 }
4795 RoleAttribute::CreateDB | RoleAttribute::NoCreateDB => {
4796 bail_never_supported!(
4797 "CREATEDB attribute",
4798 "sql/create-role/#details",
4799 "Use system privileges instead."
4800 );
4801 }
4802 RoleAttribute::CreateRole | RoleAttribute::NoCreateRole => {
4803 bail_never_supported!(
4804 "CREATEROLE attribute",
4805 "sql/create-role/#details",
4806 "Use system privileges instead."
4807 );
4808 }
4809 RoleAttribute::Password(_) if planned_attributes.password.is_some() => {
4810 sql_bail!("conflicting or redundant options");
4811 }
4812
4813 RoleAttribute::Inherit => planned_attributes.inherit = Some(true),
4814 RoleAttribute::NoInherit => planned_attributes.inherit = Some(false),
4815 RoleAttribute::Password(password) => {
4816 if let Some(password) = password {
4817 planned_attributes.password = Some(password.into());
4818 planned_attributes.scram_iterations =
4819 Some(scx.catalog.system_vars().scram_iterations())
4820 } else {
4821 planned_attributes.nopassword = Some(true);
4822 }
4823 }
4824 RoleAttribute::SuperUser => {
4825 if planned_attributes.superuser == Some(false) {
4826 sql_bail!("conflicting or redundant options");
4827 }
4828 planned_attributes.superuser = Some(true);
4829 }
4830 RoleAttribute::NoSuperUser => {
4831 if planned_attributes.superuser == Some(true) {
4832 sql_bail!("conflicting or redundant options");
4833 }
4834 planned_attributes.superuser = Some(false);
4835 }
4836 RoleAttribute::Login => {
4837 if planned_attributes.login == Some(false) {
4838 sql_bail!("conflicting or redundant options");
4839 }
4840 planned_attributes.login = Some(true);
4841 }
4842 RoleAttribute::NoLogin => {
4843 if planned_attributes.login == Some(true) {
4844 sql_bail!("conflicting or redundant options");
4845 }
4846 planned_attributes.login = Some(false);
4847 }
4848 }
4849 }
4850 if planned_attributes.inherit == Some(false) {
4851 bail_unsupported!("non inherit roles");
4852 }
4853
4854 Ok(planned_attributes)
4855}
4856
4857#[derive(Debug)]
4858pub enum PlannedRoleVariable {
4859 Set { name: String, value: VariableValue },
4860 Reset { name: String },
4861}
4862
4863impl PlannedRoleVariable {
4864 pub fn name(&self) -> &str {
4865 match self {
4866 PlannedRoleVariable::Set { name, .. } => name,
4867 PlannedRoleVariable::Reset { name } => name,
4868 }
4869 }
4870}
4871
4872fn plan_role_variable(
4873 scx: &StatementContext,
4874 variable: SetRoleVar,
4875) -> Result<PlannedRoleVariable, PlanError> {
4876 let plan = match variable {
4877 SetRoleVar::Set { name, value } => {
4878 let name = name.to_string();
4879 let value = scl::plan_set_variable_to(value)?;
4880 if let VariableValue::Values(values) = &value {
4883 vars::check_transaction_isolation_feature_flag(
4884 &name,
4885 VarInput::SqlSet(values),
4886 scx.catalog.system_vars(),
4887 )?;
4888 }
4889 PlannedRoleVariable::Set { name, value }
4890 }
4891 SetRoleVar::Reset { name } => PlannedRoleVariable::Reset {
4892 name: name.to_string(),
4893 },
4894 };
4895 Ok(plan)
4896}
4897
4898pub fn describe_create_role(
4899 _: &StatementContext,
4900 _: CreateRoleStatement,
4901) -> Result<StatementDesc, PlanError> {
4902 Ok(StatementDesc::new(None))
4903}
4904
4905pub fn plan_create_role(
4906 scx: &StatementContext,
4907 CreateRoleStatement { name, options }: CreateRoleStatement,
4908) -> Result<Plan, PlanError> {
4909 let attributes = plan_role_attributes(options, scx)?;
4910 Ok(Plan::CreateRole(CreateRolePlan {
4911 name: normalize::ident(name),
4912 attributes: attributes.into(),
4913 }))
4914}
4915
4916pub fn plan_create_network_policy(
4917 ctx: &StatementContext,
4918 CreateNetworkPolicyStatement { name, options }: CreateNetworkPolicyStatement<Aug>,
4919) -> Result<Plan, PlanError> {
4920 ctx.require_feature_flag(&vars::ENABLE_NETWORK_POLICIES)?;
4921 let policy_options: NetworkPolicyOptionExtracted = options.try_into()?;
4922
4923 let Some(rule_defs) = policy_options.rules else {
4924 sql_bail!("RULES must be specified when creating network policies.");
4925 };
4926
4927 let mut rules = vec![];
4928 for NetworkPolicyRuleDefinition { name, options } in rule_defs {
4929 let NetworkPolicyRuleOptionExtracted {
4930 seen: _,
4931 direction,
4932 action,
4933 address,
4934 } = options.try_into()?;
4935 let (direction, action, address) = match (direction, action, address) {
4936 (Some(direction), Some(action), Some(address)) => (
4937 NetworkPolicyRuleDirection::try_from(direction.as_str())?,
4938 NetworkPolicyRuleAction::try_from(action.as_str())?,
4939 PolicyAddress::try_from(address.as_str())?,
4940 ),
4941 (_, _, _) => {
4942 sql_bail!("Direction, Address, and Action must specified when creating a rule")
4943 }
4944 };
4945 rules.push(NetworkPolicyRule {
4946 name: normalize::ident(name),
4947 direction,
4948 action,
4949 address,
4950 });
4951 }
4952
4953 if rules.len()
4954 > ctx
4955 .catalog
4956 .system_vars()
4957 .max_rules_per_network_policy()
4958 .try_into()?
4959 {
4960 sql_bail!("RULES count exceeds max_rules_per_network_policy.")
4961 }
4962
4963 Ok(Plan::CreateNetworkPolicy(CreateNetworkPolicyPlan {
4964 name: normalize::ident(name),
4965 rules,
4966 }))
4967}
4968
4969pub fn plan_alter_network_policy(
4970 ctx: &StatementContext,
4971 AlterNetworkPolicyStatement { name, options }: AlterNetworkPolicyStatement<Aug>,
4972) -> Result<Plan, PlanError> {
4973 ctx.require_feature_flag(&vars::ENABLE_NETWORK_POLICIES)?;
4974
4975 let policy_options: NetworkPolicyOptionExtracted = options.try_into()?;
4976 let policy = ctx
4977 .catalog
4978 .resolve_network_policy(normalize::ident_ref(&name))?;
4979
4980 let Some(rule_defs) = policy_options.rules else {
4981 sql_bail!("RULES must be specified when creating network policies.");
4982 };
4983
4984 let mut rules = vec![];
4985 for NetworkPolicyRuleDefinition { name, options } in rule_defs {
4986 let NetworkPolicyRuleOptionExtracted {
4987 seen: _,
4988 direction,
4989 action,
4990 address,
4991 } = options.try_into()?;
4992
4993 let (direction, action, address) = match (direction, action, address) {
4994 (Some(direction), Some(action), Some(address)) => (
4995 NetworkPolicyRuleDirection::try_from(direction.as_str())?,
4996 NetworkPolicyRuleAction::try_from(action.as_str())?,
4997 PolicyAddress::try_from(address.as_str())?,
4998 ),
4999 (_, _, _) => {
5000 sql_bail!("Direction, Address, and Action must specified when creating a rule")
5001 }
5002 };
5003 rules.push(NetworkPolicyRule {
5004 name: normalize::ident(name),
5005 direction,
5006 action,
5007 address,
5008 });
5009 }
5010 if rules.len()
5011 > ctx
5012 .catalog
5013 .system_vars()
5014 .max_rules_per_network_policy()
5015 .try_into()?
5016 {
5017 sql_bail!("RULES count exceeds max_rules_per_network_policy.")
5018 }
5019
5020 Ok(Plan::AlterNetworkPolicy(AlterNetworkPolicyPlan {
5021 id: policy.id(),
5022 name: normalize::ident(name),
5023 rules,
5024 }))
5025}
5026
5027pub fn describe_create_cluster(
5028 _: &StatementContext,
5029 _: CreateClusterStatement<Aug>,
5030) -> Result<StatementDesc, PlanError> {
5031 Ok(StatementDesc::new(None))
5032}
5033
5034generate_extracted_config!(
5040 ClusterOption,
5041 (AutoScalingStrategy, ClusterAutoScalingStrategyOptionValue),
5042 (AvailabilityZones, Vec<String>),
5043 (Disk, bool),
5044 (ExperimentalArrangementCompression, bool),
5045 (IntrospectionDebugging, bool),
5046 (IntrospectionInterval, OptionalDuration),
5047 (Managed, bool),
5048 (Replicas, Vec<ReplicaDefinition<Aug>>),
5049 (ReplicationFactor, u32),
5050 (Size, String),
5051 (Schedule, ClusterScheduleOptionValue),
5052 (WorkloadClass, OptionalString)
5053);
5054
5055generate_extracted_config!(
5056 NetworkPolicyOption,
5057 (Rules, Vec<NetworkPolicyRuleDefinition<Aug>>)
5058);
5059
5060generate_extracted_config!(
5061 NetworkPolicyRuleOption,
5062 (Direction, String),
5063 (Action, String),
5064 (Address, String)
5065);
5066
5067generate_extracted_config!(ClusterAlterOption, (Wait, ClusterAlterOptionValue<Aug>));
5068
5069generate_extracted_config!(
5070 ClusterAlterUntilReadyOption,
5071 (Timeout, Duration),
5072 (OnTimeout, String)
5073);
5074
5075generate_extracted_config!(
5076 ClusterFeature,
5077 (ReoptimizeImportedViews, Option<bool>, Default(None)),
5078 (EnableEagerDeltaJoins, Option<bool>, Default(None)),
5079 (EnableNewOuterJoinLowering, Option<bool>, Default(None)),
5080 (EnableVariadicLeftJoinLowering, Option<bool>, Default(None)),
5081 (EnableLetrecFixpointAnalysis, Option<bool>, Default(None)),
5082 (EnableJoinPrioritizeArranged, Option<bool>, Default(None)),
5083 (
5084 EnableProjectionPushdownAfterRelationCse,
5085 Option<bool>,
5086 Default(None)
5087 ),
5088 (
5089 EnableUnionCancellationAfterRelationCse,
5090 Option<bool>,
5091 Default(None)
5092 )
5093);
5094
5095pub fn plan_create_cluster(
5099 scx: &StatementContext,
5100 stmt: CreateClusterStatement<Aug>,
5101) -> Result<Plan, PlanError> {
5102 let plan = plan_create_cluster_inner(scx, stmt)?;
5103
5104 if let CreateClusterVariant::Managed(_) = &plan.variant {
5106 let stmt = unplan_create_cluster(scx, plan.clone())
5107 .map_err(|e| PlanError::Replan(e.to_string()))?;
5108 let create_sql = stmt.to_ast_string_stable();
5109 let stmt = parse::parse(&create_sql)
5110 .map_err(|e| PlanError::Replan(e.to_string()))?
5111 .into_element()
5112 .ast;
5113 let (stmt, _resolved_ids) =
5114 names::resolve(scx.catalog, stmt).map_err(|e| PlanError::Replan(e.to_string()))?;
5115 let stmt = match stmt {
5116 Statement::CreateCluster(stmt) => stmt,
5117 stmt => {
5118 return Err(PlanError::Replan(format!(
5119 "replan does not match: plan={plan:?}, create_sql={create_sql:?}, stmt={stmt:?}"
5120 )));
5121 }
5122 };
5123 let replan =
5124 plan_create_cluster_inner(scx, stmt).map_err(|e| PlanError::Replan(e.to_string()))?;
5125 if plan != replan {
5126 return Err(PlanError::Replan(format!(
5127 "replan does not match: plan={plan:?}, replan={replan:?}"
5128 )));
5129 }
5130 }
5131
5132 Ok(Plan::CreateCluster(plan))
5133}
5134
5135pub fn plan_create_cluster_inner(
5136 scx: &StatementContext,
5137 CreateClusterStatement {
5138 name,
5139 options,
5140 features,
5141 if_not_exists,
5142 }: CreateClusterStatement<Aug>,
5143) -> Result<CreateClusterPlan, PlanError> {
5144 let ClusterOptionExtracted {
5145 auto_scaling_strategy,
5146 availability_zones,
5147 experimental_arrangement_compression,
5148 introspection_debugging,
5149 introspection_interval,
5150 managed,
5151 replicas,
5152 replication_factor,
5153 seen: _,
5154 size,
5155 disk,
5156 schedule,
5157 workload_class,
5158 }: ClusterOptionExtracted = options.try_into()?;
5159
5160 let managed = managed.unwrap_or_else(|| replicas.is_none());
5161
5162 if !scx.catalog.active_role_id().is_system() {
5163 if !features.is_empty() {
5164 sql_bail!("FEATURES not supported for non-system users");
5165 }
5166 if workload_class.is_some() {
5167 sql_bail!("WORKLOAD CLASS not supported for non-system users");
5168 }
5169 }
5170
5171 let schedule = schedule.unwrap_or(ClusterScheduleOptionValue::Manual);
5172 let workload_class = workload_class.and_then(|v| v.0);
5173
5174 if managed {
5175 if replicas.is_some() {
5176 sql_bail!("REPLICAS not supported for managed clusters");
5177 }
5178 let Some(size) = size else {
5179 sql_bail!("SIZE must be specified for managed clusters");
5180 };
5181
5182 if disk.is_some() {
5183 if scx.catalog.is_cluster_size_cc(&size) {
5187 sql_bail!(
5188 "DISK option not supported for modern cluster sizes because disk is always enabled"
5189 );
5190 }
5191
5192 scx.catalog
5193 .add_notice(PlanNotice::ReplicaDiskOptionDeprecated);
5194 }
5195
5196 let compute = plan_compute_replica_config(
5197 introspection_interval,
5198 introspection_debugging.unwrap_or(false),
5199 experimental_arrangement_compression.unwrap_or(false),
5200 )?;
5201
5202 let replication_factor = if matches!(schedule, ClusterScheduleOptionValue::Manual) {
5203 replication_factor.unwrap_or_else(|| {
5204 scx.catalog
5205 .system_vars()
5206 .default_cluster_replication_factor()
5207 })
5208 } else {
5209 scx.require_feature_flag(&ENABLE_CLUSTER_SCHEDULE_REFRESH)?;
5210 if replication_factor.is_some() {
5211 sql_bail!(
5212 "REPLICATION FACTOR cannot be given together with any SCHEDULE other than MANUAL"
5213 );
5214 }
5215 0
5219 };
5220 let availability_zones = availability_zones.unwrap_or_default();
5221
5222 if !availability_zones.is_empty() {
5223 scx.require_feature_flag(&vars::ENABLE_MANAGED_CLUSTER_AVAILABILITY_ZONES)?;
5224 }
5225
5226 let ClusterFeatureExtracted {
5228 reoptimize_imported_views,
5229 enable_eager_delta_joins,
5230 enable_new_outer_join_lowering,
5231 enable_variadic_left_join_lowering,
5232 enable_letrec_fixpoint_analysis,
5233 enable_join_prioritize_arranged,
5234 enable_projection_pushdown_after_relation_cse,
5235 enable_union_cancellation_after_relation_cse,
5236 seen: _,
5237 } = ClusterFeatureExtracted::try_from(features)?;
5238 let optimizer_feature_overrides = OptimizerFeatureOverrides {
5239 reoptimize_imported_views,
5240 enable_eager_delta_joins,
5241 enable_new_outer_join_lowering,
5242 enable_variadic_left_join_lowering,
5243 enable_letrec_fixpoint_analysis,
5244 enable_join_prioritize_arranged,
5245 enable_projection_pushdown_after_relation_cse,
5246 enable_union_cancellation_after_relation_cse,
5247 ..Default::default()
5248 };
5249
5250 let auto_scaling_strategy = match auto_scaling_strategy {
5253 Some(value) => {
5254 scx.require_feature_flag(&ENABLE_AUTO_SCALING_STRATEGY)?;
5255 let strategy = plan_auto_scaling_strategy(value)?;
5256 if let Some(strategy) = &strategy {
5257 let schedule_non_manual =
5258 !matches!(schedule, ClusterScheduleOptionValue::Manual);
5259 validate_auto_scaling_strategy(strategy, Some(&size), schedule_non_manual)?;
5260 }
5261 strategy
5262 }
5263 None => None,
5264 };
5265
5266 let schedule = plan_cluster_schedule(schedule)?;
5267
5268 Ok(CreateClusterPlan {
5269 name: normalize::ident(name),
5270 variant: CreateClusterVariant::Managed(CreateClusterManagedPlan {
5271 replication_factor,
5272 size,
5273 availability_zones,
5274 compute,
5275 optimizer_feature_overrides,
5276 schedule,
5277 auto_scaling_strategy,
5278 }),
5279 workload_class,
5280 if_not_exists,
5281 })
5282 } else {
5283 let Some(replica_defs) = replicas else {
5284 sql_bail!("REPLICAS must be specified for unmanaged clusters");
5285 };
5286 if auto_scaling_strategy.is_some() {
5287 sql_bail!("AUTO SCALING STRATEGY not supported for unmanaged clusters");
5288 }
5289 if availability_zones.is_some() {
5290 sql_bail!("AVAILABILITY ZONES not supported for unmanaged clusters");
5291 }
5292 if replication_factor.is_some() {
5293 sql_bail!("REPLICATION FACTOR not supported for unmanaged clusters");
5294 }
5295 if introspection_debugging.is_some() {
5296 sql_bail!("INTROSPECTION DEBUGGING not supported for unmanaged clusters");
5297 }
5298 if introspection_interval.is_some() {
5299 sql_bail!("INTROSPECTION INTERVAL not supported for unmanaged clusters");
5300 }
5301 if experimental_arrangement_compression.is_some() {
5302 sql_bail!("EXPERIMENTAL ARRANGEMENT COMPRESSION not supported for unmanaged clusters");
5303 }
5304 if size.is_some() {
5305 sql_bail!("SIZE not supported for unmanaged clusters");
5306 }
5307 if disk.is_some() {
5308 sql_bail!("DISK not supported for unmanaged clusters");
5309 }
5310 if !features.is_empty() {
5311 sql_bail!("FEATURES not supported for unmanaged clusters");
5312 }
5313 if !matches!(schedule, ClusterScheduleOptionValue::Manual) {
5314 sql_bail!(
5315 "cluster schedules other than MANUAL are not supported for unmanaged clusters"
5316 );
5317 }
5318
5319 let mut replicas = vec![];
5320 for ReplicaDefinition { name, options } in replica_defs {
5321 replicas.push((normalize::ident(name), plan_replica_config(scx, options)?));
5322 }
5323
5324 Ok(CreateClusterPlan {
5325 name: normalize::ident(name),
5326 variant: CreateClusterVariant::Unmanaged(CreateClusterUnmanagedPlan { replicas }),
5327 workload_class,
5328 if_not_exists,
5329 })
5330 }
5331}
5332
5333pub fn unplan_create_cluster(
5340 scx: &StatementContext,
5341 CreateClusterPlan {
5342 name,
5343 variant,
5344 workload_class,
5345 if_not_exists,
5346 }: CreateClusterPlan,
5347) -> Result<CreateClusterStatement<Aug>, PlanError> {
5348 match variant {
5349 CreateClusterVariant::Managed(CreateClusterManagedPlan {
5350 replication_factor,
5351 size,
5352 availability_zones,
5353 compute,
5354 optimizer_feature_overrides,
5355 schedule,
5356 auto_scaling_strategy,
5357 }) => {
5358 let schedule = unplan_cluster_schedule(schedule);
5359 let auto_scaling_strategy = auto_scaling_strategy
5360 .as_ref()
5361 .map(unplan_auto_scaling_strategy);
5362 let OptimizerFeatureOverrides {
5363 enable_reduce_mfp_fusion: _,
5364 enable_cardinality_estimates: _,
5365 persist_fast_path_limit: _,
5366 reoptimize_imported_views,
5367 enable_eager_delta_joins,
5368 enable_new_outer_join_lowering,
5369 enable_variadic_left_join_lowering,
5370 enable_letrec_fixpoint_analysis,
5371 enable_join_prioritize_arranged,
5372 enable_projection_pushdown_after_relation_cse,
5373 enable_union_cancellation_after_relation_cse,
5374 enable_less_reduce_in_eqprop: _,
5375 enable_dequadratic_eqprop_map: _,
5376 enable_eq_classes_withholding_errors: _,
5377 enable_fast_path_plan_insights: _,
5378 enable_cast_elimination: _,
5379 enable_case_literal_transform: _,
5380 enable_simplify_quantified_comparisons: _,
5381 enable_simplify_from_less_existence: _,
5382 enable_coalesce_case_transform: _,
5383 enable_will_distinct_propagation: _,
5384 enable_fixed_correlated_cte_lowering: _,
5385 } = optimizer_feature_overrides;
5386 let features_extracted = ClusterFeatureExtracted {
5388 seen: Default::default(),
5390 reoptimize_imported_views,
5391 enable_eager_delta_joins,
5392 enable_new_outer_join_lowering,
5393 enable_variadic_left_join_lowering,
5394 enable_letrec_fixpoint_analysis,
5395 enable_join_prioritize_arranged,
5396 enable_projection_pushdown_after_relation_cse,
5397 enable_union_cancellation_after_relation_cse,
5398 };
5399 let features = features_extracted.into_values(scx.catalog);
5400 let availability_zones = if availability_zones.is_empty() {
5401 None
5402 } else {
5403 Some(availability_zones)
5404 };
5405 let (introspection_interval, introspection_debugging, arrangement_compression) =
5406 unplan_compute_replica_config(compute);
5407 let replication_factor = match &schedule {
5410 ClusterScheduleOptionValue::Manual => Some(replication_factor),
5411 ClusterScheduleOptionValue::Refresh { .. } => {
5412 soft_assert_or_log!(
5419 replication_factor <= 1,
5420 "replication factor, {replication_factor:?}, must be <= 1 with a refresh schedule"
5421 );
5422 None
5423 }
5424 };
5425 let workload_class = workload_class.map(|s| OptionalString(Some(s)));
5426 let options_extracted = ClusterOptionExtracted {
5427 seen: Default::default(),
5429 auto_scaling_strategy,
5430 availability_zones,
5431 disk: None,
5432 experimental_arrangement_compression: Some(arrangement_compression),
5433 introspection_debugging: Some(introspection_debugging),
5434 introspection_interval,
5435 managed: Some(true),
5436 replicas: None,
5437 replication_factor,
5438 size: Some(size),
5439 schedule: Some(schedule),
5440 workload_class,
5441 };
5442 let options = options_extracted.into_values(scx.catalog);
5443 let name = Ident::new_unchecked(name);
5444 Ok(CreateClusterStatement {
5445 name,
5446 options,
5447 features,
5448 if_not_exists,
5449 })
5450 }
5451 CreateClusterVariant::Unmanaged(_) => {
5452 bail_unsupported!("SHOW CREATE for unmanaged clusters")
5453 }
5454 }
5455}
5456
5457generate_extracted_config!(
5458 ReplicaOption,
5459 (AvailabilityZone, String),
5460 (BilledAs, String),
5461 (ComputeAddresses, Vec<String>),
5462 (ComputectlAddresses, Vec<String>),
5463 (Disk, bool),
5464 (ExperimentalArrangementCompression, bool, Default(false)),
5465 (Internal, bool, Default(false)),
5466 (IntrospectionDebugging, bool, Default(false)),
5467 (IntrospectionInterval, OptionalDuration),
5468 (Size, String),
5469 (StorageAddresses, Vec<String>),
5470 (StoragectlAddresses, Vec<String>),
5471 (Workers, u16)
5472);
5473
5474fn plan_replica_config(
5475 scx: &StatementContext,
5476 options: Vec<ReplicaOption<Aug>>,
5477) -> Result<ReplicaConfig, PlanError> {
5478 let ReplicaOptionExtracted {
5479 availability_zone,
5480 billed_as,
5481 computectl_addresses,
5482 disk,
5483 experimental_arrangement_compression,
5484 internal,
5485 introspection_debugging,
5486 introspection_interval,
5487 size,
5488 storagectl_addresses,
5489 ..
5490 }: ReplicaOptionExtracted = options.try_into()?;
5491
5492 let compute = plan_compute_replica_config(
5493 introspection_interval,
5494 introspection_debugging,
5495 experimental_arrangement_compression,
5496 )?;
5497
5498 match (
5499 size,
5500 availability_zone,
5501 billed_as,
5502 storagectl_addresses,
5503 computectl_addresses,
5504 ) {
5505 (None, _, None, None, None) => {
5507 sql_bail!("SIZE option must be specified");
5510 }
5511 (Some(size), availability_zone, billed_as, None, None) => {
5512 if disk.is_some() {
5513 if scx.catalog.is_cluster_size_cc(&size) {
5517 sql_bail!(
5518 "DISK option not supported for modern cluster sizes because disk is always enabled"
5519 );
5520 }
5521
5522 scx.catalog
5523 .add_notice(PlanNotice::ReplicaDiskOptionDeprecated);
5524 }
5525
5526 Ok(ReplicaConfig::Orchestrated {
5527 size,
5528 availability_zone,
5529 compute,
5530 billed_as,
5531 internal,
5532 })
5533 }
5534
5535 (None, None, None, storagectl_addresses, computectl_addresses) => {
5536 scx.require_feature_flag(&vars::UNSAFE_ENABLE_UNORCHESTRATED_CLUSTER_REPLICAS)?;
5537
5538 let Some(storagectl_addrs) = storagectl_addresses else {
5542 sql_bail!("missing STORAGECTL ADDRESSES option");
5543 };
5544 let Some(computectl_addrs) = computectl_addresses else {
5545 sql_bail!("missing COMPUTECTL ADDRESSES option");
5546 };
5547
5548 if storagectl_addrs.len() != computectl_addrs.len() {
5549 sql_bail!(
5550 "COMPUTECTL ADDRESSES and STORAGECTL ADDRESSES must have the same length"
5551 );
5552 }
5553
5554 if disk.is_some() {
5555 sql_bail!("DISK can't be specified for unorchestrated clusters");
5556 }
5557
5558 Ok(ReplicaConfig::Unorchestrated {
5559 storagectl_addrs,
5560 computectl_addrs,
5561 compute,
5562 })
5563 }
5564 _ => {
5565 sql_bail!("invalid mixture of orchestrated and unorchestrated replica options");
5568 }
5569 }
5570}
5571
5572fn plan_compute_replica_config(
5576 introspection_interval: Option<OptionalDuration>,
5577 introspection_debugging: bool,
5578 arrangement_compression: bool,
5579) -> Result<ComputeReplicaConfig, PlanError> {
5580 let introspection_interval = introspection_interval
5581 .map(|OptionalDuration(i)| i)
5582 .unwrap_or(Some(DEFAULT_REPLICA_LOGGING_INTERVAL));
5583 let introspection = match introspection_interval {
5584 Some(interval) => Some(ComputeReplicaIntrospectionConfig {
5585 interval,
5586 debugging: introspection_debugging,
5587 }),
5588 None if introspection_debugging => {
5589 sql_bail!("INTROSPECTION DEBUGGING cannot be specified without INTROSPECTION INTERVAL")
5590 }
5591 None => None,
5592 };
5593 let compute = ComputeReplicaConfig {
5594 introspection,
5595 arrangement_compression,
5596 };
5597 Ok(compute)
5598}
5599
5600fn unplan_compute_replica_config(
5605 compute_replica_config: ComputeReplicaConfig,
5606) -> (Option<OptionalDuration>, bool, bool) {
5607 let ComputeReplicaConfig {
5608 introspection,
5609 arrangement_compression,
5610 } = compute_replica_config;
5611 match introspection {
5612 Some(ComputeReplicaIntrospectionConfig {
5613 debugging,
5614 interval,
5615 }) => (
5616 Some(OptionalDuration(Some(interval))),
5617 debugging,
5618 arrangement_compression,
5619 ),
5620 None => (Some(OptionalDuration(None)), false, arrangement_compression),
5621 }
5622}
5623
5624fn plan_cluster_schedule(
5628 schedule: ClusterScheduleOptionValue,
5629) -> Result<ClusterSchedule, PlanError> {
5630 Ok(match schedule {
5631 ClusterScheduleOptionValue::Manual => ClusterSchedule::Manual,
5632 ClusterScheduleOptionValue::Refresh {
5634 hydration_time_estimate: None,
5635 } => ClusterSchedule::Refresh {
5636 hydration_time_estimate: Duration::from_millis(0),
5637 },
5638 ClusterScheduleOptionValue::Refresh {
5640 hydration_time_estimate: Some(interval_value),
5641 } => {
5642 let interval = Interval::try_from_value(Value::Interval(interval_value))?;
5643 if interval.as_microseconds() < 0 {
5644 sql_bail!(
5645 "HYDRATION TIME ESTIMATE must be non-negative; got: {}",
5646 interval
5647 );
5648 }
5649 if interval.months != 0 {
5650 sql_bail!("HYDRATION TIME ESTIMATE must not involve units larger than days");
5654 }
5655 let duration = interval.duration()?;
5656 if u64::try_from(duration.as_millis()).is_err()
5657 || Interval::from_duration(&duration).is_err()
5658 {
5659 sql_bail!("HYDRATION TIME ESTIMATE too large");
5660 }
5661 ClusterSchedule::Refresh {
5662 hydration_time_estimate: duration,
5663 }
5664 }
5665 })
5666}
5667
5668fn unplan_cluster_schedule(schedule: ClusterSchedule) -> ClusterScheduleOptionValue {
5672 match schedule {
5673 ClusterSchedule::Manual => ClusterScheduleOptionValue::Manual,
5674 ClusterSchedule::Refresh {
5675 hydration_time_estimate,
5676 } => {
5677 let interval = Interval::from_duration(&hydration_time_estimate)
5678 .expect("planning ensured that this is convertible back to Interval");
5679 let interval_value = literal::unplan_interval(&interval);
5680 ClusterScheduleOptionValue::Refresh {
5681 hydration_time_estimate: Some(interval_value),
5682 }
5683 }
5684 }
5685}
5686
5687fn plan_auto_scaling_strategy(
5695 value: ClusterAutoScalingStrategyOptionValue,
5696) -> Result<Option<AutoScalingStrategy>, PlanError> {
5697 let ClusterAutoScalingStrategyOptionValue { on_hydration } = value;
5698 let Some(on_hydration) = on_hydration else {
5699 return Ok(None);
5701 };
5702
5703 let hydration_size = String::try_from_value(on_hydration.hydration_size)?;
5704
5705 let linger_duration = on_hydration
5706 .linger_duration
5707 .map(Duration::try_from_value)
5708 .transpose()?;
5709
5710 Ok(Some(AutoScalingStrategy {
5711 on_hydration: Some(OnHydration {
5712 hydration_size,
5713 linger_duration,
5714 }),
5715 }))
5716}
5717
5718fn validate_auto_scaling_strategy(
5727 strategy: &AutoScalingStrategy,
5728 cluster_size: Option<&str>,
5729 schedule_non_manual: bool,
5730) -> Result<(), PlanError> {
5731 if let (Some(on_hydration), Some(cluster_size)) = (&strategy.on_hydration, cluster_size) {
5732 if on_hydration.hydration_size == cluster_size {
5733 return Err(PlanError::HydrationSizeEqualsClusterSize {
5734 size: cluster_size.to_string(),
5735 });
5736 }
5737 }
5738 if schedule_non_manual {
5739 sql_bail!("AUTO SCALING STRATEGY cannot be combined with a SCHEDULE other than MANUAL");
5740 }
5741 Ok(())
5742}
5743
5744fn unplan_auto_scaling_strategy(
5749 strategy: &AutoScalingStrategy,
5750) -> ClusterAutoScalingStrategyOptionValue {
5751 ClusterAutoScalingStrategyOptionValue {
5752 on_hydration: strategy
5753 .on_hydration
5754 .as_ref()
5755 .map(|on_hydration| OnHydrationOptionValue {
5756 hydration_size: Value::String(on_hydration.hydration_size.clone()),
5757 linger_duration: on_hydration.linger_duration.map(|d| {
5758 let interval = Interval::from_duration(&d)
5759 .expect("planning ensured this is convertible back to Interval");
5760 Value::Interval(literal::unplan_interval(&interval))
5761 }),
5762 }),
5763 }
5764}
5765
5766pub fn describe_create_cluster_replica(
5767 _: &StatementContext,
5768 _: CreateClusterReplicaStatement<Aug>,
5769) -> Result<StatementDesc, PlanError> {
5770 Ok(StatementDesc::new(None))
5771}
5772
5773pub fn plan_create_cluster_replica(
5774 scx: &StatementContext,
5775 CreateClusterReplicaStatement {
5776 definition: ReplicaDefinition { name, options },
5777 of_cluster,
5778 if_not_exists,
5779 }: CreateClusterReplicaStatement<Aug>,
5780) -> Result<Plan, PlanError> {
5781 let cluster = scx
5782 .catalog
5783 .resolve_cluster(Some(&normalize::ident(of_cluster)))?;
5784
5785 let config = plan_replica_config(scx, options)?;
5786
5787 if let ReplicaConfig::Orchestrated { internal: true, .. } = &config {
5788 if MANAGED_REPLICA_PATTERN.is_match(name.as_str()) {
5789 return Err(PlanError::MangedReplicaName(name.into_string()));
5790 }
5791 } else {
5792 ensure_cluster_is_not_managed(scx, cluster.id())?;
5793 }
5794
5795 Ok(Plan::CreateClusterReplica(CreateClusterReplicaPlan {
5796 name: normalize::ident(name),
5797 cluster_id: cluster.id(),
5798 config,
5799 if_not_exists,
5800 }))
5801}
5802
5803pub fn describe_create_secret(
5804 _: &StatementContext,
5805 _: CreateSecretStatement<Aug>,
5806) -> Result<StatementDesc, PlanError> {
5807 Ok(StatementDesc::new(None))
5808}
5809
5810pub fn plan_create_secret(
5811 scx: &StatementContext,
5812 stmt: CreateSecretStatement<Aug>,
5813) -> Result<Plan, PlanError> {
5814 let CreateSecretStatement {
5815 name,
5816 if_not_exists,
5817 value,
5818 } = &stmt;
5819
5820 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name.to_owned())?)?;
5821 let mut create_sql_statement = stmt.clone();
5822 create_sql_statement.value = Expr::Value(Value::String("********".to_string()));
5823 let create_sql =
5824 normalize::create_statement(scx, Statement::CreateSecret(create_sql_statement))?;
5825 let secret_as = query::plan_secret_as(scx, value.clone())?;
5826
5827 let secret = Secret {
5828 create_sql,
5829 secret_as,
5830 };
5831
5832 Ok(Plan::CreateSecret(CreateSecretPlan {
5833 name,
5834 secret,
5835 if_not_exists: *if_not_exists,
5836 }))
5837}
5838
5839pub fn describe_create_connection(
5840 _: &StatementContext,
5841 _: CreateConnectionStatement<Aug>,
5842) -> Result<StatementDesc, PlanError> {
5843 Ok(StatementDesc::new(None))
5844}
5845
5846generate_extracted_config!(CreateConnectionOption, (Validate, bool));
5847
5848pub fn plan_create_connection(
5849 scx: &StatementContext,
5850 mut stmt: CreateConnectionStatement<Aug>,
5851) -> Result<Plan, PlanError> {
5852 let CreateConnectionStatement {
5853 name,
5854 connection_type,
5855 values,
5856 if_not_exists,
5857 with_options,
5858 } = stmt.clone();
5859 let connection_options_extracted = connection::ConnectionOptionExtracted::try_from(values)?;
5860 let details = connection_options_extracted.try_into_connection_details(scx, connection_type)?;
5861 let name = scx.allocate_qualified_name(normalize::unresolved_item_name(name)?)?;
5862
5863 let options = CreateConnectionOptionExtracted::try_from(with_options)?;
5864 if options.validate.is_some() {
5865 scx.require_feature_flag(&vars::ENABLE_CONNECTION_VALIDATION_SYNTAX)?;
5866 }
5867 let validate = match options.validate {
5868 Some(val) => val,
5869 None => {
5870 scx.catalog
5871 .system_vars()
5872 .enable_default_connection_validation()
5873 && details.to_connection().validate_by_default()
5874 }
5875 };
5876
5877 let full_name = scx.catalog.resolve_full_name(&name);
5879 let partial_name = PartialItemName::from(full_name.clone());
5880 if let (false, Ok(item)) = (if_not_exists, scx.catalog.resolve_item(&partial_name)) {
5881 return Err(PlanError::ItemAlreadyExists {
5882 name: full_name.to_string(),
5883 item_type: item.item_type(),
5884 });
5885 }
5886
5887 if let ConnectionDetails::Ssh { key_1, key_2, .. } = &details {
5890 stmt.values.retain(|v| {
5891 v.name != ConnectionOptionName::PublicKey1 && v.name != ConnectionOptionName::PublicKey2
5892 });
5893 stmt.values.push(ConnectionOption {
5894 name: ConnectionOptionName::PublicKey1,
5895 value: Some(WithOptionValue::Value(Value::String(key_1.public_key()))),
5896 });
5897 stmt.values.push(ConnectionOption {
5898 name: ConnectionOptionName::PublicKey2,
5899 value: Some(WithOptionValue::Value(Value::String(key_2.public_key()))),
5900 });
5901 }
5902 let create_sql = normalize::create_statement(scx, Statement::CreateConnection(stmt))?;
5903
5904 let plan = CreateConnectionPlan {
5905 name,
5906 if_not_exists,
5907 connection: crate::plan::Connection {
5908 create_sql,
5909 details,
5910 },
5911 validate,
5912 };
5913 Ok(Plan::CreateConnection(plan))
5914}
5915
5916fn plan_drop_database(
5917 scx: &StatementContext,
5918 if_exists: bool,
5919 name: &UnresolvedDatabaseName,
5920 cascade: bool,
5921) -> Result<Option<DatabaseId>, PlanError> {
5922 Ok(match resolve_database(scx, name, if_exists)? {
5923 Some(database) => {
5924 if !cascade && database.has_schemas() {
5925 sql_bail!(
5926 "database '{}' cannot be dropped with RESTRICT while it contains schemas",
5927 name,
5928 );
5929 }
5930 Some(database.id())
5931 }
5932 None => None,
5933 })
5934}
5935
5936pub fn describe_drop_objects(
5937 _: &StatementContext,
5938 _: DropObjectsStatement,
5939) -> Result<StatementDesc, PlanError> {
5940 Ok(StatementDesc::new(None))
5941}
5942
5943pub fn plan_drop_objects(
5944 scx: &mut StatementContext,
5945 DropObjectsStatement {
5946 object_type,
5947 if_exists,
5948 names,
5949 cascade,
5950 }: DropObjectsStatement,
5951) -> Result<Plan, PlanError> {
5952 if object_type == mz_sql_parser::ast::ObjectType::Func {
5953 bail_unsupported!("DROP FUNCTION");
5954 }
5955 let object_type = object_type.into();
5956
5957 let mut referenced_ids = Vec::new();
5958 for name in names {
5959 let id = match &name {
5960 UnresolvedObjectName::Cluster(name) => {
5961 plan_drop_cluster(scx, if_exists, name, cascade)?.map(ObjectId::Cluster)
5962 }
5963 UnresolvedObjectName::ClusterReplica(name) => {
5964 plan_drop_cluster_replica(scx, if_exists, name)?.map(ObjectId::ClusterReplica)
5965 }
5966 UnresolvedObjectName::Database(name) => {
5967 plan_drop_database(scx, if_exists, name, cascade)?.map(ObjectId::Database)
5968 }
5969 UnresolvedObjectName::Schema(name) => {
5970 plan_drop_schema(scx, if_exists, name, cascade)?.map(ObjectId::Schema)
5971 }
5972 UnresolvedObjectName::Role(name) => {
5973 plan_drop_role(scx, if_exists, name)?.map(ObjectId::Role)
5974 }
5975 UnresolvedObjectName::Item(name) => {
5976 plan_drop_item_name(scx, object_type, if_exists, name.clone())?.map(ObjectId::Item)
5980 }
5981 UnresolvedObjectName::NetworkPolicy(name) => {
5982 plan_drop_network_policy(scx, if_exists, name)?.map(ObjectId::NetworkPolicy)
5983 }
5984 };
5985 match id {
5986 Some(id) => referenced_ids.push(id),
5987 None => scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
5988 name: name.to_ast_string_simple(),
5989 object_type,
5990 }),
5991 }
5992 }
5993
5994 if !cascade {
5998 let dropped_items: BTreeSet<CatalogItemId> = referenced_ids
5999 .iter()
6000 .filter_map(|id| match id {
6001 ObjectId::Item(id) => Some(*id),
6002 _ => None,
6003 })
6004 .collect();
6005 for id in &dropped_items {
6006 let catalog_item = scx.catalog.get_item(id);
6007 ensure_no_blocking_dependents(scx, object_type, catalog_item, &dropped_items)?;
6008 }
6009 }
6010
6011 let drop_ids = scx.catalog.object_dependents(&referenced_ids);
6012
6013 Ok(Plan::DropObjects(DropObjectsPlan {
6014 referenced_ids,
6015 drop_ids,
6016 object_type,
6017 }))
6018}
6019
6020fn plan_drop_schema(
6021 scx: &StatementContext,
6022 if_exists: bool,
6023 name: &UnresolvedSchemaName,
6024 cascade: bool,
6025) -> Result<Option<(ResolvedDatabaseSpecifier, SchemaSpecifier)>, PlanError> {
6026 let normalized = normalize::unresolved_schema_name(name.clone())?;
6030 if normalized.database.is_none() && normalized.schema == mz_repr::namespaces::MZ_TEMP_SCHEMA {
6031 sql_bail!("cannot drop schema {name} because it is a temporary schema",)
6032 }
6033
6034 Ok(match resolve_schema(scx, name.clone(), if_exists)? {
6035 Some((database_spec, schema_spec)) => {
6036 if let ResolvedDatabaseSpecifier::Ambient = database_spec {
6037 sql_bail!(
6038 "cannot drop schema {name} because it is required by the database system",
6039 );
6040 }
6041 if let SchemaSpecifier::Temporary = schema_spec {
6042 sql_bail!("cannot drop schema {name} because it is a temporary schema",)
6043 }
6044 let schema = scx.get_schema(&database_spec, &schema_spec);
6045 if !cascade && schema.has_items() {
6046 let full_schema_name = scx.catalog.resolve_full_schema_name(schema.name());
6047 sql_bail!(
6048 "schema '{}' cannot be dropped without CASCADE while it contains objects",
6049 full_schema_name
6050 );
6051 }
6052 Some((database_spec, schema_spec))
6053 }
6054 None => None,
6055 })
6056}
6057
6058fn plan_drop_role(
6059 scx: &StatementContext,
6060 if_exists: bool,
6061 name: &Ident,
6062) -> Result<Option<RoleId>, PlanError> {
6063 match scx.catalog.resolve_role(name.as_str()) {
6064 Ok(role) => {
6065 let id = role.id();
6066 if &id == scx.catalog.active_role_id() {
6067 sql_bail!("current role cannot be dropped");
6068 }
6069 for role in scx.catalog.get_roles() {
6070 for (member_id, grantor_id) in role.membership() {
6071 if &id == grantor_id {
6072 let member_role = scx.catalog.get_role(member_id);
6073 sql_bail!(
6074 "cannot drop role {}: still depended up by membership of role {} in role {}",
6075 name.as_str(),
6076 role.name(),
6077 member_role.name()
6078 );
6079 }
6080 }
6081 }
6082 Ok(Some(role.id()))
6083 }
6084 Err(_) if if_exists => Ok(None),
6085 Err(e) => Err(e.into()),
6086 }
6087}
6088
6089fn plan_drop_cluster(
6090 scx: &StatementContext,
6091 if_exists: bool,
6092 name: &Ident,
6093 cascade: bool,
6094) -> Result<Option<ClusterId>, PlanError> {
6095 Ok(match resolve_cluster(scx, name, if_exists)? {
6096 Some(cluster) => {
6097 if !cascade && !cluster.bound_objects().is_empty() {
6098 return Err(PlanError::DependentObjectsStillExist {
6099 object_type: "cluster".to_string(),
6100 object_name: cluster.name().to_string(),
6101 dependents: Vec::new(),
6102 });
6103 }
6104 Some(cluster.id())
6105 }
6106 None => None,
6107 })
6108}
6109
6110fn plan_drop_network_policy(
6111 scx: &StatementContext,
6112 if_exists: bool,
6113 name: &Ident,
6114) -> Result<Option<NetworkPolicyId>, PlanError> {
6115 match scx.catalog.resolve_network_policy(name.as_str()) {
6116 Ok(policy) => {
6117 if scx.catalog.system_vars().default_network_policy_name() == policy.name() {
6120 Err(PlanError::NetworkPolicyInUse)
6121 } else {
6122 Ok(Some(policy.id()))
6123 }
6124 }
6125 Err(_) if if_exists => Ok(None),
6126 Err(e) => Err(e.into()),
6127 }
6128}
6129
6130fn plan_drop_cluster_replica(
6131 scx: &StatementContext,
6132 if_exists: bool,
6133 name: &QualifiedReplica,
6134) -> Result<Option<(ClusterId, ReplicaId)>, PlanError> {
6135 let cluster = resolve_cluster_replica(scx, name, if_exists)?;
6136 Ok(cluster.map(|(cluster, replica_id)| (cluster.id(), replica_id)))
6137}
6138
6139fn plan_drop_item(
6141 scx: &StatementContext,
6142 object_type: ObjectType,
6143 if_exists: bool,
6144 name: UnresolvedItemName,
6145 cascade: bool,
6146) -> Result<Option<CatalogItemId>, PlanError> {
6147 let Some(id) = plan_drop_item_name(scx, object_type, if_exists, name)? else {
6148 return Ok(None);
6149 };
6150 if !cascade {
6151 let catalog_item = scx.catalog.get_item(&id);
6152 ensure_no_blocking_dependents(scx, object_type, catalog_item, &BTreeSet::new())?;
6153 }
6154 Ok(Some(id))
6155}
6156
6157fn plan_drop_item_name(
6161 scx: &StatementContext,
6162 object_type: ObjectType,
6163 if_exists: bool,
6164 name: UnresolvedItemName,
6165) -> Result<Option<CatalogItemId>, PlanError> {
6166 let resolved = match resolve_item_or_type(scx, object_type, name, if_exists) {
6167 Ok(r) => r,
6168 Err(PlanError::MismatchedObjectType {
6170 name,
6171 is_type: ObjectType::MaterializedView,
6172 expected_type: ObjectType::View,
6173 }) => {
6174 return Err(PlanError::DropViewOnMaterializedView(name.to_string()));
6175 }
6176 e => e?,
6177 };
6178
6179 Ok(match resolved {
6180 Some(catalog_item) => {
6181 if catalog_item.id().is_system() {
6182 sql_bail!(
6183 "cannot drop {} {} because it is required by the database system",
6184 catalog_item.item_type(),
6185 scx.catalog.minimal_qualification(catalog_item.name()),
6186 );
6187 }
6188 Some(catalog_item.id())
6189 }
6190 None => None,
6191 })
6192}
6193
6194fn ensure_no_blocking_dependents(
6199 scx: &StatementContext,
6200 object_type: ObjectType,
6201 catalog_item: &dyn CatalogItem,
6202 also_dropped: &BTreeSet<CatalogItemId>,
6203) -> Result<(), PlanError> {
6204 for id in catalog_item.used_by() {
6205 if also_dropped.contains(id) {
6206 continue;
6207 }
6208 let dep = scx.catalog.get_item(id);
6209 if dependency_prevents_drop(object_type, dep) {
6210 return Err(PlanError::DependentObjectsStillExist {
6211 object_type: catalog_item.item_type().to_string(),
6212 object_name: scx
6213 .catalog
6214 .minimal_qualification(catalog_item.name())
6215 .to_string(),
6216 dependents: vec![(
6217 dep.item_type().to_string(),
6218 scx.catalog.minimal_qualification(dep.name()).to_string(),
6219 )],
6220 });
6221 }
6222 }
6223 Ok(())
6226}
6227
6228fn dependency_prevents_drop(object_type: ObjectType, dep: &dyn CatalogItem) -> bool {
6230 match object_type {
6231 ObjectType::Type => true,
6232 ObjectType::Table
6233 | ObjectType::View
6234 | ObjectType::MaterializedView
6235 | ObjectType::Source
6236 | ObjectType::Sink
6237 | ObjectType::MetricSink
6238 | ObjectType::Index
6239 | ObjectType::Role
6240 | ObjectType::Cluster
6241 | ObjectType::ClusterReplica
6242 | ObjectType::Secret
6243 | ObjectType::Connection
6244 | ObjectType::Database
6245 | ObjectType::Schema
6246 | ObjectType::Func
6247 | ObjectType::NetworkPolicy => match dep.item_type() {
6248 CatalogItemType::Func
6249 | CatalogItemType::Table
6250 | CatalogItemType::Source
6251 | CatalogItemType::View
6252 | CatalogItemType::MaterializedView
6253 | CatalogItemType::Sink
6254 | CatalogItemType::MetricSink
6255 | CatalogItemType::Type
6256 | CatalogItemType::Secret
6257 | CatalogItemType::Connection => true,
6258 CatalogItemType::Index => false,
6259 },
6260 }
6261}
6262
6263pub fn describe_alter_index_options(
6264 _: &StatementContext,
6265 _: AlterIndexStatement<Aug>,
6266) -> Result<StatementDesc, PlanError> {
6267 Ok(StatementDesc::new(None))
6268}
6269
6270pub fn describe_drop_owned(
6271 _: &StatementContext,
6272 _: DropOwnedStatement<Aug>,
6273) -> Result<StatementDesc, PlanError> {
6274 Ok(StatementDesc::new(None))
6275}
6276
6277pub fn plan_drop_owned(
6278 scx: &StatementContext,
6279 drop: DropOwnedStatement<Aug>,
6280) -> Result<Plan, PlanError> {
6281 let cascade = drop.cascade();
6282 let role_ids: BTreeSet<_> = drop.role_names.into_iter().map(|role| role.id).collect();
6283 let mut drop_ids = Vec::new();
6284 let mut privilege_revokes = Vec::new();
6285 let mut default_privilege_revokes = Vec::new();
6286
6287 fn update_privilege_revokes(
6288 object_id: SystemObjectId,
6289 privileges: &PrivilegeMap,
6290 role_ids: &BTreeSet<RoleId>,
6291 privilege_revokes: &mut Vec<(SystemObjectId, MzAclItem)>,
6292 ) {
6293 privilege_revokes.extend(iter::zip(
6294 iter::repeat(object_id),
6295 privileges
6296 .all_values()
6297 .filter(|privilege| role_ids.contains(&privilege.grantee))
6298 .cloned(),
6299 ));
6300 }
6301
6302 for replica in scx.catalog.get_cluster_replicas() {
6304 if role_ids.contains(&replica.owner_id()) {
6305 drop_ids.push((replica.cluster_id(), replica.replica_id()).into());
6306 }
6307 }
6308
6309 for cluster in scx.catalog.get_clusters() {
6311 if role_ids.contains(&cluster.owner_id()) {
6312 if !cascade {
6314 let non_owned_bound_objects: Vec<_> = cluster
6315 .bound_objects()
6316 .into_iter()
6317 .map(|item_id| scx.catalog.get_item(item_id))
6318 .filter(|item| !role_ids.contains(&item.owner_id()))
6319 .collect();
6320 if !non_owned_bound_objects.is_empty() {
6321 let names: Vec<_> = non_owned_bound_objects
6322 .into_iter()
6323 .map(|item| {
6324 (
6325 item.item_type().to_string(),
6326 scx.catalog.resolve_full_name(item.name()).to_string(),
6327 )
6328 })
6329 .collect();
6330 return Err(PlanError::DependentObjectsStillExist {
6331 object_type: "cluster".to_string(),
6332 object_name: cluster.name().to_string(),
6333 dependents: names,
6334 });
6335 }
6336 }
6337 drop_ids.push(cluster.id().into());
6338 }
6339 update_privilege_revokes(
6340 SystemObjectId::Object(cluster.id().into()),
6341 cluster.privileges(),
6342 &role_ids,
6343 &mut privilege_revokes,
6344 );
6345 }
6346
6347 for item in scx.catalog.get_items() {
6349 if role_ids.contains(&item.owner_id()) {
6350 if !cascade {
6351 let check_if_dependents_exist = |used_by: &[CatalogItemId]| {
6353 let non_owned_dependencies: Vec<_> = used_by
6354 .into_iter()
6355 .map(|item_id| scx.catalog.get_item(item_id))
6356 .filter(|item| dependency_prevents_drop(item.item_type().into(), *item))
6357 .filter(|item| !role_ids.contains(&item.owner_id()))
6358 .collect();
6359 if !non_owned_dependencies.is_empty() {
6360 let names: Vec<_> = non_owned_dependencies
6361 .into_iter()
6362 .map(|item| {
6363 let item_typ = item.item_type().to_string();
6364 let item_name =
6365 scx.catalog.resolve_full_name(item.name()).to_string();
6366 (item_typ, item_name)
6367 })
6368 .collect();
6369 Err(PlanError::DependentObjectsStillExist {
6370 object_type: item.item_type().to_string(),
6371 object_name: scx
6372 .catalog
6373 .resolve_full_name(item.name())
6374 .to_string()
6375 .to_string(),
6376 dependents: names,
6377 })
6378 } else {
6379 Ok(())
6380 }
6381 };
6382
6383 if let Some(id) = item.progress_id() {
6386 let progress_item = scx.catalog.get_item(&id);
6387 check_if_dependents_exist(progress_item.used_by())?;
6388 }
6389 check_if_dependents_exist(item.used_by())?;
6390 }
6391 drop_ids.push(item.id().into());
6392 }
6393 update_privilege_revokes(
6394 SystemObjectId::Object(item.id().into()),
6395 item.privileges(),
6396 &role_ids,
6397 &mut privilege_revokes,
6398 );
6399 }
6400
6401 for schema in scx.catalog.get_schemas() {
6403 if !schema.id().is_temporary() {
6404 if role_ids.contains(&schema.owner_id()) {
6405 if !cascade {
6406 let non_owned_dependencies: Vec<_> = schema
6407 .item_ids()
6408 .map(|item_id| scx.catalog.get_item(&item_id))
6409 .filter(|item| dependency_prevents_drop(item.item_type().into(), *item))
6410 .filter(|item| !role_ids.contains(&item.owner_id()))
6411 .collect();
6412 if !non_owned_dependencies.is_empty() {
6413 let full_schema_name = scx.catalog.resolve_full_schema_name(schema.name());
6414 sql_bail!(
6415 "schema {} cannot be dropped without CASCADE while it contains non-owned objects",
6416 full_schema_name.to_string().quoted()
6417 );
6418 }
6419 }
6420 drop_ids.push((*schema.database(), *schema.id()).into())
6421 }
6422 update_privilege_revokes(
6423 SystemObjectId::Object((*schema.database(), *schema.id()).into()),
6424 schema.privileges(),
6425 &role_ids,
6426 &mut privilege_revokes,
6427 );
6428 }
6429 }
6430
6431 for database in scx.catalog.get_databases() {
6433 if role_ids.contains(&database.owner_id()) {
6434 if !cascade {
6435 let non_owned_schemas: Vec<_> = database
6436 .schemas()
6437 .into_iter()
6438 .filter(|schema| !role_ids.contains(&schema.owner_id()))
6439 .collect();
6440 if !non_owned_schemas.is_empty() {
6441 sql_bail!(
6442 "database {} cannot be dropped without CASCADE while it contains non-owned schemas",
6443 database.name().quoted(),
6444 );
6445 }
6446 }
6447 drop_ids.push(database.id().into());
6448 }
6449 update_privilege_revokes(
6450 SystemObjectId::Object(database.id().into()),
6451 database.privileges(),
6452 &role_ids,
6453 &mut privilege_revokes,
6454 );
6455 }
6456
6457 for network_policy in scx.catalog.get_network_policies() {
6459 if role_ids.contains(&network_policy.owner_id()) {
6460 drop_ids.push(ObjectId::NetworkPolicy(network_policy.id()));
6461 }
6462 update_privilege_revokes(
6463 SystemObjectId::Object(ObjectId::NetworkPolicy(network_policy.id())),
6464 network_policy.privileges(),
6465 &role_ids,
6466 &mut privilege_revokes,
6467 );
6468 }
6469
6470 update_privilege_revokes(
6472 SystemObjectId::System,
6473 scx.catalog.get_system_privileges(),
6474 &role_ids,
6475 &mut privilege_revokes,
6476 );
6477
6478 for (default_privilege_object, default_privilege_acl_items) in
6479 scx.catalog.get_default_privileges()
6480 {
6481 for default_privilege_acl_item in default_privilege_acl_items {
6482 if role_ids.contains(&default_privilege_object.role_id)
6483 || role_ids.contains(&default_privilege_acl_item.grantee)
6484 {
6485 default_privilege_revokes.push((
6486 default_privilege_object.clone(),
6487 default_privilege_acl_item.clone(),
6488 ));
6489 }
6490 }
6491 }
6492
6493 let drop_ids = scx.catalog.object_dependents(&drop_ids);
6494
6495 let system_ids: Vec<_> = drop_ids.iter().filter(|id| id.is_system()).collect();
6496 if !system_ids.is_empty() {
6497 let mut owners = system_ids
6498 .into_iter()
6499 .filter_map(|object_id| scx.catalog.get_owner_id(object_id))
6500 .collect::<BTreeSet<_>>()
6501 .into_iter()
6502 .map(|role_id| scx.catalog.get_role(&role_id).name().quoted());
6503 sql_bail!(
6504 "cannot drop objects owned by role {} because they are required by the database system",
6505 owners.join(", "),
6506 );
6507 }
6508
6509 Ok(Plan::DropOwned(DropOwnedPlan {
6510 role_ids: role_ids.into_iter().collect(),
6511 drop_ids,
6512 privilege_revokes,
6513 default_privilege_revokes,
6514 }))
6515}
6516
6517fn plan_retain_history_option(
6518 scx: &StatementContext,
6519 retain_history: Option<OptionalDuration>,
6520) -> Result<Option<CompactionWindow>, PlanError> {
6521 if let Some(OptionalDuration(lcw)) = retain_history {
6522 Ok(Some(plan_retain_history(scx, lcw)?))
6523 } else {
6524 Ok(None)
6525 }
6526}
6527
6528fn plan_retain_history(
6534 scx: &StatementContext,
6535 lcw: Option<Duration>,
6536) -> Result<CompactionWindow, PlanError> {
6537 scx.require_feature_flag(&vars::ENABLE_LOGICAL_COMPACTION_WINDOW)?;
6538 match lcw {
6539 Some(Duration::ZERO) => Err(PlanError::InvalidOptionValue {
6544 option_name: "RETAIN HISTORY".to_string(),
6545 err: Box::new(PlanError::Unstructured(
6546 "internal error: unexpectedly zero".to_string(),
6547 )),
6548 }),
6549 Some(duration) => {
6550 if duration < DEFAULT_LOGICAL_COMPACTION_WINDOW_DURATION
6553 && scx
6554 .require_feature_flag(&vars::ENABLE_UNLIMITED_RETAIN_HISTORY)
6555 .is_err()
6556 {
6557 return Err(PlanError::RetainHistoryLow {
6558 limit: DEFAULT_LOGICAL_COMPACTION_WINDOW_DURATION,
6559 });
6560 }
6561 Ok(duration.try_into()?)
6562 }
6563 None => {
6566 if scx
6567 .require_feature_flag(&vars::ENABLE_UNLIMITED_RETAIN_HISTORY)
6568 .is_err()
6569 {
6570 Err(PlanError::RetainHistoryRequired)
6571 } else {
6572 Ok(CompactionWindow::DisableCompaction)
6573 }
6574 }
6575 }
6576}
6577
6578generate_extracted_config!(IndexOption, (RetainHistory, OptionalDuration));
6579
6580fn plan_index_options(
6581 scx: &StatementContext,
6582 with_opts: Vec<IndexOption<Aug>>,
6583) -> Result<Vec<crate::plan::IndexOption>, PlanError> {
6584 if !with_opts.is_empty() {
6585 scx.require_feature_flag(&vars::ENABLE_INDEX_OPTIONS)?;
6587 }
6588
6589 let IndexOptionExtracted { retain_history, .. }: IndexOptionExtracted = with_opts.try_into()?;
6590
6591 let mut out = Vec::with_capacity(1);
6592 if let Some(cw) = plan_retain_history_option(scx, retain_history)? {
6593 out.push(crate::plan::IndexOption::RetainHistory(cw));
6594 }
6595 Ok(out)
6596}
6597
6598generate_extracted_config!(
6599 TableOption,
6600 (PartitionBy, Vec<Ident>),
6601 (RetainHistory, OptionalDuration),
6602 (RedactedTest, String)
6603);
6604
6605fn plan_table_options(
6606 scx: &StatementContext,
6607 desc: &RelationDesc,
6608 with_opts: Vec<TableOption<Aug>>,
6609) -> Result<Vec<crate::plan::TableOption>, PlanError> {
6610 let TableOptionExtracted {
6611 partition_by,
6612 retain_history,
6613 redacted_test,
6614 ..
6615 }: TableOptionExtracted = with_opts.try_into()?;
6616
6617 if let Some(partition_by) = partition_by {
6618 scx.require_feature_flag(&ENABLE_COLLECTION_PARTITION_BY)?;
6619 check_partition_by(desc, partition_by)?;
6620 }
6621
6622 if redacted_test.is_some() {
6623 scx.require_feature_flag(&vars::ENABLE_REDACTED_TEST_OPTION)?;
6624 }
6625
6626 let mut out = Vec::with_capacity(1);
6627 if let Some(cw) = plan_retain_history_option(scx, retain_history)? {
6628 out.push(crate::plan::TableOption::RetainHistory(cw));
6629 }
6630 Ok(out)
6631}
6632
6633pub fn plan_alter_index_options(
6634 scx: &mut StatementContext,
6635 AlterIndexStatement {
6636 index_name,
6637 if_exists,
6638 action,
6639 }: AlterIndexStatement<Aug>,
6640) -> Result<Plan, PlanError> {
6641 let object_type = ObjectType::Index;
6642 match action {
6643 AlterIndexAction::ResetOptions(options) => {
6644 let mut options = options.into_iter();
6645 if let Some(opt) = options.next() {
6646 match opt {
6647 IndexOptionName::RetainHistory => {
6648 if options.next().is_some() {
6649 sql_bail!("RETAIN HISTORY must be only option");
6650 }
6651 return alter_retain_history(
6652 scx,
6653 object_type,
6654 if_exists,
6655 UnresolvedObjectName::Item(index_name),
6656 None,
6657 );
6658 }
6659 }
6660 }
6661 sql_bail!("expected option");
6662 }
6663 AlterIndexAction::SetOptions(options) => {
6664 let mut options = options.into_iter();
6665 if let Some(opt) = options.next() {
6666 match opt.name {
6667 IndexOptionName::RetainHistory => {
6668 if options.next().is_some() {
6669 sql_bail!("RETAIN HISTORY must be only option");
6670 }
6671 return alter_retain_history(
6672 scx,
6673 object_type,
6674 if_exists,
6675 UnresolvedObjectName::Item(index_name),
6676 opt.value,
6677 );
6678 }
6679 }
6680 }
6681 sql_bail!("expected option");
6682 }
6683 }
6684}
6685
6686pub fn describe_alter_cluster_set_options(
6687 _: &StatementContext,
6688 _: AlterClusterStatement<Aug>,
6689) -> Result<StatementDesc, PlanError> {
6690 Ok(StatementDesc::new(None))
6691}
6692
6693pub fn plan_alter_cluster(
6694 scx: &mut StatementContext,
6695 AlterClusterStatement {
6696 name,
6697 action,
6698 if_exists,
6699 }: AlterClusterStatement<Aug>,
6700) -> Result<Plan, PlanError> {
6701 let cluster = match resolve_cluster(scx, &name, if_exists)? {
6702 Some(entry) => entry,
6703 None => {
6704 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
6705 name: name.to_ast_string_simple(),
6706 object_type: ObjectType::Cluster,
6707 });
6708
6709 return Ok(Plan::AlterNoop(AlterNoopPlan {
6710 object_type: ObjectType::Cluster,
6711 }));
6712 }
6713 };
6714
6715 let mut options: PlanClusterOption = Default::default();
6716 let mut alter_strategy: AlterClusterPlanStrategy = AlterClusterPlanStrategy::None;
6717
6718 match action {
6719 AlterClusterAction::SetOptions {
6720 options: set_options,
6721 with_options,
6722 } => {
6723 let ClusterOptionExtracted {
6724 auto_scaling_strategy,
6725 availability_zones,
6726 experimental_arrangement_compression,
6727 introspection_debugging,
6728 introspection_interval,
6729 managed,
6730 replicas: replica_defs,
6731 replication_factor,
6732 seen: _,
6733 size,
6734 disk,
6735 schedule,
6736 workload_class,
6737 }: ClusterOptionExtracted = set_options.try_into()?;
6738
6739 if !scx.catalog.active_role_id().is_system() {
6740 if workload_class.is_some() {
6741 sql_bail!("WORKLOAD CLASS not supported for non-system users");
6742 }
6743 }
6744
6745 match managed.unwrap_or_else(|| cluster.is_managed()) {
6746 true => {
6747 let alter_strategy_extracted =
6748 ClusterAlterOptionExtracted::try_from(with_options)?;
6749 alter_strategy = AlterClusterPlanStrategy::try_from(alter_strategy_extracted)?;
6750
6751 if !matches!(alter_strategy, AlterClusterPlanStrategy::None)
6755 && size.is_none()
6756 && availability_zones.is_none()
6757 && introspection_debugging.is_none()
6758 && introspection_interval.is_none()
6759 && experimental_arrangement_compression.is_none()
6760 {
6761 sql_bail!(
6762 "WAIT can only be used together with a SIZE, AVAILABILITY ZONES, \
6763 INTROSPECTION, or EXPERIMENTAL ARRANGEMENT COMPRESSION change"
6764 );
6765 }
6766
6767 if replica_defs.is_some() {
6768 sql_bail!("REPLICAS not supported for managed clusters");
6769 }
6770 if schedule.is_some()
6771 && !matches!(schedule, Some(ClusterScheduleOptionValue::Manual))
6772 {
6773 scx.require_feature_flag(&ENABLE_CLUSTER_SCHEDULE_REFRESH)?;
6774
6775 if replication_factor.is_none()
6784 && cluster.replication_factor().is_some_and(|rf| rf > 1)
6785 {
6786 sql_bail!(
6787 "SCHEDULE cannot be set to anything other than MANUAL while the \
6788 cluster's REPLICATION FACTOR is greater than 1; \
6789 set the REPLICATION FACTOR to 1 first"
6790 );
6791 }
6792 }
6793
6794 if replication_factor.is_some() {
6795 if schedule.is_some()
6796 && !matches!(schedule, Some(ClusterScheduleOptionValue::Manual))
6797 {
6798 sql_bail!(
6799 "REPLICATION FACTOR cannot be given together with any SCHEDULE other than MANUAL"
6800 );
6801 }
6802 if let Some(current_schedule) = cluster.schedule() {
6803 if !matches!(current_schedule, ClusterSchedule::Manual) {
6804 sql_bail!(
6805 "REPLICATION FACTOR cannot be set if the cluster SCHEDULE is anything other than MANUAL"
6806 );
6807 }
6808 }
6809 }
6810
6811 if let Some(value) = auto_scaling_strategy {
6812 scx.require_feature_flag(&ENABLE_AUTO_SCALING_STRATEGY)?;
6813 let strategy = plan_auto_scaling_strategy(value)?;
6814 options.auto_scaling_strategy = AlterOptionParameter::Set(strategy);
6815 }
6816
6817 let effective_strategy = match &options.auto_scaling_strategy {
6821 AlterOptionParameter::Set(s) => s.clone(),
6822 AlterOptionParameter::Reset => None,
6823 AlterOptionParameter::Unchanged => cluster.auto_scaling_strategy().cloned(),
6824 };
6825 if let Some(effective_strategy) = &effective_strategy {
6826 let effective_size = size.as_deref().or_else(|| cluster.managed_size());
6827 let schedule_non_manual = match &schedule {
6828 Some(s) => !matches!(s, ClusterScheduleOptionValue::Manual),
6829 None => cluster
6830 .schedule()
6831 .is_some_and(|s| !matches!(s, ClusterSchedule::Manual)),
6832 };
6833 validate_auto_scaling_strategy(
6834 effective_strategy,
6835 effective_size,
6836 schedule_non_manual,
6837 )?;
6838 }
6839 }
6840 false => {
6841 if !with_options.is_empty() {
6842 sql_bail!("ALTER... WITH not supported for unmanaged clusters");
6843 }
6844 if auto_scaling_strategy.is_some() {
6845 sql_bail!("AUTO SCALING STRATEGY not supported for unmanaged clusters");
6846 }
6847 if availability_zones.is_some() {
6848 sql_bail!("AVAILABILITY ZONES not supported for unmanaged clusters");
6849 }
6850 if replication_factor.is_some() {
6851 sql_bail!("REPLICATION FACTOR not supported for unmanaged clusters");
6852 }
6853 if introspection_debugging.is_some() {
6854 sql_bail!("INTROSPECTION DEBUGGING not supported for unmanaged clusters");
6855 }
6856 if introspection_interval.is_some() {
6857 sql_bail!("INTROSPECTION INTERVAL not supported for unmanaged clusters");
6858 }
6859 if experimental_arrangement_compression.is_some() {
6860 sql_bail!(
6861 "EXPERIMENTAL ARRANGEMENT COMPRESSION not supported for unmanaged clusters"
6862 );
6863 }
6864 if size.is_some() {
6865 sql_bail!("SIZE not supported for unmanaged clusters");
6866 }
6867 if disk.is_some() {
6868 sql_bail!("DISK not supported for unmanaged clusters");
6869 }
6870 if schedule.is_some()
6871 && !matches!(schedule, Some(ClusterScheduleOptionValue::Manual))
6872 {
6873 sql_bail!(
6874 "cluster schedules other than MANUAL are not supported for unmanaged clusters"
6875 );
6876 }
6877 if let Some(current_schedule) = cluster.schedule() {
6878 if !matches!(current_schedule, ClusterSchedule::Manual)
6879 && schedule.is_none()
6880 {
6881 sql_bail!(
6882 "when switching a cluster to unmanaged, if the managed \
6883 cluster's SCHEDULE is anything other than MANUAL, you have to \
6884 explicitly set the SCHEDULE to MANUAL"
6885 );
6886 }
6887 }
6888 }
6889 }
6890
6891 let mut replicas = vec![];
6892 for ReplicaDefinition { name, options } in
6893 replica_defs.into_iter().flat_map(Vec::into_iter)
6894 {
6895 replicas.push((normalize::ident(name), plan_replica_config(scx, options)?));
6896 }
6897
6898 if let Some(managed) = managed {
6899 options.managed = AlterOptionParameter::Set(managed);
6900 }
6901 if let Some(replication_factor) = replication_factor {
6902 options.replication_factor = AlterOptionParameter::Set(replication_factor);
6903 } else if schedule
6904 .as_ref()
6905 .is_some_and(|s| !matches!(s, ClusterScheduleOptionValue::Manual))
6906 && managed != Some(true)
6907 {
6908 options.replication_factor = AlterOptionParameter::Set(0);
6921 }
6922 if let Some(size) = &size {
6923 options.size = AlterOptionParameter::Set(size.clone());
6924 }
6925 if let Some(availability_zones) = availability_zones {
6926 options.availability_zones = AlterOptionParameter::Set(availability_zones);
6927 }
6928 if let Some(introspection_debugging) = introspection_debugging {
6929 options.introspection_debugging =
6930 AlterOptionParameter::Set(introspection_debugging);
6931 }
6932 if let Some(introspection_interval) = introspection_interval {
6933 options.introspection_interval = AlterOptionParameter::Set(introspection_interval);
6934 }
6935 if let Some(experimental_arrangement_compression) = experimental_arrangement_compression
6936 {
6937 options.arrangement_compression =
6938 AlterOptionParameter::Set(experimental_arrangement_compression);
6939 }
6940 if disk.is_some() {
6941 let size = match size.as_deref() {
6945 Some(s) => s,
6946 None => cluster
6947 .managed_size()
6948 .ok_or_else(|| sql_err!("cluster is not managed"))?,
6949 };
6950 if scx.catalog.is_cluster_size_cc(size) {
6951 sql_bail!(
6952 "DISK option not supported for modern cluster sizes because disk is always enabled"
6953 );
6954 }
6955
6956 scx.catalog
6957 .add_notice(PlanNotice::ReplicaDiskOptionDeprecated);
6958 }
6959 if !replicas.is_empty() {
6960 options.replicas = AlterOptionParameter::Set(replicas);
6961 }
6962 if let Some(schedule) = schedule {
6963 options.schedule = AlterOptionParameter::Set(plan_cluster_schedule(schedule)?);
6964 }
6965 if let Some(workload_class) = workload_class {
6966 options.workload_class = AlterOptionParameter::Set(workload_class.0);
6967 }
6968 }
6969 AlterClusterAction::ResetOptions(reset_options) => {
6970 use AlterOptionParameter::Reset;
6971 use ClusterOptionName::*;
6972
6973 if !scx.catalog.active_role_id().is_system() {
6974 if reset_options.contains(&WorkloadClass) {
6975 sql_bail!("WORKLOAD CLASS not supported for non-system users");
6976 }
6977 }
6978
6979 for option in reset_options {
6985 match option {
6986 AutoScalingStrategy => options.auto_scaling_strategy = Reset,
6987 AvailabilityZones => options.availability_zones = Reset,
6988 Disk => scx
6989 .catalog
6990 .add_notice(PlanNotice::ReplicaDiskOptionDeprecated),
6991 IntrospectionInterval => options.introspection_interval = Reset,
6992 IntrospectionDebugging => options.introspection_debugging = Reset,
6993 ExperimentalArrangementCompression => options.arrangement_compression = Reset,
6994 Managed => options.managed = Reset,
6995 Replicas => options.replicas = Reset,
6996 ReplicationFactor => options.replication_factor = Reset,
6997 Size => options.size = Reset,
6998 Schedule => options.schedule = Reset,
6999 WorkloadClass => options.workload_class = Reset,
7000 }
7001 }
7002 }
7003 }
7004 Ok(Plan::AlterCluster(AlterClusterPlan {
7005 id: cluster.id(),
7006 name: cluster.name().to_string(),
7007 options,
7008 strategy: alter_strategy,
7009 }))
7010}
7011
7012pub fn describe_alter_set_cluster(
7013 _: &StatementContext,
7014 _: AlterSetClusterStatement<Aug>,
7015) -> Result<StatementDesc, PlanError> {
7016 Ok(StatementDesc::new(None))
7017}
7018
7019pub fn plan_alter_item_set_cluster(
7020 scx: &StatementContext,
7021 AlterSetClusterStatement {
7022 if_exists,
7023 set_cluster: in_cluster_name,
7024 name,
7025 object_type,
7026 }: AlterSetClusterStatement<Aug>,
7027) -> Result<Plan, PlanError> {
7028 scx.require_feature_flag(&vars::ENABLE_ALTER_SET_CLUSTER)?;
7029
7030 let object_type = object_type.into();
7031
7032 match object_type {
7034 ObjectType::MaterializedView => {}
7035 ObjectType::Index | ObjectType::Sink | ObjectType::MetricSink | ObjectType::Source => {
7036 bail_unsupported!(29606, format!("ALTER {object_type} SET CLUSTER"))
7037 }
7038 ObjectType::Table
7039 | ObjectType::View
7040 | ObjectType::Type
7041 | ObjectType::Role
7042 | ObjectType::Cluster
7043 | ObjectType::ClusterReplica
7044 | ObjectType::Secret
7045 | ObjectType::Connection
7046 | ObjectType::Database
7047 | ObjectType::Schema
7048 | ObjectType::Func
7049 | ObjectType::NetworkPolicy => {
7050 bail_never_supported!(
7051 format!("ALTER {object_type} SET CLUSTER"),
7052 "sql/alter-set-cluster/",
7053 format!("{object_type} has no associated cluster")
7054 )
7055 }
7056 }
7057
7058 let in_cluster = scx.catalog.get_cluster(in_cluster_name.id);
7059
7060 match resolve_item_or_type(scx, object_type, name.clone(), if_exists)? {
7061 Some(entry) => {
7062 let current_cluster = entry.cluster_id();
7063 let Some(current_cluster) = current_cluster else {
7064 sql_bail!("No cluster associated with {name}");
7065 };
7066
7067 if current_cluster == in_cluster.id() {
7068 Ok(Plan::AlterNoop(AlterNoopPlan { object_type }))
7069 } else {
7070 Ok(Plan::AlterSetCluster(AlterSetClusterPlan {
7071 id: entry.id(),
7072 set_cluster: in_cluster.id(),
7073 }))
7074 }
7075 }
7076 None => {
7077 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7078 name: name.to_ast_string_simple(),
7079 object_type,
7080 });
7081
7082 Ok(Plan::AlterNoop(AlterNoopPlan { object_type }))
7083 }
7084 }
7085}
7086
7087pub fn describe_alter_object_rename(
7088 _: &StatementContext,
7089 _: AlterObjectRenameStatement,
7090) -> Result<StatementDesc, PlanError> {
7091 Ok(StatementDesc::new(None))
7092}
7093
7094pub fn plan_alter_object_rename(
7095 scx: &mut StatementContext,
7096 AlterObjectRenameStatement {
7097 name,
7098 object_type,
7099 to_item_name,
7100 if_exists,
7101 }: AlterObjectRenameStatement,
7102) -> Result<Plan, PlanError> {
7103 let object_type = object_type.into();
7104 match (object_type, name) {
7105 (
7106 ObjectType::View
7107 | ObjectType::MaterializedView
7108 | ObjectType::Table
7109 | ObjectType::Source
7110 | ObjectType::Index
7111 | ObjectType::Sink
7112 | ObjectType::Secret
7113 | ObjectType::Connection,
7114 UnresolvedObjectName::Item(name),
7115 ) => plan_alter_item_rename(scx, object_type, name, to_item_name, if_exists),
7116 (ObjectType::Cluster, UnresolvedObjectName::Cluster(name)) => {
7117 plan_alter_cluster_rename(scx, object_type, name, to_item_name, if_exists)
7118 }
7119 (ObjectType::ClusterReplica, UnresolvedObjectName::ClusterReplica(name)) => {
7120 plan_alter_cluster_replica_rename(scx, object_type, name, to_item_name, if_exists)
7121 }
7122 (ObjectType::Schema, UnresolvedObjectName::Schema(name)) => {
7123 plan_alter_schema_rename(scx, name, to_item_name, if_exists)
7124 }
7125 (object_type, name) => {
7126 bail_internal!("invalid object type '{object_type}' for ALTER RENAME with name {name}")
7129 }
7130 }
7131}
7132
7133pub fn plan_alter_schema_rename(
7134 scx: &mut StatementContext,
7135 name: UnresolvedSchemaName,
7136 to_schema_name: Ident,
7137 if_exists: bool,
7138) -> Result<Plan, PlanError> {
7139 let normalized = normalize::unresolved_schema_name(name.clone())?;
7143 if normalized.database.is_none() && normalized.schema == mz_repr::namespaces::MZ_TEMP_SCHEMA {
7144 sql_bail!(
7145 "cannot rename schemas in the ambient database: {:?}",
7146 mz_repr::namespaces::MZ_TEMP_SCHEMA
7147 );
7148 }
7149
7150 let Some((db_spec, schema_spec)) = resolve_schema(scx, name.clone(), if_exists)? else {
7151 let object_type = ObjectType::Schema;
7152 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7153 name: name.to_ast_string_simple(),
7154 object_type,
7155 });
7156 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
7157 };
7158
7159 if scx
7161 .resolve_schema_in_database(&db_spec, &to_schema_name)
7162 .is_ok()
7163 {
7164 return Err(PlanError::Catalog(CatalogError::SchemaAlreadyExists(
7165 to_schema_name.clone().into_string(),
7166 )));
7167 }
7168
7169 let schema = scx.catalog.get_schema(&db_spec, &schema_spec);
7171 if schema.id().is_system() {
7172 bail_never_supported!(format!("renaming the {} schema", schema.name().schema))
7173 }
7174
7175 Ok(Plan::AlterSchemaRename(AlterSchemaRenamePlan {
7176 cur_schema_spec: (db_spec, schema_spec),
7177 new_schema_name: to_schema_name.into_string(),
7178 }))
7179}
7180
7181pub fn plan_alter_schema_swap<F>(
7182 scx: &mut StatementContext,
7183 name_a: UnresolvedSchemaName,
7184 name_b: Ident,
7185 if_exists: bool,
7186 gen_temp_suffix: F,
7187) -> Result<Plan, PlanError>
7188where
7189 F: Fn(&dyn Fn(&str) -> bool) -> Result<String, PlanError>,
7190{
7191 let normalized_a = normalize::unresolved_schema_name(name_a.clone())?;
7195 if normalized_a.database.is_none() && normalized_a.schema == mz_repr::namespaces::MZ_TEMP_SCHEMA
7196 {
7197 sql_bail!("cannot swap schemas that are in the ambient database");
7198 }
7199 let name_b_str = normalize::ident_ref(&name_b);
7201 if name_b_str == mz_repr::namespaces::MZ_TEMP_SCHEMA {
7202 sql_bail!("cannot swap schemas that are in the ambient database");
7203 }
7204
7205 let schema_a = match scx.resolve_schema(name_a.clone()) {
7206 Ok(schema) => schema,
7207 Err(_) if if_exists => {
7208 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7209 name: name_a.to_ast_string_simple(),
7210 object_type: ObjectType::Schema,
7211 });
7212 return Ok(Plan::AlterNoop(AlterNoopPlan {
7213 object_type: ObjectType::Schema,
7214 }));
7215 }
7216 Err(e) => return Err(e),
7217 };
7218
7219 let db_spec = schema_a.database().clone();
7220 if matches!(db_spec, ResolvedDatabaseSpecifier::Ambient) {
7221 sql_bail!("cannot swap schemas that are in the ambient database");
7222 };
7223 let schema_b = scx.resolve_schema_in_database(&db_spec, &name_b)?;
7224
7225 if schema_a.id().is_system() || schema_b.id().is_system() {
7227 bail_never_supported!("swapping a system schema".to_string())
7228 }
7229
7230 const SCHEMA_SWAP_PREFIX: &str = "mz_schema_swap_";
7234 let check = |temp_suffix: &str| {
7235 let mut temp_name = ident!(SCHEMA_SWAP_PREFIX);
7236 temp_name.append_lossy(temp_suffix);
7237 scx.resolve_schema_in_database(&db_spec, &temp_name)
7238 .is_err()
7239 };
7240 let temp_suffix = gen_temp_suffix(&check)?;
7241 let name_temp = format!("{SCHEMA_SWAP_PREFIX}{temp_suffix}");
7242
7243 Ok(Plan::AlterSchemaSwap(AlterSchemaSwapPlan {
7244 schema_a_spec: (*schema_a.database(), *schema_a.id()),
7245 schema_a_name: schema_a.name().schema.to_string(),
7246 schema_b_spec: (*schema_b.database(), *schema_b.id()),
7247 schema_b_name: schema_b.name().schema.to_string(),
7248 name_temp,
7249 }))
7250}
7251
7252pub fn plan_alter_item_rename(
7253 scx: &mut StatementContext,
7254 object_type: ObjectType,
7255 name: UnresolvedItemName,
7256 to_item_name: Ident,
7257 if_exists: bool,
7258) -> Result<Plan, PlanError> {
7259 let resolved = match resolve_item_or_type(scx, object_type, name.clone(), if_exists) {
7260 Ok(r) => r,
7261 Err(PlanError::MismatchedObjectType {
7263 name,
7264 is_type: ObjectType::MaterializedView,
7265 expected_type: ObjectType::View,
7266 }) => {
7267 return Err(PlanError::AlterViewOnMaterializedView(name.to_string()));
7268 }
7269 e => e?,
7270 };
7271
7272 match resolved {
7273 Some(entry) => {
7274 let full_name = scx.catalog.resolve_full_name(entry.name());
7275 let item_type = entry.item_type();
7276
7277 let proposed_name = QualifiedItemName {
7278 qualifiers: entry.name().qualifiers.clone(),
7279 item: to_item_name.clone().into_string(),
7280 };
7281
7282 let conflicting_type_exists;
7286 let conflicting_item_exists;
7287 if item_type == CatalogItemType::Type {
7288 conflicting_type_exists = scx.catalog.get_type_by_name(&proposed_name).is_some();
7289 conflicting_item_exists = scx
7290 .catalog
7291 .get_item_by_name(&proposed_name)
7292 .map(|item| item.item_type().conflicts_with_type())
7293 .unwrap_or(false);
7294 } else {
7295 conflicting_type_exists = item_type.conflicts_with_type()
7296 && scx.catalog.get_type_by_name(&proposed_name).is_some();
7297 conflicting_item_exists = scx.catalog.get_item_by_name(&proposed_name).is_some();
7298 };
7299 if conflicting_type_exists || conflicting_item_exists {
7300 sql_bail!("catalog item '{}' already exists", to_item_name);
7301 }
7302
7303 Ok(Plan::AlterItemRename(AlterItemRenamePlan {
7304 id: entry.id(),
7305 current_full_name: full_name,
7306 to_name: normalize::ident(to_item_name),
7307 object_type,
7308 }))
7309 }
7310 None => {
7311 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7312 name: name.to_ast_string_simple(),
7313 object_type,
7314 });
7315
7316 Ok(Plan::AlterNoop(AlterNoopPlan { object_type }))
7317 }
7318 }
7319}
7320
7321pub fn plan_alter_cluster_rename(
7322 scx: &mut StatementContext,
7323 object_type: ObjectType,
7324 name: Ident,
7325 to_name: Ident,
7326 if_exists: bool,
7327) -> Result<Plan, PlanError> {
7328 match resolve_cluster(scx, &name, if_exists)? {
7329 Some(entry) => Ok(Plan::AlterClusterRename(AlterClusterRenamePlan {
7330 id: entry.id(),
7331 name: entry.name().to_string(),
7332 to_name: ident(to_name),
7333 })),
7334 None => {
7335 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7336 name: name.to_ast_string_simple(),
7337 object_type,
7338 });
7339
7340 Ok(Plan::AlterNoop(AlterNoopPlan { object_type }))
7341 }
7342 }
7343}
7344
7345pub fn plan_alter_cluster_swap<F>(
7346 scx: &mut StatementContext,
7347 name_a: Ident,
7348 name_b: Ident,
7349 if_exists: bool,
7350 gen_temp_suffix: F,
7351) -> Result<Plan, PlanError>
7352where
7353 F: Fn(&dyn Fn(&str) -> bool) -> Result<String, PlanError>,
7354{
7355 let cluster_a = match scx.resolve_cluster(Some(&name_a)) {
7356 Ok(cluster) => cluster,
7357 Err(_) if if_exists => {
7358 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7359 name: name_a.to_ast_string_simple(),
7360 object_type: ObjectType::Cluster,
7361 });
7362 return Ok(Plan::AlterNoop(AlterNoopPlan {
7363 object_type: ObjectType::Cluster,
7364 }));
7365 }
7366 Err(e) => return Err(e),
7367 };
7368 let cluster_b = scx.resolve_cluster(Some(&name_b))?;
7369
7370 const CLUSTER_SWAP_PREFIX: &str = "mz_cluster_swap_";
7371 let check = |temp_suffix: &str| {
7372 let mut temp_name = ident!(CLUSTER_SWAP_PREFIX);
7373 temp_name.append_lossy(temp_suffix);
7374 match scx.catalog.resolve_cluster(Some(temp_name.as_str())) {
7375 Err(CatalogError::UnknownCluster(_)) => true,
7377 Ok(_) | Err(_) => false,
7379 }
7380 };
7381 let temp_suffix = gen_temp_suffix(&check)?;
7382 let name_temp = format!("{CLUSTER_SWAP_PREFIX}{temp_suffix}");
7383
7384 Ok(Plan::AlterClusterSwap(AlterClusterSwapPlan {
7385 id_a: cluster_a.id(),
7386 id_b: cluster_b.id(),
7387 name_a: name_a.into_string(),
7388 name_b: name_b.into_string(),
7389 name_temp,
7390 }))
7391}
7392
7393pub fn plan_alter_cluster_replica_rename(
7394 scx: &mut StatementContext,
7395 object_type: ObjectType,
7396 name: QualifiedReplica,
7397 to_item_name: Ident,
7398 if_exists: bool,
7399) -> Result<Plan, PlanError> {
7400 match resolve_cluster_replica(scx, &name, if_exists)? {
7401 Some((cluster, replica)) => {
7402 ensure_cluster_is_not_managed(scx, cluster.id())?;
7403 Ok(Plan::AlterClusterReplicaRename(
7404 AlterClusterReplicaRenamePlan {
7405 cluster_id: cluster.id(),
7406 replica_id: replica,
7407 name: QualifiedReplica {
7408 cluster: Ident::new(cluster.name())?,
7409 replica: name.replica,
7410 },
7411 to_name: normalize::ident(to_item_name),
7412 },
7413 ))
7414 }
7415 None => {
7416 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7417 name: name.to_ast_string_simple(),
7418 object_type,
7419 });
7420
7421 Ok(Plan::AlterNoop(AlterNoopPlan { object_type }))
7422 }
7423 }
7424}
7425
7426pub fn describe_alter_object_swap(
7427 _: &StatementContext,
7428 _: AlterObjectSwapStatement,
7429) -> Result<StatementDesc, PlanError> {
7430 Ok(StatementDesc::new(None))
7431}
7432
7433pub fn plan_alter_object_swap(
7434 scx: &mut StatementContext,
7435 stmt: AlterObjectSwapStatement,
7436) -> Result<Plan, PlanError> {
7437 scx.require_feature_flag(&vars::ENABLE_ALTER_SWAP)?;
7438
7439 let AlterObjectSwapStatement {
7440 object_type,
7441 if_exists,
7442 name_a,
7443 name_b,
7444 } = stmt;
7445 let object_type = object_type.into();
7446
7447 let gen_temp_suffix = |check_fn: &dyn Fn(&str) -> bool| {
7449 let mut attempts = 0;
7450 let name_temp = loop {
7451 attempts += 1;
7452 if attempts > 10 {
7453 tracing::warn!("Unable to generate temp id for swapping");
7454 sql_bail!("unable to swap!");
7455 }
7456
7457 let short_id = mz_ore::id_gen::temp_id();
7459 if check_fn(&short_id) {
7460 break short_id;
7461 }
7462 };
7463
7464 Ok(name_temp)
7465 };
7466
7467 match (object_type, name_a, name_b) {
7468 (ObjectType::Schema, UnresolvedObjectName::Schema(name_a), name_b) => {
7469 plan_alter_schema_swap(scx, name_a, name_b, if_exists, gen_temp_suffix)
7470 }
7471 (ObjectType::Cluster, UnresolvedObjectName::Cluster(name_a), name_b) => {
7472 plan_alter_cluster_swap(scx, name_a, name_b, if_exists, gen_temp_suffix)
7473 }
7474 (ObjectType::Schema | ObjectType::Cluster, _, _) => {
7475 bail_internal!("name type does not match object type for ALTER SWAP")
7476 }
7477 (
7478 ObjectType::Table
7479 | ObjectType::View
7480 | ObjectType::MaterializedView
7481 | ObjectType::Source
7482 | ObjectType::Sink
7483 | ObjectType::MetricSink
7484 | ObjectType::Index
7485 | ObjectType::Type
7486 | ObjectType::Role
7487 | ObjectType::ClusterReplica
7488 | ObjectType::Secret
7489 | ObjectType::Connection
7490 | ObjectType::Database
7491 | ObjectType::Func
7492 | ObjectType::NetworkPolicy,
7493 _,
7494 _,
7495 ) => Err(PlanError::Unsupported {
7496 feature: format!("ALTER {object_type} .. SWAP WITH ..."),
7497 discussion_no: None,
7498 }),
7499 }
7500}
7501
7502pub fn describe_alter_retain_history(
7503 _: &StatementContext,
7504 _: AlterRetainHistoryStatement<Aug>,
7505) -> Result<StatementDesc, PlanError> {
7506 Ok(StatementDesc::new(None))
7507}
7508
7509pub fn plan_alter_retain_history(
7510 scx: &StatementContext,
7511 AlterRetainHistoryStatement {
7512 object_type,
7513 if_exists,
7514 name,
7515 history,
7516 }: AlterRetainHistoryStatement<Aug>,
7517) -> Result<Plan, PlanError> {
7518 alter_retain_history(scx, object_type.into(), if_exists, name, history)
7519}
7520
7521fn alter_retain_history(
7522 scx: &StatementContext,
7523 object_type: ObjectType,
7524 if_exists: bool,
7525 name: UnresolvedObjectName,
7526 history: Option<WithOptionValue<Aug>>,
7527) -> Result<Plan, PlanError> {
7528 let name = match (object_type, name) {
7529 (
7530 ObjectType::View
7532 | ObjectType::MaterializedView
7533 | ObjectType::Table
7534 | ObjectType::Source
7535 | ObjectType::Index,
7536 UnresolvedObjectName::Item(name),
7537 ) => name,
7538 (object_type, _) => {
7539 bail_unsupported!(format!("RETAIN HISTORY on {object_type}"))
7540 }
7541 };
7542 match resolve_item_or_type(scx, object_type, name.clone(), if_exists)? {
7543 Some(entry) => {
7544 let full_name = scx.catalog.resolve_full_name(entry.name());
7545 let item_type = entry.item_type();
7546
7547 if object_type == ObjectType::View && item_type == CatalogItemType::MaterializedView {
7549 return Err(PlanError::AlterViewOnMaterializedView(
7550 full_name.to_string(),
7551 ));
7552 } else if object_type == ObjectType::View {
7553 sql_bail!("{object_type} does not support RETAIN HISTORY")
7554 } else if object_type != item_type {
7555 sql_bail!(
7556 "\"{}\" is a {} not a {}",
7557 full_name,
7558 entry.item_type(),
7559 format!("{object_type}").to_lowercase()
7560 )
7561 }
7562
7563 let (value, lcw) = match &history {
7565 Some(WithOptionValue::RetainHistoryFor(value)) => {
7566 let window = OptionalDuration::try_from_value(value.clone())?;
7567 (Some(value.clone()), window.0)
7568 }
7569 None => (None, Some(DEFAULT_LOGICAL_COMPACTION_WINDOW_DURATION)),
7571 _ => sql_bail!("unexpected value type for RETAIN HISTORY"),
7572 };
7573 let window = plan_retain_history(scx, lcw)?;
7574
7575 Ok(Plan::AlterRetainHistory(AlterRetainHistoryPlan {
7576 id: entry.id(),
7577 value,
7578 window,
7579 object_type,
7580 }))
7581 }
7582 None => {
7583 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7584 name: name.to_ast_string_simple(),
7585 object_type,
7586 });
7587
7588 Ok(Plan::AlterNoop(AlterNoopPlan { object_type }))
7589 }
7590 }
7591}
7592
7593fn alter_source_timestamp_interval(
7594 scx: &StatementContext,
7595 if_exists: bool,
7596 source_name: UnresolvedItemName,
7597 value: Option<WithOptionValue<Aug>>,
7598) -> Result<Plan, PlanError> {
7599 let object_type = ObjectType::Source;
7600 match resolve_item_or_type(scx, object_type, source_name.clone(), if_exists)? {
7601 Some(entry) => {
7602 let full_name = scx.catalog.resolve_full_name(entry.name());
7603 if entry.item_type() != CatalogItemType::Source {
7604 sql_bail!(
7605 "\"{}\" is a {} not a {}",
7606 full_name,
7607 entry.item_type(),
7608 format!("{object_type}").to_lowercase()
7609 )
7610 }
7611
7612 match value {
7613 Some(val) => {
7614 let val = match val {
7615 WithOptionValue::Value(v) => v,
7616 _ => sql_bail!("TIMESTAMP INTERVAL requires an interval value"),
7617 };
7618 let duration = Duration::try_from_value(val.clone())?;
7619
7620 let min = scx.catalog.system_vars().min_timestamp_interval();
7621 let max = scx.catalog.system_vars().max_timestamp_interval();
7622 if duration < min || duration > max {
7623 return Err(PlanError::InvalidTimestampInterval {
7624 min,
7625 max,
7626 requested: duration,
7627 });
7628 }
7629
7630 Ok(Plan::AlterSourceTimestampInterval(
7631 AlterSourceTimestampIntervalPlan {
7632 id: entry.id(),
7633 value: Some(val),
7634 interval: duration,
7635 },
7636 ))
7637 }
7638 None => {
7639 let interval = scx.catalog.system_vars().default_timestamp_interval();
7640 Ok(Plan::AlterSourceTimestampInterval(
7641 AlterSourceTimestampIntervalPlan {
7642 id: entry.id(),
7643 value: None,
7644 interval,
7645 },
7646 ))
7647 }
7648 }
7649 }
7650 None => {
7651 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7652 name: source_name.to_ast_string_simple(),
7653 object_type,
7654 });
7655
7656 Ok(Plan::AlterNoop(AlterNoopPlan { object_type }))
7657 }
7658 }
7659}
7660
7661pub fn describe_alter_secret_options(
7662 _: &StatementContext,
7663 _: AlterSecretStatement<Aug>,
7664) -> Result<StatementDesc, PlanError> {
7665 Ok(StatementDesc::new(None))
7666}
7667
7668pub fn plan_alter_secret(
7669 scx: &mut StatementContext,
7670 stmt: AlterSecretStatement<Aug>,
7671) -> Result<Plan, PlanError> {
7672 let AlterSecretStatement {
7673 name,
7674 if_exists,
7675 value,
7676 } = stmt;
7677 let object_type = ObjectType::Secret;
7678 let id = match resolve_item_or_type(scx, object_type, name.clone(), if_exists)? {
7679 Some(entry) => entry.id(),
7680 None => {
7681 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7682 name: name.to_string(),
7683 object_type,
7684 });
7685
7686 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
7687 }
7688 };
7689
7690 let secret_as = query::plan_secret_as(scx, value)?;
7691
7692 Ok(Plan::AlterSecret(AlterSecretPlan { id, secret_as }))
7693}
7694
7695pub fn describe_alter_connection(
7696 _: &StatementContext,
7697 _: AlterConnectionStatement<Aug>,
7698) -> Result<StatementDesc, PlanError> {
7699 Ok(StatementDesc::new(None))
7700}
7701
7702generate_extracted_config!(AlterConnectionOption, (Validate, bool));
7703
7704pub fn plan_alter_connection(
7705 scx: &StatementContext,
7706 stmt: AlterConnectionStatement<Aug>,
7707) -> Result<Plan, PlanError> {
7708 let AlterConnectionStatement {
7709 name,
7710 if_exists,
7711 actions,
7712 with_options,
7713 } = stmt;
7714 let conn_name = normalize::unresolved_item_name(name)?;
7715 let entry = match scx.catalog.resolve_item(&conn_name) {
7716 Ok(entry) => entry,
7717 Err(_) if if_exists => {
7718 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7719 name: conn_name.to_string(),
7720 object_type: ObjectType::Connection,
7721 });
7722
7723 return Ok(Plan::AlterNoop(AlterNoopPlan {
7724 object_type: ObjectType::Connection,
7725 }));
7726 }
7727 Err(e) => return Err(e.into()),
7728 };
7729
7730 let connection = entry.connection()?;
7731
7732 if actions
7733 .iter()
7734 .any(|action| matches!(action, AlterConnectionAction::RotateKeys))
7735 {
7736 if actions.len() > 1 {
7737 sql_bail!("cannot specify any other actions alongside ALTER CONNECTION...ROTATE KEYS");
7738 }
7739
7740 if !with_options.is_empty() {
7741 sql_bail!(
7742 "ALTER CONNECTION...ROTATE KEYS does not support WITH ({})",
7743 with_options
7744 .iter()
7745 .map(|o| o.to_ast_string_simple())
7746 .join(", ")
7747 );
7748 }
7749
7750 if !matches!(connection, Connection::Ssh(_)) {
7751 sql_bail!(
7752 "{} is not an SSH connection",
7753 scx.catalog.resolve_full_name(entry.name())
7754 )
7755 }
7756
7757 return Ok(Plan::AlterConnection(AlterConnectionPlan {
7758 id: entry.id(),
7759 action: crate::plan::AlterConnectionAction::RotateKeys,
7760 }));
7761 }
7762
7763 let options = AlterConnectionOptionExtracted::try_from(with_options)?;
7764 if options.validate.is_some() {
7765 scx.require_feature_flag(&vars::ENABLE_CONNECTION_VALIDATION_SYNTAX)?;
7766 }
7767
7768 let validate = match options.validate {
7769 Some(val) => val,
7770 None => {
7771 scx.catalog
7772 .system_vars()
7773 .enable_default_connection_validation()
7774 && connection.validate_by_default()
7775 }
7776 };
7777
7778 let connection_type = match connection {
7779 Connection::Aws(_) => CreateConnectionType::Aws,
7780 Connection::AwsPrivatelink(_) => CreateConnectionType::AwsPrivatelink,
7781 Connection::Gcp(_) => CreateConnectionType::Gcp,
7782 Connection::Kafka(_) => CreateConnectionType::Kafka,
7783 Connection::Csr(_) => CreateConnectionType::Csr,
7784 Connection::GlueSchemaRegistry(_) => CreateConnectionType::GlueSchemaRegistry,
7785 Connection::Postgres(_) => CreateConnectionType::Postgres,
7786 Connection::Ssh(_) => CreateConnectionType::Ssh,
7787 Connection::MySql(_) => CreateConnectionType::MySql,
7788 Connection::SqlServer(_) => CreateConnectionType::SqlServer,
7789 Connection::IcebergCatalog(_) => CreateConnectionType::IcebergCatalog,
7790 };
7791
7792 let specified_options: BTreeSet<_> = actions
7794 .iter()
7795 .map(|action: &AlterConnectionAction<Aug>| match action {
7796 AlterConnectionAction::SetOption(option) => Ok(option.name.clone()),
7797 AlterConnectionAction::DropOption(name) => Ok(name.clone()),
7798 AlterConnectionAction::RotateKeys => {
7799 Err(internal_err!("RotateKeys is handled separately above"))
7800 }
7801 })
7802 .collect::<Result<_, PlanError>>()?;
7803
7804 for invalid in INALTERABLE_OPTIONS {
7805 if specified_options.contains(invalid) {
7806 sql_bail!("cannot ALTER {} option {}", connection_type, invalid);
7807 }
7808 }
7809
7810 connection::validate_options_per_connection_type(connection_type, specified_options)?;
7811
7812 let mut set_options_vec: Vec<_> = Vec::new();
7814 let mut drop_options: BTreeSet<_> = BTreeSet::new();
7815 for action in actions {
7816 match action {
7817 AlterConnectionAction::SetOption(option) => set_options_vec.push(option),
7818 AlterConnectionAction::DropOption(name) => {
7819 drop_options.insert(name);
7820 }
7821 AlterConnectionAction::RotateKeys => {
7822 bail_internal!("RotateKeys is handled separately above")
7823 }
7824 }
7825 }
7826
7827 let set_options: BTreeMap<_, _> = set_options_vec
7828 .clone()
7829 .into_iter()
7830 .map(|option| (option.name, option.value))
7831 .collect();
7832
7833 let connection_options_extracted =
7837 connection::ConnectionOptionExtracted::try_from(set_options_vec)?;
7838
7839 let duplicates: Vec<_> = connection_options_extracted
7840 .seen
7841 .intersection(&drop_options)
7842 .collect();
7843
7844 if !duplicates.is_empty() {
7845 sql_bail!(
7846 "cannot both SET and DROP/RESET options {}",
7847 duplicates
7848 .iter()
7849 .map(|option| option.to_string())
7850 .join(", ")
7851 )
7852 }
7853
7854 for mutually_exclusive_options in MUTUALLY_EXCLUSIVE_SETS {
7855 let set_options_count = mutually_exclusive_options
7856 .iter()
7857 .filter(|o| set_options.contains_key(o))
7858 .count();
7859 let drop_options_count = mutually_exclusive_options
7860 .iter()
7861 .filter(|o| drop_options.contains(o))
7862 .count();
7863
7864 if set_options_count > 0 && drop_options_count > 0 {
7866 sql_bail!(
7867 "cannot both SET and DROP/RESET mutually exclusive {} options {}",
7868 connection_type,
7869 mutually_exclusive_options
7870 .iter()
7871 .map(|option| option.to_string())
7872 .join(", ")
7873 )
7874 }
7875
7876 if set_options_count > 0 || drop_options_count > 0 {
7881 drop_options.extend(mutually_exclusive_options.iter().cloned());
7882 }
7883
7884 }
7887
7888 Ok(Plan::AlterConnection(AlterConnectionPlan {
7889 id: entry.id(),
7890 action: crate::plan::AlterConnectionAction::AlterOptions {
7891 set_options,
7892 drop_options,
7893 validate,
7894 },
7895 }))
7896}
7897
7898pub fn describe_alter_sink(
7899 _: &StatementContext,
7900 _: AlterSinkStatement<Aug>,
7901) -> Result<StatementDesc, PlanError> {
7902 Ok(StatementDesc::new(None))
7903}
7904
7905pub fn plan_alter_sink(
7906 scx: &mut StatementContext,
7907 stmt: AlterSinkStatement<Aug>,
7908) -> Result<Plan, PlanError> {
7909 let AlterSinkStatement {
7910 sink_name,
7911 if_exists,
7912 action,
7913 } = stmt;
7914
7915 let object_type = ObjectType::Sink;
7916 let item = resolve_item_or_type(scx, object_type, sink_name.clone(), if_exists)?;
7917
7918 let Some(item) = item else {
7919 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
7920 name: sink_name.to_string(),
7921 object_type,
7922 });
7923
7924 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
7925 };
7926 let item = item.at_version(RelationVersionSelector::Latest);
7928
7929 let create_sql = item.create_sql();
7931 let stmts = mz_sql_parser::parser::parse_statements(create_sql)?;
7932 let [stmt]: [StatementParseResult; 1] = stmts
7933 .try_into()
7934 .map_err(|_| internal_err!("create SQL of sink was not exactly one statement"))?;
7935 let Statement::CreateSink(stmt) = stmt.ast else {
7936 bail_internal!("create SQL of sink is not a CREATE SINK statement");
7937 };
7938 let (mut stmt, _) = crate::names::resolve(scx.catalog, stmt)?;
7939
7940 let mut set_options = vec![];
7942 let mut reset_options = vec![];
7943 match action {
7944 AlterSinkAction::ChangeRelation(new_from) => {
7945 stmt.from = new_from;
7946 }
7947 AlterSinkAction::SetOptions(options) => {
7948 for option in &options {
7949 match &option.name {
7950 CreateSinkOptionName::CommitInterval => {}
7951 name => bail_unsupported!(format!(
7952 "ALTER SINK ... SET ({})",
7953 name.to_ast_string_simple()
7954 )),
7955 }
7956 }
7957 if options.iter().all(|o| stmt.with_options.contains(o)) {
7973 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
7974 }
7975 set_options = options;
7976 }
7977 AlterSinkAction::ResetOptions(names) => {
7978 for name in &names {
7979 match name {
7980 CreateSinkOptionName::CommitInterval => {}
7981 name => bail_unsupported!(format!(
7982 "ALTER SINK ... RESET ({})",
7983 name.to_ast_string_simple()
7984 )),
7985 }
7986 if !stmt.with_options.iter().any(|o| o.name == *name) {
7989 sql_bail!(
7990 "cannot RESET {}: option is not set",
7991 name.to_ast_string_simple()
7992 );
7993 }
7994 }
7995 reset_options = names;
7996 }
7997 }
7998 crate::plan::apply_sink_option_edits(&mut stmt.with_options, &set_options, &reset_options);
7999
8000 let Plan::CreateSink(mut plan) = plan_sink(scx, stmt)? else {
8002 bail_internal!("plan_sink did not produce a CreateSink plan");
8003 };
8004
8005 plan.sink.version += 1;
8006
8007 Ok(Plan::AlterSink(AlterSinkPlan {
8008 item_id: item.id(),
8009 global_id: item.global_id(),
8010 sink: plan.sink,
8011 with_snapshot: plan.with_snapshot,
8012 in_cluster: plan.in_cluster,
8013 set_options,
8014 reset_options,
8015 }))
8016}
8017
8018pub fn describe_alter_source(
8019 _: &StatementContext,
8020 _: AlterSourceStatement<Aug>,
8021) -> Result<StatementDesc, PlanError> {
8022 Ok(StatementDesc::new(None))
8024}
8025
8026generate_extracted_config!(
8027 AlterSourceAddSubsourceOption,
8028 (TextColumns, Vec::<UnresolvedItemName>, Default(vec![])),
8029 (ExcludeColumns, Vec::<UnresolvedItemName>, Default(vec![])),
8030 (Details, String)
8031);
8032
8033pub fn plan_alter_source(
8034 scx: &mut StatementContext,
8035 stmt: AlterSourceStatement<Aug>,
8036) -> Result<Plan, PlanError> {
8037 let AlterSourceStatement {
8038 source_name,
8039 if_exists,
8040 action,
8041 } = stmt;
8042 let object_type = ObjectType::Source;
8043
8044 if resolve_item_or_type(scx, object_type, source_name.clone(), if_exists)?.is_none() {
8045 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
8046 name: source_name.to_string(),
8047 object_type,
8048 });
8049
8050 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
8051 }
8052
8053 match action {
8054 AlterSourceAction::SetOptions(options) => {
8055 let mut options = options.into_iter();
8056 let option = options
8057 .next()
8058 .ok_or_else(|| sql_err!("ALTER SOURCE SET requires at least one option"))?;
8059 if option.name == CreateSourceOptionName::RetainHistory {
8060 if options.next().is_some() {
8061 sql_bail!("RETAIN HISTORY must be only option");
8062 }
8063 return alter_retain_history(
8064 scx,
8065 object_type,
8066 if_exists,
8067 UnresolvedObjectName::Item(source_name),
8068 option.value,
8069 );
8070 }
8071 if option.name == CreateSourceOptionName::TimestampInterval {
8072 if options.next().is_some() {
8073 sql_bail!("TIMESTAMP INTERVAL must be only option");
8074 }
8075 return alter_source_timestamp_interval(scx, if_exists, source_name, option.value);
8076 }
8077 sql_bail!(
8080 "Cannot modify the {} of a SOURCE.",
8081 option.name.to_ast_string_simple()
8082 );
8083 }
8084 AlterSourceAction::ResetOptions(reset) => {
8085 let mut options = reset.into_iter();
8086 let option = options
8087 .next()
8088 .ok_or_else(|| sql_err!("ALTER SOURCE RESET requires at least one option"))?;
8089 if option == CreateSourceOptionName::RetainHistory {
8090 if options.next().is_some() {
8091 sql_bail!("RETAIN HISTORY must be only option");
8092 }
8093 return alter_retain_history(
8094 scx,
8095 object_type,
8096 if_exists,
8097 UnresolvedObjectName::Item(source_name),
8098 None,
8099 );
8100 }
8101 if option == CreateSourceOptionName::TimestampInterval {
8102 if options.next().is_some() {
8103 sql_bail!("TIMESTAMP INTERVAL must be only option");
8104 }
8105 return alter_source_timestamp_interval(scx, if_exists, source_name, None);
8106 }
8107 sql_bail!(
8108 "Cannot modify the {} of a SOURCE.",
8109 option.to_ast_string_simple()
8110 );
8111 }
8112 AlterSourceAction::DropSubsources { .. } => {
8113 sql_bail!("ALTER SOURCE...DROP SUBSOURCE no longer supported; use DROP SOURCE")
8114 }
8115 AlterSourceAction::AddSubsources { .. } => {
8116 sql_bail!("ALTER SOURCE...ADD SUBSOURCE must be purified before planning")
8117 }
8118 AlterSourceAction::RefreshReferences => {
8119 sql_bail!("ALTER SOURCE...REFRESH REFERENCES must be purified before planning")
8120 }
8121 };
8122}
8123
8124pub fn describe_alter_system_set(
8125 _: &StatementContext,
8126 _: AlterSystemSetStatement,
8127) -> Result<StatementDesc, PlanError> {
8128 Ok(StatementDesc::new(None))
8129}
8130
8131pub fn plan_alter_system_set(
8132 _: &StatementContext,
8133 AlterSystemSetStatement { name, to }: AlterSystemSetStatement,
8134) -> Result<Plan, PlanError> {
8135 let name = name.to_string();
8136 Ok(Plan::AlterSystemSet(AlterSystemSetPlan {
8137 name,
8138 value: scl::plan_set_variable_to(to)?,
8139 }))
8140}
8141
8142pub fn describe_alter_system_reset(
8143 _: &StatementContext,
8144 _: AlterSystemResetStatement,
8145) -> Result<StatementDesc, PlanError> {
8146 Ok(StatementDesc::new(None))
8147}
8148
8149pub fn plan_alter_system_reset(
8150 _: &StatementContext,
8151 AlterSystemResetStatement { name }: AlterSystemResetStatement,
8152) -> Result<Plan, PlanError> {
8153 let name = name.to_string();
8154 Ok(Plan::AlterSystemReset(AlterSystemResetPlan { name }))
8155}
8156
8157pub fn describe_alter_system_reset_all(
8158 _: &StatementContext,
8159 _: AlterSystemResetAllStatement,
8160) -> Result<StatementDesc, PlanError> {
8161 Ok(StatementDesc::new(None))
8162}
8163
8164pub fn plan_alter_system_reset_all(
8165 _: &StatementContext,
8166 _: AlterSystemResetAllStatement,
8167) -> Result<Plan, PlanError> {
8168 Ok(Plan::AlterSystemResetAll(AlterSystemResetAllPlan {}))
8169}
8170
8171pub fn describe_alter_role(
8172 _: &StatementContext,
8173 _: AlterRoleStatement<Aug>,
8174) -> Result<StatementDesc, PlanError> {
8175 Ok(StatementDesc::new(None))
8176}
8177
8178pub fn plan_alter_role(
8179 scx: &StatementContext,
8180 AlterRoleStatement { name, option }: AlterRoleStatement<Aug>,
8181) -> Result<Plan, PlanError> {
8182 let option = match option {
8183 AlterRoleOption::Attributes(attrs) => {
8184 let attrs = plan_role_attributes(attrs, scx)?;
8185 PlannedAlterRoleOption::Attributes(attrs)
8186 }
8187 AlterRoleOption::Variable(variable) => {
8188 let var = plan_role_variable(scx, variable)?;
8189 PlannedAlterRoleOption::Variable(var)
8190 }
8191 };
8192
8193 Ok(Plan::AlterRole(AlterRolePlan {
8194 id: name.id,
8195 name: name.name,
8196 option,
8197 }))
8198}
8199
8200pub fn describe_alter_table_add_column(
8201 _: &StatementContext,
8202 _: AlterTableAddColumnStatement<Aug>,
8203) -> Result<StatementDesc, PlanError> {
8204 Ok(StatementDesc::new(None))
8205}
8206
8207pub fn plan_alter_table_add_column(
8208 scx: &StatementContext,
8209 stmt: AlterTableAddColumnStatement<Aug>,
8210) -> Result<Plan, PlanError> {
8211 let AlterTableAddColumnStatement {
8212 if_exists,
8213 name,
8214 if_col_not_exist,
8215 column_name,
8216 data_type,
8217 } = stmt;
8218 let object_type = ObjectType::Table;
8219
8220 scx.require_feature_flag(&vars::ENABLE_ALTER_TABLE_ADD_COLUMN)?;
8221
8222 let (relation_id, item_name, desc) =
8223 match resolve_item_or_type(scx, object_type, name.clone(), if_exists)? {
8224 Some(item) => {
8225 let item_name = scx.catalog.resolve_full_name(item.name());
8227 let item = item.at_version(RelationVersionSelector::Latest);
8228 let desc = item
8229 .relation_desc()
8230 .ok_or_else(|| sql_err!("item does not have a relation description"))?
8231 .into_owned();
8232 (item.id(), item_name, desc)
8233 }
8234 None => {
8235 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
8236 name: name.to_ast_string_simple(),
8237 object_type,
8238 });
8239 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
8240 }
8241 };
8242
8243 let column_name = ColumnName::from(column_name.as_str());
8244 if desc.get_by_name(&column_name).is_some() {
8245 if if_col_not_exist {
8246 scx.catalog.add_notice(PlanNotice::ColumnAlreadyExists {
8247 column_name: column_name.to_string(),
8248 object_name: item_name.item,
8249 });
8250 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
8251 } else {
8252 return Err(PlanError::ColumnAlreadyExists {
8253 column_name,
8254 object_name: item_name.item,
8255 });
8256 }
8257 }
8258
8259 let scalar_type = scalar_type_from_sql(scx, &data_type)?;
8260 let column_type = scalar_type.nullable(true);
8262 let raw_sql_type = mz_sql_parser::parser::parse_data_type(&data_type.to_ast_string_stable())?;
8264
8265 Ok(Plan::AlterTableAddColumn(AlterTablePlan {
8266 relation_id,
8267 column_name,
8268 column_type,
8269 raw_sql_type,
8270 }))
8271}
8272
8273pub fn describe_alter_materialized_view_apply_replacement(
8274 _: &StatementContext,
8275 _: AlterMaterializedViewApplyReplacementStatement,
8276) -> Result<StatementDesc, PlanError> {
8277 Ok(StatementDesc::new(None))
8278}
8279
8280pub fn plan_alter_materialized_view_apply_replacement(
8281 scx: &StatementContext,
8282 stmt: AlterMaterializedViewApplyReplacementStatement,
8283) -> Result<Plan, PlanError> {
8284 let AlterMaterializedViewApplyReplacementStatement {
8285 if_exists,
8286 name,
8287 replacement_name,
8288 } = stmt;
8289
8290 scx.require_feature_flag(&vars::ENABLE_REPLACEMENT_MATERIALIZED_VIEWS)?;
8291
8292 let object_type = ObjectType::MaterializedView;
8293 let Some(mv) = resolve_item_or_type(scx, object_type, name.clone(), if_exists)? else {
8294 scx.catalog.add_notice(PlanNotice::ObjectDoesNotExist {
8295 name: name.to_ast_string_simple(),
8296 object_type,
8297 });
8298 return Ok(Plan::AlterNoop(AlterNoopPlan { object_type }));
8299 };
8300
8301 let replacement = resolve_item_or_type(scx, object_type, replacement_name, false)?
8302 .ok_or_else(|| sql_err!("replacement materialized view does not exist"))?;
8303
8304 if replacement.replacement_target() != Some(mv.id()) {
8305 return Err(PlanError::InvalidReplacement {
8306 item_type: mv.item_type(),
8307 item_name: scx.catalog.minimal_qualification(mv.name()),
8308 replacement_type: replacement.item_type(),
8309 replacement_name: scx.catalog.minimal_qualification(replacement.name()),
8310 });
8311 }
8312
8313 Ok(Plan::AlterMaterializedViewApplyReplacement(
8314 AlterMaterializedViewApplyReplacementPlan {
8315 id: mv.id(),
8316 replacement_id: replacement.id(),
8317 },
8318 ))
8319}
8320
8321pub fn describe_comment(
8322 _: &StatementContext,
8323 _: CommentStatement<Aug>,
8324) -> Result<StatementDesc, PlanError> {
8325 Ok(StatementDesc::new(None))
8326}
8327
8328pub fn plan_comment(
8329 scx: &mut StatementContext,
8330 stmt: CommentStatement<Aug>,
8331) -> Result<Plan, PlanError> {
8332 const MAX_COMMENT_LENGTH: usize = 1024;
8333
8334 let CommentStatement { object, comment } = stmt;
8335
8336 if let Some(c) = &comment {
8338 if c.len() > 1024 {
8339 return Err(PlanError::CommentTooLong {
8340 length: c.len(),
8341 max_size: MAX_COMMENT_LENGTH,
8342 });
8343 }
8344 }
8345
8346 let (object_id, column_pos) = match &object {
8347 com_ty @ CommentObjectType::Table { name }
8348 | com_ty @ CommentObjectType::View { name }
8349 | com_ty @ CommentObjectType::MaterializedView { name }
8350 | com_ty @ CommentObjectType::Index { name }
8351 | com_ty @ CommentObjectType::Func { name }
8352 | com_ty @ CommentObjectType::Connection { name }
8353 | com_ty @ CommentObjectType::Source { name }
8354 | com_ty @ CommentObjectType::Sink { name }
8355 | com_ty @ CommentObjectType::Secret { name } => {
8356 let item = scx.get_item_by_resolved_name(name)?;
8357 match (com_ty, item.item_type()) {
8358 (CommentObjectType::Table { .. }, CatalogItemType::Table) => {
8359 (CommentObjectId::Table(item.id()), None)
8360 }
8361 (CommentObjectType::View { .. }, CatalogItemType::View) => {
8362 (CommentObjectId::View(item.id()), None)
8363 }
8364 (CommentObjectType::MaterializedView { .. }, CatalogItemType::MaterializedView) => {
8365 (CommentObjectId::MaterializedView(item.id()), None)
8366 }
8367 (CommentObjectType::Index { .. }, CatalogItemType::Index) => {
8368 (CommentObjectId::Index(item.id()), None)
8369 }
8370 (CommentObjectType::Func { .. }, CatalogItemType::Func) => {
8371 (CommentObjectId::Func(item.id()), None)
8372 }
8373 (CommentObjectType::Connection { .. }, CatalogItemType::Connection) => {
8374 (CommentObjectId::Connection(item.id()), None)
8375 }
8376 (CommentObjectType::Source { .. }, CatalogItemType::Source) => {
8377 (CommentObjectId::Source(item.id()), None)
8378 }
8379 (CommentObjectType::Sink { .. }, CatalogItemType::Sink) => {
8380 (CommentObjectId::Sink(item.id()), None)
8381 }
8382 (CommentObjectType::Secret { .. }, CatalogItemType::Secret) => {
8383 (CommentObjectId::Secret(item.id()), None)
8384 }
8385 (com_ty, cat_ty) => {
8386 let expected_type = match com_ty {
8387 CommentObjectType::Table { .. } => ObjectType::Table,
8388 CommentObjectType::View { .. } => ObjectType::View,
8389 CommentObjectType::MaterializedView { .. } => ObjectType::MaterializedView,
8390 CommentObjectType::Index { .. } => ObjectType::Index,
8391 CommentObjectType::Func { .. } => ObjectType::Func,
8392 CommentObjectType::Connection { .. } => ObjectType::Connection,
8393 CommentObjectType::Source { .. } => ObjectType::Source,
8394 CommentObjectType::Sink { .. } => ObjectType::Sink,
8395 CommentObjectType::Secret { .. } => ObjectType::Secret,
8396 _ => sql_bail!("cannot comment on this object type"),
8397 };
8398
8399 return Err(PlanError::InvalidObjectType {
8400 expected_type: SystemObjectType::Object(expected_type),
8401 actual_type: SystemObjectType::Object(cat_ty.into()),
8402 object_name: item.name().item.clone(),
8403 });
8404 }
8405 }
8406 }
8407 CommentObjectType::Type { ty } => match ty {
8408 ResolvedDataType::AnonymousList(_) | ResolvedDataType::AnonymousMap { .. } => {
8409 sql_bail!("cannot comment on anonymous list or map type");
8410 }
8411 ResolvedDataType::Named { id, modifiers, .. } => {
8412 if !modifiers.is_empty() {
8413 sql_bail!("cannot comment on type with modifiers");
8414 }
8415 (CommentObjectId::Type(*id), None)
8416 }
8417 ResolvedDataType::Error => bail_internal!("unresolved data type"),
8418 },
8419 CommentObjectType::Column { name } => {
8420 let (item, pos) = scx.get_column_by_resolved_name(name)?;
8421 match item.item_type() {
8422 CatalogItemType::Table => (CommentObjectId::Table(item.id()), Some(pos + 1)),
8423 CatalogItemType::Source => (CommentObjectId::Source(item.id()), Some(pos + 1)),
8424 CatalogItemType::View => (CommentObjectId::View(item.id()), Some(pos + 1)),
8425 CatalogItemType::MaterializedView => {
8426 (CommentObjectId::MaterializedView(item.id()), Some(pos + 1))
8427 }
8428 CatalogItemType::Type => (CommentObjectId::Type(item.id()), Some(pos + 1)),
8429 r => {
8430 return Err(PlanError::Unsupported {
8431 feature: format!("Specifying comments on a column of {r}"),
8432 discussion_no: None,
8433 });
8434 }
8435 }
8436 }
8437 CommentObjectType::Role { name } => (CommentObjectId::Role(name.id), None),
8438 CommentObjectType::Database { name } => {
8439 (CommentObjectId::Database(*name.database_id()), None)
8440 }
8441 CommentObjectType::Schema { name } => {
8442 if matches!(name.schema_spec(), SchemaSpecifier::Temporary) {
8446 sql_bail!(
8447 "cannot comment on schema {} because it is a temporary schema",
8448 mz_repr::namespaces::MZ_TEMP_SCHEMA
8449 );
8450 }
8451 (
8452 CommentObjectId::Schema((*name.database_spec(), *name.schema_spec())),
8453 None,
8454 )
8455 }
8456 CommentObjectType::Cluster { name } => (CommentObjectId::Cluster(name.id), None),
8457 CommentObjectType::ClusterReplica { name } => {
8458 let replica = scx.catalog.resolve_cluster_replica(name)?;
8459 (
8460 CommentObjectId::ClusterReplica((replica.cluster_id(), replica.replica_id())),
8461 None,
8462 )
8463 }
8464 CommentObjectType::NetworkPolicy { name } => {
8465 (CommentObjectId::NetworkPolicy(name.id), None)
8466 }
8467 };
8468
8469 if let Some(p) = column_pos {
8475 i32::try_from(p).map_err(|_| PlanError::TooManyColumns {
8476 max_num_columns: MAX_NUM_COLUMNS,
8477 req_num_columns: p,
8478 })?;
8479 }
8480
8481 Ok(Plan::Comment(CommentPlan {
8482 object_id,
8483 sub_component: column_pos,
8484 comment,
8485 }))
8486}
8487
8488pub(crate) fn resolve_cluster<'a>(
8489 scx: &'a StatementContext,
8490 name: &'a Ident,
8491 if_exists: bool,
8492) -> Result<Option<&'a dyn CatalogCluster<'a>>, PlanError> {
8493 match scx.resolve_cluster(Some(name)) {
8494 Ok(cluster) => Ok(Some(cluster)),
8495 Err(_) if if_exists => Ok(None),
8496 Err(e) => Err(e),
8497 }
8498}
8499
8500pub(crate) fn resolve_cluster_replica<'a>(
8501 scx: &'a StatementContext,
8502 name: &QualifiedReplica,
8503 if_exists: bool,
8504) -> Result<Option<(&'a dyn CatalogCluster<'a>, ReplicaId)>, PlanError> {
8505 match scx.resolve_cluster(Some(&name.cluster)) {
8506 Ok(cluster) => match cluster.replica_ids().get(name.replica.as_str()) {
8507 Some(replica_id) => Ok(Some((cluster, *replica_id))),
8508 None if if_exists => Ok(None),
8509 None => Err(sql_err!(
8510 "CLUSTER {} has no CLUSTER REPLICA named {}",
8511 cluster.name(),
8512 name.replica.as_str().quoted(),
8513 )),
8514 },
8515 Err(_) if if_exists => Ok(None),
8516 Err(e) => Err(e),
8517 }
8518}
8519
8520pub(crate) fn resolve_database<'a>(
8521 scx: &'a StatementContext,
8522 name: &'a UnresolvedDatabaseName,
8523 if_exists: bool,
8524) -> Result<Option<&'a dyn CatalogDatabase>, PlanError> {
8525 match scx.resolve_database(name) {
8526 Ok(database) => Ok(Some(database)),
8527 Err(_) if if_exists => Ok(None),
8528 Err(e) => Err(e),
8529 }
8530}
8531
8532pub(crate) fn resolve_schema<'a>(
8533 scx: &'a StatementContext,
8534 name: UnresolvedSchemaName,
8535 if_exists: bool,
8536) -> Result<Option<(ResolvedDatabaseSpecifier, SchemaSpecifier)>, PlanError> {
8537 match scx.resolve_schema(name) {
8538 Ok(schema) => Ok(Some((schema.database().clone(), schema.id().clone()))),
8539 Err(_) if if_exists => Ok(None),
8540 Err(e) => Err(e),
8541 }
8542}
8543
8544pub(crate) fn resolve_network_policy<'a>(
8545 scx: &'a StatementContext,
8546 name: Ident,
8547 if_exists: bool,
8548) -> Result<Option<ResolvedNetworkPolicyName>, PlanError> {
8549 match scx
8550 .catalog
8551 .resolve_network_policy(normalize::ident_ref(&name))
8552 {
8553 Ok(policy) => Ok(Some(ResolvedNetworkPolicyName {
8554 id: policy.id(),
8555 name: policy.name().to_string(),
8556 })),
8557 Err(_) if if_exists => Ok(None),
8558 Err(e) => Err(e.into()),
8559 }
8560}
8561
8562pub(crate) fn resolve_item_or_type<'a>(
8563 scx: &'a StatementContext,
8564 object_type: ObjectType,
8565 name: UnresolvedItemName,
8566 if_exists: bool,
8567) -> Result<Option<&'a dyn CatalogItem>, PlanError> {
8568 let name = normalize::unresolved_item_name(name)?;
8569 let catalog_item = match object_type {
8570 ObjectType::Type => scx.catalog.resolve_type(&name),
8571 ObjectType::Table
8572 | ObjectType::View
8573 | ObjectType::MaterializedView
8574 | ObjectType::Source
8575 | ObjectType::Sink
8576 | ObjectType::MetricSink
8577 | ObjectType::Index
8578 | ObjectType::Role
8579 | ObjectType::Cluster
8580 | ObjectType::ClusterReplica
8581 | ObjectType::Secret
8582 | ObjectType::Connection
8583 | ObjectType::Database
8584 | ObjectType::Schema
8585 | ObjectType::Func
8586 | ObjectType::NetworkPolicy => scx.catalog.resolve_item(&name),
8587 };
8588
8589 match catalog_item {
8590 Ok(item) => {
8591 let is_type = ObjectType::from(item.item_type());
8592 if object_type == is_type {
8593 Ok(Some(item))
8594 } else {
8595 Err(PlanError::MismatchedObjectType {
8596 name: scx.catalog.minimal_qualification(item.name()),
8597 is_type,
8598 expected_type: object_type,
8599 })
8600 }
8601 }
8602 Err(_) if if_exists => Ok(None),
8603 Err(e) => Err(e.into()),
8604 }
8605}
8606
8607fn ensure_cluster_is_not_managed(
8609 scx: &StatementContext,
8610 cluster_id: ClusterId,
8611) -> Result<(), PlanError> {
8612 let cluster = scx.catalog.get_cluster(cluster_id);
8613 if cluster.is_managed() {
8614 Err(PlanError::ManagedCluster {
8615 cluster_name: cluster.name().to_string(),
8616 })
8617 } else {
8618 Ok(())
8619 }
8620}
8621
8622#[cfg(test)]
8623mod tests {
8624 use super::{
8625 METRIC_SINK_CURATED_PREFIX_MARKER, METRIC_SINK_PREFIX_MARKER,
8626 validate_user_metric_sink_prefix,
8627 };
8628
8629 #[mz_ore::test]
8630 fn user_metric_sink_prefix_cannot_reach_the_curated_prefix() {
8631 assert!(validate_user_metric_sink_prefix(METRIC_SINK_CURATED_PREFIX_MARKER).is_err());
8633 assert!(validate_user_metric_sink_prefix("mz_metric_sink_curated_x").is_err());
8634 assert!(validate_user_metric_sink_prefix(METRIC_SINK_PREFIX_MARKER).is_err());
8636 assert!(validate_user_metric_sink_prefix("mz_metric_sink_cur").is_err());
8637 assert!(validate_user_metric_sink_prefix("myapp_").is_err());
8639 assert!(validate_user_metric_sink_prefix("mz_metric_sink_myapp_").is_ok());
8642 assert!(validate_user_metric_sink_prefix("mz_metric_sink_curatex_").is_ok());
8643 assert!(validate_user_metric_sink_prefix("mz_metric_sink_my_curated_").is_ok());
8644 }
8645}