fix(data-transfer): handle missing owner role policy and align web sync behavior

* Feat: Add PostgreSQL transfer ownership preview and resolution

- Add ownership policy support for data transfers: preserve, skip, and reassign missing owners
- Preview missing PostgreSQL roles before creating target schema objects
- Show ownership confirmation dialog when target roles are missing
- Allow users to skip ownership changes or reassign missing owners to the target connection user
- Add Tauri and HTTP API support for ownership preview
- Apply ownership statements according to selected policy during PostgreSQL transfers
- Add Chinese localization for ownership confirmation UI

* Feat: Add ownership confirmation localization

- Add ownership confirmation messages for supported desktop locales
- Include skip and confirm labels for ownership handling
- Update transfer test requests to use Preserve ownership policy
This commit is contained in:
Moe. 2026-07-09 00:35:34 +08:00 committed by GitHub
parent 135d4e15db
commit a07f115407
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
18 changed files with 508 additions and 33 deletions

View File

@ -55,6 +55,11 @@ const transferMode = ref<TransferMode>("append");
const targetTableNameCase = ref<TransferTableNameCase>("preserve");
const batchSize = ref(1000);
const isSubmitting = ref(false);
const ownershipDialogOpen = ref(false);
const ownershipMissingOwners = ref<string[]>([]);
const ownershipTargetOwner = ref("");
const pendingOwnershipRequest = ref<api.TransferRequest | null>(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) {
</DialogFooter>
</DialogContent>
</Dialog>
<Dialog v-model:open="ownershipDialogOpen">
<DialogContent class="sm:max-w-[520px]" @interact-outside.prevent>
<DialogHeader>
<DialogTitle>{{ t("transfer.ownershipTitle") }}</DialogTitle>
</DialogHeader>
<div class="space-y-3 text-sm">
<p class="text-muted-foreground">
{{ t("transfer.ownershipMessage", { owners: ownershipMissingOwners.join(", ") }) }}
</p>
<div class="rounded-md border bg-muted/30 px-3 py-2 text-xs text-muted-foreground">
{{ t("transfer.ownershipSkipDetails") }}
</div>
<div class="rounded-md border bg-muted/30 px-3 py-2 text-xs text-muted-foreground">
{{ t("transfer.ownershipTargetOwner", { owner: ownershipTargetOwner }) }}
</div>
</div>
<DialogFooter class="gap-2">
<Button variant="outline" size="sm" @click="resolveOwnershipDecision(null)">
{{ t("transfer.cancel") }}
</Button>
<Button variant="secondary" size="sm" @click="resolveOwnershipDecision('skip')">
{{ t("transfer.ownershipSkip") }}
</Button>
<Button size="sm" @click="resolveOwnershipDecision('reassignMissing')">
{{ t("transfer.ownershipConfirm") }}
</Button>
</DialogFooter>
</DialogContent>
</Dialog>
</template>

View File

@ -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",

View File

@ -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",

View File

@ -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",

View File

@ -2172,6 +2172,12 @@ export default withEnglishFallback({
editConfig: "設定を編集",
retry: "リトライ",
targetTableBusy: "別の転送がすでにターゲットテーブルに書き込み中です: {tables}",
ownershipTitle: "所有者の確認",
ownershipMessage: "対象データベースに次の所有者ユーザーまたはロールがありません: {owners}。元の所有者を適用すると失敗する可能性があります。",
ownershipSkipDetails: "スキップすると、オブジェクト所有者の変更ステートメントは実行されません。テーブルデータとその他の構造の同期は続行されます。",
ownershipTargetOwner: "確認すると、所有者が存在しないオブジェクトは対象接続ユーザーに再割り当てされます: {owner}。",
ownershipSkip: "スキップ",
ownershipConfirm: "確認",
selectConnection: "接続を選択",
selectDatabase: "データベースを選択",
selectSchema: "スキーマを選択",

View File

@ -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",

View File

@ -2242,6 +2242,12 @@ export default withEnglishFallback({
editConfig: "编辑配置",
retry: "重新传输",
targetTableBusy: "已有传输任务正在写入目标表:{tables}",
ownershipTitle: "归属用户确认",
ownershipMessage: "目标端缺少以下归属用户或角色:{owners}。按原归属用户应用可能失败。",
ownershipSkipDetails: "选择跳过后,将不执行对象归属变更语句,表数据和其它结构同步会继续执行。",
ownershipTargetOwner: "选择确认后,缺失归属用户对应的对象会重归属至目标连接用户:{owner}。",
ownershipSkip: "跳过",
ownershipConfirm: "确认",
selectConnection: "选择连接",
selectDatabase: "选择数据库",
selectSchema: "选择模式",

View File

@ -2075,6 +2075,12 @@ export default withEnglishFallback({
editConfig: "編輯配置",
retry: "重試",
targetTableBusy: "另一個傳輸正在寫入目標資料表:{tables}",
ownershipTitle: "歸屬使用者確認",
ownershipMessage: "目標端缺少以下歸屬使用者或角色:{owners}。依原歸屬使用者套用可能失敗。",
ownershipSkipDetails: "選擇跳過後,將不執行物件歸屬變更語句,資料表資料和其他結構同步會繼續執行。",
ownershipTargetOwner: "選擇確認後,缺失歸屬使用者對應的物件會重歸屬至目標連線使用者:{owner}。",
ownershipSkip: "跳過",
ownershipConfirm: "確認",
selectConnection: "選擇連線",
selectDatabase: "選擇資料庫",
selectSchema: "選擇結構描述",

View File

@ -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,

View File

@ -75,6 +75,7 @@ import type {
SqlFileProgress,
TransferRequest,
TransferProgress,
TransferOwnershipPreview,
TableImportPreviewRequest,
TableImportPreview,
TableImportRequest,
@ -1319,6 +1320,10 @@ export async function cancelTransfer(transferId: string): Promise<void> {
return post("/api/transfer/cancel", { transferId });
}
export async function previewTransferOwnership(request: TransferRequest): Promise<TransferOwnershipPreview> {
return post("/api/transfer/ownership-preview", { request });
}
export interface SortTablesByFkOptions {
connectionId: string;
database: string;

View File

@ -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<void> {
return invoke("cancel_transfer", { transferId });
}
export async function previewTransferOwnership(request: TransferRequest): Promise<TransferOwnershipPreview> {
return invoke("preview_transfer_ownership", { request });
}
export interface SortTablesByFkOptions {
connectionId: string;
database: string;

View File

@ -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<String>,
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<String> {
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<serde_json::Value>>) -> Vec<St
rows.into_iter().filter_map(|row| json_string_cell(&row, 0)).filter(|stmt| !stmt.trim().is_empty()).collect()
}
fn result_rows_to_postgres_ownership_statements(rows: Vec<Vec<serde_json::Value>>) -> Vec<PostgresOwnershipStatement> {
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<Vec<String>, String> {
) -> Result<Vec<PostgresOwnershipStatement>, String> {
let table_list = tables.iter().map(|table| quote_string_literal(table)).collect::<Vec<_>>().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<String> {
let mut roles = statements.iter().map(|statement| statement.owner.clone()).collect::<Vec<_>>();
roles.sort();
roles.dedup();
roles
}
async fn get_postgres_current_user(state: &AppState, target_pool_key: &str) -> Result<String, String> {
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<HashSet<String>, String> {
if roles.is_empty() {
return Ok(HashSet::new());
}
let role_list = roles.iter().map(|role| quote_string_literal(role)).collect::<Vec<_>>().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<TransferOwnershipPreview, String> {
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::<Vec<_>>();
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,
}
}

View File

@ -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,
};

View File

@ -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(

View File

@ -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))

View File

@ -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<Arc<WebState>>,
Json(body): Json<StartTransferRequest>,
@ -95,6 +101,64 @@ pub async fn start_transfer(
tables
});
let mut failed_tables: Vec<String> = 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<Arc<WebState>>,
Json(body): Json<PreviewTransferOwnershipRequest>,
) -> Result<Json<dbx_core::transfer::TransferOwnershipPreview>, 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<Arc<WebState>>,
Path(transfer_id): Path<String>,

View File

@ -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<AppState>>,
request: TransferRequest,
) -> Result<TransferOwnershipPreview, String> {
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;

View File

@ -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,