feat(doris): expand native compatibility support

This commit is contained in:
t8y2 2026-07-24 14:10:06 +08:00
parent 1497ae1009
commit 7a0dbf3fb7
No known key found for this signature in database
10 changed files with 922 additions and 682 deletions

View File

@ -5877,29 +5877,35 @@ function openExternalUrl(url: string) {
</label>
</div>
<div v-if="supportsGenericUrlParams" class="grid grid-cols-4 items-center gap-4">
<div v-if="supportsGenericUrlParams" class="grid grid-cols-4 items-start gap-4">
<Label :class="connectionLabelClass">{{ t("connection.urlParams") }}</Label>
<Input
v-model="form.url_params"
class="col-span-3"
:placeholder="
form.db_type === 'mysql'
? 'charset=utf8mb4'
: form.db_type === 'saphana'
? 'databaseName=TENANT_DB'
: form.db_type === 'clickhouse'
? 'secure=true'
: form.db_type === 'bigquery'
? 'OAuthType=0;OAuthServiceAcctEmail=svc@project.iam.gserviceaccount.com;OAuthPvtKeyPath=/path/key.json'
: form.db_type === 'informix'
? 'CLIENT_LOCALE=en_US.utf8;DB_LOCALE=en_US.utf8'
: form.db_type === 'spark'
? 'catalog=paimon_catalog'
: form.db_type === 'cassandra'
? 'localdatacenter=dc1'
: 'sslmode=prefer'
"
/>
<div class="col-span-3 space-y-1.5">
<Input
v-model="form.url_params"
:placeholder="
form.db_type === 'mysql'
? 'charset=utf8mb4'
: form.db_type === 'doris' || form.db_type === 'starrocks'
? 'sessionVariables=query_timeout=60'
: form.db_type === 'saphana'
? 'databaseName=TENANT_DB'
: form.db_type === 'clickhouse'
? 'secure=true'
: form.db_type === 'bigquery'
? 'OAuthType=0;OAuthServiceAcctEmail=svc@project.iam.gserviceaccount.com;OAuthPvtKeyPath=/path/key.json'
: form.db_type === 'informix'
? 'CLIENT_LOCALE=en_US.utf8;DB_LOCALE=en_US.utf8'
: form.db_type === 'spark'
? 'catalog=paimon_catalog'
: form.db_type === 'cassandra'
? 'localdatacenter=dc1'
: 'sslmode=prefer'
"
/>
<p v-if="form.db_type === 'mysql' || form.db_type === 'doris' || form.db_type === 'starrocks'" class="text-xs leading-5 text-muted-foreground">
{{ t("connection.localInfilePathHint") }}
</p>
</div>
</div>
<template v-if="isPrestoSqlConnection">

View File

@ -228,6 +228,7 @@ export default {
driverName: "Driver Name",
driverNamePlaceholder: "Vendor or environment name",
urlParams: "URL Params",
localInfilePathHint: "For LOAD DATA LOCAL INFILE, whitelist each file explicitly with localInfilePath=/absolute/path/file.csv. Repeat the parameter for additional files.",
d1AccountId: "Account ID",
d1DatabaseId: "Database ID",
d1ApiToken: "API Token",

View File

@ -230,6 +230,7 @@ export default withEnglishFallback({
driverName: "驱动名称",
driverNamePlaceholder: "厂商或环境名称",
urlParams: "URL 参数",
localInfilePathHint: "使用 LOAD DATA LOCAL INFILE 时,请通过 localInfilePath=/绝对路径/file.csv 明确授权文件;多个文件需重复填写该参数。",
d1AccountId: "Account ID",
d1DatabaseId: "Database ID",
d1ApiToken: "API Token",

View File

@ -0,0 +1,5 @@
pub use super::mysql_compatible::{
get_columns_show_from as get_catalog_columns, list_catalog_indexes, list_catalogs,
list_databases_show_from as list_catalog_databases, list_indexes_with_ddl_fallback as list_indexes,
list_tables_show_from as list_catalog_tables, show_create_table_ddl_from as get_catalog_table_ddl,
};

View File

@ -3,6 +3,7 @@ pub mod clickhouse_driver;
pub mod cloudflare_d1;
pub use cloudflare_d1 as cloudflare_d1_driver;
pub mod document_result;
pub mod doris;
pub mod duckdb_driver;
pub mod duckdb_sql;
#[cfg(feature = "duckdb-bundled")]
@ -19,6 +20,7 @@ pub mod influxdb_driver;
pub mod manticoresearch;
pub mod mongo_driver;
pub mod mysql;
pub mod mysql_compatible;
pub mod ob_oracle;
pub mod postgres;
pub mod proxy_tunnel;
@ -28,6 +30,7 @@ pub mod rqlite_driver;
pub mod sqlite;
pub mod sqlserver;
pub mod ssh_tunnel;
pub mod starrocks;
pub mod transport_layer_tunnel;
pub mod turso_driver;
pub mod vector_driver;

View File

@ -53,7 +53,7 @@ fn quote_value(s: &str) -> String {
format!("'{}'", s.replace('\\', "\\\\").replace('\'', "\\'"))
}
fn quote_identifier(s: &str) -> String {
pub(super) fn quote_identifier(s: &str) -> String {
format!("`{}`", s.replace('`', "``"))
}
@ -76,30 +76,30 @@ where
/// 字节转 String合法 UTF-8绝大多数场景时直接复用入参缓冲零拷贝
/// 仅在非法序列时退化为 lossy 替换。from_utf8_lossy(&b).to_string() 即使
/// 对合法输入也会多一次分配+拷贝。
fn bytes_to_string_lossy(bytes: Vec<u8>) -> String {
pub(super) fn bytes_to_string_lossy(bytes: Vec<u8>) -> String {
String::from_utf8(bytes).unwrap_or_else(|err| String::from_utf8_lossy(err.as_bytes()).into_owned())
}
fn get_str(row: &mysql_async::Row, idx: usize) -> String {
pub(super) fn get_str(row: &mysql_async::Row, idx: usize) -> String {
row_get::<String, _>(row, idx)
.or_else(|| row_get::<Vec<u8>, _>(row, idx).map(bytes_to_string_lossy))
.unwrap_or_default()
}
fn get_str_by_name(row: &mysql_async::Row, name: &str) -> String {
pub(super) fn get_str_by_name(row: &mysql_async::Row, name: &str) -> String {
row_get::<String, _>(row, name)
.or_else(|| row_get::<Vec<u8>, _>(row, name).map(bytes_to_string_lossy))
.unwrap_or_default()
}
fn get_opt_str(row: &mysql_async::Row, name: &str) -> Option<String> {
pub(super) fn get_opt_str(row: &mysql_async::Row, name: &str) -> Option<String> {
row_get::<String, _>(row, name).or_else(|| row_get::<Vec<u8>, _>(row, name).map(bytes_to_string_lossy))
}
/// First non-empty string value among the named columns (e.g. Doris `CatalogName`
/// vs StarRocks `Catalog`). Returns an empty string when none of the columns
/// are present or all are empty.
fn first_nonempty_str_by_name(row: &mysql_async::Row, names: &[&str]) -> String {
pub(super) fn first_nonempty_str_by_name(row: &mysql_async::Row, names: &[&str]) -> String {
for name in names {
let value = get_str_by_name(row, name);
if !value.is_empty() {
@ -775,6 +775,7 @@ fn create_pool(
setup_mode: MySqlSetupMode,
) -> Result<MySqlPool, String> {
let tls_url = mysql_tls_url(url)?;
let local_infile_paths = mysql_local_infile_paths(&tls_url.url);
let opts =
mysql_async::Opts::from_url(&mysql_async_url(&tls_url.url)).map_err(|e| format!("Invalid MySQL URL: {e}"))?;
let tcp_host = mysql_async_tcp_host(opts.ip_or_hostname()).to_string();
@ -809,6 +810,11 @@ fn create_pool(
if let Some(ssl_opts) = mysql_ssl_opts(base_ssl_opts, url, ca_cert_path, &tls_url.files)? {
builder = builder.ssl_opts(ssl_opts);
}
if !local_infile_paths.is_empty() {
// LOCAL INFILE lets the server request a client-side file. Restrict it
// to paths explicitly supplied by the user instead of enabling arbitrary reads.
builder = builder.local_infile_handler(Some(mysql_async::WhiteListFsHandler::new(local_infile_paths)));
}
Ok(MySqlPool::new(builder))
}
@ -952,6 +958,9 @@ fn mysql_setup_queries_for_database_with_mode(
if let Some(time_zone) = mysql_connection_time_zone(url) {
queries.push(format!("SET time_zone = {}", quote_value(&time_zone)));
}
if let Some(session_variables) = mysql_connection_session_variables(url) {
queries.push(session_variables);
}
queries.push(format!("SET NAMES {charset}"));
// MySQL defaults group_concat_max_len to 1024, which silently truncates
// GROUP_CONCAT results. Skip it for MySQL protocol-compatible databases
@ -1071,6 +1080,114 @@ fn mysql_connection_catalog(url: &str) -> Option<String> {
})
}
fn mysql_connection_session_variables(url: &str) -> Option<String> {
let (_, query) = url.split_once('?')?;
let query = query.split('#').next().unwrap_or(query);
let value = query.split('&').find_map(|segment| {
let (key, value) = segment.split_once('=')?;
percent_decode_str(key).decode_utf8().ok().filter(|key| key.eq_ignore_ascii_case("sessionVariables"))?;
percent_decode_str(value).decode_utf8().ok().map(|value| value.into_owned())
})?;
let assignments = split_mysql_session_variables(&value);
if assignments.is_empty() {
return None;
}
// Match Connector/J: separators inside strings or expressions are preserved,
// while system variables receive SESSION and user variables keep their @ prefix.
Some(format!(
"SET {}",
assignments
.into_iter()
.map(|assignment| {
if assignment.starts_with('@') {
assignment
} else {
format!("SESSION {assignment}")
}
})
.collect::<Vec<_>>()
.join(",")
))
}
fn mysql_local_infile_paths(url: &str) -> Vec<PathBuf> {
let Some((_, query)) = url.split_once('?') else {
return Vec::new();
};
let query = query.split('#').next().unwrap_or(query);
query
.split('&')
.filter_map(|segment| {
let (key, value) = segment.split_once('=')?;
percent_decode_str(key).decode_utf8().ok().filter(|key| key.eq_ignore_ascii_case("localInfilePath"))?;
let path = percent_decode_str(value).decode_utf8().ok()?.trim().to_string();
(!path.is_empty()).then(|| PathBuf::from(path))
})
.collect()
}
fn split_mysql_session_variables(value: &str) -> Vec<String> {
let chars: Vec<char> = value.chars().collect();
let mut assignments = Vec::new();
let mut current = String::new();
let mut quote = None;
let mut escaped = false;
let mut parenthesis_depth = 0usize;
let mut index = 0usize;
while index < chars.len() {
let ch = chars[index];
if let Some(active_quote) = quote {
current.push(ch);
if escaped {
escaped = false;
} else if ch == '\\' {
escaped = true;
} else if ch == active_quote {
if chars.get(index + 1) == Some(&active_quote) {
current.push(active_quote);
index += 1;
} else {
quote = None;
}
}
index += 1;
continue;
}
match ch {
'\'' | '"' => {
quote = Some(ch);
current.push(ch);
}
'(' => {
parenthesis_depth += 1;
current.push(ch);
}
')' => {
parenthesis_depth = parenthesis_depth.saturating_sub(1);
current.push(ch);
}
',' | ';' if parenthesis_depth == 0 => {
let assignment = current.trim();
if !assignment.is_empty() {
assignments.push(assignment.to_string());
}
current.clear();
}
_ => current.push(ch),
}
index += 1;
}
let assignment = current.trim();
if !assignment.is_empty() {
assignments.push(assignment.to_string());
}
assignments
}
fn is_safe_mysql_charset_name(value: &str) -> bool {
!value.is_empty() && value.bytes().all(|byte| byte.is_ascii_alphanumeric() || byte == b'_')
}
@ -1372,6 +1489,10 @@ fn is_jdbc_param(key: &str) -> bool {
| "resultsetsizethreshold"
| "nettimeoutforstreamingresults"
| "useusageadvisor"
| "uselocalsessionstate"
| "rewritebatchedstatements"
| "prepstmtcachesqllimit"
| "prepstmtcachesize"
)
}
@ -1390,6 +1511,8 @@ fn is_dbx_handled_mysql_url_param(key: &str) -> bool {
| "connectiontimezone"
| "servertimezone"
| "forceconnectiontimezonetosession"
| "sessionvariables"
| "localinfilepath"
)
}
@ -1558,7 +1681,7 @@ pub async fn list_databases_show(pool: &MySqlPool) -> Result<Vec<DatabaseInfo>,
Ok(database_infos_from_names(rows.iter().map(|row| get_str(row, 0)), true))
}
fn database_infos_from_names(
pub(super) fn database_infos_from_names(
names: impl IntoIterator<Item = String>,
include_catalogless_when_blank: bool,
) -> Vec<DatabaseInfo> {
@ -2693,7 +2816,7 @@ fn normalize_mysql_column_charset_metadata(columns: &mut [ColumnInfo], table_col
/// Example: "主键" → UTF-8 bytes [E4 B8 BB E9 94 AE]
/// → each byte → CP1252 char → UTF-8 re-encoded → garbled text
/// → reversal: map each char back to its CP1252 byte, decode as UTF-8
fn fix_potential_double_encoding(s: &str) -> String {
pub(super) fn fix_potential_double_encoding(s: &str) -> String {
// Map each character to its CP1252 byte value
let mut bytes = Vec::with_capacity(s.len());
for c in s.chars() {
@ -3006,7 +3129,7 @@ fn parse_usize_token(sql: &str, i: &mut usize) -> Option<usize> {
std::str::from_utf8(&bytes[start..*i]).ok()?.parse().ok()
}
fn mysql_keyword_at(sql: &str, i: usize, keyword: &str) -> bool {
pub(super) fn mysql_keyword_at(sql: &str, i: usize, keyword: &str) -> bool {
let end = i + keyword.len();
if end > sql.len() {
return false;
@ -3020,11 +3143,11 @@ fn mysql_keyword_at(sql: &str, i: usize, keyword: &str) -> bool {
&& (end == sql.len() || !is_mysql_identifier_byte(sql.as_bytes()[end]))
}
fn is_mysql_identifier_byte(byte: u8) -> bool {
pub(super) fn is_mysql_identifier_byte(byte: u8) -> bool {
byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'$')
}
fn skip_mysql_quoted(sql: &str, start: usize, quote: u8) -> usize {
pub(super) fn skip_mysql_quoted(sql: &str, start: usize, quote: u8) -> usize {
let bytes = sql.as_bytes();
let mut i = start + 1;
while i < bytes.len() {
@ -3667,40 +3790,6 @@ pub async fn list_indexes(pool: &MySqlPool, database: &str, table: &str) -> Resu
.collect())
}
pub async fn list_doris_family_indexes(
pool: &MySqlPool,
database: &str,
table: &str,
) -> Result<Vec<IndexInfo>, String> {
let statistics_result = list_indexes(pool, database, table).await;
let mut indexes = match &statistics_result {
Ok(indexes) => indexes.clone(),
Err(err) => {
log::debug!(
"Falling back to SHOW CREATE TABLE for Doris-family indexes on `{database}`.`{table}` after information_schema.STATISTICS failed: {err}"
);
Vec::new()
}
};
match show_create_table_ddl(pool, database, table).await {
Ok(ddl) => {
merge_index_infos(&mut indexes, doris_indexes_from_create_table_ddl(&ddl));
Ok(indexes)
}
Err(ddl_err) => {
if indexes.is_empty() {
if let Err(statistics_err) = statistics_result {
return Err(format!(
"{statistics_err}; SHOW CREATE TABLE fallback failed for Doris-family indexes: {ddl_err}"
));
}
}
Ok(indexes)
}
}
}
pub async fn show_create_table_ddl(pool: &MySqlPool, database: &str, table: &str) -> Result<String, String> {
let sql = format!("SHOW CREATE TABLE {}", quote_table_ref(database, table));
let mut conn = get_conn_with_health_check(pool).await?;
@ -3713,453 +3802,6 @@ pub async fn show_create_table_ddl(pool: &MySqlPool, database: &str, table: &str
.ok_or_else(|| "Failed to read DDL".to_string())
}
// ---------------------------------------------------------------------------
// Doris / StarRocks multi-catalog support.
//
// These engines expose external catalogs (iceberg, hive, jdbc, ...) alongside
// the native `internal` catalog via `SHOW CATALOGS`. The functions below address
// objects in a specific catalog using 3-part qualified names
// (`<catalog>.<database>.<table>`), which the engines accept directly without
// needing to `SWITCH` the session catalog.
// ---------------------------------------------------------------------------
/// Build a 2-part qualified identifier `` `<catalog>`.`<database>` ``.
fn doris_catalog_database_ref(catalog: &str, database: &str) -> String {
format!("{}.{}", quote_identifier(catalog), quote_identifier(database))
}
/// Build a 3-part qualified identifier `` `<catalog>`.`<database>`.`<table>` ``.
fn doris_catalog_table_ref(catalog: &str, database: &str, table: &str) -> String {
format!("{}.{}.{}", quote_identifier(catalog), quote_identifier(database), quote_identifier(table))
}
/// `SHOW CATALOGS` → list of catalogs visible to the current user.
///
/// Column layouts differ between engines: Doris exposes `CatalogName` (with
/// `CatalogId`/`IsCurrent`/`CreateTime`/`LastUpdateTime`), while StarRocks
/// exposes `Catalog` (only `Type`/`Comment`, no `IsCurrent`). The name is read
/// from either column; missing trailing columns degrade gracefully to
/// empty/None. The built-in catalog is named `internal` in Doris and
/// `default_catalog` in StarRocks (both with `Type=internal`); detection is
/// type-based (see `CatalogInfo::is_internal`), not name-based.
pub async fn list_doris_catalogs(pool: &MySqlPool) -> Result<Vec<crate::db::CatalogInfo>, String> {
let mut conn = get_conn_with_timeout(pool, super::connection_timeout()).await?;
let result = conn.query_iter("SHOW CATALOGS").await.map_err(|e| e.to_string())?;
let rows: Vec<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
let catalogs: Vec<crate::db::CatalogInfo> = rows
.iter()
.filter_map(|row| {
// Doris column is `CatalogName`; StarRocks column is `Catalog`.
let name = first_nonempty_str_by_name(row, &["CatalogName", "Catalog"]).trim().to_string();
if name.is_empty() {
return None;
}
let catalog_type = get_str_by_name(row, "Type").trim().to_string();
let is_current = {
let value = get_str_by_name(row, "IsCurrent").trim().to_ascii_lowercase();
!value.is_empty() && value != "no" && value != "false" && value != "0"
};
let comment = get_opt_str(row, "Comment").map(|s| s.trim().to_string()).filter(|s| !s.is_empty());
Some(crate::db::CatalogInfo { name, catalog_type, is_current, comment })
})
.collect();
Ok(normalize_doris_catalogs(catalogs))
}
/// Sort with the built-in catalog first, then the rest alphabetically by name.
/// The built-in catalog is identified by `CatalogInfo::is_internal` (type-based)
/// rather than by name, so StarRocks `default_catalog` sorts first just like
/// Doris `internal`. No synthetic catalog is injected: `SHOW CATALOGS` always
/// lists the built-in catalog on both engines, and a single-catalog result is
/// handled by the flat-sidebar fallback in the caller.
fn normalize_doris_catalogs(mut catalogs: Vec<crate::db::CatalogInfo>) -> Vec<crate::db::CatalogInfo> {
catalogs.sort_by(|a, b| match (a.is_internal(), b.is_internal()) {
(true, false) => std::cmp::Ordering::Less,
(false, true) => std::cmp::Ordering::Greater,
_ => a.name.cmp(&b.name),
});
catalogs
}
/// `SHOW DATABASES FROM <catalog>` → databases in the given catalog.
pub async fn list_databases_show_from(pool: &MySqlPool, catalog: &str) -> Result<Vec<DatabaseInfo>, String> {
let mut conn = get_conn_with_timeout(pool, super::connection_timeout()).await?;
let sql = format!("SHOW DATABASES FROM {}", quote_identifier(catalog));
let result = conn.query_iter(&sql).await.map_err(|e| e.to_string())?;
let rows: Vec<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
Ok(database_infos_from_names(rows.iter().map(|row| get_str(row, 0)), false))
}
/// `SHOW TABLES FROM <catalog>.<database>` → tables in an external catalog.
///
/// External catalogs do not support `SHOW TABLE STATUS`, so comments/status are
/// not fetched (the caller only needs names + types for browsing).
pub async fn list_tables_show_from(pool: &MySqlPool, catalog: &str, database: &str) -> Result<Vec<TableInfo>, String> {
let sql = format!("SHOW TABLES FROM {}", doris_catalog_database_ref(catalog, 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<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
let mut tables: Vec<TableInfo> = rows
.iter()
.filter_map(|row| {
let name = get_str(row, 0).trim().to_string();
if name.is_empty() {
return None;
}
// SHOW FULL TABLES exposes a type column; plain SHOW TABLES does not.
let table_type = get_str(row, 1);
Some(TableInfo {
name,
table_type: if table_type.trim().is_empty() { "TABLE".to_string() } else { table_type },
comment: None,
parent_schema: None,
parent_name: None,
})
})
.collect();
tables.sort_by(|a, b| a.name.cmp(&b.name));
Ok(tables)
}
/// `SHOW COLUMNS FROM <catalog>.<database>.<table>` → columns of an external
/// catalog table. Falls back to `DESCRIBE` if `SHOW COLUMNS` is rejected.
pub async fn get_columns_show_from(
pool: &MySqlPool,
catalog: &str,
database: &str,
table: &str,
) -> Result<Vec<ColumnInfo>, String> {
let qualified = doris_catalog_table_ref(catalog, database, table);
let full_sql = format!("SHOW FULL COLUMNS FROM {qualified}");
let plain_sql = format!("SHOW COLUMNS FROM {qualified}");
let describe_sql = format!("DESCRIBE {qualified}");
let mut conn = get_conn_with_health_check(pool).await?;
let rows: Vec<mysql_async::Row> = match conn.query_iter(&full_sql).await {
Ok(result) => result.collect_and_drop().await.map_err(|e| e.to_string())?,
Err(_) => match conn.query_iter(&plain_sql).await {
Ok(result) => result.collect_and_drop().await.map_err(|e| e.to_string())?,
Err(_) => {
let result = conn.query_iter(&describe_sql).await.map_err(|e| e.to_string())?;
result.collect_and_drop().await.map_err(|e| e.to_string())?
}
},
};
Ok(rows
.iter()
.filter_map(|row| {
let name = get_str_by_name(row, "Field").trim().to_string();
if name.is_empty() {
return None;
}
let key = get_str_by_name(row, "Key");
let collation = get_opt_str(row, "Collation").filter(|s| !s.is_empty());
Some(ColumnInfo {
name,
data_type: get_str_by_name(row, "Type"),
is_nullable: get_str_by_name(row, "Null").eq_ignore_ascii_case("YES"),
column_default: get_opt_str(row, "Default"),
is_primary_key: key.eq_ignore_ascii_case("PRI"),
extra: get_opt_str(row, "Extra"),
comment: get_opt_str(row, "Comment")
.map(|s| fix_potential_double_encoding(&s))
.filter(|s| !s.is_empty()),
numeric_precision: None,
numeric_scale: None,
character_maximum_length: None,
enum_values: None,
character_set: collation
.as_deref()
.and_then(|c| c.split_once('_').map(|(charset, _)| charset.to_string()))
.filter(|s| !s.is_empty()),
collation,
})
})
.collect())
}
/// `SHOW CREATE TABLE <catalog>.<database>.<table>` → DDL for an external
/// catalog table.
pub async fn show_create_table_ddl_from(
pool: &MySqlPool,
catalog: &str,
database: &str,
table: &str,
) -> Result<String, String> {
let sql = format!("SHOW CREATE TABLE {}", doris_catalog_table_ref(catalog, database, table));
let mut conn = get_conn_with_health_check(pool).await?;
let result = conn.query_iter(&sql).await.map_err(|e| e.to_string())?;
let rows: Vec<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
let row = rows.first().ok_or("DDL not found")?;
row.get_opt::<String, usize>(1)
.and_then(|result| result.ok())
.or_else(|| row.get_opt::<Vec<u8>, usize>(1).and_then(|result| result.ok()).map(bytes_to_string_lossy))
.ok_or_else(|| "Failed to read DDL".to_string())
}
/// Best-effort index listing for an external catalog table. External catalogs
/// generally do not expose MySQL-style index metadata via `information_schema`
/// (that view is scoped to the internal catalog), so indexes are derived from
/// `SHOW CREATE TABLE` parsing. Returns empty on failure (graceful degradation
/// — indexes are informational for external tables).
pub async fn list_doris_catalog_indexes(
pool: &MySqlPool,
catalog: &str,
database: &str,
table: &str,
) -> Result<Vec<IndexInfo>, String> {
let ddl = show_create_table_ddl_from(pool, catalog, database, table).await?;
Ok(doris_indexes_from_create_table_ddl(&ddl))
}
fn merge_index_infos(target: &mut Vec<IndexInfo>, parsed: Vec<IndexInfo>) {
let mut seen_names: HashSet<String> = target.iter().map(|index| index.name.to_ascii_lowercase()).collect();
for index in parsed {
if index.columns.is_empty() {
continue;
}
if seen_names.contains(&index.name.to_ascii_lowercase())
|| target.iter().any(|existing| {
existing.is_unique == index.is_unique
&& existing.is_primary == index.is_primary
&& existing.columns == index.columns
})
{
continue;
}
seen_names.insert(index.name.to_ascii_lowercase());
target.push(index);
}
}
fn doris_indexes_from_create_table_ddl(ddl: &str) -> Vec<IndexInfo> {
let mut indexes = Vec::new();
for raw_line in ddl.lines() {
let line = trim_ddl_definition_line(raw_line);
if line.is_empty() {
continue;
}
let upper = line.to_ascii_uppercase();
if upper.starts_with("PRIMARY KEY") {
if let Some(index) = doris_table_key_index("PRIMARY", line, true, true, "PRIMARY KEY") {
indexes.push(index);
}
} else if upper.starts_with("UNIQUE KEY") {
if let Some(index) = doris_table_key_index("UNIQUE KEY", line, true, false, "UNIQUE KEY") {
indexes.push(index);
}
} else if upper.starts_with("INDEX ") {
if let Some(index) = doris_secondary_index(line) {
indexes.push(index);
}
}
}
indexes
}
fn trim_ddl_definition_line(line: &str) -> &str {
let mut trimmed = line.trim();
if let Some(rest) = trimmed.strip_prefix(',') {
trimmed = rest.trim_start();
}
while let Some(rest) = trimmed.strip_suffix(',') {
trimmed = rest.trim_end();
}
trimmed
}
fn doris_table_key_index(
name: &str,
line: &str,
is_unique: bool,
is_primary: bool,
index_type: &str,
) -> Option<IndexInfo> {
let columns = parse_mysql_index_columns(first_parenthesized_content(line)?);
if columns.is_empty() {
return None;
}
Some(IndexInfo {
name: name.to_string(),
columns,
is_unique,
is_primary,
filter: None,
index_type: Some(index_type.to_string()),
included_columns: None,
comment: None,
})
}
fn doris_secondary_index(line: &str) -> Option<IndexInfo> {
let (_, rest) = split_keyword_prefix(line, "INDEX")?;
let (name, after_name) = read_mysql_identifier(rest.trim_start())?;
let columns = parse_mysql_index_columns(first_parenthesized_content(after_name)?);
if columns.is_empty() {
return None;
}
Some(IndexInfo {
name,
columns,
is_unique: false,
is_primary: false,
filter: None,
index_type: mysql_keyword_argument(after_name, "USING").or_else(|| Some("INDEX".to_string())),
included_columns: None,
comment: mysql_quoted_string_argument(after_name, "COMMENT"),
})
}
fn split_keyword_prefix<'a>(line: &'a str, keyword: &str) -> Option<(&'a str, &'a str)> {
if line.len() < keyword.len() || !line[..keyword.len()].eq_ignore_ascii_case(keyword) {
return None;
}
let rest = &line[keyword.len()..];
if !rest.is_empty() && is_mysql_identifier_byte(rest.as_bytes()[0]) {
return None;
}
Some((&line[..keyword.len()], rest))
}
fn read_mysql_identifier(input: &str) -> Option<(String, &str)> {
let input = input.trim_start();
if input.is_empty() {
return None;
}
let bytes = input.as_bytes();
if bytes[0] == b'`' {
let mut i = 1;
let mut value = String::new();
while i < bytes.len() {
if bytes[i] == b'`' {
if i + 1 < bytes.len() && bytes[i + 1] == b'`' {
value.push('`');
i += 2;
continue;
}
return Some((value, &input[i + 1..]));
}
let ch = input[i..].chars().next()?;
value.push(ch);
i += ch.len_utf8();
}
return None;
}
let end = input.find(|ch: char| ch.is_whitespace() || matches!(ch, '(' | ')' | ',')).unwrap_or(input.len());
if end == 0 {
return None;
}
Some((input[..end].to_string(), &input[end..]))
}
fn first_parenthesized_content(input: &str) -> Option<&str> {
let bytes = input.as_bytes();
let mut depth = 0usize;
let mut start = None;
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
b'(' => {
if depth == 0 {
start = Some(i + 1);
}
depth += 1;
}
b')' if depth > 0 => {
depth -= 1;
if depth == 0 {
return start.map(|start| &input[start..i]);
}
}
_ => {}
}
i += 1;
}
None
}
fn split_top_level_csv(input: &str) -> Vec<&str> {
let bytes = input.as_bytes();
let mut parts = Vec::new();
let mut depth = 0usize;
let mut start = 0usize;
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
b'(' => depth += 1,
b')' if depth > 0 => depth -= 1,
b',' if depth == 0 => {
parts.push(input[start..i].trim());
start = i + 1;
}
_ => {}
}
i += 1;
}
parts.push(input[start..].trim());
parts
}
fn parse_mysql_index_columns(input: &str) -> Vec<String> {
split_top_level_csv(input)
.into_iter()
.filter_map(|part| read_mysql_identifier(part).map(|(column, _)| column))
.filter(|column| !column.is_empty())
.collect()
}
fn mysql_keyword_argument(input: &str, keyword: &str) -> Option<String> {
let bytes = input.as_bytes();
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
_ if mysql_keyword_at(input, i, keyword) => {
return read_mysql_identifier(&input[i + keyword.len()..]).map(|(value, _)| value);
}
_ => i += 1,
}
}
None
}
fn mysql_quoted_string_argument(input: &str, keyword: &str) -> Option<String> {
let bytes = input.as_bytes();
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
_ if mysql_keyword_at(input, i, keyword) => {
let rest = input[i + keyword.len()..].trim_start();
if rest.as_bytes().first().copied() != Some(b'\'') {
return None;
}
let end = skip_mysql_quoted(rest, 0, b'\'');
if end <= 1 || end > rest.len() {
return None;
}
return Some(rest[1..end - 1].replace("\\'", "'").replace("''", "'"));
}
_ => i += 1,
}
}
None
}
pub async fn list_foreign_keys(pool: &MySqlPool, database: &str, table: &str) -> Result<Vec<ForeignKeyInfo>, String> {
let column_sql = format!(
"SELECT CONSTRAINT_NAME, COLUMN_NAME, REFERENCED_TABLE_SCHEMA, \
@ -4915,57 +4557,6 @@ mod tests {
assert_eq!(parse_mysql_enum_values("varchar(255)"), None);
}
#[test]
fn doris_create_table_ddl_indexes_include_unique_key_and_inverted_indexes() {
let ddl = r#"
CREATE TABLE `bfm_org` (
`org_id` bigint NULL,
`ORG_CODE` varchar(255) NULL,
`ORG_NAME` varchar(255) NULL,
INDEX org_id_idx (`org_id`) USING INVERTED,
INDEX org_code_idx (`ORG_CODE`) USING INVERTED,
INDEX org_name_idx (`ORG_NAME`) USING INVERTED
) ENGINE=OLAP
UNIQUE KEY(`org_id`)
COMMENT ''
DISTRIBUTED BY HASH(`org_id`) BUCKETS 4
"#;
let indexes = doris_indexes_from_create_table_ddl(ddl);
assert_eq!(indexes.len(), 4);
assert_eq!(indexes[0].name, "org_id_idx");
assert_eq!(indexes[0].columns, vec!["org_id"]);
assert!(!indexes[0].is_unique);
assert_eq!(indexes[0].index_type.as_deref(), Some("INVERTED"));
assert_eq!(indexes[3].name, "UNIQUE KEY");
assert_eq!(indexes[3].columns, vec!["org_id"]);
assert!(indexes[3].is_unique);
assert!(!indexes[3].is_primary);
assert_eq!(indexes[3].index_type.as_deref(), Some("UNIQUE KEY"));
}
#[test]
fn doris_create_table_ddl_index_parser_handles_quoted_names_and_comments() {
let ddl = r#"
CREATE TABLE `search_test` (
`name``part` varchar(64) NULL,
INDEX `idx``name` (`name``part`) USING NGRAM_BF COMMENT 'name''s index'
) ENGINE=OLAP
UNIQUE KEY(`tenant_id`, `name``part`)
"#;
let indexes = doris_indexes_from_create_table_ddl(ddl);
assert_eq!(indexes.len(), 2);
assert_eq!(indexes[0].name, "idx`name");
assert_eq!(indexes[0].columns, vec!["name`part"]);
assert_eq!(indexes[0].index_type.as_deref(), Some("NGRAM_BF"));
assert_eq!(indexes[0].comment.as_deref(), Some("name's index"));
assert_eq!(indexes[1].columns, vec!["tenant_id", "name`part"]);
assert!(indexes[1].is_unique);
}
#[test]
fn mysql_largeint_uses_lossless_integer_decoding() {
assert!(is_mysql_lossless_integer_type("LARGEINT"));
@ -5514,6 +5105,32 @@ UNIQUE KEY(`tenant_id`, `name``part`)
assert_eq!(mysql_async_url(url).as_ref(), "mysql://host:3306/db?require_ssl=true");
}
#[test]
fn mysql_async_url_accepts_reported_doris_jdbc_params() {
let url = "mysql://host:9030/db?useLocalSessionState=true&rewriteBatchedStatements=true&prepStmtCacheSqlLimit=2048&prepStmtCacheSize=250&sessionVariables=query_timeout%3D60";
assert_eq!(mysql_async_url(url).as_ref(), "mysql://host:9030/db");
}
#[test]
fn mysql_local_infile_paths_are_explicit_and_removed_from_driver_url() {
let url = "mysql://host:9030/db?localInfilePath=%2Ftmp%2Fone.csv&require_ssl=true&localinfilepath=C%3A%5Cdata%5Ctwo.csv";
assert_eq!(
mysql_local_infile_paths(url),
vec![PathBuf::from("/tmp/one.csv"), PathBuf::from(r"C:\data\two.csv")]
);
assert_eq!(mysql_async_url(url).as_ref(), "mysql://host:9030/db?require_ssl=true");
}
#[test]
fn mysql_local_infile_paths_ignore_empty_or_unrelated_values() {
let url = "mysql://host:9030/db?localInfilePath=&charset=utf8mb4&other=%2Ftmp%2Fignored.csv";
assert!(mysql_local_infile_paths(url).is_empty());
assert_eq!(mysql_async_url(url).as_ref(), "mysql://host:9030/db?other=%2Ftmp%2Fignored.csv");
}
#[test]
fn mysql_async_url_normalizes_cleartext_password_auth_alias() {
let url = "mysql://host:3306/db?allowCleartextPasswords=true&charset=utf8mb4";
@ -5603,6 +5220,30 @@ UNIQUE KEY(`tenant_id`, `name``part`)
);
}
#[test]
fn mysql_setup_queries_apply_connector_j_session_variables() {
assert_eq!(
mysql_setup_queries(
"mysql://host:9030/db?sessionVariables=query_timeout%3D60%2Csql_mode%3D%27STRICT%2CTRADITIONAL%27%3B%40trace_id%3Dconcat%28%27a%2Cb%27%2C%27c%27%29",
&[],
),
vec![
"USE `db`",
"SET SESSION query_timeout=60,SESSION sql_mode='STRICT,TRADITIONAL',@trace_id=concat('a,b','c')",
"SET NAMES utf8mb4",
"SET SESSION group_concat_max_len = 1048576",
]
);
}
#[test]
fn mysql_setup_queries_ignore_empty_session_variables() {
assert_eq!(
mysql_setup_queries("mysql://host:9030/db?sessionVariables=%20%2C%20%3B%20", &[]),
vec!["USE `db`", "SET NAMES utf8mb4", "SET SESSION group_concat_max_len = 1048576"]
);
}
#[test]
fn mysql_setup_queries_apply_explicit_time_zone() {
assert_eq!(
@ -5720,89 +5361,4 @@ UNIQUE KEY(`tenant_id`, `name``part`)
vec!["USE `db`", "SET NAMES utf8mb4", "SET SESSION group_concat_max_len = 1048576"]
);
}
fn doris_catalog_info(name: &str, catalog_type: &str, is_current: bool) -> crate::db::CatalogInfo {
crate::db::CatalogInfo {
name: name.to_string(),
catalog_type: catalog_type.to_string(),
is_current,
comment: None,
}
}
#[test]
fn doris_catalog_database_ref_backtick_qualifies_two_parts() {
assert_eq!(doris_catalog_database_ref("iceberg_catalog", "sales"), "`iceberg_catalog`.`sales`");
}
#[test]
fn doris_catalog_table_ref_backtick_qualifies_three_parts() {
assert_eq!(doris_catalog_table_ref("iceberg_catalog", "sales", "orders"), "`iceberg_catalog`.`sales`.`orders`");
}
#[test]
fn doris_catalog_refs_escape_embedded_backticks() {
assert_eq!(doris_catalog_database_ref("a`b", "c`d"), "`a``b`.`c``d`");
assert_eq!(doris_catalog_table_ref("a`b", "c`d", "e`f"), "`a``b`.`c``d`.`e``f`");
}
#[test]
fn normalize_doris_catalogs_does_not_inject_missing_internal() {
// SHOW CATALOGS always lists the built-in catalog on both engines, so a
// missing internal catalog is not synthesized — the caller's flat-sidebar
// fallback handles a single/empty result instead.
let catalogs = vec![
doris_catalog_info("iceberg_catalog", "iceberg", true),
doris_catalog_info("hive_catalog", "hive", false),
];
let normalized = normalize_doris_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["hive_catalog", "iceberg_catalog"]);
assert!(!normalized.iter().any(|c| c.is_internal()));
}
#[test]
fn normalize_doris_catalogs_keeps_existing_internal_first() {
let catalogs = vec![
doris_catalog_info("iceberg_catalog", "iceberg", false),
doris_catalog_info("internal", "internal", true),
doris_catalog_info("hive_catalog", "hive", false),
];
let normalized = normalize_doris_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["internal", "hive_catalog", "iceberg_catalog"]);
}
#[test]
fn normalize_doris_catalogs_sorts_starrocks_default_catalog_first() {
// StarRocks names its built-in catalog `default_catalog` (Type=Internal);
// detection is type-based, so it sorts first just like Doris `internal`.
let catalogs = vec![
doris_catalog_info("hive_catalog", "hive", false),
doris_catalog_info("default_catalog", "Internal", true),
];
let normalized = normalize_doris_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["default_catalog", "hive_catalog"]);
assert!(normalized[0].is_internal());
assert!(!normalized[1].is_internal());
}
#[test]
fn normalize_doris_catalogs_detects_internal_by_type_not_name() {
// A catalog literally named `internal` but with an external type is NOT
// the built-in catalog and must not sort first.
let catalogs =
vec![doris_catalog_info("internal", "iceberg", false), doris_catalog_info("hive_catalog", "hive", false)];
let normalized = normalize_doris_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["hive_catalog", "internal"]);
assert!(!normalized.iter().any(|c| c.is_internal()));
}
#[test]
fn normalize_doris_catalogs_handles_empty_input() {
let normalized = normalize_doris_catalogs(Vec::new());
assert!(normalized.is_empty());
}
}

View File

@ -0,0 +1,623 @@
use mysql_async::prelude::*;
use std::collections::HashSet;
use crate::types::{ColumnInfo, DatabaseInfo, IndexInfo, TableInfo};
use super::mysql::{
bytes_to_string_lossy, database_infos_from_names, first_nonempty_str_by_name, fix_potential_double_encoding,
get_conn_with_health_check, get_conn_with_timeout, get_opt_str, get_str, get_str_by_name, is_mysql_identifier_byte,
list_indexes, mysql_keyword_at, quote_identifier, show_create_table_ddl, skip_mysql_quoted, MySqlPool,
};
// Doris and StarRocks reuse the MySQL wire protocol, but catalog addressing and DDL/index
// metadata semantics are MySQL-compatible distributed behavior and stay isolated in this module.
pub async fn list_indexes_with_ddl_fallback(
pool: &MySqlPool,
database: &str,
table: &str,
) -> Result<Vec<IndexInfo>, String> {
let statistics_result = list_indexes(pool, database, table).await;
let mut indexes = match &statistics_result {
Ok(indexes) => indexes.clone(),
Err(err) => {
log::debug!(
"Falling back to SHOW CREATE TABLE for MySQL-compatible distributed indexes on `{database}`.`{table}` after information_schema.STATISTICS failed: {err}"
);
Vec::new()
}
};
match show_create_table_ddl(pool, database, table).await {
Ok(ddl) => {
merge_index_infos(&mut indexes, indexes_from_create_table_ddl(&ddl));
Ok(indexes)
}
Err(ddl_err) => {
if indexes.is_empty() {
if let Err(statistics_err) = statistics_result {
return Err(format!(
"{statistics_err}; SHOW CREATE TABLE fallback failed for MySQL-compatible distributed indexes: {ddl_err}"
));
}
}
Ok(indexes)
}
}
}
// ---------------------------------------------------------------------------
// Doris / StarRocks multi-catalog support.
//
// These engines expose external catalogs (iceberg, hive, jdbc, ...) alongside
// the native `internal` catalog via `SHOW CATALOGS`. The functions below address
// objects in a specific catalog using 3-part qualified names
// (`<catalog>.<database>.<table>`), which the engines accept directly without
// needing to `SWITCH` the session catalog.
// ---------------------------------------------------------------------------
/// Build a 2-part qualified identifier `` `<catalog>`.`<database>` ``.
fn catalog_database_ref(catalog: &str, database: &str) -> String {
format!("{}.{}", quote_identifier(catalog), quote_identifier(database))
}
/// Build a 3-part qualified identifier `` `<catalog>`.`<database>`.`<table>` ``.
fn catalog_table_ref(catalog: &str, database: &str, table: &str) -> String {
format!("{}.{}.{}", quote_identifier(catalog), quote_identifier(database), quote_identifier(table))
}
/// `SHOW CATALOGS` → list of catalogs visible to the current user.
///
/// Column layouts differ between engines: Doris exposes `CatalogName` (with
/// `CatalogId`/`IsCurrent`/`CreateTime`/`LastUpdateTime`), while StarRocks
/// exposes `Catalog` (only `Type`/`Comment`, no `IsCurrent`). The name is read
/// from either column; missing trailing columns degrade gracefully to
/// empty/None. The built-in catalog is named `internal` in Doris and
/// `default_catalog` in StarRocks (both with `Type=internal`); detection is
/// type-based (see `CatalogInfo::is_internal`), not name-based.
pub async fn list_catalogs(pool: &MySqlPool) -> Result<Vec<crate::db::CatalogInfo>, String> {
let mut conn = get_conn_with_timeout(pool, super::connection_timeout()).await?;
let result = conn.query_iter("SHOW CATALOGS").await.map_err(|e| e.to_string())?;
let rows: Vec<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
let catalogs: Vec<crate::db::CatalogInfo> = rows
.iter()
.filter_map(|row| {
// Doris column is `CatalogName`; StarRocks column is `Catalog`.
let name = first_nonempty_str_by_name(row, &["CatalogName", "Catalog"]).trim().to_string();
if name.is_empty() {
return None;
}
let catalog_type = get_str_by_name(row, "Type").trim().to_string();
let is_current = {
let value = get_str_by_name(row, "IsCurrent").trim().to_ascii_lowercase();
!value.is_empty() && value != "no" && value != "false" && value != "0"
};
let comment = get_opt_str(row, "Comment").map(|s| s.trim().to_string()).filter(|s| !s.is_empty());
Some(crate::db::CatalogInfo { name, catalog_type, is_current, comment })
})
.collect();
Ok(normalize_catalogs(catalogs))
}
/// Sort with the built-in catalog first, then the rest alphabetically by name.
/// The built-in catalog is identified by `CatalogInfo::is_internal` (type-based)
/// rather than by name, so StarRocks `default_catalog` sorts first just like
/// Doris `internal`. No synthetic catalog is injected: `SHOW CATALOGS` always
/// lists the built-in catalog on both engines, and a single-catalog result is
/// handled by the flat-sidebar fallback in the caller.
fn normalize_catalogs(mut catalogs: Vec<crate::db::CatalogInfo>) -> Vec<crate::db::CatalogInfo> {
catalogs.sort_by(|a, b| match (a.is_internal(), b.is_internal()) {
(true, false) => std::cmp::Ordering::Less,
(false, true) => std::cmp::Ordering::Greater,
_ => a.name.cmp(&b.name),
});
catalogs
}
/// `SHOW DATABASES FROM <catalog>` → databases in the given catalog.
pub async fn list_databases_show_from(pool: &MySqlPool, catalog: &str) -> Result<Vec<DatabaseInfo>, String> {
let mut conn = get_conn_with_timeout(pool, super::connection_timeout()).await?;
let sql = format!("SHOW DATABASES FROM {}", quote_identifier(catalog));
let result = conn.query_iter(&sql).await.map_err(|e| e.to_string())?;
let rows: Vec<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
Ok(database_infos_from_names(rows.iter().map(|row| get_str(row, 0)), false))
}
/// `SHOW TABLES FROM <catalog>.<database>` → tables in an external catalog.
///
/// External catalogs do not support `SHOW TABLE STATUS`, so comments/status are
/// not fetched (the caller only needs names + types for browsing).
pub async fn list_tables_show_from(pool: &MySqlPool, catalog: &str, database: &str) -> Result<Vec<TableInfo>, String> {
let sql = format!("SHOW TABLES FROM {}", catalog_database_ref(catalog, 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<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
let mut tables: Vec<TableInfo> = rows
.iter()
.filter_map(|row| {
let name = get_str(row, 0).trim().to_string();
if name.is_empty() {
return None;
}
// SHOW FULL TABLES exposes a type column; plain SHOW TABLES does not.
let table_type = get_str(row, 1);
Some(TableInfo {
name,
table_type: if table_type.trim().is_empty() { "TABLE".to_string() } else { table_type },
comment: None,
parent_schema: None,
parent_name: None,
})
})
.collect();
tables.sort_by(|a, b| a.name.cmp(&b.name));
Ok(tables)
}
/// `SHOW COLUMNS FROM <catalog>.<database>.<table>` → columns of an external
/// catalog table. Falls back to `DESCRIBE` if `SHOW COLUMNS` is rejected.
pub async fn get_columns_show_from(
pool: &MySqlPool,
catalog: &str,
database: &str,
table: &str,
) -> Result<Vec<ColumnInfo>, String> {
let qualified = catalog_table_ref(catalog, database, table);
let full_sql = format!("SHOW FULL COLUMNS FROM {qualified}");
let plain_sql = format!("SHOW COLUMNS FROM {qualified}");
let describe_sql = format!("DESCRIBE {qualified}");
let mut conn = get_conn_with_health_check(pool).await?;
let rows: Vec<mysql_async::Row> = match conn.query_iter(&full_sql).await {
Ok(result) => result.collect_and_drop().await.map_err(|e| e.to_string())?,
Err(_) => match conn.query_iter(&plain_sql).await {
Ok(result) => result.collect_and_drop().await.map_err(|e| e.to_string())?,
Err(_) => {
let result = conn.query_iter(&describe_sql).await.map_err(|e| e.to_string())?;
result.collect_and_drop().await.map_err(|e| e.to_string())?
}
},
};
Ok(rows
.iter()
.filter_map(|row| {
let name = get_str_by_name(row, "Field").trim().to_string();
if name.is_empty() {
return None;
}
let key = get_str_by_name(row, "Key");
let collation = get_opt_str(row, "Collation").filter(|s| !s.is_empty());
Some(ColumnInfo {
name,
data_type: get_str_by_name(row, "Type"),
is_nullable: get_str_by_name(row, "Null").eq_ignore_ascii_case("YES"),
column_default: get_opt_str(row, "Default"),
is_primary_key: key.eq_ignore_ascii_case("PRI"),
extra: get_opt_str(row, "Extra"),
comment: get_opt_str(row, "Comment")
.map(|s| fix_potential_double_encoding(&s))
.filter(|s| !s.is_empty()),
numeric_precision: None,
numeric_scale: None,
character_maximum_length: None,
enum_values: None,
character_set: collation
.as_deref()
.and_then(|c| c.split_once('_').map(|(charset, _)| charset.to_string()))
.filter(|s| !s.is_empty()),
collation,
})
})
.collect())
}
/// `SHOW CREATE TABLE <catalog>.<database>.<table>` → DDL for an external
/// catalog table.
pub async fn show_create_table_ddl_from(
pool: &MySqlPool,
catalog: &str,
database: &str,
table: &str,
) -> Result<String, String> {
let sql = format!("SHOW CREATE TABLE {}", catalog_table_ref(catalog, database, table));
let mut conn = get_conn_with_health_check(pool).await?;
let result = conn.query_iter(&sql).await.map_err(|e| e.to_string())?;
let rows: Vec<mysql_async::Row> = result.collect_and_drop().await.map_err(|e| e.to_string())?;
let row = rows.first().ok_or("DDL not found")?;
row.get_opt::<String, usize>(1)
.and_then(|result| result.ok())
.or_else(|| row.get_opt::<Vec<u8>, usize>(1).and_then(|result| result.ok()).map(bytes_to_string_lossy))
.ok_or_else(|| "Failed to read DDL".to_string())
}
/// Best-effort index listing for an external catalog table. External catalogs
/// generally do not expose MySQL-style index metadata via `information_schema`
/// (that view is scoped to the internal catalog), so indexes are derived from
/// `SHOW CREATE TABLE` parsing. Returns empty on failure (graceful degradation
/// — indexes are informational for external tables).
pub async fn list_catalog_indexes(
pool: &MySqlPool,
catalog: &str,
database: &str,
table: &str,
) -> Result<Vec<IndexInfo>, String> {
let ddl = show_create_table_ddl_from(pool, catalog, database, table).await?;
Ok(indexes_from_create_table_ddl(&ddl))
}
fn merge_index_infos(target: &mut Vec<IndexInfo>, parsed: Vec<IndexInfo>) {
let mut seen_names: HashSet<String> = target.iter().map(|index| index.name.to_ascii_lowercase()).collect();
for index in parsed {
if index.columns.is_empty() {
continue;
}
if seen_names.contains(&index.name.to_ascii_lowercase())
|| target.iter().any(|existing| {
existing.is_unique == index.is_unique
&& existing.is_primary == index.is_primary
&& existing.columns == index.columns
})
{
continue;
}
seen_names.insert(index.name.to_ascii_lowercase());
target.push(index);
}
}
fn indexes_from_create_table_ddl(ddl: &str) -> Vec<IndexInfo> {
let mut indexes = Vec::new();
for raw_line in ddl.lines() {
let line = trim_ddl_definition_line(raw_line);
if line.is_empty() {
continue;
}
let upper = line.to_ascii_uppercase();
if upper.starts_with("PRIMARY KEY") {
if let Some(index) = table_key_index("PRIMARY", line, true, true, "PRIMARY KEY") {
indexes.push(index);
}
} else if upper.starts_with("UNIQUE KEY") {
if let Some(index) = table_key_index("UNIQUE KEY", line, true, false, "UNIQUE KEY") {
indexes.push(index);
}
} else if upper.starts_with("INDEX ") {
if let Some(index) = secondary_index(line) {
indexes.push(index);
}
}
}
indexes
}
fn trim_ddl_definition_line(line: &str) -> &str {
let mut trimmed = line.trim();
if let Some(rest) = trimmed.strip_prefix(',') {
trimmed = rest.trim_start();
}
while let Some(rest) = trimmed.strip_suffix(',') {
trimmed = rest.trim_end();
}
trimmed
}
fn table_key_index(name: &str, line: &str, is_unique: bool, is_primary: bool, index_type: &str) -> Option<IndexInfo> {
let columns = parse_mysql_index_columns(first_parenthesized_content(line)?);
if columns.is_empty() {
return None;
}
Some(IndexInfo {
name: name.to_string(),
columns,
is_unique,
is_primary,
filter: None,
index_type: Some(index_type.to_string()),
included_columns: None,
comment: None,
})
}
fn secondary_index(line: &str) -> Option<IndexInfo> {
let (_, rest) = split_keyword_prefix(line, "INDEX")?;
let (name, after_name) = read_mysql_identifier(rest.trim_start())?;
let columns = parse_mysql_index_columns(first_parenthesized_content(after_name)?);
if columns.is_empty() {
return None;
}
Some(IndexInfo {
name,
columns,
is_unique: false,
is_primary: false,
filter: None,
index_type: mysql_keyword_argument(after_name, "USING").or_else(|| Some("INDEX".to_string())),
included_columns: None,
comment: mysql_quoted_string_argument(after_name, "COMMENT"),
})
}
fn split_keyword_prefix<'a>(line: &'a str, keyword: &str) -> Option<(&'a str, &'a str)> {
if line.len() < keyword.len() || !line[..keyword.len()].eq_ignore_ascii_case(keyword) {
return None;
}
let rest = &line[keyword.len()..];
if !rest.is_empty() && is_mysql_identifier_byte(rest.as_bytes()[0]) {
return None;
}
Some((&line[..keyword.len()], rest))
}
fn read_mysql_identifier(input: &str) -> Option<(String, &str)> {
let input = input.trim_start();
if input.is_empty() {
return None;
}
let bytes = input.as_bytes();
if bytes[0] == b'`' {
let mut i = 1;
let mut value = String::new();
while i < bytes.len() {
if bytes[i] == b'`' {
if i + 1 < bytes.len() && bytes[i + 1] == b'`' {
value.push('`');
i += 2;
continue;
}
return Some((value, &input[i + 1..]));
}
let ch = input[i..].chars().next()?;
value.push(ch);
i += ch.len_utf8();
}
return None;
}
let end = input.find(|ch: char| ch.is_whitespace() || matches!(ch, '(' | ')' | ',')).unwrap_or(input.len());
if end == 0 {
return None;
}
Some((input[..end].to_string(), &input[end..]))
}
fn first_parenthesized_content(input: &str) -> Option<&str> {
let bytes = input.as_bytes();
let mut depth = 0usize;
let mut start = None;
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
b'(' => {
if depth == 0 {
start = Some(i + 1);
}
depth += 1;
}
b')' if depth > 0 => {
depth -= 1;
if depth == 0 {
return start.map(|start| &input[start..i]);
}
}
_ => {}
}
i += 1;
}
None
}
fn split_top_level_csv(input: &str) -> Vec<&str> {
let bytes = input.as_bytes();
let mut parts = Vec::new();
let mut depth = 0usize;
let mut start = 0usize;
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
b'(' => depth += 1,
b')' if depth > 0 => depth -= 1,
b',' if depth == 0 => {
parts.push(input[start..i].trim());
start = i + 1;
}
_ => {}
}
i += 1;
}
parts.push(input[start..].trim());
parts
}
fn parse_mysql_index_columns(input: &str) -> Vec<String> {
split_top_level_csv(input)
.into_iter()
.filter_map(|part| read_mysql_identifier(part).map(|(column, _)| column))
.filter(|column| !column.is_empty())
.collect()
}
fn mysql_keyword_argument(input: &str, keyword: &str) -> Option<String> {
let bytes = input.as_bytes();
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
_ if mysql_keyword_at(input, i, keyword) => {
return read_mysql_identifier(&input[i + keyword.len()..]).map(|(value, _)| value);
}
_ => i += 1,
}
}
None
}
fn mysql_quoted_string_argument(input: &str, keyword: &str) -> Option<String> {
let bytes = input.as_bytes();
let mut i = 0usize;
while i < bytes.len() {
match bytes[i] {
b'\'' | b'"' | b'`' => {
i = skip_mysql_quoted(input, i, bytes[i]);
continue;
}
_ if mysql_keyword_at(input, i, keyword) => {
let rest = input[i + keyword.len()..].trim_start();
if rest.as_bytes().first().copied() != Some(b'\'') {
return None;
}
let end = skip_mysql_quoted(rest, 0, b'\'');
if end <= 1 || end > rest.len() {
return None;
}
return Some(rest[1..end - 1].replace("\\'", "'").replace("''", "'"));
}
_ => i += 1,
}
}
None
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn doris_create_table_ddl_indexes_include_unique_key_and_inverted_indexes() {
let ddl = r#"
CREATE TABLE `bfm_org` (
`org_id` bigint NULL,
`ORG_CODE` varchar(255) NULL,
`ORG_NAME` varchar(255) NULL,
INDEX org_id_idx (`org_id`) USING INVERTED,
INDEX org_code_idx (`ORG_CODE`) USING INVERTED,
INDEX org_name_idx (`ORG_NAME`) USING INVERTED
) ENGINE=OLAP
UNIQUE KEY(`org_id`)
COMMENT ''
DISTRIBUTED BY HASH(`org_id`) BUCKETS 4
"#;
let indexes = indexes_from_create_table_ddl(ddl);
assert_eq!(indexes.len(), 4);
assert_eq!(indexes[0].name, "org_id_idx");
assert_eq!(indexes[0].columns, vec!["org_id"]);
assert!(!indexes[0].is_unique);
assert_eq!(indexes[0].index_type.as_deref(), Some("INVERTED"));
assert_eq!(indexes[3].name, "UNIQUE KEY");
assert_eq!(indexes[3].columns, vec!["org_id"]);
assert!(indexes[3].is_unique);
assert!(!indexes[3].is_primary);
assert_eq!(indexes[3].index_type.as_deref(), Some("UNIQUE KEY"));
}
#[test]
fn doris_create_table_ddl_index_parser_handles_quoted_names_and_comments() {
let ddl = r#"
CREATE TABLE `search_test` (
`name``part` varchar(64) NULL,
INDEX `idx``name` (`name``part`) USING NGRAM_BF COMMENT 'name''s index'
) ENGINE=OLAP
UNIQUE KEY(`tenant_id`, `name``part`)
"#;
let indexes = indexes_from_create_table_ddl(ddl);
assert_eq!(indexes.len(), 2);
assert_eq!(indexes[0].name, "idx`name");
assert_eq!(indexes[0].columns, vec!["name`part"]);
assert_eq!(indexes[0].index_type.as_deref(), Some("NGRAM_BF"));
assert_eq!(indexes[0].comment.as_deref(), Some("name's index"));
assert_eq!(indexes[1].columns, vec!["tenant_id", "name`part"]);
assert!(indexes[1].is_unique);
}
fn catalog_info(name: &str, catalog_type: &str, is_current: bool) -> crate::db::CatalogInfo {
crate::db::CatalogInfo {
name: name.to_string(),
catalog_type: catalog_type.to_string(),
is_current,
comment: None,
}
}
#[test]
fn catalog_database_ref_backtick_qualifies_two_parts() {
assert_eq!(catalog_database_ref("iceberg_catalog", "sales"), "`iceberg_catalog`.`sales`");
}
#[test]
fn catalog_table_ref_backtick_qualifies_three_parts() {
assert_eq!(catalog_table_ref("iceberg_catalog", "sales", "orders"), "`iceberg_catalog`.`sales`.`orders`");
}
#[test]
fn doris_catalog_refs_escape_embedded_backticks() {
assert_eq!(catalog_database_ref("a`b", "c`d"), "`a``b`.`c``d`");
assert_eq!(catalog_table_ref("a`b", "c`d", "e`f"), "`a``b`.`c``d`.`e``f`");
}
#[test]
fn normalize_catalogs_does_not_inject_missing_internal() {
// SHOW CATALOGS always lists the built-in catalog on both engines, so a
// missing internal catalog is not synthesized — the caller's flat-sidebar
// fallback handles a single/empty result instead.
let catalogs =
vec![catalog_info("iceberg_catalog", "iceberg", true), catalog_info("hive_catalog", "hive", false)];
let normalized = normalize_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["hive_catalog", "iceberg_catalog"]);
assert!(!normalized.iter().any(|c| c.is_internal()));
}
#[test]
fn normalize_catalogs_keeps_existing_internal_first() {
let catalogs = vec![
catalog_info("iceberg_catalog", "iceberg", false),
catalog_info("internal", "internal", true),
catalog_info("hive_catalog", "hive", false),
];
let normalized = normalize_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["internal", "hive_catalog", "iceberg_catalog"]);
}
#[test]
fn normalize_catalogs_sorts_starrocks_default_catalog_first() {
// StarRocks names its built-in catalog `default_catalog` (Type=Internal);
// detection is type-based, so it sorts first just like Doris `internal`.
let catalogs =
vec![catalog_info("hive_catalog", "hive", false), catalog_info("default_catalog", "Internal", true)];
let normalized = normalize_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["default_catalog", "hive_catalog"]);
assert!(normalized[0].is_internal());
assert!(!normalized[1].is_internal());
}
#[test]
fn normalize_catalogs_detects_internal_by_type_not_name() {
// A catalog literally named `internal` but with an external type is NOT
// the built-in catalog and must not sort first.
let catalogs = vec![catalog_info("internal", "iceberg", false), catalog_info("hive_catalog", "hive", false)];
let normalized = normalize_catalogs(catalogs);
let names: Vec<&str> = normalized.iter().map(|c| c.name.as_str()).collect();
assert_eq!(names, vec!["hive_catalog", "internal"]);
assert!(!normalized.iter().any(|c| c.is_internal()));
}
#[test]
fn normalize_catalogs_handles_empty_input() {
let normalized = normalize_catalogs(Vec::new());
assert!(normalized.is_empty());
}
}

View File

@ -0,0 +1,5 @@
pub use super::mysql_compatible::{
get_columns_show_from as get_catalog_columns, list_catalog_indexes, list_catalogs,
list_databases_show_from as list_catalog_databases, list_indexes_with_ddl_fallback as list_indexes,
list_tables_show_from as list_catalog_tables, show_create_table_ddl_from as get_catalog_table_ddl,
};

View File

@ -785,9 +785,14 @@ pub async fn list_sqlserver_linked_server_tables_core(
/// flat-sidebar fallback then renders the standard database list.
pub async fn list_doris_catalogs_core(state: &AppState, connection_id: &str) -> Result<Vec<db::CatalogInfo>, String> {
let pool_key = state.get_or_create_pool(connection_id, None).await?;
let db_config = connection_config(state, connection_id).await;
let connections = state.connections.read().await;
if let Some(PoolKind::Mysql(p, _)) = connections.get(&pool_key) {
return db::mysql::list_doris_catalogs(p).await;
return if db_config.as_ref().is_some_and(is_starrocks_config) {
db::starrocks::list_catalogs(p).await
} else {
db::doris::list_catalogs(p).await
};
}
Ok(vec![])
}
@ -809,7 +814,11 @@ pub async fn list_doris_catalog_databases_core(
let PoolKind::Mysql(p, _) = pool else {
return Ok(vec![]);
};
let databases = db::mysql::list_databases_show_from(p, catalog).await;
let databases = if db_config.as_ref().is_some_and(is_starrocks_config) {
db::starrocks::list_catalog_databases(p, catalog).await
} else {
db::doris::list_catalog_databases(p, catalog).await
};
// External catalogs may reject `SHOW DATABASES FROM <catalog>` when the user
// lacks permission — surface as an empty list rather than an error.
let databases = match databases {
@ -842,14 +851,18 @@ pub async fn list_doris_catalog_tables_core(
object_types: Option<&[String]>,
) -> Result<Vec<db::TableInfo>, String> {
let pool_key = state.get_or_create_pool(connection_id, None).await?;
let db_config = connection_config(state, connection_id).await;
let connections = state.connections.read().await;
let pool = connections.get(&pool_key).ok_or("Pool not found")?;
let PoolKind::Mysql(p, _) = pool else {
return Ok(vec![]);
};
db::mysql::list_tables_show_from(p, catalog, database)
.await
.map(|tables| filter_table_infos(tables, filter, limit, offset, object_types))
let tables = if db_config.as_ref().is_some_and(is_starrocks_config) {
db::starrocks::list_catalog_tables(p, catalog, database).await
} else {
db::doris::list_catalog_tables(p, catalog, database).await
}?;
Ok(filter_table_infos(tables, filter, limit, offset, object_types))
}
/// Columns of an external catalog table via `SHOW COLUMNS FROM <catalog>.<db>.<table>`.
@ -861,12 +874,18 @@ pub async fn get_doris_catalog_columns_core(
table: &str,
) -> Result<Vec<db::ColumnInfo>, String> {
let pool_key = state.get_or_create_pool(connection_id, None).await?;
let db_config = connection_config(state, connection_id).await;
let connections = state.connections.read().await;
let pool = connections.get(&pool_key).ok_or("Pool not found")?;
let PoolKind::Mysql(p, _) = pool else {
return Ok(vec![]);
};
db::mysql::get_columns_show_from(p, catalog, database, table).await.map(deduplicate_column_infos)
let columns = if db_config.as_ref().is_some_and(is_starrocks_config) {
db::starrocks::get_catalog_columns(p, catalog, database, table).await
} else {
db::doris::get_catalog_columns(p, catalog, database, table).await
}?;
Ok(deduplicate_column_infos(columns))
}
/// DDL for an external catalog table via `SHOW CREATE TABLE <catalog>.<db>.<table>`.
@ -878,12 +897,17 @@ pub async fn get_doris_catalog_table_ddl_core(
table: &str,
) -> Result<String, String> {
let pool_key = state.get_or_create_pool(connection_id, None).await?;
let db_config = connection_config(state, connection_id).await;
let connections = state.connections.read().await;
let pool = connections.get(&pool_key).ok_or("Pool not found")?;
let PoolKind::Mysql(p, _) = pool else {
return Err("DDL not supported for this connection".to_string());
};
db::mysql::show_create_table_ddl_from(p, catalog, database, table).await
if db_config.as_ref().is_some_and(is_starrocks_config) {
db::starrocks::get_catalog_table_ddl(p, catalog, database, table).await
} else {
db::doris::get_catalog_table_ddl(p, catalog, database, table).await
}
}
/// Best-effort index listing for an external catalog table (derived from DDL).
@ -895,12 +919,17 @@ pub async fn list_doris_catalog_indexes_core(
table: &str,
) -> Result<Vec<db::IndexInfo>, String> {
let pool_key = state.get_or_create_pool(connection_id, None).await?;
let db_config = connection_config(state, connection_id).await;
let connections = state.connections.read().await;
let pool = connections.get(&pool_key).ok_or("Pool not found")?;
let PoolKind::Mysql(p, _) = pool else {
return Ok(vec![]);
};
db::mysql::list_doris_catalog_indexes(p, catalog, database, table).await
if db_config.as_ref().is_some_and(is_starrocks_config) {
db::starrocks::list_catalog_indexes(p, catalog, database, table).await
} else {
db::doris::list_catalog_indexes(p, catalog, database, table).await
}
}
/// Table comment for an external catalog table. Doris does not reliably expose
@ -4972,8 +5001,10 @@ pub async fn list_indexes_core(
}
if *mode == MysqlMode::OceanBaseOracle {
db::ob_oracle::list_indexes(p, schema, table).await
} else if db_config.as_ref().is_some_and(is_doris_family_config) {
db::mysql::list_doris_family_indexes(p, mysql_table_metadata_catalog(database, schema), table).await
} else if db_config.as_ref().is_some_and(is_starrocks_config) {
db::starrocks::list_indexes(p, mysql_table_metadata_catalog(database, schema), table).await
} else if db_config.as_ref().is_some_and(is_doris_config) {
db::doris::list_indexes(p, mysql_table_metadata_catalog(database, schema), table).await
} else {
db::mysql::list_indexes(p, mysql_table_metadata_catalog(database, schema), table).await
}
@ -5558,6 +5589,10 @@ fn is_doris_family_config(config: &ConnectionConfig) -> bool {
|| matches!(config.driver_profile.as_deref(), Some("doris" | "selectdb" | "starrocks" | "manticoresearch"))
}
fn is_doris_config(config: &ConnectionConfig) -> bool {
config.db_type == DatabaseType::Doris || matches!(config.driver_profile.as_deref(), Some("doris" | "selectdb"))
}
fn is_starrocks_config(config: &ConnectionConfig) -> bool {
config.db_type == DatabaseType::StarRocks || matches!(config.driver_profile.as_deref(), Some("starrocks"))
}

View File

@ -46,9 +46,7 @@ pub(in crate::schema) async fn list_tables(
let db = if schema.is_empty() { database } else { schema };
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_questdb_config) => db::questdb::list_tables(p, schema).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,
@ -59,9 +57,9 @@ pub(in crate::schema) async fn list_tables(
PoolKind::Elasticsearch(client) => {
db::elasticsearch_driver::list_indices(client).await.map(|names| collection_names_to_tables(names, "INDEX"))
}
PoolKind::VectorDb(client) => {
db::vector_driver::list_collections(client).await.map(|infos| collection_names_to_tables(infos.into_iter().map(|i| i.name).collect(), "COLLECTION"))
}
PoolKind::VectorDb(client) => db::vector_driver::list_collections(client)
.await
.map(|infos| collection_names_to_tables(infos.into_iter().map(|i| i.name).collect(), "COLLECTION")),
_ => Ok(vec![]),
}
}
@ -82,9 +80,9 @@ pub(in crate::schema) async fn list_objects(
PoolKind::Mysql(p, _) if config.is_some_and(is_doris_family_config) => {
db::mysql::list_table_objects_show(p, database).await.map(Some)
}
PoolKind::Mysql(p, _) => db::mysql::list_objects(p, database, None, None, None)
.await
.map(|result| Some(result.objects)),
PoolKind::Mysql(p, _) => {
db::mysql::list_objects(p, database, None, None, None).await.map(|result| Some(result.objects))
}
PoolKind::Postgres(p) if config.is_some_and(is_questdb_config) => {
db::questdb::list_objects(p, schema).await.map(Some)
}
@ -158,8 +156,14 @@ pub(in crate::schema) async fn list_indexes(
PoolKind::Mysql(p, mode) if *mode == MysqlMode::OceanBaseOracle => {
db::ob_oracle::list_indexes(p, schema, table).await
}
PoolKind::Mysql(p, _) if config.is_some_and(is_doris_family_config) => {
db::mysql::list_doris_family_indexes(p, database, table).await
PoolKind::Mysql(p, _) if config.is_some_and(is_starrocks_config) => {
db::starrocks::list_indexes(p, database, table).await
}
PoolKind::Mysql(p, _) if config.is_some_and(is_doris_config) => {
db::doris::list_indexes(p, database, table).await
}
PoolKind::Mysql(p, _) if config.is_some_and(is_manticoresearch_config) => {
db::mysql_compatible::list_indexes_with_ddl_fallback(p, database, table).await
}
PoolKind::Mysql(p, _) => db::mysql::list_indexes(p, schema, table).await,
PoolKind::Postgres(p) if config.is_some_and(is_questdb_config) => {
@ -262,11 +266,8 @@ pub(in crate::schema) async fn object_source(
}
PoolKind::Sqlite(pool) => {
let source = super::super::first_string_cell(
db::sqlite::execute_query(
pool,
&super::super::sqlite_object_source_sql(schema, name, object_type),
)
.await?,
db::sqlite::execute_query(pool, &super::super::sqlite_object_source_sql(schema, name, object_type))
.await?,
)?;
Ok(Some(source))
}
@ -320,15 +321,20 @@ fn is_doris_family_config(config: &ConnectionConfig) -> bool {
|| matches!(config.driver_profile.as_deref(), Some("doris" | "selectdb" | "starrocks" | "manticoresearch"))
}
fn is_doris_config(config: &ConnectionConfig) -> bool {
config.db_type == DatabaseType::Doris || matches!(config.driver_profile.as_deref(), Some("doris" | "selectdb"))
}
fn is_starrocks_config(config: &ConnectionConfig) -> bool {
config.db_type == DatabaseType::StarRocks || matches!(config.driver_profile.as_deref(), Some("starrocks"))
}
fn is_manticoresearch_config(config: &ConnectionConfig) -> bool {
matches!(config.db_type, DatabaseType::ManticoreSearch)
|| matches!(config.driver_profile.as_deref(), Some("manticoresearch"))
}
fn mysql_show_metadata_database_for_config<'a>(
config: Option<&ConnectionConfig>,
database: &'a str,
) -> &'a str {
fn mysql_show_metadata_database_for_config<'a>(config: Option<&ConnectionConfig>, database: &'a str) -> &'a str {
if config.is_some_and(is_manticoresearch_config) {
""
} else {
@ -352,6 +358,5 @@ fn is_mysql_system_database(name: &str) -> bool {
}
fn is_questdb_config(config: &ConnectionConfig) -> bool {
matches!(config.db_type, DatabaseType::Questdb)
|| matches!(config.driver_profile.as_deref(), Some("questdb"))
matches!(config.db_type, DatabaseType::Questdb) || matches!(config.driver_profile.as_deref(), Some("questdb"))
}