diff --git a/crates/dbx-core/src/db/mysql.rs b/crates/dbx-core/src/db/mysql.rs index 2b0b25b15..13087def2 100644 --- a/crates/dbx-core/src/db/mysql.rs +++ b/crates/dbx-core/src/db/mysql.rs @@ -535,14 +535,28 @@ fn mysql_ssl_opts( fn mysql_setup_queries(url: &str) -> Vec { let charset = mysql_connection_charset(url).unwrap_or("utf8mb4"); + let catalog = mysql_connection_catalog(url); + let database = mysql_connection_database(url); let mut queries = Vec::new(); - if let Some(database) = mysql_connection_database(url) { - queries.push(format!("USE {}", quote_identifier(&database))); + if let Some(database) = database.as_deref() { + queries.push(format!("USE {}", quote_identifier(database))); } if let Some(time_zone) = mysql_connection_time_zone(url) { queries.push(format!("SET time_zone = {}", quote_value(&time_zone))); } queries.push(format!("SET NAMES {charset}")); + // StarRocks/Doris expose external storage (Paimon, Hive, ...) through a + // catalog. `SET catalog` must run *before* `USE ` (the database + // lives in the external catalog and is unknown to the default one). + // mysql_async drains the setup list back-to-front (Vec::pop), so push it + // last to make it execute first. The handshake does not send the database + // as schema (see `mysql_async_url`, which strips the path when a catalog is + // configured), so the connection establishes in the default catalog and + // this setup query is what switches it. The pool re-runs these queries + // after every connection reset, so the catalog stays current. + if let Some(catalog) = catalog.as_deref() { + queries.push(format!("SET catalog = {}", quote_identifier(catalog))); + } queries } @@ -622,6 +636,26 @@ fn mysql_connection_database(url: &str) -> Option { percent_decode_str(database).decode_utf8().ok().map(|value| value.into_owned()) } +/// Extracts an opt-in `catalog=` URL parameter. dbx strips it from the +/// URL before handing it to mysql_async (see `is_dbx_handled_mysql_url_param`) +/// and instead emits `SET catalog = ` during connection setup. This is +/// how StarRocks/Doris connections reach an external catalog such as Paimon. +fn mysql_connection_catalog(url: &str) -> Option { + let (_, query) = url.split_once('?')?; + let query = query.split('#').next().unwrap_or(query); + query.split('&').find_map(|segment| { + let (key, value) = segment.split_once('=')?; + if !key.eq_ignore_ascii_case("catalog") { + return None; + } + let value = value.trim(); + if value.is_empty() { + return None; + } + percent_decode_str(value).decode_utf8().ok().map(|value| value.into_owned()) + }) +} + fn is_safe_mysql_charset_name(value: &str) -> bool { !value.is_empty() && value.bytes().all(|byte| byte.is_ascii_alphanumeric() || byte == b'_') } @@ -853,6 +887,7 @@ fn is_dbx_handled_mysql_url_param(key: &str) -> bool { matches!( key.to_ascii_lowercase().as_str(), "charset" + | "catalog" | "time_zone" | "time-zone" | "timezone" @@ -866,6 +901,20 @@ fn is_dbx_handled_mysql_url_param(key: &str) -> bool { ) } +/// Strips the database path from a `mysql://[user[:pass]@]host[:port][/path]` +/// URL, returning only the scheme and authority. Used so mysql_async does not +/// send the database as the schema during the MySQL handshake (StarRocks would +/// reject an external-catalog database before `SET catalog` runs in setup). +fn strip_mysql_url_path(base: &str) -> &str { + let Some(rest) = base.strip_prefix("mysql://") else { + return base; + }; + match rest.find('/') { + Some(idx) => &base[.."mysql://".len() + idx], + None => base, + } +} + fn mysql_async_url(url: &str) -> Cow<'_, str> { let Some((base, query)) = url.split_once('?') else { return Cow::Borrowed(url); @@ -874,6 +923,7 @@ fn mysql_async_url(url: &str) -> Cow<'_, str> { let original_count = query.split('&').filter(|segment| !segment.trim().is_empty()).count(); let mut filtered: Vec = Vec::new(); let mut changed = false; + let mut has_catalog = false; for segment in query.split('&') { let segment = segment.trim(); if segment.is_empty() { @@ -885,6 +935,9 @@ fn mysql_async_url(url: &str) -> Cow<'_, str> { filtered.push(segment.to_string()); continue; }; + if key.eq_ignore_ascii_case("catalog") { + has_catalog = true; + } if is_dbx_handled_mysql_url_param(key) { changed = true; continue; @@ -914,7 +967,12 @@ fn mysql_async_url(url: &str) -> Cow<'_, str> { filtered.push(segment.to_string()); } - if !changed && filtered.len() == original_count { + // When a catalog is configured, the database in the URL path must not be + // sent as the schema during the MySQL handshake. Strip the path so mysql_async + // connects without a default schema; the database is selected via setup queries. + let base = if has_catalog { strip_mysql_url_path(base) } else { base }; + + if !changed && filtered.len() == original_count && !has_catalog { Cow::Borrowed(url) } else if filtered.is_empty() { Cow::Owned(base.to_string()) @@ -2779,6 +2837,29 @@ mod tests { assert_eq!(mysql_async_url(url).as_ref(), "mysql://host:3306/db?require_ssl=true"); } + #[test] + fn mysql_async_url_strips_database_path_when_catalog_present() { + // With a catalog configured, the database path must not reach mysql_async + // (it would be sent as the handshake schema and rejected before SET catalog). + assert_eq!( + mysql_async_url("mysql://root:secret@host:3306/clip?catalog=paimon_catalog").as_ref(), + "mysql://root:secret@host:3306" + ); + assert_eq!( + mysql_async_url("mysql://host:3306/clip?catalog=paimon_catalog&require_ssl=true").as_ref(), + "mysql://host:3306?require_ssl=true" + ); + } + + #[test] + fn mysql_async_url_keeps_database_path_when_catalog_absent() { + assert_eq!( + mysql_async_url("mysql://host:3306/clip?require_ssl=true").as_ref(), + "mysql://host:3306/clip?require_ssl=true" + ); + assert_eq!(mysql_async_url("mysql://host:3306/clip").as_ref(), "mysql://host:3306/clip"); + } + #[test] fn ssl_fallback_does_not_disable_required_tls() { assert_eq!(ssl_fallback_url("mysql://host:3306/db?require_ssl=true&charset=utf8mb4"), None); @@ -2845,4 +2926,35 @@ mod tests { vec!["USE `db`", "SET NAMES utf8mb4"] ); } + + #[test] + fn mysql_setup_queries_switch_catalog_when_present() { + // `SET catalog` is pushed last so mysql_async's back-to-front setup + // execution (Vec::pop) runs it before `USE `. + assert_eq!( + mysql_setup_queries("mysql://host:3306/clip?catalog=paimon_catalog"), + vec!["USE `clip`", "SET NAMES utf8mb4", "SET catalog = `paimon_catalog`"] + ); + } + + #[test] + fn mysql_setup_queries_switch_catalog_without_database() { + assert_eq!( + mysql_setup_queries("mysql://host:3306/?catalog=paimon_catalog"), + vec!["SET NAMES utf8mb4", "SET catalog = `paimon_catalog`"] + ); + } + + #[test] + fn mysql_setup_queries_decodes_catalog_parameter() { + assert_eq!( + mysql_setup_queries("mysql://host:3306/db?catalog=my%5Fcatalog"), + vec!["USE `db`", "SET NAMES utf8mb4", "SET catalog = `my_catalog`"] + ); + } + + #[test] + fn mysql_setup_queries_omits_catalog_when_absent() { + assert_eq!(mysql_setup_queries("mysql://host:3306/db?charset=utf8mb4"), vec!["USE `db`", "SET NAMES utf8mb4"]); + } }