diff --git a/crates/dbx-core/src/schema.rs b/crates/dbx-core/src/schema.rs index 54c61a247..144a59a13 100644 --- a/crates/dbx-core/src/schema.rs +++ b/crates/dbx-core/src/schema.rs @@ -3382,6 +3382,18 @@ mod ddl_tests { assert!(ddl.contains("COMMENT ON COLUMN \"public\".\"users\".\"display_name\" IS 'User''s display name';")); } + #[test] + fn postgres_table_ddl_includes_generated_identity() { + let mut id = column("id", "integer"); + id.is_nullable = false; + id.is_primary_key = true; + id.extra = Some("generated by default as identity".to_string()); + + let ddl = render_postgres_table_ddl("public", "users", &[id], &[], &[]); + + assert!(ddl.contains("\"id\" integer generated by default as identity NOT NULL"), "ddl: {ddl}"); + } + #[test] fn sqlserver_table_ddl_includes_column_comments() { let mut display_name = column("display]name", "nvarchar(100)"); @@ -3469,11 +3481,21 @@ pub fn render_postgres_table_ddl( .iter() .map(|c| { let mut line = format!(" {} {}", pg_ident(&c.name), c.data_type); + let generated_clause = c + .extra + .as_deref() + .map(str::trim) + .filter(|value| !value.is_empty() && value.to_ascii_lowercase().starts_with("generated ")); + if let Some(extra) = generated_clause { + line.push_str(&format!(" {extra}")); + } if !c.is_nullable { line.push_str(" NOT NULL"); } - if let Some(ref def) = c.column_default { - line.push_str(&format!(" DEFAULT {def}")); + if generated_clause.is_none() { + if let Some(ref def) = c.column_default { + line.push_str(&format!(" DEFAULT {def}")); + } } line }) diff --git a/crates/dbx-core/src/transfer.rs b/crates/dbx-core/src/transfer.rs index 4e7bb511f..00a95947c 100644 --- a/crates/dbx-core/src/transfer.rs +++ b/crates/dbx-core/src/transfer.rs @@ -9,6 +9,7 @@ use crate::db::mongo_driver::MongoDocumentResult; use crate::models::connection::DatabaseType; use crate::object_source_sql::{build_executable_object_source_statements, EditableObjectSourceSqlInput}; use crate::query::{agent_execute_query_params, should_discard_pool_after_error, QueryExecutionOptions}; +use crate::sql::split_sql_statements; #[cfg(feature = "duckdb-bundled")] use crate::sql::starts_with_executable_sql_keyword; use crate::sql_dialect::{qualified_transfer_table, quote_transfer_identifier}; @@ -142,6 +143,17 @@ fn is_postgres_sequence_default(default_value: Option<&str>) -> bool { default_value.is_some_and(|value| value.to_ascii_lowercase().contains("nextval(")) } +fn is_postgres_generated_extra(extra: Option<&str>) -> bool { + extra.is_some_and(|value| value.trim().to_ascii_lowercase().starts_with("generated ")) +} + +fn is_postgres_identity_extra(extra: Option<&str>) -> bool { + extra.is_some_and(|value| { + let normalized = value.trim().to_ascii_lowercase(); + normalized.starts_with("generated ") && normalized.contains(" identity") + }) +} + fn rewrite_postgres_schema_qualified_references(input: &str, source_schema: &str, target_schema: &str) -> String { if source_schema.trim().is_empty() || source_schema == target_schema { return input.to_string(); @@ -183,6 +195,11 @@ fn postgres_default_clause( if !is_postgres_compat_transfer(source_db, target_db) { return None; } + if is_postgres_generated_extra(column.extra.as_deref()) { + if let Some(extra) = column.extra.as_deref().map(str::trim).filter(|value| !value.is_empty()) { + return Some(extra.to_string()); + } + } let default_value = column.column_default.as_deref()?.trim(); if default_value.is_empty() { return None; @@ -423,8 +440,13 @@ fn generate_postgres_index_ddl(indexes: &[db::IndexInfo], table: &str, schema: & statements } -fn generate_postgres_foreign_key_ddl(foreign_keys: &[db::ForeignKeyInfo], table: &str, schema: &str) -> Vec { - let full_table = qualified_table(table, schema, &DatabaseType::Postgres); +fn generate_postgres_foreign_key_ddl( + foreign_keys: &[db::ForeignKeyInfo], + table: &str, + source_schema: &str, + target_schema: &str, +) -> Vec { + let full_table = qualified_table(table, target_schema, &DatabaseType::Postgres); let mut grouped: HashMap<&str, Vec<&db::ForeignKeyInfo>> = HashMap::new(); let mut order = Vec::new(); @@ -450,7 +472,11 @@ fn generate_postgres_foreign_key_ddl(foreign_keys: &[db::ForeignKeyInfo], table: .map(|foreign_key| quote_identifier(&foreign_key.ref_column, &DatabaseType::Postgres)) .collect::>() .join(", "); - let referenced_schema = group[0].ref_schema.as_deref().unwrap_or(schema); + let referenced_schema = match group[0].ref_schema.as_deref() { + Some(ref_schema) if ref_schema == source_schema => target_schema, + Some(ref_schema) => ref_schema, + None => target_schema, + }; let referenced_table = qualified_table(&group[0].ref_table, referenced_schema, &DatabaseType::Postgres); statements.push(format!( "ALTER TABLE {full_table} ADD CONSTRAINT {} FOREIGN KEY ({columns}) REFERENCES {referenced_table} ({ref_columns})", @@ -465,7 +491,10 @@ fn generate_postgres_sequence_sync_sql(columns: &[db::ColumnInfo], table: &str, let full_table = qualified_table(table, schema, &DatabaseType::Postgres); columns .iter() - .filter(|column| is_postgres_sequence_default(column.column_default.as_deref())) + .filter(|column| { + is_postgres_sequence_default(column.column_default.as_deref()) + || is_postgres_identity_extra(column.extra.as_deref()) + }) .map(|column| { let quoted_column = quote_identifier(&column.name, &DatabaseType::Postgres); format!( @@ -582,7 +611,7 @@ fn generate_postgres_materialized_view_ddls(view: &PostgresMaterializedViewSourc ] } -fn rewrite_postgres_routine_schema(source: &str, target_schema: &str) -> Option { +fn rewrite_postgres_routine_schema(source: &str, source_schema: &str, target_schema: &str) -> Option { let re = Regex::new( r#"(?is)^(\s*CREATE\s+(?:OR\s+REPLACE\s+)?(?:(?:NON)?EDITIONABLE\s+)?(?:FUNCTION|PROCEDURE)\s+)((?:"(?:""|[^"])+"|[A-Za-z_][\w$]*)(?:\s*\.\s*(?:"(?:""|[^"])+"|[A-Za-z_][\w$]*))?)"#, ) @@ -602,7 +631,8 @@ fn rewrite_postgres_routine_schema(source: &str, target_schema: &str) -> Option< quote_identifier(target_schema, &DatabaseType::Postgres), quote_identifier(name, &DatabaseType::Postgres) ); - Some(format!("{}{}{}{}", &source[..full.start()], prefix, replacement, &source[full.end()..])) + let rewritten = format!("{}{}{}{}", &source[..full.start()], prefix, replacement, &source[full.end()..]); + Some(rewrite_postgres_schema_qualified_references(&rewritten, source_schema, target_schema)) } fn rewrite_postgres_trigger_table_schema( @@ -624,10 +654,11 @@ fn rewrite_postgres_trigger_table_schema( ]; for pattern in candidate_patterns { if source.contains(&pattern) { - return source.replacen(&pattern, &format!(" ON {qualified_target_table} "), 1); + let rewritten = source.replacen(&pattern, &format!(" ON {qualified_target_table} "), 1); + return rewrite_postgres_schema_qualified_references(&rewritten, source_schema, target_schema); } } - source.to_string() + rewrite_postgres_schema_qualified_references(source, source_schema, target_schema) } pub fn escape_value(val: &serde_json::Value, db_type: &DatabaseType) -> String { @@ -1468,11 +1499,25 @@ fn can_reuse_source_table_ddl( ) -> bool { preserves_target_table_name && !matches!(target_db_type, DatabaseType::ClickHouse) - && ((source_db_type == target_db_type) + && (source_db_type == target_db_type || (is_mysql_family_target(source_db_type) && is_mysql_family_target(target_db_type)) || (is_postgres_family_target(source_db_type) && is_postgres_family_target(target_db_type))) } +fn rewrite_transfer_source_table_ddl( + sql: &str, + source_schema: &str, + target_schema: &str, + source_db_type: &DatabaseType, + target_db_type: &DatabaseType, +) -> String { + if is_postgres_family_target(source_db_type) && is_postgres_family_target(target_db_type) { + rewrite_postgres_schema_qualified_references(sql, source_schema, target_schema) + } else { + sql.to_string() + } +} + #[allow(clippy::too_many_arguments)] fn generate_transfer_write_sql( mode: &TransferMode, @@ -1900,6 +1945,63 @@ pub async fn execute_on_pool(state: &AppState, pool_key: &str, sql: &str) -> Res execute_on_pool_with_max_rows(state, pool_key, sql, None).await } +async fn execute_transfer_ddl_on_pool( + state: &AppState, + pool_key: &str, + sql: &str, + db_type: &DatabaseType, +) -> Result<(), String> { + for statement in transfer_ddl_statements(sql, db_type) { + execute_on_pool(state, pool_key, &statement).await?; + } + Ok(()) +} + +fn transfer_ddl_statements(sql: &str, db_type: &DatabaseType) -> Vec { + if matches!(db_type, DatabaseType::Postgres) { + let statements = split_sql_statements(sql); + if statements.is_empty() { + vec![sql.trim().to_string()] + } else { + statements + .into_iter() + .map(|statement| sanitize_postgres_transfer_ddl_statement(&statement)) + .filter(|statement| !is_postgres_post_table_index_statement(statement)) + .collect() + } + } else { + vec![sql.to_string()] + } +} + +fn sanitize_postgres_transfer_ddl_statement(statement: &str) -> String { + if !statement.trim_start().to_ascii_uppercase().starts_with("CREATE TABLE ") { + return statement.to_string(); + } + + let mut lines: Vec = Vec::new(); + for line in statement.lines() { + if line.to_ascii_uppercase().contains(" FOREIGN KEY ") { + if let Some(previous) = lines.last_mut() { + let trimmed_len = previous.trim_end_matches(char::is_whitespace).len(); + if previous[..trimmed_len].ends_with(',') { + previous.truncate(trimmed_len - 1); + } + } + continue; + } + lines.push(line.to_string()); + } + lines.join("\n") +} + +fn is_postgres_post_table_index_statement(statement: &str) -> bool { + let normalized = statement.trim_start().to_ascii_uppercase(); + normalized.starts_with("CREATE INDEX ") + || normalized.starts_with("CREATE UNIQUE INDEX ") + || normalized.starts_with("COMMENT ON INDEX ") +} + pub async fn execute_on_pool_with_max_rows( state: &AppState, pool_key: &str, @@ -2512,11 +2614,11 @@ async fn get_postgres_grant_statements_for_transfer( relation_grants AS ( \ SELECT format( \ 'GRANT %s ON %s %I.%I TO %s%s', \ - string_agg(a.privilege_type, ', ' ORDER BY a.privilege_type), \ + string_agg(privilege_type, ', ' ORDER BY privilege_type), \ CASE WHEN relkind = 'S' THEN 'SEQUENCE' ELSE 'TABLE' END, \ {target_schema}, relname, \ - CASE WHEN a.grantee = 0 THEN 'PUBLIC' ELSE quote_ident(grantee.rolname) END, \ - CASE WHEN bool_or(a.is_grantable) THEN ' WITH GRANT OPTION' ELSE '' END \ + CASE WHEN grantee = 0 THEN 'PUBLIC' ELSE quote_ident(rolname) END, \ + CASE WHEN bool_or(is_grantable) THEN ' WITH GRANT OPTION' ELSE '' END \ ) AS stmt \ FROM ( \ SELECT c.relname, c.relkind, a.grantee, a.privilege_type, a.is_grantable, grantee.rolname \ @@ -2531,11 +2633,11 @@ async fn get_postgres_grant_statements_for_transfer( routine_grants AS ( \ SELECT format( \ 'GRANT %s ON %s %I.%I(%s) TO %s%s', \ - string_agg(a.privilege_type, ', ' ORDER BY a.privilege_type), \ + string_agg(privilege_type, ', ' ORDER BY privilege_type), \ CASE WHEN prokind = 'p' THEN 'PROCEDURE' ELSE 'FUNCTION' END, \ {target_schema}, proname, identity_args, \ - CASE WHEN a.grantee = 0 THEN 'PUBLIC' ELSE quote_ident(grantee.rolname) END, \ - CASE WHEN bool_or(a.is_grantable) THEN ' WITH GRANT OPTION' ELSE '' END \ + CASE WHEN grantee = 0 THEN 'PUBLIC' ELSE quote_ident(rolname) END, \ + CASE WHEN bool_or(is_grantable) THEN ' WITH GRANT OPTION' ELSE '' END \ ) AS stmt \ FROM ( \ SELECT p.proname, p.prokind, pg_get_function_identity_arguments(p.oid) AS identity_args, a.grantee, a.privilege_type, a.is_grantable, grantee.rolname \ @@ -3027,7 +3129,7 @@ where let can_reuse_source_ddl = can_reuse_source_table_ddl(source_db_type, target_db_type, preserves_target_table_name); let ddl = if can_reuse_source_ddl { - crate::schema::get_table_ddl_core( + let source_ddl = crate::schema::get_table_ddl_core( &state, &request.source_connection_id, &request.source_database, @@ -3046,7 +3148,14 @@ where source_db_type, table_comment.as_deref(), ) - }) + }); + rewrite_transfer_source_table_ddl( + &source_ddl, + &request.source_schema, + &request.target_schema, + source_db_type, + target_db_type, + ) } else { generate_create_table_ddl( &columns, @@ -3059,7 +3168,7 @@ where ) }; log::info!("[transfer] creating target table: {}", ddl.chars().take(200).collect::()); - let table_exists = match execute_on_pool(state, target_pool_key, &ddl).await { + let table_exists = match execute_transfer_ddl_on_pool(state, target_pool_key, &ddl, target_db_type).await { Ok(_) => true, Err(e) => { let err_lower = e.to_lowercase(); @@ -3212,8 +3321,12 @@ where .await .map_err(|e| format!("Failed to create PostgreSQL index for {target_table}: {e}"))?; } - for statement in generate_postgres_foreign_key_ddl(&source_foreign_keys, &target_table, &request.target_schema) - { + for statement in generate_postgres_foreign_key_ddl( + &source_foreign_keys, + &target_table, + &request.source_schema, + &request.target_schema, + ) { execute_on_pool(state, target_pool_key, &statement) .await .map_err(|e| format!("Failed to create PostgreSQL foreign key for {target_table}: {e}"))?; @@ -3400,7 +3513,7 @@ where let rewritten_source = match object.object_type { db::ObjectSourceKind::View | db::ObjectSourceKind::MaterializedView => object.source.clone(), db::ObjectSourceKind::Procedure | db::ObjectSourceKind::Function => { - rewrite_postgres_routine_schema(&object.source, &request.target_schema) + rewrite_postgres_routine_schema(&object.source, &request.source_schema, &request.target_schema) .unwrap_or_else(|| object.source.clone()) } db::ObjectSourceKind::Sequence | db::ObjectSourceKind::Package | db::ObjectSourceKind::PackageBody => { @@ -3822,6 +3935,52 @@ mod tests { assert!(stmts[2].contains("COMMENT ON COLUMN \"public\".\"items\".\"name\" IS '名称'")); } + #[test] + fn postgres_transfer_ddl_splits_reused_multi_statement_table_ddl() { + let ddl = + "CREATE TABLE \"public\".\"items\" (\"id\" integer);\nCOMMENT ON TABLE \"public\".\"items\" IS 'items';"; + + let statements = transfer_ddl_statements(ddl, &DatabaseType::Postgres); + + assert_eq!( + statements, + vec![ + "CREATE TABLE \"public\".\"items\" (\"id\" integer)".to_string(), + "COMMENT ON TABLE \"public\".\"items\" IS 'items'".to_string(), + ] + ); + } + + #[test] + fn postgres_transfer_ddl_skips_reused_index_statements() { + let ddl = "CREATE TABLE \"public\".\"items\" (\"id\" integer);\n\ + CREATE INDEX \"items_lower_idx\" ON \"public\".\"items\" USING btree (\"lower(name)\");\n\ + COMMENT ON INDEX \"public\".\"items_lower_idx\" IS 'lookup';"; + + let statements = transfer_ddl_statements(ddl, &DatabaseType::Postgres); + + assert_eq!(statements, vec!["CREATE TABLE \"public\".\"items\" (\"id\" integer)".to_string()]); + } + + #[test] + fn postgres_transfer_ddl_removes_inline_foreign_keys_from_reused_table_ddl() { + let ddl = "CREATE TABLE \"public\".\"audit_logs\" (\n \"id\" integer,\n \"user_id\" integer,\n CONSTRAINT \"audit_logs_user_id_fkey\" FOREIGN KEY (\"user_id\") REFERENCES \"users\"(\"id\")\n);"; + + let statements = transfer_ddl_statements(ddl, &DatabaseType::Postgres); + + assert_eq!( + statements, + vec!["CREATE TABLE \"public\".\"audit_logs\" (\n \"id\" integer,\n \"user_id\" integer\n)".to_string()] + ); + } + + #[test] + fn non_postgres_transfer_ddl_keeps_statement_text_intact() { + let ddl = "CREATE TABLE `items` (`id` int);\nALTER TABLE `items` COMMENT = 'items';"; + + assert_eq!(transfer_ddl_statements(ddl, &DatabaseType::Mysql), vec![ddl.to_string()]); + } + #[test] fn clickhouse_comment_ddl_uses_alter_table() { let cols = vec![db::ColumnInfo { comment: Some("日志消息".to_string()), ..test_column("message", "text") }]; @@ -3895,12 +4054,25 @@ mod tests { } #[test] - fn clickhouse_transfer_does_not_reuse_source_table_ddl() { + fn transfer_reuses_source_table_ddl_only_when_target_shape_matches() { assert!(!can_reuse_source_table_ddl(&DatabaseType::ClickHouse, &DatabaseType::ClickHouse, true)); assert!(can_reuse_source_table_ddl(&DatabaseType::Postgres, &DatabaseType::Postgres, true)); assert!(!can_reuse_source_table_ddl(&DatabaseType::Postgres, &DatabaseType::Postgres, false)); } + #[test] + fn postgres_transfer_reused_table_ddl_rewrites_target_schema() { + let ddl = + "CREATE TABLE \"src\".\"items\" (\"id\" integer);\nCOMMENT ON COLUMN \"src\".\"items\".\"id\" IS 'id';"; + + let rewritten = + rewrite_transfer_source_table_ddl(ddl, "src", "dst", &DatabaseType::Postgres, &DatabaseType::Postgres); + + assert!(rewritten.contains("CREATE TABLE \"dst\".\"items\"")); + assert!(rewritten.contains("COMMENT ON COLUMN \"dst\".\"items\".\"id\"")); + assert!(!rewritten.contains("\"src\".\"items\"")); + } + #[test] fn hive_create_table_uses_hive_friendly_columns() { let cols = vec![ @@ -4211,7 +4383,7 @@ mod tests { ]; let index_sql = generate_postgres_index_ddl(&indexes, "users", "public"); - let foreign_key_sql = generate_postgres_foreign_key_ddl(&foreign_keys, "orders", "public"); + let foreign_key_sql = generate_postgres_foreign_key_ddl(&foreign_keys, "orders", "public", "archive"); assert_eq!( index_sql, @@ -4223,7 +4395,7 @@ mod tests { assert_eq!( foreign_key_sql, vec![ - "ALTER TABLE \"public\".\"orders\" ADD CONSTRAINT \"orders_user_id_fkey\" FOREIGN KEY (\"user_id\", \"tenant_id\") REFERENCES \"public\".\"users\" (\"id\", \"tenant_id\")".to_string() + "ALTER TABLE \"archive\".\"orders\" ADD CONSTRAINT \"orders_user_id_fkey\" FOREIGN KEY (\"user_id\", \"tenant_id\") REFERENCES \"archive\".\"users\" (\"id\", \"tenant_id\")".to_string() ] ); } @@ -4231,18 +4403,44 @@ mod tests { #[test] fn postgres_sequence_sync_sql_uses_table_max_values() { let sql = generate_postgres_sequence_sync_sql( - &[db::ColumnInfo { - name: "id".to_string(), - data_type: "integer".to_string(), - is_nullable: false, - column_default: Some("nextval('public.users_id_seq'::regclass)".to_string()), - is_primary_key: true, - extra: None, - comment: None, - numeric_precision: None, - numeric_scale: None, - character_maximum_length: None, - }], + &[ + db::ColumnInfo { + name: "id".to_string(), + data_type: "integer".to_string(), + is_nullable: false, + column_default: Some("nextval('public.users_id_seq'::regclass)".to_string()), + is_primary_key: true, + extra: None, + comment: None, + numeric_precision: None, + numeric_scale: None, + character_maximum_length: None, + }, + db::ColumnInfo { + name: "identity_id".to_string(), + data_type: "integer".to_string(), + is_nullable: false, + column_default: None, + is_primary_key: false, + extra: Some("generated by default as identity".to_string()), + comment: None, + numeric_precision: None, + numeric_scale: None, + character_maximum_length: None, + }, + db::ColumnInfo { + name: "computed_id".to_string(), + data_type: "integer".to_string(), + is_nullable: false, + column_default: None, + is_primary_key: false, + extra: Some("generated always as (identity_id + 1) stored".to_string()), + comment: None, + numeric_precision: None, + numeric_scale: None, + character_maximum_length: None, + }, + ], "users", "public", ); @@ -4250,7 +4448,8 @@ mod tests { assert_eq!( sql, vec![ - "SELECT setval(pg_get_serial_sequence('\"public\".\"users\"', 'id'), GREATEST(COALESCE(MAX(\"id\"), 0), 1), MAX(\"id\") IS NOT NULL) FROM \"public\".\"users\"".to_string() + "SELECT setval(pg_get_serial_sequence('\"public\".\"users\"', 'id'), GREATEST(COALESCE(MAX(\"id\"), 0), 1), MAX(\"id\") IS NOT NULL) FROM \"public\".\"users\"".to_string(), + "SELECT setval(pg_get_serial_sequence('\"public\".\"users\"', 'identity_id'), GREATEST(COALESCE(MAX(\"identity_id\"), 0), 1), MAX(\"identity_id\") IS NOT NULL) FROM \"public\".\"users\"".to_string() ] ); } @@ -4258,12 +4457,14 @@ mod tests { #[test] fn postgres_routine_schema_rewrite_targets_destination_schema() { let rewritten = rewrite_postgres_routine_schema( - "CREATE OR REPLACE FUNCTION public.bump_counter(id integer)\nRETURNS integer\nLANGUAGE plpgsql\nAS $$ BEGIN RETURN id + 1; END; $$", + "CREATE OR REPLACE FUNCTION public.bump_counter(id integer)\nRETURNS integer\nLANGUAGE plpgsql\nAS $$ BEGIN INSERT INTO public.audit_logs(user_id) VALUES (id); RETURN id + 1; END; $$", + "public", "archive", ) .unwrap(); assert!(rewritten.starts_with("CREATE OR REPLACE FUNCTION \"archive\".\"bump_counter\"(")); + assert!(rewritten.contains("INSERT INTO \"archive\".audit_logs")); } #[test] @@ -4276,6 +4477,7 @@ mod tests { ); assert!(rewritten.contains(" ON \"archive\".\"users\" ")); + assert!(rewritten.contains("EXECUTE FUNCTION \"archive\".bump_counter()")); } #[test] @@ -4730,4 +4932,27 @@ mod tests { assert!(ddl.contains("GENERATED BY DEFAULT AS IDENTITY"), "ddl: {ddl}"); } + + #[test] + fn postgres_create_table_preserves_identity_from_column_extra() { + let cols = vec![db::ColumnInfo { + data_type: "integer".to_string(), + extra: Some("generated by default as identity".to_string()), + is_primary_key: true, + is_nullable: false, + ..test_column("id", "integer") + }]; + + let ddl = generate_create_table_ddl( + &cols, + "t", + "public", + "public", + &DatabaseType::Postgres, + &DatabaseType::Postgres, + None, + ); + + assert!(ddl.contains("\"id\" integer generated by default as identity NOT NULL"), "ddl: {ddl}"); + } } diff --git a/crates/dbx-core/tests/live_postgres_transfer.rs b/crates/dbx-core/tests/live_postgres_transfer.rs index 92647cadf..208e95805 100644 --- a/crates/dbx-core/tests/live_postgres_transfer.rs +++ b/crates/dbx-core/tests/live_postgres_transfer.rs @@ -114,6 +114,10 @@ async fn live_postgres_transfer_preserves_data_and_schema_objects() { "CREATE INDEX \"users_display_name_idx\" ON \"{}\".\"users\" USING btree (lower(display_name))", source_schema ), + format!( + "COMMENT ON COLUMN \"{}\".\"users\".\"display_name\" IS 'Display name used in transfer test'", + source_schema + ), format!("COMMENT ON INDEX \"{}\".\"users_display_name_idx\" IS 'lookup index'", source_schema), format!( "CREATE OR REPLACE FUNCTION \"{}\".\"log_user_insert\"() RETURNS trigger LANGUAGE plpgsql AS $$ \ @@ -274,6 +278,21 @@ async fn live_postgres_transfer_preserves_data_and_schema_objects() { .await, json!("email_text") ); + assert_eq!( + query_scalar( + &target_pool, + &format!( + "SELECT col_description(c.oid, a.attnum) \ + FROM pg_catalog.pg_class c \ + JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace \ + JOIN pg_catalog.pg_attribute a ON a.attrelid = c.oid \ + WHERE n.nspname = '{}' AND c.relname = 'users' AND a.attname = 'display_name'", + target_schema + ) + ) + .await, + json!("Display name used in transfer test") + ); assert_eq!( query_scalar(&target_pool, &format!("SELECT count(*) FROM \"{}\".\"active_users\"", target_schema)).await, json!(1)