Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ members = [
"packages/runtime",
"packages/schema-diff",
"packages/data-sync",
"packages/migration-common",
"packages/data-transfer",
"packages/drivers/*",
]
Expand All @@ -23,6 +24,7 @@ datazen-driver-api = { path = "packages/driver-api" }
datazen-platform-api = { path = "packages/platform-api" }
datazen-application = { path = "packages/application" }
datazen-schema-diff = { path = "packages/schema-diff" }
datazen-migration-common = { path = "packages/migration-common" }
datazen-data-sync = { path = "packages/data-sync" }
datazen-data-transfer = { path = "packages/data-transfer" }

Expand Down
2 changes: 2 additions & 0 deletions docs/architecture/backend/data-sync.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

> P5 桌面 JobRuntime 接入已完成;迁移三件套独立测试与冒烟测试通过。Source of truth: `packages/data-sync/`、`src-tauri/src/commands/sync/` 和 `src/windows/data-sync/`。

共享边界值规范化、比较与端点配对算法位于 `packages/migration-common/`;SQL 名称处理工具位于 `packages/driver-api/src/sql_identifiers.rs`。Data Sync 与 Data Transfer 直接依赖这些公共实现,Data Transfer 不依赖 Data Sync 领域包。Job 生命周期继续复用 `packages/runtime/src/job/`。

Data Sync 用于**同族数据库的行级差异同步**。它与 Schema Diff、Data Transfer 是三个独立执行模型。

| 能力 | 用途 |
Expand Down
2 changes: 2 additions & 0 deletions docs/architecture/backend/data-transfer.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

> 当前实现说明;P5 桌面 JobRuntime 与数据迁移三件套于 2026-10-09 完成并通过独立测试、冒烟测试。Data Transfer 的实现事实见 `packages/data-transfer/`、`src-tauri/src/commands/data_transfer/`、`packages/runtime/src/job/` 与 `src-tauri/src/store/app_db/jobs/`;后续服务端与跨进程恢复边界见[迁移任务详细设计](../platform/data-migration-jobs.md)。

共享边界值规范化、比较与端点配对算法位于 `packages/migration-common/`;SQL 名称处理工具位于 `packages/driver-api/src/sql_identifiers.rs`。Data Sync 与 Data Transfer 直接依赖这些公共实现,Data Transfer 不依赖 Data Sync 领域包。Job 生命周期继续复用 `packages/runtime/src/job/`。

## 1. 职责与执行边界

Data Transfer 搬运结构和数据,支持异构数据库、表列映射,以及 SQL 文件目的地。它不承担 Schema Diff 的差异 DAG,也不使用 Data Sync 的同行 ChangeSet。三者可以共用驱动与类型模型,但执行入口、检查和恢复证据分别成立。
Expand Down
4 changes: 4 additions & 0 deletions docs/architecture/platform/data-migration-jobs.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,10 @@

三个领域包和对应桌面 Job adapter 已由 P5 落地。`src-tauri/src/lib.rs` 将领域包导出给 Host;Driver 方言、DDL renderer、类型适配和专属测试留在 driver 包。领域引擎不引用 Tauri、HTTP、窗口 Store 或 driver 实现库类型。

公共代码的实际归属:三套引擎共享 `runtime` 的 JobHandler/JobRuntime 和 `platform-api` 的任务 DTO;SQL 标识符引用与旧 family-based 表名限定工具位于 `driver-api::sql_identifiers`。`packages/migration-common` 只提供记录集边界值规范化/比较和端点 category/family 配对规则,不依赖领域引擎、runtime 或平台传输。Data Sync 与 Data Transfer 直接依赖它,Data Transfer 不再依赖 Data Sync。Data Sync 的 `recordset_bounds`、`sync_pairing` 和 `sql` 工具旧路径保留为转导出,只有一份实现。

领域语义仍分别拥有:Schema Diff 的结构差异与操作依赖图、Data Sync 的同族门闸与 ChangeSet、Data Transfer 的异构 IR 转换与续传策略。共享配对规则只分类端点;各领域继续执行自己的能力限制(例如 Transfer 拒绝缺少适配器的 Redis 配对)。`qualify_relation_sql` 保留既有 family-based 兼容语义;驱动自有 SQL 重写继续使用 `DatabaseDriver::qualified_sql`,本次不改变驱动 trait、协议版本或执行行为。

Runtime 承担接受/认领、预算、资源申请、子 execution、事件、取消和 cleanup;领域 handler 承担计划校验、分阶段算法及恢复核验。Application 服务校验授权、输入与计划消费;前端只提交稳定目标、审阅选择及幂等令牌,不提交执行 SQL、原始检查点或 live handle 作为恢复资格。

