From 785c4098d1b5cd1c4ae273d7ffe60155de88f04e Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Wed, 22 Jul 2026 01:14:18 +0800 Subject: [PATCH] fix(starrocks): classify materialized views --- crates/dbx-core/src/db/mysql.rs | 204 ++++++++++++++++++++++++++++++-- crates/dbx-core/src/schema.rs | 59 +++++++++ 2 files changed, 253 insertions(+), 10 deletions(-) diff --git a/crates/dbx-core/src/db/mysql.rs b/crates/dbx-core/src/db/mysql.rs index 755325fe6..1b3841b89 100644 --- a/crates/dbx-core/src/db/mysql.rs +++ b/crates/dbx-core/src/db/mysql.rs @@ -2070,6 +2070,66 @@ pub async fn list_tables_show(pool: &MySqlPool, database: &str) -> Result String { + format!( + "SELECT TABLE_NAME FROM information_schema.materialized_views WHERE TABLE_SCHEMA = {}", + quote_value(database) + ) +} + +async fn list_starrocks_materialized_view_names(pool: &MySqlPool, database: &str) -> Result, String> { + let sql = starrocks_materialized_views_sql(database); + let mut conn = get_conn_with_timeout(pool, super::connection_timeout()).await?; + let result = conn.query_iter(&sql).await.map_err(|e| e.to_string())?; + let rows: Vec = result.collect_and_drop().await.map_err(|e| e.to_string())?; + Ok(rows + .iter() + .filter_map(|row| { + let name = get_str_by_name(row, "TABLE_NAME").trim().to_string(); + (!name.is_empty()).then_some(name) + }) + .collect()) +} + +fn classify_starrocks_materialized_views( + tables: &mut [TableInfo], + materialized_view_names: Result, String>, + database: &str, +) { + let materialized_view_names = match materialized_view_names { + Ok(names) => names, + Err(err) => { + // Older StarRocks versions and restricted accounts may not expose this + // information_schema view; keep the base SHOW TABLES result usable. + log::warn!("Skipping materialized view classification for StarRocks database `{database}`: {err}"); + return; + } + }; + + for table in tables { + if table.table_type.eq_ignore_ascii_case("VIEW") && materialized_view_names.contains(&table.name) { + table.table_type = "MATERIALIZED_VIEW".to_string(); + } + } +} + +async fn list_starrocks_tables_with_status( + pool: &MySqlPool, + database: &str, +) -> Result<(Vec, HashMap), String> { + let (tables, materialized_view_names) = tokio::join!( + list_tables_show_with_status(pool, database), + list_starrocks_materialized_view_names(pool, database) + ); + let (mut tables, status) = tables?; + classify_starrocks_materialized_views(&mut tables, materialized_view_names, database); + Ok((tables, status)) +} + +pub async fn list_starrocks_tables(pool: &MySqlPool, database: &str) -> Result, String> { + list_starrocks_tables_with_status(pool, database).await.map(|(tables, _)| tables) +} + fn requested_object_type(object_types: Option<&[String]>, object_type: &str) -> bool { object_types.is_none_or(|types| { types.is_empty() || types.iter().any(|candidate| candidate.eq_ignore_ascii_case(object_type)) @@ -2282,13 +2342,49 @@ pub async fn list_table_objects_show(pool: &MySqlPool, database: &str) -> Result let (tables, routines) = tokio::join!(list_tables_show_with_status(pool, database), list_routine_objects(pool, database)); let (tables, status) = tables?; - let mut objects: Vec = tables + let mut objects = table_infos_to_objects(tables, &status, database); + + match routines { + Ok(routines) => objects.extend(routines), + Err(err) => log::warn!("Skipping routines for database `{}` in object browser: {}", database, err), + } + + Ok(objects) +} + +pub async fn list_starrocks_table_objects(pool: &MySqlPool, database: &str) -> Result, String> { + let (tables, routines) = + tokio::join!(list_starrocks_tables_with_status(pool, database), list_routine_objects(pool, database)); + let (tables, status) = tables?; + let mut objects = table_infos_to_objects(tables, &status, database); + + match routines { + Ok(routines) => objects.extend(routines), + Err(err) => log::warn!("Skipping routines for database `{}` in object browser: {}", database, err), + } + + Ok(objects) +} + +fn table_infos_to_objects( + tables: Vec, + status: &HashMap, + database: &str, +) -> Vec { + tables .into_iter() .map(|table| { let meta = status.get(&table.name); ObjectInfo { name: table.name, - object_type: if table.table_type.eq_ignore_ascii_case("VIEW") { "VIEW" } else { "TABLE" }.to_string(), + object_type: if table.table_type.eq_ignore_ascii_case("MATERIALIZED_VIEW") { + "MATERIALIZED_VIEW" + } else if table.table_type.eq_ignore_ascii_case("VIEW") { + "VIEW" + } else { + "TABLE" + } + .to_string(), schema: Some(database.to_string()), valid: None, signature: None, @@ -2299,14 +2395,7 @@ pub async fn list_table_objects_show(pool: &MySqlPool, database: &str) -> Result parent_name: table.parent_name, } }) - .collect(); - - match routines { - Ok(routines) => objects.extend(routines), - Err(err) => log::warn!("Skipping routines for database `{}` in object browser: {}", database, err), - } - - Ok(objects) + .collect() } async fn list_routine_objects(pool: &MySqlPool, database: &str) -> Result, String> { @@ -4190,6 +4279,101 @@ mod tests { assert_eq!(filtered.iter().map(|table| table.name.as_str()).collect::>(), vec!["t_0001"]); } + #[test] + fn starrocks_materialized_views_are_classified_without_duplicating_tables() { + let mut tables = vec![ + TableInfo { + name: "orders".to_string(), + table_type: "BASE TABLE".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + TableInfo { + name: "orders_view".to_string(), + table_type: "VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + TableInfo { + name: "orders_mv".to_string(), + table_type: "VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + ]; + let materialized_views = HashSet::from(["orders_mv".to_string(), "orders_mv".to_string()]); + + classify_starrocks_materialized_views(&mut tables, Ok(materialized_views), "analytics"); + + assert_eq!(tables.len(), 3); + assert_eq!( + tables.iter().map(|table| (table.name.as_str(), table.table_type.as_str())).collect::>(), + vec![("orders", "BASE TABLE"), ("orders_view", "VIEW"), ("orders_mv", "MATERIALIZED_VIEW")] + ); + } + + #[test] + fn starrocks_materialized_view_lookup_failure_keeps_base_types() { + let mut tables = vec![TableInfo { + name: "orders_mv".to_string(), + table_type: "VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }]; + + classify_starrocks_materialized_views(&mut tables, Err("permission denied".to_string()), "analytics"); + + assert_eq!(tables[0].table_type, "VIEW"); + } + + #[test] + fn starrocks_materialized_view_query_is_scoped_to_database() { + let sql = starrocks_materialized_views_sql("tenant's analytics"); + + assert_eq!( + sql, + "SELECT TABLE_NAME FROM information_schema.materialized_views WHERE TABLE_SCHEMA = 'tenant\\'s analytics'" + ); + } + + #[test] + fn starrocks_object_conversion_preserves_table_view_and_materialized_view_types() { + let tables = vec![ + TableInfo { + name: "orders".to_string(), + table_type: "BASE TABLE".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + TableInfo { + name: "orders_view".to_string(), + table_type: "VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + TableInfo { + name: "orders_mv".to_string(), + table_type: "MATERIALIZED_VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + ]; + + let objects = table_infos_to_objects(tables, &HashMap::new(), "analytics"); + + assert_eq!( + objects.iter().map(|object| (object.name.as_str(), object.object_type.as_str())).collect::>(), + vec![("orders", "TABLE"), ("orders_view", "VIEW"), ("orders_mv", "MATERIALIZED_VIEW")] + ); + } + #[test] fn mysql_table_comment_sql_targets_single_table() { let sql = table_comment_sql("app", "users"); diff --git a/crates/dbx-core/src/schema.rs b/crates/dbx-core/src/schema.rs index 4c3d7a294..9a81c2edf 100644 --- a/crates/dbx-core/src/schema.rs +++ b/crates/dbx-core/src/schema.rs @@ -2023,6 +2023,11 @@ async fn list_tables_once( let pool = connections.get(&pool_key).ok_or("Pool not found")?; match pool { + PoolKind::Mysql(p, _) if db_config.as_ref().is_some_and(is_starrocks_config) => { + db::mysql::list_starrocks_tables(p, database) + .await + .map(|tables| filter_table_infos(tables, filter, limit, offset, object_types)) + } PoolKind::Mysql(p, _) if db_config.as_ref().is_some_and(is_doris_family_config) => { db::mysql::list_tables_show(p, database) .await @@ -2781,6 +2786,39 @@ mod tests { assert_eq!(filtered[0].name, "active_users"); } + #[test] + fn filter_table_infos_pages_starrocks_materialized_views_independently() { + let tables = vec![ + test_table_info("orders"), + super::db::TableInfo { + name: "orders_view".to_string(), + table_type: "VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + super::db::TableInfo { + name: "daily_orders_mv".to_string(), + table_type: "MATERIALIZED_VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + super::db::TableInfo { + name: "monthly_orders_mv".to_string(), + table_type: "MATERIALIZED_VIEW".to_string(), + comment: None, + parent_schema: None, + parent_name: None, + }, + ]; + let object_types = vec!["MATERIALIZED_VIEW".to_string()]; + + let filtered = filter_table_infos(tables, Some("orders"), Some(1), Some(1), Some(&object_types)); + + assert_eq!(filtered.into_iter().map(|table| table.name).collect::>(), vec!["monthly_orders_mv"]); + } + #[test] fn filter_object_infos_filters_object_type_before_offset_and_limit() { let objects = vec![ @@ -2797,6 +2835,21 @@ mod tests { assert_eq!(filtered[0].name, "fetch_name"); } + #[test] + fn filter_object_infos_pages_starrocks_materialized_views_independently() { + let objects = vec![ + test_object_info("orders", "TABLE"), + test_object_info("orders_view", "VIEW"), + test_object_info("daily_orders_mv", "MATERIALIZED_VIEW"), + test_object_info("monthly_orders_mv", "MATERIALIZED_VIEW"), + ]; + let object_types = vec!["MATERIALIZED_VIEW".to_string()]; + + let filtered = filter_object_infos(objects, Some("orders"), Some(1), Some(1), Some(&object_types)); + + assert_eq!(filtered.into_iter().map(|object| object.name).collect::>(), vec!["monthly_orders_mv"]); + } + #[test] fn filter_object_infos_matches_comments() { let mut order_view = test_object_info("order_view", "VIEW"); @@ -4038,6 +4091,8 @@ async fn list_objects_once( db::ob_oracle::list_objects(p, schema).await.map(unpaged_object_list) } else if db_config.as_ref().is_some_and(is_manticoresearch_config) { db::manticoresearch::list_objects(p, database).await.map(unpaged_object_list) + } else if db_config.as_ref().is_some_and(is_starrocks_config) { + db::mysql::list_starrocks_table_objects(p, database).await.map(unpaged_object_list) } else if db_config.as_ref().is_some_and(is_doris_family_config) { db::mysql::list_table_objects_show(p, database).await.map(unpaged_object_list) } else { @@ -4985,6 +5040,10 @@ fn is_doris_family_config(config: &ConnectionConfig) -> bool { || matches!(config.driver_profile.as_deref(), Some("doris" | "selectdb" | "starrocks" | "manticoresearch")) } +fn is_starrocks_config(config: &ConnectionConfig) -> bool { + config.db_type == DatabaseType::StarRocks || matches!(config.driver_profile.as_deref(), Some("starrocks")) +} + /// Doris-family engines that support multi-catalog federation (`SHOW CATALOGS`). /// Manticore Search is excluded — it shares the MySQL code path but has no /// catalog concept.