Skip to main content

mz_catalog_decode/
create_sql.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Decoders for catalog `create_sql` strings.
11//!
12//! Each function parses a persisted `create_sql` and returns a JSON value with
13//! the fields a builtin catalog view projects. Errors are human-readable
14//! strings.
15
16use 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
31/// Parses `sql`, which must contain exactly one statement.
32fn 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
67/// Extracts the string form of a `WithOptionValue`, matching how the planner
68/// coerces option values. The blanket `TryFromValue<WithOptionValue<T>>` impl in
69/// `src/sql/src/plan/with_options.rs` accepts a quoted string, a bare identifier,
70/// and a 1-part unresolved item name, all yielding a string. Any other variant
71/// yields None.
72///
73/// `AstDisplay` does not re-quote bare identifiers or 1-part names, so those
74/// forms persist unquoted in `create_sql`. Matching only `Value::String` here
75/// would miss them and silently fall back to a default, so the catalog-raw
76/// parsers below route their string options through this helper.
77fn 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
88/// Parses a catalog `create_sql` string into a JSON object.
89///
90/// The returned JSON does not fully reflect the parsed SQL and instead contains only fields
91/// required by current callers.
92pub 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            // PostgreSQL appends a semicolon in `pg_views.definition`, we
109            // do the same for compatibility's sake.
110            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            // An old-syntax kafka source ingests into its own relation, so an
165            // omitted ENVELOPE means the default ENVELOPE NONE and the pre-MV
166            // packer reported 'none'. A new-syntax source (no progress
167            // subsource, hence no EXPOSE PROGRESS AS in create_sql) ingests
168            // nothing itself. Its envelopes live on the per-table exports, so
169            // its own envelope_type stays absent (SQL NULL), matching released
170            // behavior. `progress_subsource.is_some()` is the planner's own
171            // old-vs-new discriminator (see `OldSyntaxIngestion` in
172            // plan_create_source). Non-kafka sources carry no envelope either.
173            // See the `mz_sources.envelope_type` column.
174            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                        // Debezium sources with a single format spec implicitly use
195                        // the same format for both key and value.
196                        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        // Everything the mz_sinks, mz_kafka_sinks and mz_iceberg_sinks
239        // views read. The Rust side of each value lives in `Sink` and
240        // `StorageSinkConnection`, so those and this have to move together.
241        //
242        // NOTE: we bail below if a sink has no resolved `IN CLUSTER`, no
243        // `TOPIC` on kafka, or no `NAMESPACE`/`TABLE` on iceberg. Planning
244        // guarantees all four. But if one ever slipped through it would
245        // take down every view built on this function, not just the sink
246        // ones.
247        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                        // A key format only survives if the sink has a
280                        // `KEY`. Without one `kafka_sink_builder` throws
281                        // away the key half of a key/value spec, and does
282                        // not copy a bare spec over to the key either.
283                        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                        // The deprecated combined `format`. Only avro/avro
301                        // and json/json collapse to a single name.
302                        // Everything else, text/text and bytes/bytes
303                        // included, gets the composite form.
304                        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                    // The catalog connection, not the optional AWS one.
326                    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                    // Iceberg spells the envelope `MODE`, and has no format
347                    // columns at all.
348                    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        // NOTE: every statement that creates a catalog item needs an arm above. These
382        // catalog views run this over every item row before their type filter drops the
383        // unwanted rows, so one unclassified `create_sql` takes out `mz_objects`,
384        // `mz_indexes`, and every sibling view at once. The match is exhaustive to make
385        // that a compile error here, not a runtime failure.
386        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
455/// Extracts the catalog item references from a catalog `create_sql` string as a JSONB object.
456///
457/// - `ids`: array of catalog id strings from `[<id> AS <name>]` references.
458/// - `named_funcs`: array of function references such as "pg_catalog"."max".
459/// - `named_types`: array of type references such as "pg_catalog"."int4", exclusive from the ones
460///   in `ids`.
461/// - `named_relations`: array of relation references, exclusive from the ones in `ids`.
462pub 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/// Minimal decoder for `ProtoPostgresSourcePublicationDetails`. The
489/// canonical proto lives in `mz-storage-types`, which depends on this
490/// crate via `mz-expr`, so we redeclare the two tags we read here. Upstream tag
491/// renumbers slip past silently. The `mz_postgres_sources` lockdown
492/// SLTs catch them.
493#[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
501/// Extracts postgres source publication details (slot, timeline_id) from a
502/// catalog `create_sql`. Returns:
503///
504/// - jsonb `{"slot": <text>, "timeline_id": <u64 | null>}` for `CREATE SOURCE ... FROM POSTGRES
505///   CONNECTION ... (DETAILS = ...)` statements.
506/// - jsonb `null` for any other statement.
507///
508/// Errors if the statement fails to parse, is a postgres source without
509/// a `DETAILS` option, or if the `DETAILS` value can't be hex- and
510/// proto-decoded.
511pub 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
545/// Extracts kafka source configuration (topic, group id prefix, connection
546/// id) from a catalog `create_sql`. Returns:
547///
548/// - jsonb `{"topic": <text>, "group_id_prefix": <text | null>, "connection_id": <text>}` for
549///   `CREATE SOURCE ... FROM KAFKA CONNECTION ... (TOPIC = ..., [GROUP ID PREFIX = ...])`
550///   statements.
551/// - jsonb `null` for any other statement.
552///
553/// Errors if the statement fails to parse, is a kafka source without a
554/// `TOPIC` option, or references an unresolved connection name (i.e. one
555/// that hasn't been through purification).
556pub 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
594/// Extracts source-export (source table) metadata from a catalog `create_sql`.
595///
596/// Returns, for a `CREATE TABLE ... FROM SOURCE` or a non-progress
597/// `CREATE SUBSOURCE ... OF SOURCE ...` statement:
598///
599/// ```json
600/// {
601///   "source_id": "<parent source item id>",
602///   "external_reference": ["part1", "part2", ...],
603///   "envelope_type": <text | null>,
604///   "key_format": <text | null>,
605///   "value_format": <text | null>
606/// }
607/// ```
608///
609/// `envelope_type`, `key_format`, and `value_format` are always null for a
610/// `CREATE SUBSOURCE` (the postgres/mysql/sql-server exports that use the old
611/// subsource syntax carry neither format nor envelope). They may also be null
612/// for a `CREATE TABLE ... FROM SOURCE` that omits FORMAT/ENVELOPE.
613///
614/// Returns jsonb `null` for progress subsources and for any statement that is
615/// not a source export. The caller distinguishes the four source-table views
616/// by joining `source_id` against `mz_sources` and filtering on the parent's
617/// type, so this helper stays connection-type agnostic.
618///
619/// Errors if the statement fails to parse, references an unresolved item name,
620/// or is a non-progress subsource missing its OF SOURCE or EXTERNAL REFERENCE.
621///
622/// The `key_format`/`value_format` derivation mirrors the runtime
623/// `DataSourceDesc::formats()` that the removed `pack_kafka_source_tables_update`
624/// packer read. A bare FORMAT only carries a key when it resolves to an
625/// encoding that has one, which among bare formats is only Avro or Protobuf
626/// read from a Confluent Schema Registry whose purified seed carries a key
627/// schema. A KEY FORMAT ... VALUE FORMAT ... spec always carries both.
628pub fn source_export_details(a: &str) -> Result<serde_json::Value, String> {
629    // A bare FORMAT resolves to an encoding with a key only for Avro or
630    // Protobuf read from a schema registry whose purified seed carries a key
631    // schema. Every other bare format is value-only.
632    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            // Progress subsources track ingestion progress and are not
711            // source tables. They have no external reference.
712            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
751/// Extracts connection-detail metadata from a catalog `create_sql`.
752///
753/// Returns a per-connection-type object with the fields that the
754/// `mz_kafka_connections`, `mz_ssh_tunnel_connections`, and `mz_aws_connections`
755/// builtin views need. For everything else (other connection types, including
756/// aws-privatelink whose only detail is context-derived, and non-connection
757/// statements) it returns jsonb `null`, so callers filter on `IS NOT NULL` and
758/// gate on the connection type separately (via
759/// `parse_catalog_create_sql(...)->>'connection_type'`, the way `mz_connections`
760/// already does).
761///
762/// The shape per type:
763///
764/// ```json
765/// // kafka
766/// { "brokers": ["host:port", ...], "progress_topic": <text | null> }
767/// // ssh-tunnel
768/// { "public_key_1": "<text>", "public_key_2": "<text>" }
769/// // aws
770/// {
771///   "auth_kind": "credentials" | "assume-role",
772///   "endpoint": <text | null>, "region": <text | null>,
773///   "access_key_id": <text | null>, "access_key_id_secret_id": <text | null>,
774///   "secret_access_key_secret_id": <text | null>,
775///   "session_token": <text | null>, "session_token_secret_id": <text | null>,
776///   "assume_role_arn": <text | null>, "assume_role_session_name": <text | null>
777/// }
778/// ```
779///
780/// `progress_topic` is null when the connection does not set an explicit
781/// `PROGRESS TOPIC`. The default (`_materialize-progress-<env>-<conn_id>`) is
782/// reconstructed by the view, not here, because it needs the environment id and
783/// the connection's own id. Values derived only from environment context
784/// (AWS principal, external id, trust policy, privatelink principal) are also
785/// left to the view. This keeps the helper a pure function of the `create_sql`.
786///
787/// For aws, an option is either an inline value or a secret reference. Inline
788/// values land in `access_key_id`/`session_token`; a secret reference lands in
789/// the matching `*_secret_id` as the referenced secret's catalog item id (the
790/// persisted `create_sql` stores resolved references as `[uNNN AS name]`).
791/// `auth_kind` is `assume-role` when `ASSUME ROLE ARN` is present, else
792/// `credentials`, matching the `AwsAuth` variant the removed packer read.
793///
794/// Errors if the statement fails to parse.
795pub fn connection_details(a: &str) -> Result<serde_json::Value, String> {
796    // The persisted `create_sql` stores an inline broker as a single `BROKER`
797    // option and a broker list as a `BROKERS (...)` sequence. Either way we
798    // only need the addresses, which are present regardless of the tunnel
799    // (direct, SSH, or PrivateLink).
800    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    // The catalog id of the secret a `SECRET ...` option references. Resolved
832    // references persist as `RawItemName::Id`, so an unresolved name yields
833    // None (the same treatment [`item_id`] gives item names).
834    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                // Planning coerces an empty ENDPOINT to None (see
873                // `src/sql/src/plan/statement/ddl/connection.rs`), so the
874                // removed packer wrote NULL for `ENDPOINT = ''`. Match that.
875                "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    /// Encode the two proto fields our decoder cares about, using the same
901    /// tag numbering as the canonical proto.
902    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    // --- parse_postgres_source_details ---------------------------------------
935
936    #[mz_ore::test]
937    #[cfg_attr(miri, ignore)] // error: unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
938    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        // Pre-2024 sources have no timeline_id field. The decoder must
948        // surface that as JSON null, not error.
949        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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
968    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        // Valid hex, garbage bytes. Prost decoding fails on unexpected wire
997        // format.
998        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    // --- parse_kafka_source_details ------------------------------------------
1007
1008    #[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        // Planning accepts a bare identifier for TOPIC / GROUP ID PREFIX, and it
1057        // persists unquoted in create_sql. Matching only quoted strings would
1058        // drop TOPIC and error the whole mz_kafka_source_tables view.
1059        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        // A bare-name connection reference never happens after purification,
1071        // but the decoder must reject it explicitly rather than silently
1072        // dropping the connection_id.
1073        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    // --- parse_source_export_details -----------------------------------------
1084
1085    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        // Postgres/mysql/sql-server tables carry a multi-part external
1096        // reference and no format or envelope.
1097        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        // A bare non-registry FORMAT is value-only: no key format.
1114        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        // Omitting ENVELOPE persists as absent in create_sql, so this
1131        // source-type-agnostic helper reports null. The mz_kafka_source_tables
1132        // view is responsible for defaulting kafka's null envelope to 'none'.
1133        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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1159    fn export_table_kafka_bare_avro_seed_with_key() {
1160        // A bare Avro CSR format whose seed carries a key schema resolves to
1161        // an encoding with a key, so key_format mirrors value_format. This is
1162        // the upsert/debezium path.
1163        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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1184    fn export_table_kafka_bare_avro_seed_without_key() {
1185        // A bare Avro CSR seed with only a value schema is value-only.
1186        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        // Old-syntax subsource: external reference lives in a WITH option, and
1208        // there is never a format or envelope.
1209        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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1235    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    // --- parse_connection_details --------------------------------------------
1253
1254    #[mz_ore::test]
1255    fn connection_kafka_single_broker_default_progress() {
1256        // No explicit PROGRESS TOPIC: the helper leaves it null and the view
1257        // reconstructs the default.
1258        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        // Inline ACCESS KEY ID, secret SECRET ACCESS KEY. Assume-role columns
1317        // stay null and auth_kind is credentials.
1318        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        // Every credential provided as a secret reference lands in the matching
1342        // *_secret_id column as the referenced secret's catalog id.
1343        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        // Assume-role sets auth_kind and the assume-role columns; credential
1369        // columns stay null.
1370        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        // A bare identifier PROGRESS TOPIC persists unquoted in create_sql. The
1394        // helper must surface it, else the view substitutes the default topic.
1395        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        // Planning accepts bare identifiers for these options and persists them
1405        // unquoted. The helper must surface them, not fall back to NULL.
1406        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        // Planning coerces ENDPOINT = '' to None, so the packer wrote NULL. The
1419        // view must match, not report an empty string.
1420        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        // A connection type without a detail view (postgres) yields null.
1429        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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1437    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    // --- parse_catalog_create_sql envelope_type ------------------------------
1444
1445    #[mz_ore::test]
1446    fn catalog_kafka_old_syntax_omitted_envelope_defaults_none() {
1447        // An old-syntax kafka source (carries EXPOSE PROGRESS AS) ingests into
1448        // its own relation, so an omitted ENVELOPE means the default NONE, which
1449        // the pre-MV packer reported as 'none'.
1450        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        // A new-syntax kafka source has no progress subsource (no EXPOSE PROGRESS
1473        // AS). It ingests nothing itself. Envelopes live on the per-table exports,
1474        // so its own envelope_type stays absent (SQL NULL).
1475        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        // Non-kafka sources carry no envelope, so envelope_type stays absent
1485        // (SQL NULL), not 'none'.
1486        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    // --- parse_catalog_create_sql, CreateSink arm ----------------------------
1493
1494    const AVRO_FORMAT: &str = "FORMAT AVRO USING CONFLUENT SCHEMA REGISTRY CONNECTION \
1495         [u12 AS \"materialize\".\"public\".\"csr_conn\"]";
1496
1497    /// A persisted kafka-sink `create_sql`: resolved names, and a `TOPIC` that
1498    /// planning guarantees.
1499    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        // A bare format applies to the key too once the sink has a KEY, which
1544        // is what makes the deprecated `format` column collapse to `avro`.
1545        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        // Only avro/avro and json/json collapse, so a keyed text sink reports
1566        // the composite form even though both halves are `text`.
1567        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        // `kafka_sink_builder` ignores the key half of the format spec when the
1595        // sink has no KEY, so neither `key_format` nor the composite `format`
1596        // may reflect it.
1597        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                // The catalog connection, never the AWS connection (u21).
1626                "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        // `append` is reachable only through an iceberg sink's MODE.
1639        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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1655    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    // --- parse_catalog_create_sql --------------------------------------------
1662
1663    /// `type` for a `create_sql`, or the error message if parsing failed.
1664    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    /// `definition` for a `CREATE VIEW` whose query is `query`.
1682    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_tables` and `mz_views` select rows by
1693    /// `parse_catalog_create_sql(...)->>'type'`, and the function runs over
1694    /// every `Item` row in the catalog, so a statement kind that changes its
1695    /// reported type silently gains or loses rows in those relations. Pin the
1696    /// type of every kind the catalog can hold.
1697    #[mz_ore::test]
1698    #[cfg_attr(miri, ignore)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1699    fn catalog_item_type_per_statement_kind() {
1700        let cases = [
1701            (
1702                "CREATE TABLE \"materialize\".\"public\".\"t\" (a int4)",
1703                "table",
1704            ),
1705            // A table created from a source, and a webhook table, are both
1706            // `table`, so both land in mz_tables.
1707            (
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_views.definition` is produced here. It used to be produced by
1778    /// `pack_view_update` in the adapter, so the exact rendering is a
1779    /// compatibility surface: `pg_views.definition` reads it.
1780    #[mz_ore::test]
1781    #[cfg_attr(miri, ignore)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1782    fn catalog_view_definition() {
1783        // Identifiers and function names come back fully quoted, literals
1784        // untouched, and PostgreSQL's trailing semicolon is appended.
1785        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        // Identifiers needing quotes, an embedded double quote, non-ASCII, an
1799        // embedded single quote in a literal, and ORDER BY all survive.
1800        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    /// The rendering must be a fixed point: `pg_views` consumers re-issue
1813    /// `definition` as the body of a new view, so a second pass through the
1814    /// parser has to produce the identical string. The trailing `;` is part of
1815    /// what gets re-parsed.
1816    #[mz_ore::test]
1817    #[cfg_attr(miri, ignore)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1818    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_tables.source_id` comes from this key. A table with no source must
1836    /// omit it entirely, so the MV's `->>'source_id'` yields SQL NULL.
1837    #[mz_ore::test]
1838    #[cfg_attr(miri, ignore)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1839    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    /// Every error here is fatal to the whole of `mz_tables`/`mz_views`, not to
1861    /// one row: the MVs call this function inside their `WHERE` clause, so an
1862    /// item the parser rejects makes the relation unreadable for everyone.
1863    #[mz_ore::test]
1864    #[cfg_attr(miri, ignore)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1865    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        // A statement that is not a CREATE of a catalog item, e.g. if a future
1879        // change persists something else in an Item record.
1880        assert_eq!(
1881            item_type("SELECT 1"),
1882            Err("not a CREATE item statement".to_string())
1883        );
1884        // Catalog `create_sql` always names items by id. An unresolved name
1885        // means the record was written wrong.
1886        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    // --- parse_catalog_item_references ----------------------------------------
1896
1897    #[mz_ore::test]
1898    #[cfg_attr(miri, ignore)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1899    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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1917    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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1929    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)] // error: unsupported operation: can't call foreign function `rust_psm_stack_pointer` on OS `linux`
1940    fn item_references_name_based_relation_surfaces() {
1941        // Stored SQL always prints relations as id references. A name-based
1942        // relation must surface in "named_relations" rather than vanish.
1943        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}