## 2. 公共执行模型
Expand Down
5 changes: 5 additions & 0 deletions docs/architecture/platform/shared-boundaries-and-ports.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,7 @@ flowchart TD
PA[packages/platform-api]
DAPI[packages/driver-api]
DOM[schema-diff / data-sync / data-transfer]
MC[packages/migration-common]
DRV[packages/drivers/*]

TA --> APP
Expand All @@ -124,6 +125,9 @@ flowchart TD
RT --> DAPI
RT -.注册 JobHandler.-> DOM
DOM --> PA
DOM --> DAPI
DOM -.Sync / Transfer.-> MC
MC --> DAPI
DRV --> DAPI
FE -.只依赖契约.-> PA
```
Expand All @@ -139,6 +143,7 @@ flowchart TD
| F-05 | `packages/application`、`packages/runtime`、`packages/platform-api` 不出现 `react`、`@tauri-apps/api` 等前端标识(Rust 侧通过 crate 名与 `build.rs` 依赖检查,前端侧由 §7 的字符串扫描补齐) | UI 运行时进入后端依赖图 |
| F-06 | `server` 的 normal + build 依赖闭包不含 `tauri*` 与 `datazen` 宿主 crate | server 无法独立构建 |
| F-07 | `packages/backend-client` 不含 `@tauri-apps/` 前缀、`fetch(`、`XMLHttpRequest` 字面量 | 传输无关契约被具体传输污染 |
| F-08 | `packages/migration-common` 的 normal + build 工作区依赖只允许自身与 driver-api;不含 `tauri`、`axum`、`actix-web`、`warp`、`tonic`、`react` | 公共算法反向依赖领域引擎、runtime 或平台传输,重新形成耦合 |

**F-01 的作用域同时覆盖两侧:manifest 声明边与解析闭包。** 上面"crate 依赖闭包"若按字面只理解为 `cargo metadata` 的 `resolve` 图,是**不够的**:那个图是 **feature-resolved 的**,只含当前 feature 解析下真正被链接的边。无人启用的 feature 背后的 optional 依赖根本不在 `resolve.nodes[].deps[]` 里,任何闭包遍历都看不见它。因此 F-01 也约束 `Cargo.toml` 的**声明边**——声明了、但当前 feature 解析不链接的边,同样算违反。两侧任一出现即违规。

Expand Down
1 change: 1 addition & 0 deletions packages/data-sync/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ license = "GPL-3.0-or-later"

[dependencies]
datazen-driver-api = { workspace = true }
datazen-migration-common = { workspace = true }
datazen-platform-api = { workspace = true }
datazen-runtime = { path = "../runtime" }
async-trait = "0.1"
Expand Down
2 changes: 1 addition & 1 deletion packages/data-sync/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ pub mod model;
pub mod pairing;
pub mod profile;
mod recordset;
pub mod recordset_bounds;
pub use datazen_migration_common::recordset_bounds;
pub mod session;
pub mod sql;
pub mod state;
Expand Down
62 changes: 4 additions & 58 deletions packages/data-sync/src/sql.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,64 +30,10 @@ pub struct SqlStatement {
pub identity_insert: Option<IdentityInsertTarget>,
}

pub fn quote_ident_sql(name: &str, quote: char) -> String {
if quote == '[' {
return format!("[{}]", name.replace(']', "]]"));
}
let doubled = name.replace(quote, &format!("{quote}{quote}"));
format!("{quote}{doubled}{quote}")
}

/// Qualify `table` as `schema.table` when `schema` is non-empty (PostgreSQL etc.).
pub fn qualify_table_sql(schema: Option<&str>, table: &str, quote: char) -> String {
match schema.map(str::trim).filter(|s| !s.is_empty()) {
Some(schema) => format!(
"{}.{}",
quote_ident_sql(schema, quote),
quote_ident_sql(table, quote)
),
None => quote_ident_sql(table, quote),
}
}

/// Qualify a table reference for DML/SELECT without switching the session catalog.
///
/// - MySQL/MariaDB/ClickHouse: `` `database`.`table` `` when `database` is set.
/// - SQL Server: `[database].[schema].[table]` with either qualifier set.
/// - PostgreSQL and similar: `"schema"."table"` when `schema` is set.
/// - Otherwise: bare `table`.
pub fn qualify_relation_sql(
family: &str,
database: Option<&str>,
schema: Option<&str>,
table: &str,
quote: char,
) -> String {
let family = family.to_ascii_lowercase();
if matches!(family.as_str(), "mysql" | "mariadb" | "clickhouse") {
return match database.map(str::trim).filter(|s| !s.is_empty()) {
Some(db) => format!(
"{}.{}",
quote_ident_sql(db, quote),
quote_ident_sql(table, quote)
),
None => quote_ident_sql(table, quote),
};
}
if family == "sqlserver" {
return [
database.map(str::trim).filter(|s| !s.is_empty()),
schema.map(str::trim).filter(|s| !s.is_empty()),
Some(table),
]
.into_iter()
.flatten()
.map(|part| quote_ident_sql(part, quote))
.collect::<Vec<_>>()
.join(".");
}
qualify_table_sql(schema, table, quote)
}
// Compatibility exports; implementation belongs to the driver SQL contract.
pub use datazen_driver_api::sql_identifiers::{
qualify_relation_sql, qualify_table_sql, quote_ident_sql,
};

pub fn qualify_table_ident<Q>(schema: Option<&str>, table: &str, quote_ident: Q) -> String
where
Expand Down
77 changes: 2 additions & 75 deletions packages/data-sync/src/sync_pairing.rs
Original file line number Diff line number Diff line change
@@ -1,79 +1,6 @@
//! Sync pairing policy: same dialect family → Direct; cross SQL → IR; cross category → forbidden.
//!
//! Category/family rules are driven by driver-api taxonomy (`sync_category_of` /
//! `sync_family_of`); frontend Transfer UI reads `DatabaseTypeMeta` instead.
//! Compatibility path for the shared migration endpoint pairing policy.

use datazen_driver_api::{sync_category_of, sync_family_of, SyncCategory};

/// Resolved sync path for a source/target database type pair.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SyncPairing {
Direct { family: String },
Ir,
Unsupported { reason: String },
}

impl SyncPairing {
pub fn path_label(&self) -> &'static str {
match self {
Self::Direct { .. } => "direct",
Self::Ir => "ir",
Self::Unsupported { .. } => "unsupported",
}
}
}

/// Normalize a database type id to its sync dialect family.
pub fn normalize_sync_family(raw: &str) -> String {
sync_family_of(raw)
}

/// Map a database type id to a sync category.
pub fn sync_category(raw: &str) -> SyncCategory {
sync_category_of(raw)
}

/// Classify how a source/target pair should sync.
pub fn resolve_sync_pairing(source: &str, target: &str) -> SyncPairing {
let src_cat = sync_category(source);
let tgt_cat = sync_category(target);

if src_cat != tgt_cat {
return SyncPairing::Unsupported {
reason: format!(
"Sync between {} ({src_cat}) and {} ({tgt_cat}) is not supported",
source, target
),
};
}

match src_cat {
SyncCategory::Other => SyncPairing::Unsupported {
reason: format!("Sync is not supported for database type '{source}'"),
},
SyncCategory::Sql | SyncCategory::Document | SyncCategory::Kv => {
let src_family = normalize_sync_family(source);
let tgt_family = normalize_sync_family(target);
if src_family == tgt_family {
SyncPairing::Direct { family: src_family }
} else if src_cat == SyncCategory::Sql {
SyncPairing::Ir
} else {
SyncPairing::Unsupported {
reason: format!("Sync between {source} and {target} is not supported"),
}
}
}
}
}

/// Fail fast when a pair is forbidden; returns the resolved pairing otherwise.
pub fn enforce_sync_pairing(source: &str, target: &str) -> Result<SyncPairing, String> {
match resolve_sync_pairing(source, target) {
SyncPairing::Unsupported { reason } => Err(reason),
ok @ (SyncPairing::Direct { .. } | SyncPairing::Ir) => Ok(ok),
}
}
pub use datazen_migration_common::pairing::*;

#[cfg(test)]
mod tests {
Expand Down
2 changes: 1 addition & 1 deletion packages/data-transfer/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ webdriver = []

[dependencies]
datazen-driver-api = { workspace = true }
datazen-data-sync = { path = "../data-sync" }
datazen-migration-common = { workspace = true }
datazen-runtime = { path = "../runtime" }
datazen-platform-api = { workspace = true }
async-trait = "0.1"
Expand Down
2 changes: 1 addition & 1 deletion packages/data-transfer/src/execute.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use datazen_driver_api::TableSchema;

use crate::transfer::adapter::SyncTargetAdapter;
use crate::transfer::ir::IRType;
use datazen_data_sync::sql::{qualify_relation_sql, quote_ident_sql};
use datazen_driver_api::sql_identifiers::{qualify_relation_sql, quote_ident_sql};
use datazen_driver_api::{ConnectionHandle, DatabaseDriver, Value};

use super::error::TransferError;
Expand Down
2 changes: 1 addition & 1 deletion packages/data-transfer/src/filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,8 @@

use crate::error::TransferError;
use base64::{engine::general_purpose::STANDARD as BASE64, Engine};
use datazen_data_sync::sql::quote_ident_sql;
use datazen_driver_api::filters::{FilterCondition, FilterOperator};
use datazen_driver_api::sql_identifiers::quote_ident_sql;
use datazen_driver_api::{TableSchema, Value};
use serde::{Deserialize, Serialize};

Expand Down
4 changes: 3 additions & 1 deletion packages/data-transfer/src/job/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -163,7 +163,9 @@ pub async fn execute_bounded_table(
"SELECT {} FROM {}",
projection
.iter()
.map(|column| { datazen_data_sync::sql::quote_ident_sql(column, context.source_quote) })
.map(|column| {
datazen_driver_api::sql_identifiers::quote_ident_sql(column, context.source_quote)
})
.collect::<Vec<_>>()
.join(", "),
context.source_table_ref
Expand Down
12 changes: 8 additions & 4 deletions packages/data-transfer/src/recordset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,15 @@

use datazen_driver_api::TableSchema;

use datazen_data_sync::sql::quote_ident_sql;
use datazen_driver_api::sql_identifiers::quote_ident_sql;
use datazen_driver_api::Value;

use super::error::TransferError;
use super::filter::SourceFilter;
use super::model::{TransferRecordset, TransferRecordsetBound, TransferRecordsetTupleBound};
use datazen_data_sync::recordset_bounds::{canonical_bound_value, compare_bound_keys, BoundKey};
use datazen_migration_common::recordset_bounds::{
canonical_bound_value, compare_bound_keys, BoundKey,
};

#[derive(Debug, Clone)]
pub struct ResolvedRecordset {
Expand Down Expand Up @@ -90,7 +92,7 @@ pub fn resolve_recordset(
"source primary-key column '{key}' is nullable; tuple ranges require non-null keys"
)));
}
datazen_data_sync::recordset_bounds::ensure_supported_bound_type(
datazen_migration_common::recordset_bounds::ensure_supported_bound_type(
&column.data_type,
key,
)
Expand Down Expand Up @@ -409,7 +411,9 @@ pub fn preview_summary(
.iter()
.find(|column| column.name == *key)
.is_some_and(|column| {
datazen_data_sync::recordset_bounds::is_text_bound_type(&column.data_type)
datazen_migration_common::recordset_bounds::is_text_bound_type(
&column.data_type,
)
})
});
if has_text_key {
Expand Down
2 changes: 1 addition & 1 deletion packages/data-transfer/src/resume.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ use std::sync::Arc;

use datazen_driver_api::TableSchema;

use datazen_data_sync::sql::quote_ident_sql;
use datazen_driver_api::sql_identifiers::quote_ident_sql;
use datazen_driver_api::{ConnectionHandle, DatabaseDriver, TransactionHandle, Value};

use super::error::TransferError;
Expand Down
4 changes: 2 additions & 2 deletions packages/data-transfer/src/resume/fingerprint.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,13 +69,13 @@ fn primary_key_for_paging(
)));
}
if resumable {
datazen_data_sync::recordset_bounds::ensure_supported_bound_type(
datazen_migration_common::recordset_bounds::ensure_supported_bound_type(
&column.data_type,
key,
)
.map_err(TransferError::validation)
.map_err(|error| TransferError::unsupported(error.to_string()))?;
if datazen_data_sync::recordset_bounds::is_text_bound_type(&column.data_type) {
if datazen_migration_common::recordset_bounds::is_text_bound_type(&column.data_type) {
return Err(TransferError::unsupported(format!(
"primary-key column '{key}' uses text ordering whose collation cannot be verified for resume"
)));
Expand Down
2 changes: 1 addition & 1 deletion packages/data-transfer/src/sql_file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ use super::model::{
};
use crate::transfer::adapter::{SyncSourceAdapter, SyncTargetAdapter};
use crate::transfer::ir::{IRDefault, IRType};
use datazen_data_sync::sql::{qualify_relation_sql, quote_ident_sql};
use datazen_driver_api::sql_identifiers::{qualify_relation_sql, quote_ident_sql};
use datazen_driver_api::{ConnectionHandle, Value};

pub use super::sql_structure::build_structure_plan;
Expand Down
2 changes: 1 addition & 1 deletion packages/data-transfer/src/transfer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ pub mod adapter_registry;
pub mod adapters;
pub mod ddl;
pub mod full_types;
pub use datazen_data_sync::sync_pairing as pairing;
pub use datazen_migration_common::pairing;

pub use datazen_driver_api::sync::{
BoxedSyncAdapter, IRColumn, IRDefault, IRTable, IRType, SyncAdapterFactory, SyncSourceAdapter,
Expand Down
Loading
Loading