diff --git a/apps/desktop/src/components/transfer/DataTransferDialog.vue b/apps/desktop/src/components/transfer/DataTransferDialog.vue index e15e64ca2..64ddb04b5 100644 --- a/apps/desktop/src/components/transfer/DataTransferDialog.vue +++ b/apps/desktop/src/components/transfer/DataTransferDialog.vue @@ -55,6 +55,11 @@ const transferMode = ref("append"); const targetTableNameCase = ref("preserve"); const batchSize = ref(1000); const isSubmitting = ref(false); +const ownershipDialogOpen = ref(false); +const ownershipMissingOwners = ref([]); +const ownershipTargetOwner = ref(""); +const pendingOwnershipRequest = ref(null); +const pendingOwnershipRefresh = ref<{ targetConnection: string; targetDatabase: string; targetSchema: string; shouldRefreshTargetTree: boolean } | null>(null); const filteredTables = computed(() => { const q = tableSearch.value.toLowerCase(); @@ -247,6 +252,11 @@ function resetState() { targetTableNameCase.value = "preserve"; batchSize.value = 1000; isSubmitting.value = false; + ownershipDialogOpen.value = false; + ownershipMissingOwners.value = []; + ownershipTargetOwner.value = ""; + pendingOwnershipRequest.value = null; + pendingOwnershipRefresh.value = null; } async function startTransfer() { @@ -272,10 +282,39 @@ async function startTransfer() { createTable: createTable.value, mode: transferMode.value, targetTableNameCase: targetTableNameCase.value, + ownershipPolicy: "preserve", batchSize: batchSize.value, }; - startDataTransferTask(request, `${sourceDatabaseName} → ${targetDatabaseName}`, { + if (createTable.value) { + try { + const preview = await api.previewTransferOwnership(request); + if (preview.missingOwners.length > 0) { + ownershipMissingOwners.value = preview.missingOwners; + ownershipTargetOwner.value = preview.targetOwner; + pendingOwnershipRequest.value = request; + pendingOwnershipRefresh.value = { + targetConnection, + targetDatabase: targetDatabaseName, + targetSchema: effectiveTargetSchema, + shouldRefreshTargetTree, + }; + ownershipDialogOpen.value = true; + isSubmitting.value = false; + return; + } + } catch { + isSubmitting.value = false; + return; + } + } + + runTransfer(request, targetConnection, targetDatabaseName, effectiveTargetSchema, shouldRefreshTargetTree); +} + +function runTransfer(request: api.TransferRequest, targetConnection: string, targetDatabaseName: string, effectiveTargetSchema: string, shouldRefreshTargetTree: boolean) { + isSubmitting.value = true; + startDataTransferTask(request, `${request.sourceDatabase} → ${targetDatabaseName}`, { formatOverlapError: (tables) => t("transfer.targetTableBusy", { tables: tables.join(", ") }), onDone: async () => { if (shouldRefreshTargetTree) { @@ -287,6 +326,21 @@ async function startTransfer() { resetState(); } +function resolveOwnershipDecision(policy: api.TransferOwnershipPolicy | null) { + const request = pendingOwnershipRequest.value; + const refresh = pendingOwnershipRefresh.value; + pendingOwnershipRequest.value = null; + pendingOwnershipRefresh.value = null; + ownershipDialogOpen.value = false; + ownershipMissingOwners.value = []; + ownershipTargetOwner.value = ""; + if (!policy || !request || !refresh) { + isSubmitting.value = false; + return; + } + runTransfer({ ...request, ownershipPolicy: policy }, refresh.targetConnection, refresh.targetDatabase, refresh.targetSchema, refresh.shouldRefreshTargetTree); +} + function getConnectionName(id: string) { return store.connections.find((c) => c.id === id)?.name ?? id; } @@ -514,4 +568,34 @@ function getConnectionName(id: string) { + + + + + {{ t("transfer.ownershipTitle") }} + +
+

+ {{ t("transfer.ownershipMessage", { owners: ownershipMissingOwners.join(", ") }) }} +

+
+ {{ t("transfer.ownershipSkipDetails") }} +
+
+ {{ t("transfer.ownershipTargetOwner", { owner: ownershipTargetOwner }) }} +
+
+ + + + + +
+
diff --git a/apps/desktop/src/i18n/locales/en.ts b/apps/desktop/src/i18n/locales/en.ts index 61f6a5df8..bde1635e5 100644 --- a/apps/desktop/src/i18n/locales/en.ts +++ b/apps/desktop/src/i18n/locales/en.ts @@ -2242,6 +2242,12 @@ export default { editConfig: "Edit config", retry: "Retry", targetTableBusy: "Another transfer is already writing target table(s): {tables}", + ownershipTitle: "Confirm ownership", + ownershipMessage: "The target database is missing these owner users or roles: {owners}. Applying the original ownership may fail.", + ownershipSkipDetails: "Skip will not run object ownership change statements. Table data and other structure synchronization will continue.", + ownershipTargetOwner: "Confirm will reassign objects with missing owners to the target connection user: {owner}.", + ownershipSkip: "Skip", + ownershipConfirm: "Confirm", selectConnection: "Select connection", selectDatabase: "Select database", selectSchema: "Select schema", diff --git a/apps/desktop/src/i18n/locales/es.ts b/apps/desktop/src/i18n/locales/es.ts index 9db27e4f5..2d086a8b0 100644 --- a/apps/desktop/src/i18n/locales/es.ts +++ b/apps/desktop/src/i18n/locales/es.ts @@ -2174,6 +2174,12 @@ export default withEnglishFallback({ editConfig: "Editar config", retry: "Reintentar", targetTableBusy: "Otra transferencia ya está escribiendo en la(s) tabla(s) de destino: {tables}", + ownershipTitle: "Confirmar propietario", + ownershipMessage: "La base de datos de destino no tiene estos usuarios o roles propietarios: {owners}. Aplicar el propietario original puede fallar.", + ownershipSkipDetails: "Omitir no ejecutará las sentencias de cambio de propietario de objetos. Los datos de tablas y otras estructuras seguirán sincronizándose.", + ownershipTargetOwner: "Confirmar reasignará los objetos con propietarios ausentes al usuario de conexión de destino: {owner}.", + ownershipSkip: "Omitir", + ownershipConfirm: "Confirmar", selectConnection: "Seleccionar conexión", selectDatabase: "Seleccionar base de datos", selectSchema: "Seleccionar esquema", diff --git a/apps/desktop/src/i18n/locales/it.ts b/apps/desktop/src/i18n/locales/it.ts index 17c67e7c0..9de710d06 100644 --- a/apps/desktop/src/i18n/locales/it.ts +++ b/apps/desktop/src/i18n/locales/it.ts @@ -2172,6 +2172,12 @@ export default withEnglishFallback({ editConfig: "Modifica config", retry: "Riprova", targetTableBusy: "Un altro trasferimento sta gia scrivendo le tabelle di destinazione: {tables}", + ownershipTitle: "Conferma proprietario", + ownershipMessage: "Nel database di destinazione mancano questi utenti o ruoli proprietari: {owners}. Applicare il proprietario originale potrebbe non riuscire.", + ownershipSkipDetails: "Salta non eseguirà le istruzioni di modifica del proprietario degli oggetti. I dati delle tabelle e le altre strutture continueranno a essere sincronizzati.", + ownershipTargetOwner: "Conferma riassegnerà gli oggetti con proprietari mancanti all'utente della connessione di destinazione: {owner}.", + ownershipSkip: "Salta", + ownershipConfirm: "Conferma", selectConnection: "Seleziona connessione", selectDatabase: "Seleziona database", selectSchema: "Seleziona schema", diff --git a/apps/desktop/src/i18n/locales/ja.ts b/apps/desktop/src/i18n/locales/ja.ts index a14aebf35..dd83f573d 100644 --- a/apps/desktop/src/i18n/locales/ja.ts +++ b/apps/desktop/src/i18n/locales/ja.ts @@ -2172,6 +2172,12 @@ export default withEnglishFallback({ editConfig: "設定を編集", retry: "リトライ", targetTableBusy: "別の転送がすでにターゲットテーブルに書き込み中です: {tables}", + ownershipTitle: "所有者の確認", + ownershipMessage: "対象データベースに次の所有者ユーザーまたはロールがありません: {owners}。元の所有者を適用すると失敗する可能性があります。", + ownershipSkipDetails: "スキップすると、オブジェクト所有者の変更ステートメントは実行されません。テーブルデータとその他の構造の同期は続行されます。", + ownershipTargetOwner: "確認すると、所有者が存在しないオブジェクトは対象接続ユーザーに再割り当てされます: {owner}。", + ownershipSkip: "スキップ", + ownershipConfirm: "確認", selectConnection: "接続を選択", selectDatabase: "データベースを選択", selectSchema: "スキーマを選択", diff --git a/apps/desktop/src/i18n/locales/pt-BR.ts b/apps/desktop/src/i18n/locales/pt-BR.ts index f483cb80a..8c870b166 100644 --- a/apps/desktop/src/i18n/locales/pt-BR.ts +++ b/apps/desktop/src/i18n/locales/pt-BR.ts @@ -2173,6 +2173,12 @@ export default withEnglishFallback({ editConfig: "Editar configuração", retry: "Tentar novamente", targetTableBusy: "Outra transferência já está gravando na(s) tabela(s) de destino: {tables}", + ownershipTitle: "Confirmar proprietário", + ownershipMessage: "O banco de dados de destino não tem estes usuários ou roles proprietários: {owners}. Aplicar o proprietário original pode falhar.", + ownershipSkipDetails: "Pular não executará as instruções de alteração de proprietário dos objetos. Os dados das tabelas e outras estruturas continuarão sendo sincronizados.", + ownershipTargetOwner: "Confirmar reatribuirá os objetos com proprietários ausentes ao usuário da conexão de destino: {owner}.", + ownershipSkip: "Pular", + ownershipConfirm: "Confirmar", selectConnection: "Selecionar conexão", selectDatabase: "Selecionar banco de dados", selectSchema: "Selecionar schema", diff --git a/apps/desktop/src/i18n/locales/zh-CN.ts b/apps/desktop/src/i18n/locales/zh-CN.ts index 093d1591a..4bfb351ff 100644 --- a/apps/desktop/src/i18n/locales/zh-CN.ts +++ b/apps/desktop/src/i18n/locales/zh-CN.ts @@ -2242,6 +2242,12 @@ export default withEnglishFallback({ editConfig: "编辑配置", retry: "重新传输", targetTableBusy: "已有传输任务正在写入目标表:{tables}", + ownershipTitle: "归属用户确认", + ownershipMessage: "目标端缺少以下归属用户或角色:{owners}。按原归属用户应用可能失败。", + ownershipSkipDetails: "选择跳过后,将不执行对象归属变更语句,表数据和其它结构同步会继续执行。", + ownershipTargetOwner: "选择确认后,缺失归属用户对应的对象会重归属至目标连接用户:{owner}。", + ownershipSkip: "跳过", + ownershipConfirm: "确认", selectConnection: "选择连接", selectDatabase: "选择数据库", selectSchema: "选择模式", diff --git a/apps/desktop/src/i18n/locales/zh-TW.ts b/apps/desktop/src/i18n/locales/zh-TW.ts index 2d84fd45f..990fbe5a4 100644 --- a/apps/desktop/src/i18n/locales/zh-TW.ts +++ b/apps/desktop/src/i18n/locales/zh-TW.ts @@ -2075,6 +2075,12 @@ export default withEnglishFallback({ editConfig: "編輯配置", retry: "重試", targetTableBusy: "另一個傳輸正在寫入目標資料表:{tables}", + ownershipTitle: "歸屬使用者確認", + ownershipMessage: "目標端缺少以下歸屬使用者或角色:{owners}。依原歸屬使用者套用可能失敗。", + ownershipSkipDetails: "選擇跳過後,將不執行物件歸屬變更語句,資料表資料和其他結構同步會繼續執行。", + ownershipTargetOwner: "選擇確認後,缺失歸屬使用者對應的物件會重歸屬至目標連線使用者:{owner}。", + ownershipSkip: "跳過", + ownershipConfirm: "確認", selectConnection: "選擇連線", selectDatabase: "選擇資料庫", selectSchema: "選擇結構描述", diff --git a/apps/desktop/src/lib/backend/api.ts b/apps/desktop/src/lib/backend/api.ts index 766c6e822..afc7f306c 100644 --- a/apps/desktop/src/lib/backend/api.ts +++ b/apps/desktop/src/lib/backend/api.ts @@ -278,6 +278,7 @@ export const nacosRawRequest = forward("nacosRawRequest"); // Data Transfer export const startTransfer = forward("startTransfer"); export const cancelTransfer = forward("cancelTransfer"); +export const previewTransferOwnership = forward("previewTransferOwnership"); export const sortTablesByFkDependency = forward("sortTablesByFkDependency"); // Table File Import @@ -511,6 +512,8 @@ export type { TransferProgress, TransferMode, TransferTableNameCase, + TransferOwnershipPolicy, + TransferOwnershipPreview, TableImportMode, TableImportStatus, TableImportSourceFormat, diff --git a/apps/desktop/src/lib/backend/http.ts b/apps/desktop/src/lib/backend/http.ts index 8a9a1db4c..2987a001c 100644 --- a/apps/desktop/src/lib/backend/http.ts +++ b/apps/desktop/src/lib/backend/http.ts @@ -75,6 +75,7 @@ import type { SqlFileProgress, TransferRequest, TransferProgress, + TransferOwnershipPreview, TableImportPreviewRequest, TableImportPreview, TableImportRequest, @@ -1319,6 +1320,10 @@ export async function cancelTransfer(transferId: string): Promise { return post("/api/transfer/cancel", { transferId }); } +export async function previewTransferOwnership(request: TransferRequest): Promise { + return post("/api/transfer/ownership-preview", { request }); +} + export interface SortTablesByFkOptions { connectionId: string; database: string; diff --git a/apps/desktop/src/lib/backend/tauri.ts b/apps/desktop/src/lib/backend/tauri.ts index 5e120ad1e..fbcadeb45 100644 --- a/apps/desktop/src/lib/backend/tauri.ts +++ b/apps/desktop/src/lib/backend/tauri.ts @@ -1889,6 +1889,7 @@ export async function listenSqlFileProgress(handler: (progress: SqlFileProgress) // --- Data Transfer --- export type TransferMode = "append" | "overwrite" | "upsert"; export type TransferTableNameCase = "preserve" | "lower" | "upper"; +export type TransferOwnershipPolicy = "preserve" | "skip" | "reassignMissing"; export interface TransferRequest { transferId: string; @@ -1902,9 +1903,15 @@ export interface TransferRequest { createTable: boolean; mode: TransferMode; targetTableNameCase: TransferTableNameCase; + ownershipPolicy?: TransferOwnershipPolicy; batchSize: number; } +export interface TransferOwnershipPreview { + missingOwners: string[]; + targetOwner: string; +} + export interface TransferProgress { transferId: string; table: string; @@ -1943,6 +1950,10 @@ export async function cancelTransfer(transferId: string): Promise { return invoke("cancel_transfer", { transferId }); } +export async function previewTransferOwnership(request: TransferRequest): Promise { + return invoke("preview_transfer_ownership", { request }); +} + export interface SortTablesByFkOptions { connectionId: string; database: string; diff --git a/crates/dbx-core/src/transfer.rs b/crates/dbx-core/src/transfer.rs index 3b9cd19b7..c2724f716 100644 --- a/crates/dbx-core/src/transfer.rs +++ b/crates/dbx-core/src/transfer.rs @@ -41,6 +41,15 @@ pub enum TransferTableNameCase { Upper, } +#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq)] +#[serde(rename_all = "camelCase")] +pub enum TransferOwnershipPolicy { + #[default] + Preserve, + Skip, + ReassignMissing, +} + #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct TransferRequest { @@ -57,9 +66,18 @@ pub struct TransferRequest { pub mode: TransferMode, #[serde(default)] pub target_table_name_case: TransferTableNameCase, + #[serde(default)] + pub ownership_policy: TransferOwnershipPolicy, pub batch_size: usize, } +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct TransferOwnershipPreview { + pub missing_owners: Vec, + pub target_owner: String, +} + impl TransferRequest { pub fn target_table_name(&self, source_table: &str) -> String { match self.target_table_name_case { @@ -795,6 +813,12 @@ struct PostgresMaterializedViewSource { source: String, } +#[derive(Debug, Clone)] +struct PostgresOwnershipStatement { + sql_prefix: String, + owner: String, +} + fn json_string_cell(row: &[serde_json::Value], index: usize) -> Option { row.get(index).and_then(|value| value.as_str().map(str::to_string)) } @@ -803,6 +827,20 @@ fn result_rows_to_string_statements(rows: Vec>) -> Vec>) -> Vec { + rows.into_iter() + .filter_map(|row| { + let sql_prefix = json_string_cell(&row, 0)?; + let owner = json_string_cell(&row, 1)?; + if sql_prefix.trim().is_empty() || owner.trim().is_empty() { + None + } else { + Some(PostgresOwnershipStatement { sql_prefix, owner }) + } + }) + .collect() +} + fn ensure_sql_statement_terminated(sql: &str) -> String { let trimmed = sql.trim(); if trimmed.ends_with(';') { @@ -3123,50 +3161,120 @@ async fn get_postgres_ownership_statements_for_transfer( source_schema: &str, target_schema: &str, tables: &[String], -) -> Result, String> { +) -> Result, String> { let table_list = tables.iter().map(|table| quote_string_literal(table)).collect::>().join(", "); let table_filter = if tables.is_empty() { "FALSE".to_string() } else { format!("c.relname IN ({table_list})") }; let sql = format!( "WITH relation_owners AS ( \ SELECT CASE c.relkind \ - WHEN 'm' THEN format('ALTER MATERIALIZED VIEW %I.%I OWNER TO %I', {target_schema}, c.relname, pg_get_userbyid(c.relowner)) \ - WHEN 'v' THEN format('ALTER VIEW %I.%I OWNER TO %I', {target_schema}, c.relname, pg_get_userbyid(c.relowner)) \ - WHEN 'f' THEN format('ALTER FOREIGN TABLE %I.%I OWNER TO %I', {target_schema}, c.relname, pg_get_userbyid(c.relowner)) \ - WHEN 'S' THEN format('ALTER SEQUENCE %I.%I OWNER TO %I', {target_schema}, c.relname, pg_get_userbyid(c.relowner)) \ - ELSE format('ALTER TABLE %I.%I OWNER TO %I', {target_schema}, c.relname, pg_get_userbyid(c.relowner)) \ - END AS stmt \ + WHEN 'm' THEN format('ALTER MATERIALIZED VIEW %I.%I OWNER TO ', {target_schema}, c.relname) \ + WHEN 'v' THEN format('ALTER VIEW %I.%I OWNER TO ', {target_schema}, c.relname) \ + WHEN 'f' THEN format('ALTER FOREIGN TABLE %I.%I OWNER TO ', {target_schema}, c.relname) \ + WHEN 'S' THEN format('ALTER SEQUENCE %I.%I OWNER TO ', {target_schema}, c.relname) \ + ELSE format('ALTER TABLE %I.%I OWNER TO ', {target_schema}, c.relname) \ + END AS stmt_prefix, \ + pg_get_userbyid(c.relowner) AS owner_name \ FROM pg_catalog.pg_class c \ JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace \ WHERE n.nspname = {source_schema} AND (c.relkind IN ('v','m') OR ({table_filter} AND c.relkind IN ('r','p','f','S'))) \ ), \ routine_owners AS ( \ - SELECT format('ALTER %s %I.%I(%s) OWNER TO %I', \ + SELECT format('ALTER %s %I.%I(%s) OWNER TO ', \ CASE p.prokind WHEN 'p' THEN 'PROCEDURE' ELSE 'FUNCTION' END, \ - {target_schema}, p.proname, pg_get_function_identity_arguments(p.oid), pg_get_userbyid(p.proowner)) AS stmt \ + {target_schema}, p.proname, pg_get_function_identity_arguments(p.oid)) AS stmt_prefix, \ + pg_get_userbyid(p.proowner) AS owner_name \ FROM pg_catalog.pg_proc p \ JOIN pg_catalog.pg_namespace n ON n.oid = p.pronamespace \ WHERE n.nspname = {source_schema} AND p.prokind IN ('p','f') \ ), \ type_owners AS ( \ - SELECT format('ALTER %s %I.%I OWNER TO %I', \ + SELECT format('ALTER %s %I.%I OWNER TO ', \ CASE t.typtype WHEN 'd' THEN 'DOMAIN' ELSE 'TYPE' END, \ - {target_schema}, t.typname, pg_get_userbyid(t.typowner)) AS stmt \ + {target_schema}, t.typname) AS stmt_prefix, \ + pg_get_userbyid(t.typowner) AS owner_name \ FROM pg_catalog.pg_type t \ JOIN pg_catalog.pg_namespace n ON n.oid = t.typnamespace \ WHERE n.nspname = {source_schema} AND t.typtype IN ('e','d') \ ) \ - SELECT stmt FROM ( \ - SELECT format('ALTER SCHEMA %I OWNER TO %I', {target_schema}, pg_get_userbyid(n.nspowner)) AS stmt \ + SELECT stmt_prefix, owner_name FROM ( \ + SELECT format('ALTER SCHEMA %I OWNER TO ', {target_schema}) AS stmt_prefix, \ + pg_get_userbyid(n.nspowner) AS owner_name \ FROM pg_catalog.pg_namespace n WHERE n.nspname = {source_schema} \ - UNION ALL SELECT stmt FROM relation_owners \ - UNION ALL SELECT stmt FROM routine_owners \ - UNION ALL SELECT stmt FROM type_owners \ - ) statements", + UNION ALL SELECT stmt_prefix, owner_name FROM relation_owners \ + UNION ALL SELECT stmt_prefix, owner_name FROM routine_owners \ + UNION ALL SELECT stmt_prefix, owner_name FROM type_owners \ + ) statements \ + WHERE stmt_prefix IS NOT NULL AND owner_name IS NOT NULL", source_schema = quote_string_literal(source_schema), target_schema = quote_string_literal(target_schema), table_filter = table_filter, ); - Ok(result_rows_to_string_statements(execute_on_pool(state, pool_key, &sql).await?.rows)) + Ok(result_rows_to_postgres_ownership_statements(execute_on_pool(state, pool_key, &sql).await?.rows)) +} + +fn distinct_postgres_ownership_roles(statements: &[PostgresOwnershipStatement]) -> Vec { + let mut roles = statements.iter().map(|statement| statement.owner.clone()).collect::>(); + roles.sort(); + roles.dedup(); + roles +} + +async fn get_postgres_current_user(state: &AppState, target_pool_key: &str) -> Result { + let rows = execute_on_pool(state, target_pool_key, "SELECT current_user").await?.rows; + rows.first() + .and_then(|row| json_string_cell(row, 0)) + .filter(|user| !user.trim().is_empty()) + .ok_or_else(|| "Failed to read target PostgreSQL current user".to_string()) +} + +async fn get_existing_postgres_roles( + state: &AppState, + target_pool_key: &str, + roles: &[String], +) -> Result, String> { + if roles.is_empty() { + return Ok(HashSet::new()); + } + let role_list = roles.iter().map(|role| quote_string_literal(role)).collect::>().join(", "); + let sql = format!("SELECT rolname FROM pg_catalog.pg_roles WHERE rolname IN ({role_list})"); + let rows = execute_on_pool(state, target_pool_key, &sql).await?.rows; + Ok(rows.into_iter().filter_map(|row| json_string_cell(&row, 0)).collect()) +} + +fn build_postgres_ownership_statement(statement: &PostgresOwnershipStatement, owner: &str) -> String { + format!("{}{}", statement.sql_prefix, quote_identifier(owner, &DatabaseType::Postgres)) +} + +pub async fn preview_transfer_ownership( + state: &AppState, + request: &TransferRequest, + source_db_type: &DatabaseType, + target_db_type: &DatabaseType, + source_pool_key: &str, + target_pool_key: &str, +) -> Result { + if !request.create_table || !is_postgres_compat_transfer(source_db_type, target_db_type) { + return Ok(TransferOwnershipPreview { missing_owners: Vec::new(), target_owner: String::new() }); + } + + let statements = get_postgres_ownership_statements_for_transfer( + state, + source_pool_key, + &request.source_schema, + &request.target_schema, + &request.tables, + ) + .await?; + let roles = distinct_postgres_ownership_roles(&statements); + let existing_roles = get_existing_postgres_roles(state, target_pool_key, &roles).await?; + let missing_owners = roles.into_iter().filter(|role| !existing_roles.contains(role)).collect::>(); + let target_owner = if missing_owners.is_empty() { + String::new() + } else { + get_postgres_current_user(state, target_pool_key).await? + }; + + Ok(TransferOwnershipPreview { missing_owners, target_owner }) } async fn get_postgres_grant_statements_for_transfer( @@ -4114,14 +4222,31 @@ where &request.tables, ) .await?; - let ownership_statements = get_postgres_ownership_statements_for_transfer( - state, - source_pool_key, - &request.source_schema, - &request.target_schema, - &request.tables, - ) - .await?; + let ownership_statements = if matches!(request.ownership_policy, TransferOwnershipPolicy::Skip) { + Vec::new() + } else { + get_postgres_ownership_statements_for_transfer( + state, + source_pool_key, + &request.source_schema, + &request.target_schema, + &request.tables, + ) + .await? + }; + let ownership_existing_roles = if matches!(request.ownership_policy, TransferOwnershipPolicy::ReassignMissing) { + let roles = distinct_postgres_ownership_roles(&ownership_statements); + get_existing_postgres_roles(state, target_pool_key, &roles).await? + } else { + HashSet::new() + }; + let ownership_target_user = if matches!(request.ownership_policy, TransferOwnershipPolicy::ReassignMissing) + && !ownership_statements.is_empty() + { + Some(get_postgres_current_user(state, target_pool_key).await?) + } else { + None + }; let grant_statements = get_postgres_grant_statements_for_transfer( state, source_pool_key, @@ -4286,7 +4411,17 @@ where status: TransferStatus::Running, error: None, }); - execute_on_pool(state, target_pool_key, &statement) + let ownership_owner = if matches!(request.ownership_policy, TransferOwnershipPolicy::ReassignMissing) + && !ownership_existing_roles.contains(&statement.owner) + { + ownership_target_user + .as_deref() + .ok_or_else(|| "Failed to read target PostgreSQL current user".to_string())? + } else { + &statement.owner + }; + let ownership_sql = build_postgres_ownership_statement(&statement, ownership_owner); + execute_on_pool(state, target_pool_key, &ownership_sql) .await .map_err(|e| format!("Failed to apply PostgreSQL ownership statement: {e}"))?; } @@ -4431,6 +4566,7 @@ mod tests { create_table: true, mode: TransferMode::Append, target_table_name_case: TransferTableNameCase::Preserve, + ownership_policy: TransferOwnershipPolicy::Preserve, batch_size: 1000, } } diff --git a/crates/dbx-core/tests/live_postgres_transfer.rs b/crates/dbx-core/tests/live_postgres_transfer.rs index c2847c4c1..b25d881d3 100644 --- a/crates/dbx-core/tests/live_postgres_transfer.rs +++ b/crates/dbx-core/tests/live_postgres_transfer.rs @@ -4,7 +4,7 @@ use dbx_core::models::connection::{ConnectionConfig, DatabaseType}; use dbx_core::storage::Storage; use dbx_core::transfer::{ get_db_type, transfer_postgres_schema_dependencies, transfer_postgres_schema_objects, transfer_table, TransferMode, - TransferRequest, TransferTableNameCase, + TransferOwnershipPolicy, TransferRequest, TransferTableNameCase, }; use serde_json::json; @@ -211,6 +211,7 @@ async fn live_postgres_transfer_preserves_data_and_schema_objects() { create_table: true, mode: TransferMode::Append, target_table_name_case: TransferTableNameCase::Preserve, + ownership_policy: TransferOwnershipPolicy::Preserve, batch_size: 100, }; @@ -491,6 +492,7 @@ async fn live_postgres_transfer_skips_create_ddl_for_existing_target_table() { create_table: true, mode: TransferMode::Append, target_table_name_case: TransferTableNameCase::Preserve, + ownership_policy: TransferOwnershipPolicy::Preserve, batch_size: 100, }; diff --git a/crates/dbx-core/tests/live_sqlserver_completion.rs b/crates/dbx-core/tests/live_sqlserver_completion.rs index 80762d387..7ec787428 100644 --- a/crates/dbx-core/tests/live_sqlserver_completion.rs +++ b/crates/dbx-core/tests/live_sqlserver_completion.rs @@ -401,6 +401,7 @@ async fn live_sqlserver_transfer_table_skips_rowversion_insert_column() { create_table: true, mode: dbx_core::transfer::TransferMode::Append, target_table_name_case: dbx_core::transfer::TransferTableNameCase::Upper, + ownership_policy: dbx_core::transfer::TransferOwnershipPolicy::Preserve, batch_size: 100, }; let result = dbx_core::transfer::transfer_table( diff --git a/crates/dbx-web/src/main.rs b/crates/dbx-web/src/main.rs index 9e0c1cd56..e5bad506a 100644 --- a/crates/dbx-web/src/main.rs +++ b/crates/dbx-web/src/main.rs @@ -507,6 +507,7 @@ async fn main() { .route("/ai/models", post(routes::ai::ai_list_models)) // Transfer .route("/transfer/start", post(routes::transfer::start_transfer)) + .route("/transfer/ownership-preview", post(routes::transfer::preview_transfer_ownership)) .route("/transfer/progress/{transferId}", get(routes::transfer::transfer_progress)) .route("/transfer/cancel", post(routes::transfer::cancel_transfer)) .route("/transfer/sort-tables-by-fk", post(routes::transfer::sort_tables_by_fk_dependency)) diff --git a/crates/dbx-web/src/routes/transfer.rs b/crates/dbx-web/src/routes/transfer.rs index c77d371a8..00d12b0b4 100644 --- a/crates/dbx-web/src/routes/transfer.rs +++ b/crates/dbx-web/src/routes/transfer.rs @@ -22,6 +22,12 @@ pub struct CancelTransferRequest { pub transfer_id: String, } +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct PreviewTransferOwnershipRequest { + pub request: TransferRequest, +} + pub async fn start_transfer( State(state): State>, Json(body): Json, @@ -95,6 +101,64 @@ pub async fn start_transfer( tables }); let mut failed_tables: Vec = Vec::new(); + + if matches!(source_db_type, dbx_core::models::connection::DatabaseType::Postgres) + && matches!(target_db_type, dbx_core::models::connection::DatabaseType::Postgres) + { + let tx_clone = tx.clone(); + match transfer::transfer_postgres_schema_dependencies( + &app, + &req, + &source_pool_key, + &target_pool_key, + |progress| { + if let Ok(json) = serde_json::to_string(&progress) { + let _ = tx_clone.send(json); + } + }, + ) + .await + { + Ok(()) => {} + Err(e) if e == "Cancelled" => { + let progress = transfer::TransferProgress { + transfer_id: req.transfer_id.clone(), + table: "schema dependencies".to_string(), + table_index: 0, + total_tables: tables.len(), + rows_transferred: 0, + total_rows: None, + status: TransferStatus::Cancelled, + error: None, + }; + if let Ok(json) = serde_json::to_string(&progress) { + let _ = tx.send(json); + } + transfer::clear_cancelled(&req.transfer_id).await; + state_clone.remove_sse_channel(&req.transfer_id).await; + return; + } + Err(e) => { + let progress = transfer::TransferProgress { + transfer_id: req.transfer_id.clone(), + table: "schema dependencies".to_string(), + table_index: 0, + total_tables: tables.len(), + rows_transferred: 0, + total_rows: None, + status: TransferStatus::Error, + error: Some(e), + }; + if let Ok(json) = serde_json::to_string(&progress) { + let _ = tx.send(json); + } + transfer::clear_cancelled(&req.transfer_id).await; + state_clone.remove_sse_channel(&req.transfer_id).await; + return; + } + } + } + for (i, table) in tables.iter().enumerate() { if transfer::is_cancelled(&req.transfer_id).await { let progress = transfer::TransferProgress { @@ -138,14 +202,14 @@ pub async fn start_transfer( .await; match result { - Ok(_) => { + Ok(rows) => { let progress = transfer::TransferProgress { transfer_id: req.transfer_id.clone(), table: table.clone(), table_index: i, total_tables: tables.len(), - rows_transferred: last_rows_transferred, - total_rows: last_total_rows.or(Some(last_rows_transferred)), + rows_transferred: rows, + total_rows: last_total_rows.or(Some(rows)), status: TransferStatus::TableDone, error: None, }; @@ -154,6 +218,24 @@ pub async fn start_transfer( } } Err(e) => { + if e == "Cancelled" { + let progress = transfer::TransferProgress { + transfer_id: req.transfer_id.clone(), + table: table.clone(), + table_index: i, + total_tables: tables.len(), + rows_transferred: 0, + total_rows: None, + status: TransferStatus::Cancelled, + error: None, + }; + if let Ok(json) = serde_json::to_string(&progress) { + let _ = tx.send(json); + } + transfer::clear_cancelled(&req.transfer_id).await; + state_clone.remove_sse_channel(&req.transfer_id).await; + return; + } failed_tables.push(table.clone()); let progress = transfer::TransferProgress { transfer_id: req.transfer_id.clone(), @@ -172,6 +254,61 @@ pub async fn start_transfer( } } + if matches!(source_db_type, dbx_core::models::connection::DatabaseType::Postgres) + && matches!(target_db_type, dbx_core::models::connection::DatabaseType::Postgres) + { + let tx_clone = tx.clone(); + match transfer::transfer_postgres_schema_objects( + &app, + &req, + &source_pool_key, + &target_pool_key, + |progress| { + if let Ok(json) = serde_json::to_string(&progress) { + let _ = tx_clone.send(json); + } + }, + ) + .await + { + Ok(()) => {} + Err(e) if e == "Cancelled" => { + let progress = transfer::TransferProgress { + transfer_id: req.transfer_id.clone(), + table: "schema objects".to_string(), + table_index: tables.len(), + total_tables: tables.len(), + rows_transferred: 0, + total_rows: None, + status: TransferStatus::Cancelled, + error: None, + }; + if let Ok(json) = serde_json::to_string(&progress) { + let _ = tx.send(json); + } + transfer::clear_cancelled(&req.transfer_id).await; + state_clone.remove_sse_channel(&req.transfer_id).await; + return; + } + Err(e) => { + failed_tables.push("schema objects".to_string()); + let progress = transfer::TransferProgress { + transfer_id: req.transfer_id.clone(), + table: "schema objects".to_string(), + table_index: tables.len(), + total_tables: tables.len(), + rows_transferred: 0, + total_rows: None, + status: TransferStatus::Error, + error: Some(e), + }; + if let Ok(json) = serde_json::to_string(&progress) { + let _ = tx.send(json); + } + } + } + } + // Send done let done = transfer::TransferProgress { transfer_id: req.transfer_id.clone(), @@ -202,6 +339,31 @@ pub async fn start_transfer( Ok(Json(serde_json::json!({ "transferId": transfer_id }))) } +pub async fn preview_transfer_ownership( + State(state): State>, + Json(body): Json, +) -> Result, AppError> { + let req = body.request; + transfer::validate_transfer_target_table_names(&req).map_err(AppError)?; + let source_db_type = transfer::get_db_type(&state.app, &req.source_connection_id).await.map_err(AppError)?; + let target_db_type = transfer::get_db_type(&state.app, &req.target_connection_id).await.map_err(AppError)?; + let source_pool_key = + state.app.get_or_create_pool(&req.source_connection_id, Some(&req.source_database)).await.map_err(AppError)?; + let target_pool_key = + state.app.get_or_create_pool(&req.target_connection_id, Some(&req.target_database)).await.map_err(AppError)?; + let preview = transfer::preview_transfer_ownership( + &state.app, + &req, + &source_db_type, + &target_db_type, + &source_pool_key, + &target_pool_key, + ) + .await + .map_err(AppError)?; + Ok(Json(preview)) +} + pub async fn transfer_progress( State(state): State>, Path(transfer_id): Path, diff --git a/src-tauri/src/commands/transfer.rs b/src-tauri/src/commands/transfer.rs index 47b09b9b2..f3055813b 100644 --- a/src-tauri/src/commands/transfer.rs +++ b/src-tauri/src/commands/transfer.rs @@ -4,7 +4,9 @@ use tauri::{AppHandle, Emitter, State}; use crate::commands::connection::{ensure_connection_writable, AppState}; // Re-export types and functions used by other modules -pub use dbx_core::transfer::{get_db_type, TransferProgress, TransferRequest, TransferStatus}; +pub use dbx_core::transfer::{ + get_db_type, TransferOwnershipPreview, TransferProgress, TransferRequest, TransferStatus, +}; fn emit_progress(app: &AppHandle, progress: TransferProgress) { let _ = app.emit("transfer-progress", progress); @@ -271,6 +273,31 @@ pub async fn start_transfer( Ok(()) } +#[tauri::command] +pub async fn preview_transfer_ownership( + state: State<'_, Arc>, + request: TransferRequest, +) -> Result { + let state = state.inner().clone(); + let source_db_type = get_db_type(&state, &request.source_connection_id).await?; + let target_db_type = get_db_type(&state, &request.target_connection_id).await?; + dbx_core::transfer::validate_transfer_target_table_names(&request)?; + let source_pool_key = + state.get_or_create_pool(&request.source_connection_id, Some(&request.source_database)).await?; + let target_pool_key = + state.get_or_create_pool(&request.target_connection_id, Some(&request.target_database)).await?; + + dbx_core::transfer::preview_transfer_ownership( + &state, + &request, + &source_db_type, + &target_db_type, + &source_pool_key, + &target_pool_key, + ) + .await +} + #[tauri::command] pub async fn cancel_transfer(transfer_id: String) -> Result<(), String> { dbx_core::transfer::set_cancelled(&transfer_id).await; diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index c496e6c1a..a6b06c364 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -1105,6 +1105,7 @@ pub fn run() { commands::update::get_system_proxy_url, commands::update::download_and_install_update, commands::transfer::start_transfer, + commands::transfer::preview_transfer_ownership, commands::transfer::cancel_transfer, commands::database_export::export_database_sql, commands::database_export::cancel_database_export,