1use std::collections::BTreeMap;
17
18use mz_sql_parser::ast::display::AstDisplay;
19use mz_sql_parser::ast::item_refs::collect_item_references;
20use mz_sql_parser::ast::{
21 AstInfo, AvroSchema, ConnectionOption, ConnectionOptionName, CreateConnectionType,
22 CreateSinkConnection, CreateSubsourceOptionName, Format, FormatSpecifier,
23 IcebergSinkConfigOptionName, IcebergSinkMode, KafkaSinkConfigOptionName,
24 KafkaSourceConfigOptionName, PgConfigOptionName, ProtobufSchema, Raw, RawClusterName,
25 RawItemName, SinkEnvelope, SourceEnvelope, SourceErrorPolicy, Statement, UnresolvedItemName,
26 Value, WithOptionValue,
27};
28use prost::Message as _;
29use serde_json::json;
30
31fn parse_single_statement(sql: &str) -> Result<Statement<Raw>, String> {
33 let mut stmts = mz_sql_parser::parser::parse_statements(sql)
34 .map_err(|e| format!("failed to parse create_sql: {e}"))?;
35 match stmts.len() {
36 1 => Ok(stmts.remove(0).ast),
37 n => Err(format!("expected a single statement, found {n}")),
38 }
39}
40
41fn item_id(item: RawItemName) -> Result<String, &'static str> {
42 match item {
43 RawItemName::Id(id, _, _) => Ok(id),
44 RawItemName::Name(_) => Err("unresolved item name"),
45 }
46}
47
48fn cluster_id(in_cluster: RawClusterName) -> Result<String, &'static str> {
49 match in_cluster {
50 RawClusterName::Resolved(s) => Ok(s),
51 RawClusterName::Unresolved(_) => Err("unresolved cluster name"),
52 }
53}
54
55fn format_name<T: AstInfo>(fmt: &Format<T>) -> &'static str {
56 match fmt {
57 Format::Bytes => "bytes",
58 Format::Avro(_) => "avro",
59 Format::Protobuf(_) => "protobuf",
60 Format::Regex(_) => "regex",
61 Format::Csv { .. } => "csv",
62 Format::Json { .. } => "json",
63 Format::Text => "text",
64 }
65}
66
67fn option_string<T: AstInfo>(value: &WithOptionValue<T>) -> Option<String> {
78 match value {
79 WithOptionValue::Value(Value::String(s)) => Some(s.clone()),
80 WithOptionValue::Ident(ident) => Some(ident.clone().into_string()),
81 WithOptionValue::UnresolvedItemName(UnresolvedItemName(parts)) if parts.len() == 1 => {
82 Some(parts[0].clone().into_string())
83 }
84 _ => None,
85 }
86}
87
88pub fn item_details(a: &str) -> Result<serde_json::Value, String> {
93 let stmt = parse_single_statement(a)?;
94
95 let mut info = BTreeMap::<&str, serde_json::Value>::new();
96
97 use mz_sql_parser::ast::Statement::*;
98 let item_type = match stmt {
99 CreateSecret(_) => "secret",
100 CreateConnection(stmt) => {
101 let connection_type = stmt.connection_type.as_str();
102 info.insert("connection_type", json!(connection_type));
103
104 "connection"
105 }
106 CreateView(stmt) => {
107 let mut definition = stmt.definition.query.to_ast_string_stable();
108 definition.push(';');
111 info.insert("definition", json!(definition));
112
113 "view"
114 }
115 CreateMaterializedView(stmt) => {
116 let Some(in_cluster) = stmt.in_cluster else {
117 return Err("missing IN CLUSTER".into());
118 };
119 let cluster_id = match in_cluster {
120 RawClusterName::Unresolved(ident) => ident.into_string(),
121 RawClusterName::Resolved(s) => s,
122 };
123 info.insert("cluster_id", json!(cluster_id));
124
125 let mut definition = stmt.query.to_ast_string_stable();
126 definition.push(';');
127 info.insert("definition", json!(definition));
128
129 "materialized-view"
130 }
131 CreateTable(_) => "table",
132 CreateTableFromSource(stmt) => {
133 let source_id = item_id(stmt.source)?;
134 info.insert("source_id", json!(source_id));
135
136 "table"
137 }
138 CreateSource(stmt) => {
139 let Some(in_cluster) = stmt.in_cluster else {
140 return Err("missing IN CLUSTER".into());
141 };
142 let cluster_id = cluster_id(in_cluster)?;
143 info.insert("cluster_id", json!(cluster_id));
144
145 use mz_sql_parser::ast::CreateSourceConnection::*;
146 let (source_type, connection) = match stmt.connection {
147 Kafka { connection, .. } => ("kafka", Some(connection)),
148 Postgres { connection, .. } => ("postgres", Some(connection)),
149 MySql { connection, .. } => ("mysql", Some(connection)),
150 SqlServer { connection, .. } => ("sql-server", Some(connection)),
151 LoadGenerator { .. } => ("load-generator", None),
152 };
153 info.insert("source_type", json!(source_type));
154 if let Some(conn) = connection {
155 let conn_id = item_id(conn)?;
156 info.insert("connection_id", json!(conn_id));
157 }
158
159 let is_debezium = matches!(
160 stmt.envelope,
161 Some(mz_sql_parser::ast::SourceEnvelope::Debezium)
162 );
163
164 let envelope_type = match &stmt.envelope {
175 Some(envelope) => {
176 use mz_sql_parser::ast::SourceEnvelope::*;
177 Some(match envelope {
178 None => "none",
179 Debezium => "debezium",
180 Upsert { .. } => "upsert",
181 CdcV2 => "materialize",
182 })
183 }
184 None if source_type == "kafka" && stmt.progress_subsource.is_some() => Some("none"),
185 None => None,
186 };
187 if let Some(envelope_type) = envelope_type {
188 info.insert("envelope_type", json!(envelope_type));
189 }
190
191 if let Some(format_spec) = stmt.format {
192 match &format_spec {
193 FormatSpecifier::Bare(fmt) => {
194 if is_debezium {
197 info.insert("key_format", json!(format_name(fmt)));
198 }
199 info.insert("value_format", json!(format_name(fmt)));
200 }
201 FormatSpecifier::KeyValue { key, value } => {
202 info.insert("key_format", json!(format_name(key)));
203 info.insert("value_format", json!(format_name(value)));
204 }
205 }
206 }
207
208 "source"
209 }
210 CreateWebhookSource(stmt) => {
211 if stmt.is_table {
212 "table"
213 } else {
214 info.insert("source_type", json!("webhook"));
215 if let Some(in_cluster) = stmt.in_cluster {
216 let cluster_id = cluster_id(in_cluster)?;
217 info.insert("cluster_id", json!(cluster_id));
218 }
219 "source"
220 }
221 }
222 CreateSubsource(stmt) => {
223 use mz_sql_parser::ast::CreateSubsourceOptionName;
224 let is_progress = stmt
225 .with_options
226 .iter()
227 .any(|o| matches!(o.name, CreateSubsourceOptionName::Progress));
228 let source_type = if is_progress { "progress" } else { "subsource" };
229 info.insert("source_type", json!(source_type));
230
231 if let Some(of_source) = stmt.of_source {
232 let of_source_id = item_id(of_source)?;
233 info.insert("of_source_id", json!(of_source_id));
234 }
235
236 "subsource"
237 }
238 CreateSink(stmt) => {
248 let Some(in_cluster) = stmt.in_cluster else {
249 return Err("missing IN CLUSTER".into());
250 };
251 info.insert("cluster_id", json!(cluster_id(in_cluster)?));
252
253 match stmt.connection {
254 CreateSinkConnection::Kafka {
255 connection,
256 options,
257 key: sink_key,
258 ..
259 } => {
260 info.insert("sink_type", json!("kafka"));
261 info.insert("connection_id", json!(item_id(connection)?));
262
263 let topic = options
264 .into_iter()
265 .find(|o| o.name == KafkaSinkConfigOptionName::Topic)
266 .and_then(|o| o.value.as_ref().and_then(option_string))
267 .ok_or("kafka sink missing TOPIC")?;
268 info.insert("topic", json!(topic));
269
270 if let Some(envelope) = stmt.envelope {
271 let envelope_type = match envelope {
272 SinkEnvelope::Upsert => "upsert",
273 SinkEnvelope::Debezium => "debezium",
274 };
275 info.insert("envelope_type", json!(envelope_type));
276 }
277
278 if let Some(format_spec) = stmt.format {
279 let (key_format, value_format) = match (&format_spec, sink_key.is_some()) {
284 (FormatSpecifier::Bare(fmt), false) => (None, format_name(fmt)),
285 (FormatSpecifier::Bare(fmt), true) => {
286 (Some(format_name(fmt)), format_name(fmt))
287 }
288 (FormatSpecifier::KeyValue { value, .. }, false) => {
289 (None, format_name(value))
290 }
291 (FormatSpecifier::KeyValue { key, value }, true) => {
292 (Some(format_name(key)), format_name(value))
293 }
294 };
295 if let Some(key_format) = key_format {
296 info.insert("key_format", json!(key_format));
297 }
298 info.insert("value_format", json!(value_format));
299
300 let combined = match key_format {
305 None => value_format.to_string(),
306 Some(key_format)
307 if key_format == value_format
308 && matches!(value_format, "avro" | "json") =>
309 {
310 value_format.to_string()
311 }
312 Some(key_format) => {
313 format!("key-{key_format}-value-{value_format}")
314 }
315 };
316 info.insert("format", json!(combined));
317 }
318 }
319 CreateSinkConnection::Iceberg {
320 catalog_connection,
321 options,
322 ..
323 } => {
324 info.insert("sink_type", json!("iceberg"));
325 info.insert("connection_id", json!(item_id(catalog_connection)?));
327
328 let mut namespace = None;
329 let mut table = None;
330 for option in options {
331 match option.name {
332 IcebergSinkConfigOptionName::Namespace => {
333 namespace = option.value.as_ref().and_then(option_string)
334 }
335 IcebergSinkConfigOptionName::Table => {
336 table = option.value.as_ref().and_then(option_string)
337 }
338 }
339 }
340 info.insert(
341 "namespace",
342 json!(namespace.ok_or("iceberg sink missing NAMESPACE")?),
343 );
344 info.insert("table", json!(table.ok_or("iceberg sink missing TABLE")?));
345
346 if let Some(mode) = stmt.mode {
349 let envelope_type = match mode {
350 IcebergSinkMode::Upsert => "upsert",
351 IcebergSinkMode::Append => "append",
352 };
353 info.insert("envelope_type", json!(envelope_type));
354 }
355 }
356 }
357
358 "sink"
359 }
360 CreateMetricSink(stmt) => {
361 let Some(in_cluster) = stmt.in_cluster else {
362 return Err("missing IN CLUSTER".into());
363 };
364 let cluster_id = cluster_id(in_cluster)?;
365 info.insert("cluster_id", json!(cluster_id));
366 let from_id = item_id(stmt.from)?;
367 info.insert("from_id", json!(from_id));
368 "metric-sink"
369 }
370 CreateIndex(stmt) => {
371 let Some(in_cluster) = stmt.in_cluster else {
372 return Err("missing IN CLUSTER".into());
373 };
374 let cluster_id = cluster_id(in_cluster)?;
375 info.insert("cluster_id", json!(cluster_id));
376 let on_id = item_id(stmt.on_name)?;
377 info.insert("on_id", json!(on_id));
378 "index"
379 }
380 CreateType(_) => "type",
381 Select(_)
387 | Insert(_)
388 | Copy(_)
389 | Update(_)
390 | Delete(_)
391 | CreateDatabase(_)
392 | CreateSchema(_)
393 | CreateRole(_)
394 | CreateCluster(_)
395 | CreateClusterReplica(_)
396 | CreateNetworkPolicy(_)
397 | AlterCluster(_)
398 | AlterOwner(_)
399 | AlterObjectRename(_)
400 | AlterObjectSwap(_)
401 | AlterRetainHistory(_)
402 | AlterIndex(_)
403 | AlterSecret(_)
404 | AlterSetCluster(_)
405 | AlterSink(_)
406 | AlterSource(_)
407 | AlterSystemSet(_)
408 | AlterSystemReset(_)
409 | AlterSystemResetAll(_)
410 | AlterConnection(_)
411 | AlterNetworkPolicy(_)
412 | AlterRole(_)
413 | AlterTableAddColumn(_)
414 | AlterMaterializedViewApplyReplacement(_)
415 | Discard(_)
416 | DropObjects(_)
417 | DropOwned(_)
418 | SetVariable(_)
419 | ResetVariable(_)
420 | Show(_)
421 | StartTransaction(_)
422 | SetTransaction(_)
423 | Commit(_)
424 | Rollback(_)
425 | Subscribe(_)
426 | ExplainPlan(_)
427 | ExplainPushdown(_)
428 | ExplainTimestamp(_)
429 | ExplainSinkSchema(_)
430 | ExplainAnalyzeObject(_)
431 | ExplainAnalyzeCluster(_)
432 | Declare(_)
433 | Fetch(_)
434 | Close(_)
435 | Prepare(_)
436 | Execute(_)
437 | ExecuteUnitTest(_)
438 | Deallocate(_)
439 | Raise(_)
440 | GrantRole(_)
441 | RevokeRole(_)
442 | GrantPrivileges(_)
443 | RevokePrivileges(_)
444 | AlterDefaultPrivileges(_)
445 | ReassignOwned(_)
446 | ValidateConnection(_)
447 | Comment(_) => return Err("not a CREATE item statement".into()),
448 };
449 info.insert("type", json!(item_type));
450
451 let info = info.into_iter().map(|(k, v)| (k.to_string(), v)).collect();
452 Ok(info)
453}
454
455pub fn item_references(a: &str) -> Result<serde_json::Value, String> {
463 fn qualified(name: &UnresolvedItemName) -> serde_json::Value {
464 match &name.0[..] {
465 [.., schema, item] => json!({"schema": schema.as_str(), "name": item.as_str()}),
466 [item] => json!({"schema": serde_json::Value::Null, "name": item.as_str()}),
467 [] => json!({"schema": serde_json::Value::Null, "name": ""}),
468 }
469 }
470
471 let stmt = parse_single_statement(a)?;
472
473 let refs = collect_item_references(&stmt);
474 mz_ore::soft_assert_or_log!(
475 refs.named_array_elements.is_empty(),
476 "persisted create_sql should never carry an array type T[] \
477 and resolve its array type in `ids`: {:?}",
478 refs.named_array_elements
479 );
480 Ok(json!({
481 "ids": refs.ids.iter().collect::<Vec<_>>(),
482 "named_funcs": refs.named_funcs.iter().map(qualified).collect::<Vec<_>>(),
483 "named_types": refs.named_types.iter().map(qualified).collect::<Vec<_>>(),
484 "named_relations": refs.named_relations.iter().map(qualified).collect::<Vec<_>>(),
485 }))
486}
487
488#[derive(Clone, PartialEq, ::prost::Message)]
494struct PostgresPublicationDetailsSubset {
495 #[prost(string, tag = "2")]
496 slot: String,
497 #[prost(uint64, optional, tag = "3")]
498 timeline_id: Option<u64>,
499}
500
501pub fn postgres_source_details(a: &str) -> Result<serde_json::Value, String> {
512 let stmt = parse_single_statement(a)?;
513
514 use mz_sql_parser::ast::CreateSourceConnection;
515 use mz_sql_parser::ast::Statement::CreateSource;
516 let options = match stmt {
517 CreateSource(stmt) => match stmt.connection {
518 CreateSourceConnection::Postgres { options, .. } => options,
519 _ => return Ok(serde_json::Value::Null),
520 },
521 _ => return Ok(serde_json::Value::Null),
522 };
523
524 let details_hex = options
525 .into_iter()
526 .find(|opt| opt.name == PgConfigOptionName::Details)
527 .and_then(|opt| match opt.value {
528 Some(WithOptionValue::Value(Value::String(s))) => Some(s),
529 _ => None,
530 })
531 .ok_or("missing DETAILS option on postgres source")?;
532
533 let details_bytes =
534 hex::decode(&details_hex).map_err(|e| format!("DETAILS is not valid hex: {e}"))?;
535
536 let details = PostgresPublicationDetailsSubset::decode(&*details_bytes)
537 .map_err(|e| format!("DETAILS is not a valid publication-details proto: {e}"))?;
538
539 Ok(json!({
540 "slot": details.slot,
541 "timeline_id": details.timeline_id,
542 }))
543}
544
545pub fn kafka_source_details(a: &str) -> Result<serde_json::Value, String> {
557 let stmt = parse_single_statement(a)?;
558
559 use mz_sql_parser::ast::CreateSourceConnection;
560 use mz_sql_parser::ast::Statement::CreateSource;
561 let (connection, options) = match stmt {
562 CreateSource(stmt) => match stmt.connection {
563 CreateSourceConnection::Kafka {
564 connection,
565 options,
566 } => (connection, options),
567 _ => return Ok(serde_json::Value::Null),
568 },
569 _ => return Ok(serde_json::Value::Null),
570 };
571
572 let connection_id = item_id(connection)?;
573
574 let mut topic: Option<String> = None;
575 let mut group_id_prefix: Option<String> = None;
576 for opt in options {
577 let string_value = opt.value.as_ref().and_then(option_string);
578 match opt.name {
579 KafkaSourceConfigOptionName::Topic => topic = string_value,
580 KafkaSourceConfigOptionName::GroupIdPrefix => group_id_prefix = string_value,
581 _ => {}
582 }
583 }
584
585 let topic = topic.ok_or("missing TOPIC option on kafka source")?;
586
587 Ok(json!({
588 "topic": topic,
589 "group_id_prefix": group_id_prefix,
590 "connection_id": connection_id,
591 }))
592}
593
594pub fn source_export_details(a: &str) -> Result<serde_json::Value, String> {
629 fn bare_format_has_key<T: AstInfo>(fmt: &Format<T>) -> bool {
633 match fmt {
634 Format::Avro(AvroSchema::Csr { csr_connection }) => csr_connection
635 .seed
636 .as_ref()
637 .is_some_and(|seed| seed.key_schema.is_some()),
638 Format::Protobuf(ProtobufSchema::Csr { csr_connection }) => csr_connection
639 .seed
640 .as_ref()
641 .is_some_and(|seed| seed.key.is_some()),
642 _ => false,
643 }
644 }
645
646 fn key_value_formats<T: AstInfo>(
647 spec: &FormatSpecifier<T>,
648 ) -> (Option<&'static str>, Option<&'static str>) {
649 match spec {
650 FormatSpecifier::KeyValue { key, value } => {
651 (Some(format_name(key)), Some(format_name(value)))
652 }
653 FormatSpecifier::Bare(fmt) => {
654 let value = Some(format_name(fmt));
655 let key = bare_format_has_key(fmt).then(|| format_name(fmt));
656 (key, value)
657 }
658 }
659 }
660
661 fn envelope_name(envelope: &SourceEnvelope) -> &'static str {
662 match envelope {
663 SourceEnvelope::None => "none",
664 SourceEnvelope::Debezium => "debezium",
665 SourceEnvelope::Upsert {
666 value_decode_err_policy,
667 } => {
668 if value_decode_err_policy
669 .iter()
670 .any(|p| matches!(p, SourceErrorPolicy::Inline { .. }))
671 {
672 "upsert-value-err-inline"
673 } else {
674 "upsert"
675 }
676 }
677 SourceEnvelope::CdcV2 => "materialize",
678 }
679 }
680
681 let stmt = parse_single_statement(a)?;
682
683 use mz_sql_parser::ast::Statement::{CreateSubsource, CreateTableFromSource};
684 match stmt {
685 CreateTableFromSource(stmt) => {
686 let source_id = item_id(stmt.source)?;
687 let external_reference = stmt
688 .external_reference
689 .ok_or("missing external reference on CREATE TABLE FROM SOURCE")?
690 .0
691 .into_iter()
692 .map(|ident| ident.into_string())
693 .collect::<Vec<_>>();
694
695 let envelope_type = stmt.envelope.as_ref().map(envelope_name);
696 let (key_format, value_format) = match &stmt.format {
697 Some(spec) => key_value_formats(spec),
698 None => (None, None),
699 };
700
701 Ok(json!({
702 "source_id": source_id,
703 "external_reference": external_reference,
704 "envelope_type": envelope_type,
705 "key_format": key_format,
706 "value_format": value_format,
707 }))
708 }
709 CreateSubsource(stmt) => {
710 let is_progress = stmt
713 .with_options
714 .iter()
715 .any(|o| matches!(o.name, CreateSubsourceOptionName::Progress));
716 if is_progress {
717 return Ok(serde_json::Value::Null);
718 }
719
720 let source_id = stmt
721 .of_source
722 .ok_or("non-progress CREATE SUBSOURCE without OF SOURCE")
723 .and_then(item_id)?;
724
725 let external_reference = stmt
726 .with_options
727 .into_iter()
728 .find(|o| matches!(o.name, CreateSubsourceOptionName::ExternalReference))
729 .and_then(|o| match o.value {
730 Some(WithOptionValue::UnresolvedItemName(name)) => Some(name),
731 _ => None,
732 })
733 .ok_or("CREATE SUBSOURCE missing EXTERNAL REFERENCE option")?
734 .0
735 .into_iter()
736 .map(|ident| ident.into_string())
737 .collect::<Vec<_>>();
738
739 Ok(json!({
740 "source_id": source_id,
741 "external_reference": external_reference,
742 "envelope_type": serde_json::Value::Null,
743 "key_format": serde_json::Value::Null,
744 "value_format": serde_json::Value::Null,
745 }))
746 }
747 _ => Ok(serde_json::Value::Null),
748 }
749}
750
751pub fn connection_details(a: &str) -> Result<serde_json::Value, String> {
796 fn broker_addresses<T: AstInfo>(values: &[ConnectionOption<T>]) -> Vec<String> {
801 let mut brokers = Vec::new();
802 for opt in values {
803 match (&opt.name, &opt.value) {
804 (ConnectionOptionName::Broker, Some(WithOptionValue::ConnectionKafkaBroker(b))) => {
805 brokers.push(b.address.clone());
806 }
807 (ConnectionOptionName::Brokers, Some(WithOptionValue::Sequence(seq))) => {
808 for v in seq {
809 if let WithOptionValue::ConnectionKafkaBroker(b) = v {
810 brokers.push(b.address.clone());
811 }
812 }
813 }
814 _ => {}
815 }
816 }
817 brokers
818 }
819
820 fn string_option<T: AstInfo>(
821 values: &[ConnectionOption<T>],
822 name: ConnectionOptionName,
823 ) -> Option<String> {
824 values
825 .iter()
826 .find(|o| o.name == name)
827 .and_then(|o| o.value.as_ref())
828 .and_then(option_string)
829 }
830
831 fn secret_id_option<T: AstInfo<ItemName = RawItemName>>(
835 values: &[ConnectionOption<T>],
836 name: ConnectionOptionName,
837 ) -> Option<String> {
838 values.iter().find_map(|o| match &o.value {
839 Some(WithOptionValue::Secret(RawItemName::Id(id, _, _))) if o.name == name => {
840 Some(id.clone())
841 }
842 _ => None,
843 })
844 }
845
846 let stmt = parse_single_statement(a)?;
847
848 use mz_sql_parser::ast::Statement::CreateConnection;
849 let stmt = match stmt {
850 CreateConnection(stmt) => stmt,
851 _ => return Ok(serde_json::Value::Null),
852 };
853
854 match stmt.connection_type {
855 CreateConnectionType::Kafka => Ok(json!({
856 "brokers": broker_addresses(&stmt.values),
857 "progress_topic": string_option(&stmt.values, ConnectionOptionName::ProgressTopic),
858 })),
859 CreateConnectionType::Ssh => Ok(json!({
860 "public_key_1": string_option(&stmt.values, ConnectionOptionName::PublicKey1),
861 "public_key_2": string_option(&stmt.values, ConnectionOptionName::PublicKey2),
862 })),
863 CreateConnectionType::Aws => {
864 let assume_role_arn = string_option(&stmt.values, ConnectionOptionName::AssumeRoleArn);
865 let auth_kind = if assume_role_arn.is_some() {
866 "assume-role"
867 } else {
868 "credentials"
869 };
870 Ok(json!({
871 "auth_kind": auth_kind,
872 "endpoint": string_option(&stmt.values, ConnectionOptionName::Endpoint)
876 .filter(|s| !s.is_empty()),
877 "region": string_option(&stmt.values, ConnectionOptionName::Region),
878 "access_key_id": string_option(&stmt.values, ConnectionOptionName::AccessKeyId),
879 "access_key_id_secret_id":
880 secret_id_option(&stmt.values, ConnectionOptionName::AccessKeyId),
881 "secret_access_key_secret_id":
882 secret_id_option(&stmt.values, ConnectionOptionName::SecretAccessKey),
883 "session_token": string_option(&stmt.values, ConnectionOptionName::SessionToken),
884 "session_token_secret_id":
885 secret_id_option(&stmt.values, ConnectionOptionName::SessionToken),
886 "assume_role_arn": assume_role_arn,
887 "assume_role_session_name":
888 string_option(&stmt.values, ConnectionOptionName::AssumeRoleSessionName),
889 }))
890 }
891 _ => Ok(serde_json::Value::Null),
892 }
893}
894
895#[cfg(test)]
896mod tests {
897 use prost::Message as _;
898 use serde_json::json;
899
900 fn encode_pg_details(slot: &str, timeline_id: Option<u64>) -> String {
903 let details = super::PostgresPublicationDetailsSubset {
904 slot: slot.to_string(),
905 timeline_id,
906 };
907 hex::encode(details.encode_to_vec())
908 }
909
910 fn pg_source_sql(details_hex: &str) -> String {
911 format!(
912 "CREATE SOURCE \"materialize\".\"public\".\"pg_src\" \
913 IN CLUSTER [u42] \
914 FROM POSTGRES CONNECTION [u10 AS \"materialize\".\"public\".\"pg_conn\"] \
915 (DETAILS = '{details_hex}', PUBLICATION = 'mz_source') \
916 FOR ALL TABLES"
917 )
918 }
919
920 fn kafka_source_sql(with_prefix: bool) -> String {
921 let prefix_opt = if with_prefix {
922 ", GROUP ID PREFIX 'my-prefix-'"
923 } else {
924 ""
925 };
926 format!(
927 "CREATE SOURCE \"materialize\".\"public\".\"k_src\" \
928 IN CLUSTER [u42] \
929 FROM KAFKA CONNECTION [u11 AS \"materialize\".\"public\".\"k_conn\"] \
930 (TOPIC 'test'{prefix_opt}) FORMAT TEXT"
931 )
932 }
933
934 #[mz_ore::test]
937 #[cfg_attr(miri, ignore)] fn pg_happy_path_with_timeline() {
939 let hex = encode_pg_details("materialize_abc", Some(42));
940 let sql = pg_source_sql(&hex);
941 let out = super::postgres_source_details(&sql).expect("ok");
942 assert_eq!(out, json!({ "slot": "materialize_abc", "timeline_id": 42 }));
943 }
944
945 #[mz_ore::test]
946 fn pg_happy_path_null_timeline() {
947 let hex = encode_pg_details("materialize_legacy", None);
950 let sql = pg_source_sql(&hex);
951 let out = super::postgres_source_details(&sql).expect("ok");
952 assert_eq!(
953 out,
954 json!({ "slot": "materialize_legacy", "timeline_id": null }),
955 );
956 }
957
958 #[mz_ore::test]
959 fn pg_non_postgres_source_returns_null_jsonb() {
960 let sql = "CREATE SOURCE \"materialize\".\"public\".\"lg\" \
961 IN CLUSTER [u42] FROM LOAD GENERATOR COUNTER";
962 let out = super::postgres_source_details(sql).expect("ok");
963 assert_eq!(out, serde_json::Value::Null);
964 }
965
966 #[mz_ore::test]
967 #[cfg_attr(miri, ignore)] fn pg_non_create_source_returns_null_jsonb() {
969 let sql = "CREATE VIEW v AS SELECT 1";
970 let out = super::postgres_source_details(sql).expect("ok");
971 assert_eq!(out, serde_json::Value::Null);
972 }
973
974 #[mz_ore::test]
975 fn pg_missing_details_option_errors() {
976 let sql = "CREATE SOURCE \"materialize\".\"public\".\"pg_src\" \
977 IN CLUSTER [u42] \
978 FROM POSTGRES CONNECTION [u10 AS \"materialize\".\"public\".\"pg_conn\"] \
979 (PUBLICATION = 'mz_source') FOR ALL TABLES";
980 let err = super::postgres_source_details(sql).unwrap_err();
981 assert!(
982 err.contains("missing DETAILS"),
983 "wrong error message: {err}"
984 );
985 }
986
987 #[mz_ore::test]
988 fn pg_malformed_hex_errors() {
989 let sql = pg_source_sql("not-hex!!");
990 let err = super::postgres_source_details(&sql).unwrap_err();
991 assert!(err.contains("valid hex"), "wrong error message: {err}");
992 }
993
994 #[mz_ore::test]
995 fn pg_malformed_proto_errors() {
996 let sql = pg_source_sql("ffff");
999 let err = super::postgres_source_details(&sql).unwrap_err();
1000 assert!(
1001 err.contains("publication-details proto"),
1002 "wrong error message: {err}"
1003 );
1004 }
1005
1006 #[mz_ore::test]
1009 fn kafka_happy_path_with_prefix() {
1010 let sql = kafka_source_sql(true);
1011 let out = super::kafka_source_details(&sql).expect("ok");
1012 assert_eq!(
1013 out,
1014 json!({
1015 "topic": "test",
1016 "group_id_prefix": "my-prefix-",
1017 "connection_id": "u11",
1018 }),
1019 );
1020 }
1021
1022 #[mz_ore::test]
1023 fn kafka_happy_path_without_prefix() {
1024 let sql = kafka_source_sql(false);
1025 let out = super::kafka_source_details(&sql).expect("ok");
1026 assert_eq!(
1027 out,
1028 json!({
1029 "topic": "test",
1030 "group_id_prefix": null,
1031 "connection_id": "u11",
1032 }),
1033 );
1034 }
1035
1036 #[mz_ore::test]
1037 fn kafka_non_kafka_source_returns_null_jsonb() {
1038 let sql = "CREATE SOURCE \"materialize\".\"public\".\"lg\" \
1039 IN CLUSTER [u42] FROM LOAD GENERATOR COUNTER";
1040 let out = super::kafka_source_details(sql).expect("ok");
1041 assert_eq!(out, serde_json::Value::Null);
1042 }
1043
1044 #[mz_ore::test]
1045 fn kafka_missing_topic_errors() {
1046 let sql = "CREATE SOURCE \"materialize\".\"public\".\"k_src\" \
1047 IN CLUSTER [u42] \
1048 FROM KAFKA CONNECTION [u11 AS \"materialize\".\"public\".\"k_conn\"] \
1049 FORMAT TEXT";
1050 let err = super::kafka_source_details(sql).unwrap_err();
1051 assert!(err.contains("missing TOPIC"), "wrong error message: {err}");
1052 }
1053
1054 #[mz_ore::test]
1055 fn kafka_unquoted_topic_and_prefix() {
1056 let sql = "CREATE SOURCE \"materialize\".\"public\".\"k_src\" \
1060 IN CLUSTER [u42] \
1061 FROM KAFKA CONNECTION [u11 AS \"materialize\".\"public\".\"k_conn\"] \
1062 (TOPIC = my_topic, GROUP ID PREFIX = my_prefix) FORMAT TEXT";
1063 let out = super::kafka_source_details(sql).expect("ok");
1064 assert_eq!(out["topic"], json!("my_topic"));
1065 assert_eq!(out["group_id_prefix"], json!("my_prefix"));
1066 }
1067
1068 #[mz_ore::test]
1069 fn kafka_unresolved_connection_errors() {
1070 let sql = "CREATE SOURCE \"materialize\".\"public\".\"k_src\" \
1074 IN CLUSTER [u42] \
1075 FROM KAFKA CONNECTION k_conn (TOPIC 'test') FORMAT TEXT";
1076 let err = super::kafka_source_details(sql).unwrap_err();
1077 assert!(
1078 err.contains("unresolved item name"),
1079 "wrong error message: {err}"
1080 );
1081 }
1082
1083 fn table_from_source_sql(reference: &str, suffix: &str) -> String {
1086 format!(
1087 "CREATE TABLE \"materialize\".\"public\".\"tbl\" \
1088 FROM SOURCE [u1 AS \"materialize\".\"public\".\"src\"] \
1089 (REFERENCE = {reference}){suffix}"
1090 )
1091 }
1092
1093 #[mz_ore::test]
1094 fn export_table_postgres_style_no_format() {
1095 let sql = table_from_source_sql("\"db\".\"public\".\"t\"", "");
1098 let out = super::source_export_details(&sql).expect("ok");
1099 assert_eq!(
1100 out,
1101 json!({
1102 "source_id": "u1",
1103 "external_reference": ["db", "public", "t"],
1104 "envelope_type": null,
1105 "key_format": null,
1106 "value_format": null,
1107 }),
1108 );
1109 }
1110
1111 #[mz_ore::test]
1112 fn export_table_kafka_bare_value_only() {
1113 let sql = table_from_source_sql("\"topic\"", " FORMAT TEXT ENVELOPE NONE");
1115 let out = super::source_export_details(&sql).expect("ok");
1116 assert_eq!(
1117 out,
1118 json!({
1119 "source_id": "u1",
1120 "external_reference": ["topic"],
1121 "envelope_type": "none",
1122 "key_format": null,
1123 "value_format": "text",
1124 }),
1125 );
1126 }
1127
1128 #[mz_ore::test]
1129 fn export_table_kafka_omitted_envelope_is_null() {
1130 let sql = table_from_source_sql("\"topic\"", " FORMAT TEXT");
1134 let out = super::source_export_details(&sql).expect("ok");
1135 assert_eq!(out["envelope_type"], serde_json::Value::Null);
1136 }
1137
1138 #[mz_ore::test]
1139 fn export_table_kafka_key_value_format() {
1140 let sql = table_from_source_sql(
1141 "\"topic\"",
1142 " KEY FORMAT TEXT VALUE FORMAT TEXT ENVELOPE NONE",
1143 );
1144 let out = super::source_export_details(&sql).expect("ok");
1145 assert_eq!(
1146 out,
1147 json!({
1148 "source_id": "u1",
1149 "external_reference": ["topic"],
1150 "envelope_type": "none",
1151 "key_format": "text",
1152 "value_format": "text",
1153 }),
1154 );
1155 }
1156
1157 #[mz_ore::test]
1158 #[cfg_attr(miri, ignore)] fn export_table_kafka_bare_avro_seed_with_key() {
1160 let sql = table_from_source_sql(
1164 "\"topic\"",
1165 " FORMAT AVRO USING CONFLUENT SCHEMA REGISTRY \
1166 CONNECTION [u5 AS \"materialize\".\"public\".\"csr\"] \
1167 SEED KEY SCHEMA 'k' VALUE SCHEMA 'v' ENVELOPE UPSERT",
1168 );
1169 let out = super::source_export_details(&sql).expect("ok");
1170 assert_eq!(
1171 out,
1172 json!({
1173 "source_id": "u1",
1174 "external_reference": ["topic"],
1175 "envelope_type": "upsert",
1176 "key_format": "avro",
1177 "value_format": "avro",
1178 }),
1179 );
1180 }
1181
1182 #[mz_ore::test]
1183 #[cfg_attr(miri, ignore)] fn export_table_kafka_bare_avro_seed_without_key() {
1185 let sql = table_from_source_sql(
1187 "\"topic\"",
1188 " FORMAT AVRO USING CONFLUENT SCHEMA REGISTRY \
1189 CONNECTION [u5 AS \"materialize\".\"public\".\"csr\"] \
1190 SEED VALUE SCHEMA 'v' ENVELOPE NONE",
1191 );
1192 let out = super::source_export_details(&sql).expect("ok");
1193 assert_eq!(
1194 out,
1195 json!({
1196 "source_id": "u1",
1197 "external_reference": ["topic"],
1198 "envelope_type": "none",
1199 "key_format": null,
1200 "value_format": "avro",
1201 }),
1202 );
1203 }
1204
1205 #[mz_ore::test]
1206 fn export_subsource_non_progress() {
1207 let sql = "CREATE SUBSOURCE \"materialize\".\"public\".\"sub\" (id int4) \
1210 OF SOURCE [u1 AS \"materialize\".\"public\".\"src\"] \
1211 WITH (EXTERNAL REFERENCE = \"db\".\"public\".\"t\")";
1212 let out = super::source_export_details(sql).expect("ok");
1213 assert_eq!(
1214 out,
1215 json!({
1216 "source_id": "u1",
1217 "external_reference": ["db", "public", "t"],
1218 "envelope_type": null,
1219 "key_format": null,
1220 "value_format": null,
1221 }),
1222 );
1223 }
1224
1225 #[mz_ore::test]
1226 fn export_progress_subsource_returns_null_jsonb() {
1227 let sql = "CREATE SUBSOURCE \"materialize\".\"public\".\"progress\" (id int4) \
1228 WITH (PROGRESS)";
1229 let out = super::source_export_details(sql).expect("ok");
1230 assert_eq!(out, serde_json::Value::Null);
1231 }
1232
1233 #[mz_ore::test]
1234 #[cfg_attr(miri, ignore)] fn export_non_source_export_returns_null_jsonb() {
1236 let sql = "CREATE VIEW v AS SELECT 1";
1237 let out = super::source_export_details(sql).expect("ok");
1238 assert_eq!(out, serde_json::Value::Null);
1239 }
1240
1241 #[mz_ore::test]
1242 fn export_unresolved_source_name_errors() {
1243 let sql = "CREATE TABLE \"materialize\".\"public\".\"tbl\" \
1244 FROM SOURCE src (REFERENCE = \"topic\")";
1245 let err = super::source_export_details(sql).unwrap_err();
1246 assert!(
1247 err.contains("unresolved item name"),
1248 "wrong error message: {err}"
1249 );
1250 }
1251
1252 #[mz_ore::test]
1255 fn connection_kafka_single_broker_default_progress() {
1256 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO KAFKA \
1259 (BROKER = 'localhost:9092', SECURITY PROTOCOL = plaintext)";
1260 let out = super::connection_details(sql).expect("ok");
1261 assert_eq!(
1262 out,
1263 json!({
1264 "brokers": ["localhost:9092"],
1265 "progress_topic": null,
1266 }),
1267 );
1268 }
1269
1270 #[mz_ore::test]
1271 fn connection_kafka_explicit_progress_topic() {
1272 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO KAFKA \
1273 (BROKER = 'localhost:9092', PROGRESS TOPIC = 'override', \
1274 SECURITY PROTOCOL = plaintext)";
1275 let out = super::connection_details(sql).expect("ok");
1276 assert_eq!(
1277 out,
1278 json!({
1279 "brokers": ["localhost:9092"],
1280 "progress_topic": "override",
1281 }),
1282 );
1283 }
1284
1285 #[mz_ore::test]
1286 fn connection_kafka_broker_list() {
1287 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO KAFKA \
1288 (BROKERS ('b1:9092', 'b2:9092'), SECURITY PROTOCOL = plaintext)";
1289 let out = super::connection_details(sql).expect("ok");
1290 assert_eq!(
1291 out,
1292 json!({
1293 "brokers": ["b1:9092", "b2:9092"],
1294 "progress_topic": null,
1295 }),
1296 );
1297 }
1298
1299 #[mz_ore::test]
1300 fn connection_ssh_public_keys() {
1301 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO SSH TUNNEL \
1302 (HOST = 'ssh.example.com', PORT = 22, USER = 'mz', \
1303 PUBLIC KEY 1 = 'ssh-ed25519 AAAA', PUBLIC KEY 2 = 'ssh-ed25519 BBBB')";
1304 let out = super::connection_details(sql).expect("ok");
1305 assert_eq!(
1306 out,
1307 json!({
1308 "public_key_1": "ssh-ed25519 AAAA",
1309 "public_key_2": "ssh-ed25519 BBBB",
1310 }),
1311 );
1312 }
1313
1314 #[mz_ore::test]
1315 fn connection_aws_credentials_inline_key() {
1316 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO AWS \
1319 (ACCESS KEY ID = 'AKIAEXAMPLE', \
1320 SECRET ACCESS KEY = SECRET [u1 AS \"materialize\".\"public\".\"sk\"])";
1321 let out = super::connection_details(sql).expect("ok");
1322 assert_eq!(
1323 out,
1324 json!({
1325 "auth_kind": "credentials",
1326 "endpoint": null,
1327 "region": null,
1328 "access_key_id": "AKIAEXAMPLE",
1329 "access_key_id_secret_id": null,
1330 "secret_access_key_secret_id": "u1",
1331 "session_token": null,
1332 "session_token_secret_id": null,
1333 "assume_role_arn": null,
1334 "assume_role_session_name": null,
1335 }),
1336 );
1337 }
1338
1339 #[mz_ore::test]
1340 fn connection_aws_credentials_secret_key_and_session_token() {
1341 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO AWS \
1344 (ENDPOINT = 'http://localhost', REGION = 'us-east-1', \
1345 ACCESS KEY ID = SECRET [u1 AS \"materialize\".\"public\".\"ak\"], \
1346 SECRET ACCESS KEY = SECRET [u2 AS \"materialize\".\"public\".\"sk\"], \
1347 SESSION TOKEN = SECRET [u3 AS \"materialize\".\"public\".\"st\"])";
1348 let out = super::connection_details(sql).expect("ok");
1349 assert_eq!(
1350 out,
1351 json!({
1352 "auth_kind": "credentials",
1353 "endpoint": "http://localhost",
1354 "region": "us-east-1",
1355 "access_key_id": null,
1356 "access_key_id_secret_id": "u1",
1357 "secret_access_key_secret_id": "u2",
1358 "session_token": null,
1359 "session_token_secret_id": "u3",
1360 "assume_role_arn": null,
1361 "assume_role_session_name": null,
1362 }),
1363 );
1364 }
1365
1366 #[mz_ore::test]
1367 fn connection_aws_assume_role() {
1368 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO AWS \
1371 (ASSUME ROLE ARN 'arn:aws:iam::123:role/mz', \
1372 ASSUME ROLE SESSION NAME 'sess')";
1373 let out = super::connection_details(sql).expect("ok");
1374 assert_eq!(
1375 out,
1376 json!({
1377 "auth_kind": "assume-role",
1378 "endpoint": null,
1379 "region": null,
1380 "access_key_id": null,
1381 "access_key_id_secret_id": null,
1382 "secret_access_key_secret_id": null,
1383 "session_token": null,
1384 "session_token_secret_id": null,
1385 "assume_role_arn": "arn:aws:iam::123:role/mz",
1386 "assume_role_session_name": "sess",
1387 }),
1388 );
1389 }
1390
1391 #[mz_ore::test]
1392 fn connection_kafka_unquoted_progress_topic() {
1393 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO KAFKA \
1396 (BROKER = 'localhost:9092', PROGRESS TOPIC = my_topic, \
1397 SECURITY PROTOCOL = plaintext)";
1398 let out = super::connection_details(sql).expect("ok");
1399 assert_eq!(out["progress_topic"], json!("my_topic"));
1400 }
1401
1402 #[mz_ore::test]
1403 fn connection_aws_unquoted_option_values() {
1404 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO AWS \
1407 (ENDPOINT = localhost, REGION = useast1, \
1408 ASSUME ROLE ARN 'arn:aws:iam::123:role/mz', \
1409 ASSUME ROLE SESSION NAME = mysession)";
1410 let out = super::connection_details(sql).expect("ok");
1411 assert_eq!(out["endpoint"], json!("localhost"));
1412 assert_eq!(out["region"], json!("useast1"));
1413 assert_eq!(out["assume_role_session_name"], json!("mysession"));
1414 }
1415
1416 #[mz_ore::test]
1417 fn connection_aws_empty_endpoint_is_null() {
1418 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO AWS \
1421 (ENDPOINT = '', ASSUME ROLE ARN 'arn:aws:iam::123:role/mz')";
1422 let out = super::connection_details(sql).expect("ok");
1423 assert_eq!(out["endpoint"], serde_json::Value::Null);
1424 }
1425
1426 #[mz_ore::test]
1427 fn connection_other_type_returns_null_jsonb() {
1428 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"c\" TO POSTGRES \
1430 (HOST = 'db', DATABASE = 'postgres', USER = 'mz')";
1431 let out = super::connection_details(sql).expect("ok");
1432 assert_eq!(out, serde_json::Value::Null);
1433 }
1434
1435 #[mz_ore::test]
1436 #[cfg_attr(miri, ignore)] fn connection_non_connection_returns_null_jsonb() {
1438 let sql = "CREATE VIEW v AS SELECT 1";
1439 let out = super::connection_details(sql).expect("ok");
1440 assert_eq!(out, serde_json::Value::Null);
1441 }
1442
1443 #[mz_ore::test]
1446 fn catalog_kafka_old_syntax_omitted_envelope_defaults_none() {
1447 let sql = "CREATE SOURCE \"materialize\".\"public\".\"k\" \
1451 IN CLUSTER [u42] \
1452 FROM KAFKA CONNECTION [u11 AS \"materialize\".\"public\".\"k_conn\"] \
1453 (TOPIC 'test') FORMAT TEXT \
1454 EXPOSE PROGRESS AS [u12 AS \"materialize\".\"public\".\"k_progress\"]";
1455 let out = super::item_details(sql).expect("ok");
1456 assert_eq!(out["envelope_type"], json!("none"));
1457 }
1458
1459 #[mz_ore::test]
1460 fn catalog_kafka_old_syntax_explicit_envelope() {
1461 let sql = "CREATE SOURCE \"materialize\".\"public\".\"k\" \
1462 IN CLUSTER [u42] \
1463 FROM KAFKA CONNECTION [u11 AS \"materialize\".\"public\".\"k_conn\"] \
1464 (TOPIC 'test') FORMAT BYTES ENVELOPE UPSERT \
1465 EXPOSE PROGRESS AS [u12 AS \"materialize\".\"public\".\"k_progress\"]";
1466 let out = super::item_details(sql).expect("ok");
1467 assert_eq!(out["envelope_type"], json!("upsert"));
1468 }
1469
1470 #[mz_ore::test]
1471 fn catalog_kafka_new_syntax_source_omits_envelope_type() {
1472 let sql = "CREATE SOURCE \"materialize\".\"public\".\"k\" \
1476 IN CLUSTER [u42] \
1477 FROM KAFKA CONNECTION [u11 AS \"materialize\".\"public\".\"k_conn\"]";
1478 let out = super::item_details(sql).expect("ok");
1479 assert_eq!(out.get("envelope_type"), None);
1480 }
1481
1482 #[mz_ore::test]
1483 fn catalog_non_kafka_source_omits_envelope_type() {
1484 let sql = "CREATE SOURCE \"materialize\".\"public\".\"lg\" \
1487 IN CLUSTER [u42] FROM LOAD GENERATOR COUNTER";
1488 let out = super::item_details(sql).expect("ok");
1489 assert_eq!(out.get("envelope_type"), None);
1490 }
1491
1492 const AVRO_FORMAT: &str = "FORMAT AVRO USING CONFLUENT SCHEMA REGISTRY CONNECTION \
1495 [u12 AS \"materialize\".\"public\".\"csr_conn\"]";
1496
1497 fn kafka_sink_sql(key: Option<&str>, format: &str, envelope: &str) -> String {
1500 let key_clause = key.map(|k| format!(" KEY ({k})")).unwrap_or_default();
1501 format!(
1502 "CREATE SINK \"materialize\".\"public\".\"snk\" \
1503 IN CLUSTER [u42] \
1504 FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1505 INTO KAFKA CONNECTION [u10 AS \"materialize\".\"public\".\"k_conn\"] \
1506 (TOPIC 'sink-topic'){key_clause} {format} ENVELOPE {envelope}"
1507 )
1508 }
1509
1510 fn iceberg_sink_sql(mode: &str) -> String {
1511 format!(
1512 "CREATE SINK \"materialize\".\"public\".\"ice\" \
1513 IN CLUSTER [u42] \
1514 FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1515 INTO ICEBERG CATALOG CONNECTION [u20 AS \"materialize\".\"public\".\"cat_conn\"] \
1516 (NAMESPACE 'ns', TABLE 'tbl') \
1517 USING AWS CONNECTION [u21 AS \"materialize\".\"public\".\"aws_conn\"] \
1518 MODE {mode}"
1519 )
1520 }
1521
1522 #[mz_ore::test]
1523 fn sink_kafka_bare_format_without_key() {
1524 let sql = kafka_sink_sql(None, "FORMAT JSON", "DEBEZIUM");
1525 let out = super::item_details(&sql).expect("ok");
1526 assert_eq!(
1527 out,
1528 json!({
1529 "type": "sink",
1530 "sink_type": "kafka",
1531 "cluster_id": "u42",
1532 "connection_id": "u10",
1533 "topic": "sink-topic",
1534 "envelope_type": "debezium",
1535 "format": "json",
1536 "value_format": "json",
1537 }),
1538 );
1539 }
1540
1541 #[mz_ore::test]
1542 fn sink_kafka_bare_format_with_key_derives_key_format() {
1543 let sql = kafka_sink_sql(Some("a"), AVRO_FORMAT, "UPSERT");
1546 let out = super::item_details(&sql).expect("ok");
1547 assert_eq!(
1548 out,
1549 json!({
1550 "type": "sink",
1551 "sink_type": "kafka",
1552 "cluster_id": "u42",
1553 "connection_id": "u10",
1554 "topic": "sink-topic",
1555 "envelope_type": "upsert",
1556 "format": "avro",
1557 "key_format": "avro",
1558 "value_format": "avro",
1559 }),
1560 );
1561 }
1562
1563 #[mz_ore::test]
1564 fn sink_kafka_bare_text_format_with_key_does_not_collapse() {
1565 let sql = kafka_sink_sql(Some("a"), "FORMAT TEXT", "UPSERT");
1568 let out = super::item_details(&sql).expect("ok");
1569 assert_eq!(out["format"], json!("key-text-value-text"));
1570 assert_eq!(out["key_format"], json!("text"));
1571 assert_eq!(out["value_format"], json!("text"));
1572 }
1573
1574 #[mz_ore::test]
1575 fn sink_kafka_key_value_json_collapses() {
1576 let sql = kafka_sink_sql(Some("a"), "KEY FORMAT JSON VALUE FORMAT JSON", "UPSERT");
1577 let out = super::item_details(&sql).expect("ok");
1578 assert_eq!(out["format"], json!("json"));
1579 assert_eq!(out["key_format"], json!("json"));
1580 assert_eq!(out["value_format"], json!("json"));
1581 }
1582
1583 #[mz_ore::test]
1584 fn sink_kafka_key_value_mixed_is_composite() {
1585 let sql = kafka_sink_sql(Some("a"), "KEY FORMAT TEXT VALUE FORMAT BYTES", "UPSERT");
1586 let out = super::item_details(&sql).expect("ok");
1587 assert_eq!(out["format"], json!("key-text-value-bytes"));
1588 assert_eq!(out["key_format"], json!("text"));
1589 assert_eq!(out["value_format"], json!("bytes"));
1590 }
1591
1592 #[mz_ore::test]
1593 fn sink_kafka_key_format_without_key_is_dropped() {
1594 let sql = kafka_sink_sql(None, "KEY FORMAT JSON VALUE FORMAT TEXT", "DEBEZIUM");
1598 let out = super::item_details(&sql).expect("ok");
1599 assert_eq!(out["format"], json!("text"));
1600 assert_eq!(out["key_format"], serde_json::Value::Null);
1601 assert_eq!(out["value_format"], json!("text"));
1602 }
1603
1604 #[mz_ore::test]
1605 fn sink_kafka_missing_topic_errors() {
1606 let sql = "CREATE SINK \"materialize\".\"public\".\"snk\" \
1607 IN CLUSTER [u42] \
1608 FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1609 INTO KAFKA CONNECTION [u10 AS \"materialize\".\"public\".\"k_conn\"] \
1610 FORMAT JSON ENVELOPE DEBEZIUM";
1611 let err = super::item_details(sql).unwrap_err();
1612 assert!(err.contains("missing TOPIC"), "wrong error message: {err}");
1613 }
1614
1615 #[mz_ore::test]
1616 fn sink_iceberg_upsert_mode() {
1617 let sql = iceberg_sink_sql("UPSERT");
1618 let out = super::item_details(&sql).expect("ok");
1619 assert_eq!(
1620 out,
1621 json!({
1622 "type": "sink",
1623 "sink_type": "iceberg",
1624 "cluster_id": "u42",
1625 "connection_id": "u20",
1627 "namespace": "ns",
1628 "table": "tbl",
1629 "envelope_type": "upsert",
1630 }),
1631 );
1632 }
1633
1634 #[mz_ore::test]
1635 fn sink_iceberg_append_mode() {
1636 let sql = iceberg_sink_sql("APPEND");
1637 let out = super::item_details(&sql).expect("ok");
1638 assert_eq!(out["envelope_type"], json!("append"));
1640 }
1641
1642 #[mz_ore::test]
1643 fn sink_iceberg_missing_table_errors() {
1644 let sql = "CREATE SINK \"materialize\".\"public\".\"ice\" \
1645 IN CLUSTER [u42] \
1646 FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1647 INTO ICEBERG CATALOG CONNECTION [u20 AS \"materialize\".\"public\".\"cat_conn\"] \
1648 (NAMESPACE 'ns') MODE UPSERT";
1649 let err = super::item_details(sql).unwrap_err();
1650 assert!(err.contains("missing TABLE"), "wrong error message: {err}");
1651 }
1652
1653 #[mz_ore::test]
1654 #[cfg_attr(miri, ignore)] fn sink_arm_leaves_other_item_types_alone() {
1656 let sql = "CREATE VIEW \"materialize\".\"public\".\"v\" AS SELECT 1";
1657 let out = super::item_details(sql).expect("ok");
1658 assert_eq!(out, json!({ "type": "view", "definition": "SELECT 1;" }));
1659 }
1660
1661 fn item_type(sql: &str) -> Result<String, String> {
1665 match super::item_details(sql) {
1666 Ok(out) => match out {
1667 serde_json::Value::Object(mut m) => match m.remove("type") {
1668 Some(serde_json::Value::String(s)) => Ok(s),
1669 other => panic!("no string `type` key: {other:?}"),
1670 },
1671 other => panic!("not a JSON object: {other:?}"),
1672 },
1673 Err(msg) => Err(msg),
1674 }
1675 }
1676
1677 fn view_sql(query: &str) -> String {
1678 format!("CREATE VIEW \"materialize\".\"public\".\"v\" AS {query}")
1679 }
1680
1681 fn view_definition(query: &str) -> String {
1683 match super::item_details(&view_sql(query)).expect("ok") {
1684 serde_json::Value::Object(mut m) => match m.remove("definition") {
1685 Some(serde_json::Value::String(s)) => s,
1686 other => panic!("no string `definition` key: {other:?}"),
1687 },
1688 other => panic!("not a JSON object: {other:?}"),
1689 }
1690 }
1691
1692 #[mz_ore::test]
1698 #[cfg_attr(miri, ignore)] fn catalog_item_type_per_statement_kind() {
1700 let cases = [
1701 (
1702 "CREATE TABLE \"materialize\".\"public\".\"t\" (a int4)",
1703 "table",
1704 ),
1705 (
1708 "CREATE TABLE \"materialize\".\"public\".\"tbl\" \
1709 FROM SOURCE [u1 AS \"materialize\".\"public\".\"src\"] \
1710 (REFERENCE = \"topic\") FORMAT TEXT",
1711 "table",
1712 ),
1713 (
1714 "CREATE TABLE \"materialize\".\"public\".\"wht\" FROM WEBHOOK BODY FORMAT JSON",
1715 "table",
1716 ),
1717 (
1718 "CREATE VIEW \"materialize\".\"public\".\"v\" AS SELECT 1",
1719 "view",
1720 ),
1721 (
1722 "CREATE MATERIALIZED VIEW \"materialize\".\"public\".\"mv\" \
1723 IN CLUSTER [u1] AS SELECT 1",
1724 "materialized-view",
1725 ),
1726 (
1727 "CREATE SOURCE \"materialize\".\"public\".\"lg\" \
1728 IN CLUSTER [u1] FROM LOAD GENERATOR COUNTER",
1729 "source",
1730 ),
1731 (
1732 "CREATE SOURCE \"materialize\".\"public\".\"wh\" \
1733 IN CLUSTER [u1] FROM WEBHOOK BODY FORMAT JSON",
1734 "source",
1735 ),
1736 (
1737 "CREATE SUBSOURCE \"materialize\".\"public\".\"sub\" (id int4) \
1738 OF SOURCE [u1 AS \"materialize\".\"public\".\"src\"]",
1739 "subsource",
1740 ),
1741 (
1742 "CREATE SUBSOURCE \"materialize\".\"public\".\"progress\" (id int4) \
1743 WITH (PROGRESS)",
1744 "subsource",
1745 ),
1746 (
1747 "CREATE SINK \"materialize\".\"public\".\"snk\" IN CLUSTER [u1] \
1748 FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1749 INTO KAFKA CONNECTION [u2 AS \"materialize\".\"public\".\"c\"] \
1750 (TOPIC 'tp') FORMAT JSON ENVELOPE DEBEZIUM",
1751 "sink",
1752 ),
1753 (
1754 "CREATE INDEX \"i\" IN CLUSTER [u1] \
1755 ON [u1 AS \"materialize\".\"public\".\"t\"] (\"a\")",
1756 "index",
1757 ),
1758 (
1759 "CREATE TYPE \"materialize\".\"public\".\"ty\" AS LIST (ELEMENT TYPE = int4)",
1760 "type",
1761 ),
1762 (
1763 "CREATE SECRET \"materialize\".\"public\".\"s\" AS 'x'",
1764 "secret",
1765 ),
1766 (
1767 "CREATE CONNECTION \"materialize\".\"public\".\"c\" \
1768 TO KAFKA (BROKER 'b', SECURITY PROTOCOL PLAINTEXT)",
1769 "connection",
1770 ),
1771 ];
1772 for (sql, expected) in cases {
1773 assert_eq!(item_type(sql).as_deref(), Ok(expected), "for {sql}");
1774 }
1775 }
1776
1777 #[mz_ore::test]
1781 #[cfg_attr(miri, ignore)] fn catalog_view_definition() {
1783 assert_eq!(view_definition("SELECT 1"), "SELECT 1;");
1786 assert_eq!(
1787 view_definition("WITH c AS (SELECT 1 AS a) SELECT a FROM c"),
1788 "WITH \"c\" AS (SELECT 1 AS \"a\") SELECT \"a\" FROM \"c\";"
1789 );
1790 assert_eq!(
1791 view_definition("SELECT 1 UNION ALL SELECT 2"),
1792 "SELECT 1 UNION ALL SELECT 2;"
1793 );
1794 assert_eq!(
1795 view_definition("SELECT (SELECT max(a) FROM [u1 AS \"materialize\".\"public\".\"t\"])"),
1796 "SELECT (SELECT \"max\"(\"a\") FROM [u1 AS \"materialize\".\"public\".\"t\"]);"
1797 );
1798 assert_eq!(
1801 view_definition(
1802 "SELECT \"a b\", \"héllo\", \"q\"\"x\" \
1803 FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1804 WHERE s = 'lit''eral' AND n = 42 ORDER BY 1"
1805 ),
1806 "SELECT \"a b\", \"héllo\", \"q\"\"x\" \
1807 FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1808 WHERE \"s\" = 'lit''eral' AND \"n\" = 42 ORDER BY 1;"
1809 );
1810 }
1811
1812 #[mz_ore::test]
1817 #[cfg_attr(miri, ignore)] fn catalog_view_definition_is_idempotent() {
1819 for query in [
1820 "SELECT 1",
1821 "WITH c AS (SELECT 1 AS a) SELECT a FROM c",
1822 "SELECT 1 UNION ALL SELECT 2",
1823 "SELECT \"a b\", \"q\"\"x\" FROM [u1 AS \"materialize\".\"public\".\"t\"] \
1824 WHERE s = 'lit''eral' ORDER BY 1",
1825 ] {
1826 let once = view_definition(query);
1827 assert_eq!(
1828 view_definition(&once),
1829 once,
1830 "not a fixed point for {query}"
1831 );
1832 }
1833 }
1834
1835 #[mz_ore::test]
1838 #[cfg_attr(miri, ignore)] fn catalog_table_source_id() {
1840 let from_source = super::item_details(
1841 "CREATE TABLE \"materialize\".\"public\".\"tbl\" \
1842 FROM SOURCE [u1 AS \"materialize\".\"public\".\"src\"] \
1843 (REFERENCE = \"topic\") FORMAT TEXT",
1844 )
1845 .expect("ok");
1846 assert_eq!(from_source, json!({ "type": "table", "source_id": "u1" }));
1847
1848 for sql in [
1849 "CREATE TABLE \"materialize\".\"public\".\"t\" (a int4)",
1850 "CREATE TABLE \"materialize\".\"public\".\"wht\" FROM WEBHOOK BODY FORMAT JSON",
1851 ] {
1852 assert_eq!(
1853 super::item_details(sql).expect("ok"),
1854 json!({ "type": "table" }),
1855 "for {sql}"
1856 );
1857 }
1858 }
1859
1860 #[mz_ore::test]
1864 #[cfg_attr(miri, ignore)] fn catalog_create_sql_errors() {
1866 assert_eq!(
1867 item_type("this is not sql"),
1868 Err(
1869 "failed to parse create_sql: Expected a keyword at the beginning of a statement, \
1870 found identifier \"this\""
1871 .to_string()
1872 )
1873 );
1874 assert_eq!(
1875 item_type("CREATE TABLE t (a int4); CREATE TABLE u (b int4)"),
1876 Err("expected a single statement, found 2".to_string())
1877 );
1878 assert_eq!(
1881 item_type("SELECT 1"),
1882 Err("not a CREATE item statement".to_string())
1883 );
1884 assert_eq!(
1887 item_type(
1888 "CREATE TABLE \"materialize\".\"public\".\"tbl\" \
1889 FROM SOURCE src (REFERENCE = \"topic\")"
1890 ),
1891 Err("unresolved item name".to_string())
1892 );
1893 }
1894
1895 #[mz_ore::test]
1898 #[cfg_attr(miri, ignore)] fn item_references_view() {
1900 let sql = "CREATE VIEW \"materialize\".\"public\".\"v\" AS \
1901 SELECT \"pg_catalog\".\"abs\"(\"a\"::[s20 AS \"pg_catalog\".\"int4\"]) \
1902 FROM [u1 AS \"materialize\".\"public\".\"t\"]";
1903 let out = super::item_references(sql).expect("ok");
1904 assert_eq!(
1905 out,
1906 json!({
1907 "ids": ["s20", "u1"],
1908 "named_funcs": [{"schema": "pg_catalog", "name": "abs"}],
1909 "named_types": [],
1910 "named_relations": [],
1911 })
1912 );
1913 }
1914
1915 #[mz_ore::test]
1916 #[cfg_attr(miri, ignore)] fn item_references_connection() {
1918 let sql = "CREATE CONNECTION \"materialize\".\"public\".\"kc\" TO KAFKA (\
1919 BROKER 'kafka:9092', \
1920 SSH TUNNEL [u5 AS \"materialize\".\"public\".\"ssh\"], \
1921 SASL PASSWORD = SECRET [u7 AS \"materialize\".\"public\".\"pw\"])";
1922 let out = super::item_references(sql).expect("ok");
1923 assert_eq!(out["ids"], json!(["u5", "u7"]));
1924 assert_eq!(out["named_funcs"], json!([]));
1925 }
1926
1927 #[mz_ore::test]
1928 #[cfg_attr(miri, ignore)] fn item_references_cte_shadowing() {
1930 let sql = "CREATE VIEW \"materialize\".\"public\".\"v\" AS \
1931 WITH \"c\" AS (SELECT * FROM [u1 AS \"materialize\".\"public\".\"t\"]) \
1932 SELECT * FROM \"c\"";
1933 let out = super::item_references(sql).expect("ok");
1934 assert_eq!(out["ids"], json!(["u1"]));
1935 assert_eq!(out["named_relations"], json!([]));
1936 }
1937
1938 #[mz_ore::test]
1939 #[cfg_attr(miri, ignore)] fn item_references_name_based_relation_surfaces() {
1941 let sql = "CREATE VIEW \"materialize\".\"public\".\"v\" AS \
1944 SELECT * FROM \"mz_catalog\".\"mz_tables\"";
1945 let out = super::item_references(sql).expect("ok");
1946 assert_eq!(
1947 out["named_relations"],
1948 json!([{"schema": "mz_catalog", "name": "mz_tables"}])
1949 );
1950 }
1951
1952 #[mz_ore::test]
1953 fn item_references_parse_error() {
1954 let err = super::item_references("NOT SQL").unwrap_err();
1955 assert!(
1956 err.contains("failed to parse"),
1957 "wrong error message: {err}"
1958 );
1959 }
1960}