diff --git a/apps/desktop/public/icons/database/cloudberry.png b/apps/desktop/public/icons/database/cloudberry.png new file mode 100644 index 000000000..0e58b30bc Binary files /dev/null and b/apps/desktop/public/icons/database/cloudberry.png differ diff --git a/apps/desktop/src/components/connection/ConnectionDialog.vue b/apps/desktop/src/components/connection/ConnectionDialog.vue index 0fcedae95..8e2dbf25a 100644 --- a/apps/desktop/src/components/connection/ConnectionDialog.vue +++ b/apps/desktop/src/components/connection/ConnectionDialog.vue @@ -744,6 +744,14 @@ const driverProfiles: Record< icon: "postgres", urlParams: "", }, + cloudberry: { + type: "postgres", + port: 5432, + user: "postgres", + label: "Apache Cloudberry", + icon: "cloudberry", + urlParams: "", + }, redis: { type: "redis", port: 6379, user: "", label: "Redis", icon: "redis" }, sqlite: { type: "sqlite", port: 0, user: "", label: "SQLite", icon: "sqlite" }, rqlite: { type: "rqlite", port: 4001, user: "", label: "RQLite", icon: "rqlite" }, @@ -2108,6 +2116,7 @@ function isH2FileJdbcUrlLikePath(value: string): boolean { const iconTypeMap: Record = { mysql: "mysql", postgres: "postgres", + cloudberry: "cloudberry", sqlite: "sqlite", rqlite: "rqlite", turso: "turso", @@ -2190,6 +2199,7 @@ const iconTypeMap: Record = { const dbOptions: DbOption[] = [ { value: "postgres", label: "PostgreSQL" }, + { value: "cloudberry", label: "Apache Cloudberry" }, { value: "mysql", label: "MySQL" }, { value: "mongodb", label: "MongoDB" }, { value: "redis", label: "Redis" }, @@ -2282,7 +2292,7 @@ const dbCategoryDefinitions: Array<{ { key: "analytics", titleKey: "connection.databaseCategoryAnalytics", - optionValues: ["clickhouse", "doris", "starrocks", "databend", "selectdb", "databricks", "saphana", "teradata", "vertica", "exasol", "redshift", "snowflake", "trino", "prestosql", "hive", "spark", "bigquery", "kylin", "dremio"], + optionValues: ["cloudberry", "clickhouse", "doris", "starrocks", "databend", "selectdb", "databricks", "saphana", "teradata", "vertica", "exasol", "redshift", "snowflake", "trino", "prestosql", "hive", "spark", "bigquery", "kylin", "dremio"], }, { key: "domestic", diff --git a/apps/desktop/src/components/icons/DatabaseIcon.vue b/apps/desktop/src/components/icons/DatabaseIcon.vue index 8458f0eea..0afaffded 100644 --- a/apps/desktop/src/components/icons/DatabaseIcon.vue +++ b/apps/desktop/src/components/icons/DatabaseIcon.vue @@ -11,6 +11,7 @@ const assetIcons: Record = { mysql: "mysql", postgres: "postgres", postgresql: "postgres", + cloudberry: "cloudberry.png", sqlite: "sqlite", rqlite: "rqlite.png", turso: "turso.png", diff --git a/apps/desktop/src/lib/__tests__/connection/connectionUrl.cloudberry.spec.ts b/apps/desktop/src/lib/__tests__/connection/connectionUrl.cloudberry.spec.ts new file mode 100644 index 000000000..8a6a0c93a --- /dev/null +++ b/apps/desktop/src/lib/__tests__/connection/connectionUrl.cloudberry.spec.ts @@ -0,0 +1,37 @@ +import { describe, expect, it } from "vitest"; +import { connectionProfileForScheme, parseConnectionUrl } from "@/lib/connection/connectionUrl"; + +describe("Cloudberry connection URLs", () => { + it("parses the Cloudberry alias as a PostgreSQL-compatible profile", () => { + const parsed = parseConnectionUrl("cloudberry://analyst:secret@cb.example.com/warehouse"); + + expect(parsed).toMatchObject({ + dbType: "postgres", + driverProfile: "cloudberry", + driverLabel: "Apache Cloudberry", + host: "cb.example.com", + port: 5432, + username: "analyst", + password: "secret", + database: "warehouse", + }); + }); + + it("keeps the Cloudberry profile for standard PostgreSQL URLs", () => { + const parsed = parseConnectionUrl("postgresql://cb.example.com:6432/warehouse", "cloudberry"); + + expect(parsed.dbType).toBe("postgres"); + expect(parsed.driverProfile).toBe("cloudberry"); + expect(parsed.driverLabel).toBe("Apache Cloudberry"); + expect(parsed.port).toBe(6432); + }); + + it("exposes Cloudberry to connection deep links", () => { + expect(connectionProfileForScheme("cloudberry")).toEqual({ + type: "postgres", + profile: "cloudberry", + label: "Apache Cloudberry", + defaultPort: 5432, + }); + }); +}); diff --git a/apps/desktop/src/lib/__tests__/imports/dbeaverImport.spec.ts b/apps/desktop/src/lib/__tests__/imports/dbeaverImport.spec.ts index 730bb06e6..ad35f7a5a 100644 --- a/apps/desktop/src/lib/__tests__/imports/dbeaverImport.spec.ts +++ b/apps/desktop/src/lib/__tests__/imports/dbeaverImport.spec.ts @@ -81,3 +81,36 @@ describe("DBeaver folder import", () => { expect(layoutLabels(result.layout!, names)).toEqual([{ group: "Ad hoc", children: [{ group: "Production", children: ["Nested"] }] }]); }); }); + +describe("DBeaver Cloudberry import", () => { + it("preserves Cloudberry while reusing the PostgreSQL backend", async () => { + const connections = await parseDbeaverConnections( + payload({ + connections: { + cloudberry: { + id: "cloudberry", + name: "analytics", + provider: "cloudberry", + driver: "cloudberry-jdbc", + configuration: { + host: "cb.example.com", + port: 5432, + database: "warehouse", + user: "analyst", + }, + }, + }, + }), + ); + + expect(connections[0]).toMatchObject({ + db_type: "postgres", + driver_profile: "cloudberry", + driver_label: "Apache Cloudberry", + host: "cb.example.com", + port: 5432, + database: "warehouse", + username: "analyst", + }); + }); +}); diff --git a/apps/desktop/src/lib/connection/connectionUrl.ts b/apps/desktop/src/lib/connection/connectionUrl.ts index f70738154..9e1e64dbf 100644 --- a/apps/desktop/src/lib/connection/connectionUrl.ts +++ b/apps/desktop/src/lib/connection/connectionUrl.ts @@ -31,6 +31,7 @@ const SCHEME_PROFILES: Record = { mariadb: { type: "mysql", profile: "mariadb", label: "MariaDB", defaultPort: 3306 }, postgres: { type: "postgres", profile: "postgres", label: "PostgreSQL", defaultPort: 5432 }, postgresql: { type: "postgres", profile: "postgres", label: "PostgreSQL", defaultPort: 5432 }, + cloudberry: { type: "postgres", profile: "cloudberry", label: "Apache Cloudberry", defaultPort: 5432 }, redshift: { type: "redshift", profile: "redshift", label: "Redshift", defaultPort: 5439 }, redis: { type: "redis", profile: "redis", label: "Redis", defaultPort: 6379 }, rediss: { type: "redis", profile: "redis", label: "Redis", defaultPort: 6379 }, @@ -248,6 +249,10 @@ export function connectionProfileForScheme(scheme: string, preferredProfile?: st if ((scheme === "http" || scheme === "https") && preferredProfile) { return HTTP_SELECTED_PROFILES[preferredProfile]; } + // Cloudberry uses PostgreSQL URLs, so keep the selected product profile when parsing a pasted URL. + if ((scheme === "postgres" || scheme === "postgresql") && preferredProfile === "cloudberry") { + return SCHEME_PROFILES.cloudberry; + } return SCHEME_PROFILES[scheme]; } diff --git a/apps/desktop/src/lib/imports/dbeaverImport.ts b/apps/desktop/src/lib/imports/dbeaverImport.ts index a68631443..8cf4f4bd7 100644 --- a/apps/desktop/src/lib/imports/dbeaverImport.ts +++ b/apps/desktop/src/lib/imports/dbeaverImport.ts @@ -41,6 +41,7 @@ const profileMap: Record = { mariadb: { dbType: "mysql", profile: "mariadb", label: "MariaDB", port: 3306, user: "root" }, postgresql: { dbType: "postgres", profile: "postgres", label: "PostgreSQL", port: 5432, user: "postgres" }, postgres: { dbType: "postgres", profile: "postgres", label: "PostgreSQL", port: 5432, user: "postgres" }, + cloudberry: { dbType: "postgres", profile: "cloudberry", label: "Apache Cloudberry", port: 5432, user: "postgres" }, sqlite: { dbType: "sqlite", profile: "sqlite", label: "SQLite", port: 0, user: "" }, sqlserver: { dbType: "sqlserver", profile: "sqlserver", label: "SQL Server", port: 1433, user: "sa" }, mssql: { dbType: "sqlserver", profile: "sqlserver", label: "SQL Server", port: 1433, user: "sa" }, diff --git a/crates/dbx-core/src/db/cloudberry.rs b/crates/dbx-core/src/db/cloudberry.rs new file mode 100644 index 000000000..a85cd3ec9 --- /dev/null +++ b/crates/dbx-core/src/db/cloudberry.rs @@ -0,0 +1,335 @@ +use std::collections::HashSet; + +use deadpool_postgres::Pool; + +use super::{ObjectInfo, TableInfo}; +use crate::db; + +const CLOUD_BERRY_TABLE_DDL_SQL: &str = "SELECT pg_get_tabledef($1, $2, true)"; + +const CLOUD_BERRY_EXTERNAL_TABLES_SQL: &str = "SELECT c.relname \ + FROM pg_catalog.pg_class c \ + JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace \ + JOIN pg_catalog.pg_exttable x ON x.reloid = c.oid \ + WHERE n.nspname = $1 AND c.relname = ANY($2::text[])"; + +const CLOUD_BERRY_TABLE_MODIFIERS_SQL: &str = "SELECT COALESCE(am.amname, '')::text AS access_method, \ + COALESCE(array_to_string(c.reloptions, E'\\n'), '')::text AS reloptions, \ + COALESCE(dp.policytype::text, '')::text AS policy_type, \ + COALESCE(string_agg(a.attname, E'\\n' \ + ORDER BY array_position(dp.distkey::smallint[], a.attnum::smallint)), '')::text \ + AS distribution_columns, \ + bool_or(x.reloid IS NOT NULL) AS is_external, \ + COALESCE(fs.srvname, '')::text AS external_server, \ + ft.ftoptions AS external_options \ + FROM pg_catalog.pg_class c \ + JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace \ + LEFT JOIN pg_catalog.pg_am am ON am.oid = c.relam \ + LEFT JOIN pg_catalog.gp_distribution_policy dp ON dp.localoid = c.oid \ + LEFT JOIN pg_catalog.pg_attribute a ON a.attrelid = c.oid AND a.attnum = ANY(dp.distkey) \ + LEFT JOIN pg_catalog.pg_exttable x ON x.reloid = c.oid \ + LEFT JOIN pg_catalog.pg_foreign_table ft ON ft.ftrelid = c.oid \ + LEFT JOIN pg_catalog.pg_foreign_server fs ON fs.oid = ft.ftserver \ + WHERE n.nspname = $1 AND c.relname = $2 \ + AND c.relkind IN ('r', 'p', 'f') \ + GROUP BY c.oid, am.amname, c.reloptions, dp.policytype, dp.distkey, fs.srvname, ft.ftoptions \ + LIMIT 1"; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum DistributionPolicy { + Hash(Vec), + Random, + Replicated, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ExternalTableDefinition { + pub server: String, + pub options: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TableModifiers { + pub access_method: Option, + pub reloptions: Vec, + pub distribution: Option, + pub external: Option, +} + +pub async fn list_tables_filtered( + pool: &Pool, + schema: &str, + filter: Option<&str>, + limit: Option, + offset: Option, +) -> Result, String> { + let mut tables = db::postgres::list_tables_filtered(pool, schema, filter, limit, offset).await?; + annotate_external_tables(pool, schema, &mut tables).await; + Ok(tables) +} + +pub async fn list_objects(pool: &Pool, schema: &str) -> Result, String> { + let mut objects = db::postgres::list_objects(pool, schema).await?; + let names = objects.iter().map(|object| object.name.clone()).collect::>(); + let external_names = external_table_names(pool, schema, &names).await.unwrap_or_else(|error| { + log::debug!("[cloudberry][list_objects:external-table-fallback] error={error}"); + HashSet::new() + }); + for object in &mut objects { + if external_names.contains(&object.name) { + object.object_type = "EXTERNAL TABLE".to_string(); + } + } + Ok(objects) +} + +pub async fn table_ddl(pool: &Pool, schema: &str, table: &str) -> Result { + let client = db::postgres::checkout_postgres_client(pool, None, super::connection_timeout()).await?; + let row = client + .query_opt(CLOUD_BERRY_TABLE_DDL_SQL, &[&schema, &table]) + .await + .map_err(|error| error.to_string())? + .ok_or_else(|| format!("Cloudberry table not found: {schema}.{table}"))?; + let ddl = row.try_get::<_, Option>(0).map_err(|error| error.to_string())?.unwrap_or_default(); + normalize_ddl(ddl) +} + +pub async fn table_modifiers(pool: &Pool, schema: &str, table: &str) -> Result { + let client = db::postgres::checkout_postgres_client(pool, None, super::connection_timeout()).await?; + let row = client + .query_opt(CLOUD_BERRY_TABLE_MODIFIERS_SQL, &[&schema, &table]) + .await + .map_err(|error| error.to_string())? + .ok_or_else(|| format!("Cloudberry table not found: {schema}.{table}"))?; + + let access_method = non_empty(row.try_get::<_, String>(0).map_err(|error| error.to_string())?) + .filter(|method| !method.eq_ignore_ascii_case("heap")); + let reloptions = split_catalog_lines(row.try_get::<_, String>(1).map_err(|error| error.to_string())?); + let policy_type = row.try_get::<_, String>(2).map_err(|error| error.to_string())?; + let distribution_columns = split_catalog_lines(row.try_get::<_, String>(3).map_err(|error| error.to_string())?); + let distribution = match policy_type.as_str() { + "r" => Some(DistributionPolicy::Replicated), + "p" if distribution_columns.is_empty() => Some(DistributionPolicy::Random), + "p" => Some(DistributionPolicy::Hash(distribution_columns)), + _ => None, + }; + + let is_external = row.try_get::<_, bool>(4).map_err(|error| error.to_string())?; + let external = if is_external { + let server = row.try_get::<_, String>(5).map_err(|error| error.to_string())?; + if server.trim().is_empty() { + return Err(format!("Cloudberry external table has no foreign server: {schema}.{table}")); + } + Some(ExternalTableDefinition { + server, + options: row.try_get::<_, Option>>(6).map_err(|error| error.to_string())?.unwrap_or_default(), + }) + } else { + None + }; + + Ok(TableModifiers { access_method, reloptions, distribution, external }) +} + +pub fn append_table_modifiers(ddl: &str, modifiers: &TableModifiers) -> Result { + if let Some(external) = modifiers.external.as_ref() { + return render_external_table_ddl(ddl, external); + } + + let clauses = render_table_modifier_clauses(modifiers); + if clauses.is_empty() { + return Ok(ddl.to_string()); + } + let insertion = ddl + .find(";\n") + .or_else(|| ddl.find(';')) + .ok_or_else(|| "Cloudberry fallback DDL has no CREATE TABLE terminator".to_string())?; + let mut output = String::with_capacity(ddl.len() + clauses.len() + 2); + output.push_str(&ddl[..insertion]); + output.push('\n'); + output.push_str(&clauses); + output.push_str(&ddl[insertion..]); + Ok(output) +} + +fn render_external_table_ddl(ddl: &str, external: &ExternalTableDefinition) -> Result { + let create_table = "CREATE TABLE "; + if !ddl.starts_with(create_table) { + return Err("Cloudberry external-table fallback expected CREATE TABLE DDL".to_string()); + } + let insertion = ddl + .find(";\n") + .or_else(|| ddl.find(';')) + .ok_or_else(|| "Cloudberry fallback DDL has no CREATE TABLE terminator".to_string())?; + let mut output = String::with_capacity(ddl.len() + external.options.len() * 24 + 48); + output.push_str("CREATE FOREIGN TABLE "); + output.push_str(&ddl[create_table.len()..insertion]); + output.push_str("\nSERVER "); + output.push_str(&db::postgres::pg_quote_ident(&external.server)); + if !external.options.is_empty() { + // Cloudberry 2.x stores external tables as foreign tables. Reusing the + // server options preserves URI, format and execution-location details. + output.push_str("\nOPTIONS (\n "); + output.push_str( + &external + .options + .iter() + .map(|option| render_foreign_table_option(option)) + .collect::, _>>()? + .join(",\n "), + ); + output.push_str("\n)"); + } + output.push_str(&ddl[insertion..]); + Ok(output) +} + +fn render_foreign_table_option(option: &str) -> Result { + let (name, value) = + option.split_once('=').ok_or_else(|| format!("Invalid Cloudberry foreign-table option: {option}"))?; + Ok(format!("{} {}", db::postgres::pg_quote_ident(name), quote_sql_string(value))) +} + +fn render_table_modifier_clauses(modifiers: &TableModifiers) -> String { + let mut clauses = Vec::new(); + if let Some(access_method) = modifiers.access_method.as_deref() { + clauses.push(format!("USING {}", db::postgres::pg_quote_ident(access_method))); + } + if !modifiers.reloptions.is_empty() { + clauses.push(format!("WITH (\n {}\n)", modifiers.reloptions.join(",\n "))); + } + if let Some(distribution) = modifiers.distribution.as_ref() { + clauses.push(match distribution { + DistributionPolicy::Hash(columns) => format!( + "DISTRIBUTED BY ({})", + columns.iter().map(|column| db::postgres::pg_quote_ident(column)).collect::>().join(", ") + ), + DistributionPolicy::Random => "DISTRIBUTED RANDOMLY".to_string(), + DistributionPolicy::Replicated => "DISTRIBUTED REPLICATED".to_string(), + }); + } + clauses.join("\n") +} + +async fn annotate_external_tables(pool: &Pool, schema: &str, tables: &mut [TableInfo]) { + let names = tables.iter().map(|table| table.name.clone()).collect::>(); + let external_names = external_table_names(pool, schema, &names).await.unwrap_or_else(|error| { + log::debug!("[cloudberry][list_tables:external-table-fallback] error={error}"); + HashSet::new() + }); + for table in tables { + if external_names.contains(&table.name) { + table.table_type = "EXTERNAL TABLE".to_string(); + } + } +} + +async fn external_table_names(pool: &Pool, schema: &str, names: &[String]) -> Result, String> { + if names.is_empty() { + return Ok(HashSet::new()); + } + let client = db::postgres::checkout_postgres_client(pool, None, super::connection_timeout()).await?; + let rows = + client.query(CLOUD_BERRY_EXTERNAL_TABLES_SQL, &[&schema, &names]).await.map_err(|error| error.to_string())?; + Ok(rows.into_iter().filter_map(|row| row.try_get::<_, String>(0).ok()).collect()) +} + +fn normalize_ddl(ddl: String) -> Result { + let ddl = ddl.trim(); + if ddl.is_empty() { + return Err("Cloudberry returned an empty table DDL".to_string()); + } + if ddl.ends_with(';') { + Ok(format!("{ddl}\n")) + } else { + Ok(format!("{ddl};\n")) + } +} + +fn non_empty(value: String) -> Option { + let value = value.trim(); + (!value.is_empty()).then(|| value.to_string()) +} + +fn split_catalog_lines(value: String) -> Vec { + value.lines().map(str::trim).filter(|value| !value.is_empty()).map(str::to_string).collect() +} + +fn quote_sql_string(value: &str) -> String { + format!("'{}'", value.replace('\'', "''")) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn modifiers(distribution: Option) -> TableModifiers { + TableModifiers { access_method: None, reloptions: Vec::new(), distribution, external: None } + } + + #[test] + fn appends_hash_distribution_before_table_terminator() { + let ddl = "CREATE TABLE \"public\".\"events\" (\n \"tenant_id\" integer\n);\n"; + let rendered = append_table_modifiers( + ddl, + &modifiers(Some(DistributionPolicy::Hash(vec!["tenant_id".to_string(), "event id".to_string()]))), + ) + .unwrap(); + + assert_eq!( + rendered, + "CREATE TABLE \"public\".\"events\" (\n \"tenant_id\" integer\n)\nDISTRIBUTED BY (\"tenant_id\", \"event id\");\n" + ); + } + + #[test] + fn appends_storage_and_replicated_distribution() { + let ddl = "CREATE TABLE \"public\".\"dimensions\" (\n \"id\" integer\n);\n"; + let rendered = append_table_modifiers( + ddl, + &TableModifiers { + access_method: Some("ao_column".to_string()), + reloptions: vec!["compresstype=zstd".to_string(), "compresslevel=3".to_string()], + distribution: Some(DistributionPolicy::Replicated), + external: None, + }, + ) + .unwrap(); + + assert!(rendered.contains("USING \"ao_column\"")); + assert!(rendered.contains("WITH (\n compresstype=zstd,\n compresslevel=3\n)")); + assert!(rendered.contains("DISTRIBUTED REPLICATED;")); + } + + #[test] + fn renders_external_table_from_foreign_options() { + let ddl = "CREATE TABLE \"public\".\"external_events\" (\n \"id\" integer\n);\n"; + let rendered = append_table_modifiers( + ddl, + &TableModifiers { + external: Some(ExternalTableDefinition { + server: "gp_exttable_server".to_string(), + options: vec![ + "format=csv".to_string(), + "location_uris=file://cdw/tmp/events.csv".to_string(), + "null=".to_string(), + ], + }), + ..modifiers(None) + }, + ) + .unwrap(); + + assert!(rendered.starts_with("CREATE FOREIGN TABLE \"public\".\"external_events\"")); + assert!(rendered.contains("SERVER \"gp_exttable_server\"")); + assert!(rendered.contains("\"location_uris\" 'file://cdw/tmp/events.csv'")); + assert!(rendered.contains("\"null\" ''")); + } + + #[test] + fn ddl_query_uses_cloudberry_native_definition_function() { + assert_eq!(CLOUD_BERRY_TABLE_DDL_SQL, "SELECT pg_get_tabledef($1, $2, true)"); + assert!(CLOUD_BERRY_TABLE_MODIFIERS_SQL.contains("gp_distribution_policy")); + assert!(CLOUD_BERRY_TABLE_MODIFIERS_SQL.contains("pg_exttable")); + } +} diff --git a/crates/dbx-core/src/db/mod.rs b/crates/dbx-core/src/db/mod.rs index e28b066d0..c40a97ca4 100644 --- a/crates/dbx-core/src/db/mod.rs +++ b/crates/dbx-core/src/db/mod.rs @@ -1,5 +1,6 @@ pub mod agent_driver; pub mod clickhouse_driver; +pub mod cloudberry; pub mod cloudflare_d1; pub use cloudflare_d1 as cloudflare_d1_driver; pub mod document_result; diff --git a/crates/dbx-core/src/schema.rs b/crates/dbx-core/src/schema.rs index dcc120a87..a872ab383 100644 --- a/crates/dbx-core/src/schema.rs +++ b/crates/dbx-core/src/schema.rs @@ -2268,6 +2268,15 @@ async fn list_tables_once( .await .map(|tables| filter_table_infos(tables, filter, limit, offset, object_types)) } + PoolKind::Postgres(p) if db_config.as_ref().is_some_and(is_cloudberry_config) => { + if object_types.is_some() { + db::cloudberry::list_tables_filtered(p, schema, filter, None, None) + .await + .map(|tables| filter_table_infos(tables, filter, limit, offset, object_types)) + } else { + db::cloudberry::list_tables_filtered(p, schema, filter, limit, offset).await + } + } PoolKind::Postgres(p) => { if object_types.is_some() { db::postgres::list_tables_filtered(p, schema, filter, None, None) @@ -4383,6 +4392,9 @@ async fn list_objects_once( PoolKind::Postgres(p) if db_config.as_ref().is_some_and(is_questdb_config) => { db::questdb::list_objects(p, schema).await.map(unpaged_object_list) } + PoolKind::Postgres(p) if db_config.as_ref().is_some_and(is_cloudberry_config) => { + db::cloudberry::list_objects(p, schema).await.map(unpaged_object_list) + } PoolKind::Postgres(p) => db::postgres::list_objects(p, schema).await.map(unpaged_object_list), _ => { drop(connections); @@ -4499,6 +4511,9 @@ async fn list_completion_objects_once( PoolKind::Postgres(p) if db_config.as_ref().is_some_and(is_questdb_config) => { db::questdb::list_objects(p, schema).await.map(filter_completion_objects) } + PoolKind::Postgres(p) if db_config.as_ref().is_some_and(is_cloudberry_config) => { + db::cloudberry::list_objects(p, schema).await.map(filter_completion_objects) + } PoolKind::Postgres(p) => db::postgres::list_objects(p, schema).await.map(filter_completion_objects), PoolKind::SqlServer(_) => { drop(connections); @@ -5412,6 +5427,9 @@ pub async fn get_table_ddl_core( Err(_) => pg_ddl(p, schema, table).await, } } + PoolKind::Postgres(p) if db_config.as_ref().is_some_and(is_cloudberry_config) => { + cloudberry_ddl(p, schema, table).await + } PoolKind::Postgres(p) => pg_ddl(p, schema, table).await, PoolKind::Sqlite(p) => sqlite_ddl(p, schema, table).await, PoolKind::Rqlite(client) => db::rqlite_driver::table_ddl(client, table).await, @@ -5429,6 +5447,10 @@ fn is_opengauss_family_config(config: &ConnectionConfig) -> bool { || matches!(config.driver_profile.as_deref(), Some("opengauss" | "gaussdb")) } +fn is_cloudberry_config(config: &ConnectionConfig) -> bool { + matches!(config.driver_profile.as_deref(), Some("cloudberry")) +} + fn is_default_oracle_agent_config(config: &ConnectionConfig) -> bool { // Only the default go-oracle agent handles filtered/paged metadata; legacy profiles keep Rust fallback paging. matches!(config.db_type, DatabaseType::Oracle) @@ -6974,6 +6996,27 @@ pub async fn pg_ddl(pool: &deadpool_postgres::Pool, schema: &str, table: &str) - )) } +pub async fn cloudberry_ddl(pool: &deadpool_postgres::Pool, schema: &str, table: &str) -> Result { + match db::cloudberry::table_ddl(pool, schema, table).await { + Ok(ddl) => Ok(ddl), + Err(native_error) => { + let base_ddl = pg_ddl(pool, schema, table).await.map_err(|fallback_error| { + format!( + "Cloudberry pg_get_tabledef failed: {native_error}; PostgreSQL DDL fallback failed: {fallback_error}" + ) + })?; + let modifiers = db::cloudberry::table_modifiers(pool, schema, table).await.map_err(|fallback_error| { + format!("Cloudberry pg_get_tabledef failed: {native_error}; modifier fallback failed: {fallback_error}") + })?; + db::cloudberry::append_table_modifiers(&base_ddl, &modifiers).map_err(|fallback_error| { + format!( + "Cloudberry pg_get_tabledef failed: {native_error}; DDL rendering fallback failed: {fallback_error}" + ) + }) + } + } +} + pub fn render_postgres_table_ddl( schema: &str, table: &str, diff --git a/crates/dbx-core/src/schema/providers/native.rs b/crates/dbx-core/src/schema/providers/native.rs index 86cba7b6a..81b75226b 100644 --- a/crates/dbx-core/src/schema/providers/native.rs +++ b/crates/dbx-core/src/schema/providers/native.rs @@ -47,6 +47,9 @@ pub(in crate::schema) async fn list_tables( db::mysql::list_tables(p, db).await } PoolKind::Postgres(p) if config.is_some_and(is_questdb_config) => db::questdb::list_tables(p, schema).await, + PoolKind::Postgres(p) if config.is_some_and(is_cloudberry_config) => { + db::cloudberry::list_tables_filtered(p, schema, None, None, None).await + } PoolKind::Postgres(p) => db::postgres::list_tables(p, schema).await, PoolKind::Sqlite(p) => db::sqlite::list_tables(p, schema).await, PoolKind::Rqlite(client) => db::rqlite_driver::list_tables(client, schema).await, @@ -86,6 +89,9 @@ pub(in crate::schema) async fn list_objects( PoolKind::Postgres(p) if config.is_some_and(is_questdb_config) => { db::questdb::list_objects(p, schema).await.map(Some) } + PoolKind::Postgres(p) if config.is_some_and(is_cloudberry_config) => { + db::cloudberry::list_objects(p, schema).await.map(Some) + } PoolKind::Postgres(p) => db::postgres::list_objects(p, schema).await.map(Some), _ => Ok(None), } @@ -107,6 +113,9 @@ pub(in crate::schema) async fn list_completion_objects( PoolKind::Postgres(p) if config.is_some_and(is_questdb_config) => { db::questdb::list_objects(p, schema).await.map(Some) } + PoolKind::Postgres(p) if config.is_some_and(is_cloudberry_config) => { + db::cloudberry::list_objects(p, schema).await.map(Some) + } PoolKind::Postgres(p) => db::postgres::list_objects(p, schema).await.map(Some), _ => Ok(None), } @@ -237,6 +246,9 @@ pub(in crate::schema) async fn table_ddl( Err(_) => super::super::pg_ddl(p, schema, table).await, } } + PoolKind::Postgres(p) if config.is_some_and(is_cloudberry_config) => { + super::super::cloudberry_ddl(p, schema, table).await + } PoolKind::Postgres(p) => super::super::pg_ddl(p, schema, table).await, PoolKind::Sqlite(p) => super::super::sqlite_ddl(p, schema, table).await, PoolKind::Rqlite(client) => db::rqlite_driver::table_ddl(client, table).await, @@ -316,6 +328,10 @@ fn is_opengauss_family_config(config: &ConnectionConfig) -> bool { || matches!(config.driver_profile.as_deref(), Some("opengauss" | "gaussdb")) } +fn is_cloudberry_config(config: &ConnectionConfig) -> bool { + matches!(config.driver_profile.as_deref(), Some("cloudberry")) +} + fn is_doris_family_config(config: &ConnectionConfig) -> bool { matches!(config.db_type, DatabaseType::Doris | DatabaseType::StarRocks | DatabaseType::ManticoreSearch) || matches!(config.driver_profile.as_deref(), Some("doris" | "selectdb" | "starrocks" | "manticoresearch")) diff --git a/crates/dbx-core/tests/live_cloudberry.rs b/crates/dbx-core/tests/live_cloudberry.rs new file mode 100644 index 000000000..c2a562a73 --- /dev/null +++ b/crates/dbx-core/tests/live_cloudberry.rs @@ -0,0 +1,90 @@ +use std::time::Duration; + +use dbx_core::{db, schema}; + +#[tokio::test] +#[ignore = "requires DBX_TEST_CLOUDBERRY_URL pointing at a writable Apache Cloudberry database"] +async fn cloudberry_metadata_and_ddl_round_trip() { + let url = std::env::var("DBX_TEST_CLOUDBERRY_URL").expect("DBX_TEST_CLOUDBERRY_URL"); + let pool = db::postgres::connect(&url, Duration::from_secs(10)).await.expect("connect Cloudberry"); + let suffix = uuid::Uuid::new_v4().simple().to_string(); + let source_schema = format!("dbx_cb_source_{suffix}"); + let target_schema = format!("dbx_cb_target_{suffix}"); + let source_ident = quote_ident(&source_schema); + let target_ident = quote_ident(&target_schema); + + db::postgres::execute_batch( + &pool, + &[ + format!("CREATE SCHEMA {source_ident}"), + format!("CREATE SCHEMA {target_ident}"), + format!( + "CREATE TABLE {source_ident}.hash_events (tenant_id integer, payload text) \ + DISTRIBUTED BY (tenant_id)" + ), + format!( + "CREATE TABLE {source_ident}.random_events (id integer, payload text) \ + DISTRIBUTED RANDOMLY" + ), + format!( + "CREATE TABLE {source_ident}.replicated_dimensions (id integer, name text) \ + DISTRIBUTED REPLICATED" + ), + format!( + "CREATE TABLE {source_ident}.column_metrics (metric text, value numeric(18,4)) \ + USING ao_column WITH (compresstype=zstd, compresslevel=3) DISTRIBUTED BY (metric)" + ), + format!( + "CREATE TABLE {source_ident}.partitioned_events \ + (event_date date, tenant_id integer) PARTITION BY RANGE (event_date) \ + DISTRIBUTED BY (tenant_id)" + ), + format!( + "CREATE READABLE EXTERNAL TABLE {source_ident}.external_events (id integer, payload text) \ + LOCATION ('file://cdw/tmp/dbx-cloudberry-live-test.csv') \ + FORMAT 'CSV' (DELIMITER ',')" + ), + ], + ) + .await + .expect("create Cloudberry fixtures"); + + let exercise = async { + let tables = db::cloudberry::list_tables_filtered(&pool, &source_schema, None, None, None).await?; + let external = tables.iter().find(|table| table.name == "external_events").ok_or("missing external table")?; + assert_eq!(external.table_type, "EXTERNAL TABLE"); + + let cases = [ + ("hash_events", "DISTRIBUTED BY (\"tenant_id\")"), + ("random_events", "DISTRIBUTED RANDOMLY"), + ("replicated_dimensions", "DISTRIBUTED REPLICATED"), + ("column_metrics", "USING \"ao_column\""), + ("partitioned_events", "PARTITION BY RANGE (event_date)"), + ("external_events", "CREATE FOREIGN TABLE"), + ]; + for (table, expected) in cases { + let ddl = schema::cloudberry_ddl(&pool, &source_schema, table).await?; + assert!(ddl.contains(expected), "{table} DDL did not contain {expected}: {ddl}"); + if table == "external_events" { + assert!(ddl.contains("SERVER \"gp_exttable_server\""), "external DDL: {ddl}"); + assert!(ddl.contains("\"location_uris\" 'file://cdw/tmp/dbx-cloudberry-live-test.csv'")); + } + let target_ddl = ddl.replace(&source_ident, &target_ident); + db::postgres::execute_query(&pool, &target_ddl).await?; + } + Ok::<_, String>(()) + } + .await; + + db::postgres::execute_batch( + &pool, + &[format!("DROP SCHEMA {target_ident} CASCADE"), format!("DROP SCHEMA {source_ident} CASCADE")], + ) + .await + .expect("drop Cloudberry fixtures"); + exercise.expect("validate Cloudberry metadata and DDL"); +} + +fn quote_ident(value: &str) -> String { + format!("\"{}\"", value.replace('"', "\"\"")) +} diff --git a/docs/public/icons/database/cloudberry.png b/docs/public/icons/database/cloudberry.png new file mode 100644 index 000000000..0e58b30bc Binary files /dev/null and b/docs/public/icons/database/cloudberry.png differ