fix(transfer): support Kingbase metadata pools

This commit is contained in:
t8y2 2026-08-01 15:14:49 +08:00
parent 8282f7ae4e
commit 7a612cb747
No known key found for this signature in database
1 changed files with 65 additions and 2 deletions

View File

@ -3417,10 +3417,20 @@ pub async fn get_columns_for_transfer(
async fn get_postgres_indexes_for_transfer(
state: &AppState,
pool_key: &str,
database: &str,
schema: &str,
table: &str,
) -> Result<Vec<db::IndexInfo>, String> {
let connections = state.connections.read().await;
if let Some(PoolKind::Agent(client)) = connections.get(pool_key) {
let client = client.clone();
let database = database.to_string();
let schema = schema.to_string();
let table = table.to_string();
drop(connections);
let mut client = client.lock().await;
return client.list_indexes(&database, &schema, &table, None).await;
}
let Some(PoolKind::Postgres(pool)) = connections.get(pool_key) else {
return Err("PostgreSQL pool not found".to_string());
};
@ -3432,10 +3442,20 @@ async fn get_postgres_indexes_for_transfer(
async fn get_postgres_foreign_keys_for_transfer(
state: &AppState,
pool_key: &str,
database: &str,
schema: &str,
table: &str,
) -> Result<Vec<db::ForeignKeyInfo>, String> {
let connections = state.connections.read().await;
if let Some(PoolKind::Agent(client)) = connections.get(pool_key) {
let client = client.clone();
let database = database.to_string();
let schema = schema.to_string();
let table = table.to_string();
drop(connections);
let mut client = client.lock().await;
return client.list_foreign_keys(&database, &schema, &table, None).await;
}
let Some(PoolKind::Postgres(pool)) = connections.get(pool_key) else {
return Err("PostgreSQL pool not found".to_string());
};
@ -4572,13 +4592,27 @@ where
let source_indexes =
if request.create_table && pg_compat_transfer && preserves_target_table_name && !target_table_preexisting {
get_postgres_indexes_for_transfer(state, source_pool_key, &request.source_schema, table).await?
get_postgres_indexes_for_transfer(
state,
source_pool_key,
&request.source_database,
&request.source_schema,
table,
)
.await?
} else {
Vec::new()
};
let source_foreign_keys =
if request.create_table && pg_compat_transfer && preserves_target_table_name && !target_table_preexisting {
get_postgres_foreign_keys_for_transfer(state, source_pool_key, &request.source_schema, table).await?
get_postgres_foreign_keys_for_transfer(
state,
source_pool_key,
&request.source_database,
&request.source_schema,
table,
)
.await?
} else {
Vec::new()
};
@ -5335,6 +5369,35 @@ mod tests {
}
}
async fn test_app_state() -> (AppState, std::path::PathBuf) {
let dir = std::env::temp_dir().join(format!("dbx-transfer-test-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&dir).unwrap();
let storage = crate::storage::Storage::open(&dir.join("storage.db")).await.unwrap();
(AppState::new(storage), dir)
}
#[tokio::test]
async fn postgres_transfer_metadata_routes_agent_pools() {
let (state, dir) = test_app_state().await;
state.connections.write().await.insert(
"source:source_db".to_string(),
PoolKind::agent(crate::db::agent_driver::AgentDriverClient::test_stub()),
);
let index_error =
get_postgres_indexes_for_transfer(&state, "source:source_db", "source_db", "source_schema", "items")
.await
.unwrap_err();
let foreign_key_error =
get_postgres_foreign_keys_for_transfer(&state, "source:source_db", "source_db", "source_schema", "items")
.await
.unwrap_err();
assert!(!index_error.contains("PostgreSQL pool not found"), "index error: {index_error}");
assert!(!foreign_key_error.contains("PostgreSQL pool not found"), "foreign key error: {foreign_key_error}");
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn postgres_schema_exists_query_escapes_schema_name() {
assert_eq!(