diff --git a/agents/drivers/kingbase-go/kingbase_metadata.go b/agents/drivers/kingbase-go/kingbase_metadata.go index 5e0e916a2..51381ee39 100644 --- a/agents/drivers/kingbase-go/kingbase_metadata.go +++ b/agents/drivers/kingbase-go/kingbase_metadata.go @@ -28,6 +28,7 @@ var kingbaseDataTypes = []string{ } type kingbaseMode struct { + compatibilityMode string postgresCatalog bool mysqlCompat bool sqlServerIdentity bool @@ -108,13 +109,17 @@ type triggerInfo struct { } func detectKingbaseMode(db *sql.DB, configuredMySQL bool) kingbaseMode { - mode := kingbaseMode{mysqlCompat: configuredMySQL} if configuredMySQL { - return mode + return kingbaseMode{compatibilityMode: "mysql", mysqlCompat: true} } + mode := kingbaseMode{compatibilityMode: detectDatabaseMode(db)} mode.postgresCatalog = !catalogExists(db, "sys_catalog.sys_namespace") && catalogExists(db, "pg_catalog.pg_namespace") if !mode.postgresCatalog { - mode.mysqlCompat = detectMySQLCompatMode(db) + if mode.compatibilityMode != "" { + mode.mysqlCompat = mode.compatibilityMode == "mysql" + } else { + mode.mysqlCompat = supportsBacktickIdentifiers(db) + } mode.sqlServerIdentity = !mode.mysqlCompat && catalogExists(db, "sys.identity_columns") } return mode @@ -131,20 +136,25 @@ func catalogExists(db *sql.DB, catalog string) bool { } func detectMySQLCompatMode(db *sql.DB) bool { + if databaseMode := detectDatabaseMode(db); databaseMode != "" { + return databaseMode == "mysql" + } + return supportsBacktickIdentifiers(db) +} + +func detectDatabaseMode(db *sql.DB) string { var databaseMode string switch err := db.QueryRow("SELECT setting FROM sys_catalog.sys_settings WHERE LOWER(name) = 'database_mode'").Scan(&databaseMode); { case err == nil: // Treat database_mode as authoritative when the server exposes it. This // avoids misclassifying Oracle-compatible servers that also publish // sql_mode for MySQL syntax toggles such as ANSI_QUOTES. - return strings.EqualFold(strings.TrimSpace(databaseMode), "mysql") + return strings.ToLower(strings.TrimSpace(databaseMode)) case errors.Is(err, sql.ErrNoRows): - // Older Kingbase versions may not expose database_mode. Probe the syntax - // directly because non-MySQL modes can still expose a sql_mode setting. + return "" default: - // Ignore metadata errors and fall back to the syntax probe below. + return "" } - return supportsBacktickIdentifiers(db) } func supportsBacktickIdentifiers(db *sql.DB) bool { @@ -173,7 +183,8 @@ func (s *server) connectionInfo() (map[string]any, error) { } return map[string]any{ "database": database, "username": username, "version": version, "schema": schema, - "mysql_compat_mode": s.mode.mysqlCompat, "identifierQuote": s.identifierQuote(), + "compatibilityMode": s.mode.compatibilityMode, "mysql_compat_mode": s.mode.mysqlCompat, + "identifierQuote": s.identifierQuote(), }, nil } diff --git a/agents/drivers/kingbase-go/main_test.go b/agents/drivers/kingbase-go/main_test.go index 2cf08d50f..adfe320de 100644 --- a/agents/drivers/kingbase-go/main_test.go +++ b/agents/drivers/kingbase-go/main_test.go @@ -813,26 +813,46 @@ func TestDetectMySQLCompatModeAcceptsExplicitMySQLMode(t *testing.T) { } } +func TestDetectKingbaseModeReportsDatabaseMode(t *testing.T) { + for _, databaseMode := range []string{"oracle", "postgresql"} { + t.Run(databaseMode, func(t *testing.T) { + db := openModeDetectionDB(t, &modeDetectionDriverState{databaseMode: &databaseMode}) + + mode := detectKingbaseMode(db, false) + + if mode.compatibilityMode != databaseMode { + t.Fatalf("unexpected compatibility mode: %q", mode.compatibilityMode) + } + }) + } +} + func TestConnectionInfoReportsCompatibilityIdentifierQuote(t *testing.T) { for _, testCase := range []struct { - name string - mysqlCompat bool - expected string + name string + compatibilityMode string + mysqlCompat bool + expectedQuote string }{ - {name: "postgres compatible", expected: `"`}, - {name: "mysql compatible", mysqlCompat: true, expected: "`"}, + {name: "postgres compatible", compatibilityMode: "postgresql", expectedQuote: `"`}, + {name: "oracle compatible", compatibilityMode: "oracle", expectedQuote: `"`}, + {name: "mysql compatible", compatibilityMode: "mysql", mysqlCompat: true, expectedQuote: "`"}, } { t.Run(testCase.name, func(t *testing.T) { db := openModeDetectionDB(t, &modeDetectionDriverState{}) server := newServer() server.db = db + server.mode.compatibilityMode = testCase.compatibilityMode server.mode.mysqlCompat = testCase.mysqlCompat info, err := server.connectionInfo() if err != nil { t.Fatal(err) } - if info["identifierQuote"] != testCase.expected { + if info["compatibilityMode"] != testCase.compatibilityMode { + t.Fatalf("unexpected compatibility mode: %#v", info["compatibilityMode"]) + } + if info["identifierQuote"] != testCase.expectedQuote { t.Fatalf("unexpected identifier quote: %#v", info["identifierQuote"]) } }) diff --git a/crates/dbx-core/src/db/agent_driver.rs b/crates/dbx-core/src/db/agent_driver.rs index c66e99073..52a4652d6 100644 --- a/crates/dbx-core/src/db/agent_driver.rs +++ b/crates/dbx-core/src/db/agent_driver.rs @@ -512,6 +512,8 @@ pub struct AgentConnectionInfo { #[serde(default)] pub identifier_quote: String, #[serde(default)] + pub compatibility_mode: Option, + #[serde(default)] pub database_info: Option, } diff --git a/crates/dbx-core/src/table_import.rs b/crates/dbx-core/src/table_import.rs index a2ac47214..c4b87a6c7 100644 --- a/crates/dbx-core/src/table_import.rs +++ b/crates/dbx-core/src/table_import.rs @@ -2895,7 +2895,7 @@ fn build_import_insert_batch_from_rows_with_format( ); } let plan = compile_import_plan(columns, mappings, target_column_types)?; - build_import_insert_batch_with_plan(rows, &plan, table, schema, db_type, date_time_format) + build_import_insert_batch_with_plan(rows, &plan, table, schema, db_type, false, date_time_format) } fn build_import_insert_batch_with_plan( @@ -2904,12 +2904,13 @@ fn build_import_insert_batch_with_plan( table: &str, schema: &str, db_type: &DatabaseType, + kingbase_oracle_mode: bool, date_time_format: Option<&str>, ) -> Result, String> { if rows.is_empty() { return Ok(None); } - let mapped_rows = map_import_rows_with_plan(rows, plan, db_type, date_time_format); + let mapped_rows = map_import_rows_with_plan(rows, plan, db_type, kingbase_oracle_mode, date_time_format); let sql = generate_insert_typed(&plan.target_columns, &plan.column_types, &mapped_rows, table, schema, db_type); Ok((!sql.trim().is_empty()).then_some(ImportSqlBatch { sql, row_count: rows.len() })) } @@ -2920,12 +2921,13 @@ fn build_import_insert_batches_with_plan( table: &str, schema: &str, db_type: &DatabaseType, + kingbase_oracle_mode: bool, date_time_format: Option<&str>, ) -> Vec { if rows.is_empty() { return Vec::new(); } - let mapped_rows = map_import_rows_with_plan(rows, plan, db_type, date_time_format); + let mapped_rows = map_import_rows_with_plan(rows, plan, db_type, kingbase_oracle_mode, date_time_format); generate_insert_typed_sql_batches( &plan.target_columns, &plan.column_types, @@ -2944,6 +2946,7 @@ fn map_import_rows_with_plan( rows: &[Vec], plan: &CompiledImportPlan, db_type: &DatabaseType, + kingbase_oracle_mode: bool, date_time_format: Option<&str>, ) -> Vec> { rows.iter() @@ -2957,6 +2960,7 @@ fn map_import_rows_with_plan( &value, plan.column_types.get(target_index).and_then(|data_type| data_type.as_deref()), db_type, + kingbase_oracle_mode, date_time_format, ) }) @@ -2975,10 +2979,19 @@ fn build_import_execution_batches( table: &str, schema: &str, db_type: &DatabaseType, + kingbase_oracle_mode: bool, date_time_format: Option<&str>, ) -> Result, String> { if let Some(plan) = plan { - return Ok(build_import_insert_batches_with_plan(rows, plan, table, schema, db_type, date_time_format)); + return Ok(build_import_insert_batches_with_plan( + rows, + plan, + table, + schema, + db_type, + kingbase_oracle_mode, + date_time_format, + )); } if *db_type == DatabaseType::CloudflareD1 { return crate::db::cloudflare_d1::build_import_insert_batches( @@ -2992,7 +3005,15 @@ fn build_import_execution_batches( ); } let plan = compile_import_plan(columns, mappings, target_column_types)?; - Ok(build_import_insert_batches_with_plan(rows, &plan, table, schema, db_type, date_time_format)) + Ok(build_import_insert_batches_with_plan( + rows, + &plan, + table, + schema, + db_type, + kingbase_oracle_mode, + date_time_format, + )) } fn effective_import_batch_size(db_type: &DatabaseType, requested: usize) -> usize { @@ -3011,13 +3032,15 @@ fn normalize_import_temporal_value( value: &serde_json::Value, data_type: Option<&str>, db_type: &DatabaseType, + kingbase_oracle_mode: bool, date_time_format: Option<&str>, ) -> serde_json::Value { - let oracle_date_time = matches!(db_type, DatabaseType::Oracle | DatabaseType::OceanbaseOracle) + let date_type_preserves_time = (matches!(db_type, DatabaseType::Oracle | DatabaseType::OceanbaseOracle) + || (*db_type == DatabaseType::Kingbase && kingbase_oracle_mode)) && data_type.is_some_and(|data_type| data_type.trim().eq_ignore_ascii_case("date")); crate::temporal_format::normalize_temporal_import_value( value, - if oracle_date_time { Some("datetime") } else { data_type }, + if date_type_preserves_time { Some("datetime") } else { data_type }, date_time_format, ) } @@ -3087,9 +3110,10 @@ fn normalize_import_value( value: &serde_json::Value, data_type: Option<&str>, db_type: &DatabaseType, + kingbase_oracle_mode: bool, date_time_format: Option<&str>, ) -> serde_json::Value { - normalize_import_temporal_value(value, data_type, db_type, date_time_format) + normalize_import_temporal_value(value, data_type, db_type, kingbase_oracle_mode, date_time_format) } pub fn build_import_insert_batches( @@ -3108,6 +3132,7 @@ pub fn build_import_insert_batches( table, schema, db_type, + false, batch_size, None, ) @@ -3121,6 +3146,7 @@ fn build_import_insert_batches_with_format( table: &str, schema: &str, db_type: &DatabaseType, + kingbase_oracle_mode: bool, batch_size: usize, date_time_format: Option<&str>, ) -> Result, String> { @@ -3139,7 +3165,15 @@ fn build_import_insert_batches_with_format( let batch_size = effective_import_batch_size(db_type, batch_size); let mut batches = Vec::new(); for rows in data.rows.chunks(batch_size) { - batches.extend(build_import_insert_batches_with_plan(rows, &plan, table, schema, db_type, date_time_format)); + batches.extend(build_import_insert_batches_with_plan( + rows, + &plan, + table, + schema, + db_type, + kingbase_oracle_mode, + date_time_format, + )); } Ok(batches) } @@ -3553,7 +3587,7 @@ fn build_postgres_copy_text_batch( schema: &str, date_time_format: Option<&str>, ) -> Result<(String, Vec), String> { - let mapped_rows = map_import_rows_with_plan(rows, plan, &DatabaseType::Postgres, date_time_format); + let mapped_rows = map_import_rows_with_plan(rows, plan, &DatabaseType::Postgres, false, date_time_format); let mut data = String::new(); for row in mapped_rows { for (index, value) in row.iter().enumerate() { @@ -3723,6 +3757,7 @@ async fn execute_import_rows_batch( mode: &TableImportMode, pending_truncate: bool, allow_postgres_copy: bool, + kingbase_oracle_mode: bool, date_time_format: Option<&str>, db_write_ms: &mut u128, statement_count: &mut usize, @@ -3758,6 +3793,7 @@ async fn execute_import_rows_batch( table, schema, db_type, + kingbase_oracle_mode, date_time_format, ) .map_err(ImportRowsBatchError::before_write)?; @@ -4037,6 +4073,26 @@ pub async fn preview_table_import_file_core(file_path: &str) -> Result bool { + if *db_type != DatabaseType::Kingbase { + return false; + } + let client = { + let connections = state.connections.read().await; + match connections.get(pool_key) { + Some(PoolKind::Agent(client)) => client.clone(), + _ => return false, + } + }; + let mut agent = client.lock().await; + agent + .connection_info(Some(crate::db::connection_timeout())) + .await + .ok() + .and_then(|info| info.compatibility_mode) + .is_some_and(|mode| mode.trim().eq_ignore_ascii_case("oracle")) +} + /// Core import logic. Returns (rows_imported, total_rows). /// `progress_callback` is invoked for progress updates. pub async fn import_table_file_core( @@ -4054,6 +4110,7 @@ where let mut db_write_ms = 0u128; let mut statement_count = 0usize; let batch_size = if request.batch_size == 0 { DEFAULT_BATCH_SIZE } else { request.batch_size }; + let kingbase_oracle_mode = kingbase_oracle_compatibility_mode(state, pool_key, db_type).await; let source_format = match effective_source_format(&request.file_path, request.source_format) { Ok(format) => format, Err(error) => { @@ -4386,6 +4443,7 @@ where &request.mode, pending_truncate, allow_postgres_copy, + kingbase_oracle_mode, request.date_time_format.as_deref(), &mut db_write_ms, &mut statement_count, @@ -4686,6 +4744,7 @@ where &request.mode, pending_truncate, allow_postgres_copy, + kingbase_oracle_mode, request.date_time_format.as_deref(), &mut db_write_ms, &mut statement_count, @@ -4901,6 +4960,7 @@ where &request.mode, pending_truncate, allow_postgres_copy, + kingbase_oracle_mode, request.date_time_format.as_deref(), &mut db_write_ms, &mut statement_count, @@ -7277,6 +7337,55 @@ mod tests { ); } + fn kingbase_date_import_sql(oracle_mode: bool) -> String { + let mappings = vec![TableImportColumnMapping { + source_column: "created_at".to_string(), + target_column: "created_at".to_string(), + target_data_type: None, + }]; + let excel_date_time = + Data::DateTime(ExcelDateTime::new(45959.686111111, calamine::ExcelDateTimeType::DateTime, false)); + let imported_value = xlsx_cell_value_with_temporal_kind(&excel_date_time, Some(XlsxTemporalKind::DateTime)); + assert_eq!(imported_value, serde_json::json!("2025-10-29 16:28:00")); + let data = ParsedImportFile { + columns: vec!["created_at".to_string()], + rows: vec![vec![imported_value]], + total_rows: 1, + effective_encoding: None, + }; + + let batches = build_import_insert_batches_with_format( + &data, + &mappings, + &[("created_at".to_string(), "DATE".to_string())], + "events", + "public", + &DatabaseType::Kingbase, + oracle_mode, + 500, + None, + ) + .unwrap(); + + batches[0].sql.clone() + } + + #[test] + fn import_insert_batches_preserve_kingbase_oracle_date_time_components() { + assert_eq!( + kingbase_date_import_sql(true), + "INSERT INTO \"public\".\"events\" (\"created_at\") VALUES\n('2025-10-29 16:28:00')" + ); + } + + #[test] + fn import_insert_batches_normalize_kingbase_postgres_date() { + assert_eq!( + kingbase_date_import_sql(false), + "INSERT INTO \"public\".\"events\" (\"created_at\") VALUES\n('2025-10-29')" + ); + } + #[test] fn import_insert_batch_normalizes_oracle_date_and_timestamp_columns() { let mappings = vec![