diff --git a/crates/dbx-core/assets/database-drivers.manifest.json b/crates/dbx-core/assets/database-drivers.manifest.json index b202a6e53..e66f3b457 100644 --- a/crates/dbx-core/assets/database-drivers.manifest.json +++ b/crates/dbx-core/assets/database-drivers.manifest.json @@ -609,7 +609,7 @@ "tableDataEdit": true, "tableStructureEdit": true, "tableImport": true, - "dataTransfer": false, + "dataTransfer": true, "sqlFileExecution": true, "databaseCreate": true, "fieldLineage": false, diff --git a/crates/dbx-core/src/transfer.rs b/crates/dbx-core/src/transfer.rs index b6e0cc7f4..a3c5d3060 100644 --- a/crates/dbx-core/src/transfer.rs +++ b/crates/dbx-core/src/transfer.rs @@ -253,7 +253,13 @@ fn is_simple_identifier(value: &str) -> bool { } fn is_postgres_compat_transfer(source_db: &DatabaseType, target_db: &DatabaseType) -> bool { - matches!(source_db, DatabaseType::Postgres) && matches!(target_db, DatabaseType::Postgres) + is_postgres_transfer_dialect(source_db) && is_postgres_transfer_dialect(target_db) +} + +fn is_postgres_transfer_dialect(db_type: &DatabaseType) -> bool { + // KingbaseES supports the PostgreSQL DDL, type, and ON CONFLICT paths used by transfer; + // other PG-wire databases stay opt-in until their transfer behavior is verified. + matches!(db_type, DatabaseType::Postgres | DatabaseType::Kingbase) } fn is_postgres_integer_like_type(data_type: &str) -> bool { @@ -1479,7 +1485,7 @@ pub fn map_column_type(source_type: &str, _source_db: &DatabaseType, target_db: match base { "int" | "integer" | "int4" | "mediumint" => match target_db { - DatabaseType::Postgres => "INTEGER".into(), + target_db if is_postgres_transfer_dialect(target_db) => "INTEGER".into(), DatabaseType::Mysql => "INT".into(), DatabaseType::SqlServer => "INT".into(), _ => "INTEGER".into(), @@ -1487,26 +1493,29 @@ pub fn map_column_type(source_type: &str, _source_db: &DatabaseType, target_db: "bigint" | "int8" => "BIGINT".into(), "smallint" | "int2" => "SMALLINT".into(), "tinyint" => match target_db { - DatabaseType::Postgres => "SMALLINT".into(), + target_db if is_postgres_transfer_dialect(target_db) => "SMALLINT".into(), _ => "TINYINT".into(), }, "serial" | "bigserial" | "smallserial" => match target_db { - DatabaseType::Postgres => source_type.to_uppercase(), + target_db if is_postgres_transfer_dialect(target_db) => source_type.to_uppercase(), DatabaseType::Mysql => "BIGINT AUTO_INCREMENT".into(), _ => "INTEGER".into(), }, "float" | "float4" | "real" => match target_db { - DatabaseType::Postgres => "REAL".into(), + target_db if is_postgres_transfer_dialect(target_db) => "REAL".into(), _ => "FLOAT".into(), }, "double" | "double precision" | "float8" => match target_db { - DatabaseType::Postgres => "DOUBLE PRECISION".into(), + target_db if is_postgres_transfer_dialect(target_db) => "DOUBLE PRECISION".into(), _ => "DOUBLE".into(), }, "decimal" | "numeric" | "number" => { if t.contains('(') { match target_db { - DatabaseType::Mysql | DatabaseType::Postgres | DatabaseType::SqlServer | DatabaseType::Oracle => { + DatabaseType::Mysql | DatabaseType::SqlServer | DatabaseType::Oracle => { + format!("DECIMAL{}", &t[t.find('(').unwrap()..]) + } + target_db if is_postgres_transfer_dialect(target_db) => { format!("DECIMAL{}", &t[t.find('(').unwrap()..]) } _ => "NUMERIC".into(), @@ -1519,7 +1528,7 @@ pub fn map_column_type(source_type: &str, _source_db: &DatabaseType, target_db: if t.contains('(') { let len_part = &t[t.find('(').unwrap()..]; match target_db { - DatabaseType::Postgres => format!("VARCHAR{len_part}"), + target_db if is_postgres_transfer_dialect(target_db) => format!("VARCHAR{len_part}"), DatabaseType::Mysql => format!("VARCHAR{len_part}"), DatabaseType::SqlServer => format!("NVARCHAR{len_part}"), _ => format!("VARCHAR{len_part}"), @@ -1553,7 +1562,7 @@ pub fn map_column_type(source_type: &str, _source_db: &DatabaseType, target_db: "date" => "DATE".into(), "time" => "TIME".into(), "datetime" => match target_db { - DatabaseType::Postgres => "TIMESTAMP".into(), + target_db if is_postgres_transfer_dialect(target_db) => "TIMESTAMP".into(), DatabaseType::ClickHouse => "DateTime64(6)".into(), _ => "DATETIME".into(), }, @@ -1565,38 +1574,38 @@ pub fn map_column_type(source_type: &str, _source_db: &DatabaseType, target_db: }, "longblob" => match target_db { DatabaseType::Mysql => "LONGBLOB".into(), - DatabaseType::Postgres => "BYTEA".into(), + target_db if is_postgres_transfer_dialect(target_db) => "BYTEA".into(), DatabaseType::SqlServer => "VARBINARY(MAX)".into(), _ => "BLOB".into(), }, "mediumblob" => match target_db { DatabaseType::Mysql => "MEDIUMBLOB".into(), - DatabaseType::Postgres => "BYTEA".into(), + target_db if is_postgres_transfer_dialect(target_db) => "BYTEA".into(), DatabaseType::SqlServer => "VARBINARY(MAX)".into(), _ => "BLOB".into(), }, "blob" | "tinyblob" | "binary" | "varbinary" | "image" => match target_db { - DatabaseType::Postgres => "BYTEA".into(), + target_db if is_postgres_transfer_dialect(target_db) => "BYTEA".into(), DatabaseType::Mysql => "BLOB".into(), DatabaseType::SqlServer => "VARBINARY(MAX)".into(), _ => "BLOB".into(), }, "bytea" => match target_db { - DatabaseType::Postgres => "BYTEA".into(), + target_db if is_postgres_transfer_dialect(target_db) => "BYTEA".into(), DatabaseType::Mysql => "BLOB".into(), _ => "BLOB".into(), }, "json" | "jsonb" => match target_db { - DatabaseType::Postgres => "JSONB".into(), + target_db if is_postgres_transfer_dialect(target_db) => "JSONB".into(), DatabaseType::Mysql => "JSON".into(), _ => "TEXT".into(), }, "uuid" => match target_db { - DatabaseType::Postgres => "UUID".into(), + target_db if is_postgres_transfer_dialect(target_db) => "UUID".into(), _ => "VARCHAR(36)".into(), }, "bit" => match target_db { - DatabaseType::Postgres => "BOOLEAN".into(), + target_db if is_postgres_transfer_dialect(target_db) => "BOOLEAN".into(), _ => "BIT".into(), }, _ => "TEXT".into(), @@ -1872,7 +1881,7 @@ pub fn generate_upsert_typed( } match db_type { - DatabaseType::Postgres | DatabaseType::Sqlite | DatabaseType::DuckDb => { + db_type if is_postgres_transfer_dialect(db_type) => { let pk_list = pk_columns.iter().map(|c| quote_identifier(c, db_type)).collect::>().join(", "); let mut sql = format!("INSERT INTO {full_table} ({col_list}) VALUES\n{}", value_rows.join(",\n")); if non_pk_columns.is_empty() { @@ -2538,7 +2547,7 @@ fn transfer_create_table_created(result: Result<(), String>, error_prefix: &str) } fn transfer_ddl_statements(sql: &str, db_type: &DatabaseType) -> Vec { - if matches!(db_type, DatabaseType::Postgres) { + if is_postgres_transfer_dialect(db_type) { let statements = split_sql_statements(sql); if statements.is_empty() { vec![sql.trim().to_string()] @@ -3970,7 +3979,7 @@ where // Create table on target if requested if request.create_table { - if matches!(target_db_type, DatabaseType::Postgres) && !request.target_schema.trim().is_empty() { + if is_postgres_transfer_dialect(target_db_type) && !request.target_schema.trim().is_empty() { let create_schema_sql = format!("CREATE SCHEMA IF NOT EXISTS {}", quote_identifier(&request.target_schema, target_db_type)); execute_on_pool(state, target_pool_key, &create_schema_sql) @@ -6448,4 +6457,61 @@ SELECT 1 FROM dual"# assert!(ddl.contains("\"id\" integer generated by default as identity NOT NULL"), "ddl: {ddl}"); } + + #[test] + fn kingbase_transfer_uses_postgres_compatible_types() { + assert_eq!(map_column_type("jsonb", &DatabaseType::Postgres, &DatabaseType::Kingbase), "JSONB"); + assert_eq!(map_column_type("bytea", &DatabaseType::Postgres, &DatabaseType::Kingbase), "BYTEA"); + assert_eq!(map_column_type("uuid", &DatabaseType::Postgres, &DatabaseType::Kingbase), "UUID"); + assert_eq!(map_column_type("serial", &DatabaseType::Postgres, &DatabaseType::Kingbase), "SERIAL"); + } + + #[test] + fn kingbase_create_table_preserves_postgres_defaults() { + let cols = vec![db::ColumnInfo { + data_type: "integer".to_string(), + column_default: Some("nextval('source.items_id_seq'::regclass)".to_string()), + is_primary_key: true, + is_nullable: false, + ..test_column("id", "integer") + }]; + + let ddl = generate_create_table_ddl( + &cols, + "items", + "source", + "target", + &DatabaseType::Kingbase, + &DatabaseType::Postgres, + None, + ); + + assert!(ddl.contains("GENERATED BY DEFAULT AS IDENTITY"), "ddl: {ddl}"); + } + + #[test] + fn kingbase_upsert_uses_on_conflict() { + let sql = generate_upsert_typed( + &[String::from("id"), String::from("name")], + &[Some(String::from("integer")), Some(String::from("text"))], + &[vec![json!(1), json!("updated")]], + "items", + "public", + &DatabaseType::Kingbase, + &[String::from("id")], + ); + + assert!(sql.contains("ON CONFLICT (\"id\") DO UPDATE SET \"name\" = EXCLUDED.\"name\""), "sql: {sql}"); + } + + #[test] + fn kingbase_reused_ddl_uses_postgres_statement_sanitization() { + let ddl = r#"CREATE TABLE public.items (id integer PRIMARY KEY); +CREATE INDEX items_name_idx ON public.items (id);"#; + + let statements = transfer_ddl_statements(ddl, &DatabaseType::Kingbase); + + assert_eq!(statements.len(), 1); + assert!(statements[0].starts_with("CREATE TABLE")); + } } diff --git a/crates/dbx-core/tests/database_capabilities.rs b/crates/dbx-core/tests/database_capabilities.rs index 32bf4dd46..972eeebf9 100644 --- a/crates/dbx-core/tests/database_capabilities.rs +++ b/crates/dbx-core/tests/database_capabilities.rs @@ -326,3 +326,16 @@ fn driver_manifest_declares_expected_product_capabilities() { assert!(zookeeper.capabilities.driver_management); assert!(!zookeeper.capabilities.metadata_browse); } + +#[test] +fn kingbase_declares_data_transfer_support() { + let manifest = driver_manifest(); + let kingbase = manifest + .drivers + .iter() + .find(|driver| driver.db_type == DatabaseType::Kingbase) + .expect("Kingbase manifest entry"); + + assert!(kingbase.capabilities.table_import); + assert!(kingbase.capabilities.data_transfer); +}