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
2 changes: 1 addition & 1 deletion docs/architecture/backend/data-transfer.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ peek 不消耗计划;claim 原子将 available 转为 executing 并建立活

## 6. SQL 文件与取消

`data_transfer/sql_file.rs` 负责 SQL 文件输出;SQL 文件目的地不建立目标数据库写会话,但仍需源读取与渲染能力。文件输出失败或取消不能标成完整产物;结果必须表达已输出部分和错误。
`data_transfer/sql_file.rs` 负责 SQL 文件输出;SQL 文件目的地不建立目标数据库写会话,但仍需源读取与渲染能力。INSERT 字面量必须由目标 SQL 驱动声明的 `SqlLiteralDialect` 生成;无方言或值无法安全表示时停止生成,不能回退到默认转义。文件输出失败或取消不能标成完整产物;结果必须表达已输出部分和错误。

旧执行与兼容路径仍由 `commands/data_transfer/exec.rs` 等入口管理;P5 新路径在 `commands/data_transfer/job_api/`。新路径的 cancel intent 持久写入本机 Job repository,运行时通过 `CancelWatch` 将它传给当前 stage。该 SQLite repository 是单机桌面 host,不是团队服务跨实例协调器。Job 只存 Artifact ID 引用,不存文件字节;当前引用 TTL 为 30 天,且不保证仍可下载对应内容。

Expand Down
6 changes: 5 additions & 1 deletion docs/architecture/backend/drivers.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ Driver 基础能力包括:
- connect / test_connection / disconnect
- get_databases / get_tables / get_table_schema
- query / query_multi / query_stream
- query_with_params / execute
- query_with_params / execute / execute_with_params
- transaction
- EXPLAIN
- Driver Commands
Expand All @@ -51,6 +51,10 @@ cleanup_query_execution

只有实际声明精确取消能力的 Driver 才会被 Host 当作 cancellable;兼容默认实现不会自动获得取消能力。

Host 生成的写语句必须把 SQL 模板与 `Vec<Value>` 分开,通过 `parameter_placeholder()` 和 `execute_with_params()` 传值;没有绑定写能力的驱动返回 `Unsupported`,Host 在开启事务前拒绝该批写入。`build_update_statement()` / `build_delete_statement()` 生成的预览只含占位符,不能把值插入 SQL 文本。

SQL 筛选与 SQL 文件产物使用 `try_format_sql_literal()`。每个 SQL 驱动显式声明 `SqlLiteralDialect`,共享格式器按方言使用对会话转义模式稳定的表示(例如 MySQL 的 UTF-8 十六进制转换、PostgreSQL / DuckDB 的 dollar-quoted 文本);未声明方言时返回 `Unsupported`,不会猜测反斜杠规则。无方言的键值 / 文档驱动不应调用 SQL 字面量接口。旧的无结果 `format_sql_literal()` 和内插式 `build_*_sql()` 仅保留兼容,产品写路径不使用它们。

### 2.1 Schema 元数据目标

`get_tables`、`get_table_schema`、`get_columns` 和 `get_all_columns` 都接收显式 `database` 与可选 `schema`,驱动按传入目标读取,不依赖会话当前选中的数据库。Host 的 `list_catalog`、`read_relation_columns`、`read_relation_schema` 和 `refresh_schema_metadata` Driver Commands 以 `RelationRef { database, schema, name }` 表达目标;前端与扩展共用 Driver Command 网关,不再需要按数据库类型选择旧 Host schema IPC。
Expand Down
2 changes: 2 additions & 0 deletions docs/architecture/backend/services.md
Original file line number Diff line number Diff line change
Expand Up @@ -129,13 +129,15 @@ DataTable 的行编辑与数据导出是两条不同的写路径,不要按同

- **创建 / 复用**:都只在调用方已建的会话上执行。会话找不到时走同一条透明重建路径。
- **事务**:预览待提交改动**不**开事务,只读;真正执行改动计划时,若该会话已有显式事务就复用,否则自己开一个。语义不明确的写入(例如无法判定影响行数)一律拒绝,不做「尽力而为」。
- **写入参数**:UPDATE / DELETE 的 SQL 只包含驱动占位符;列值与主键值通过驱动绑定 API 单独传递。驱动先构造整批语句,确认占位符与绑定执行均受支持后才开启事务。预览中的 SQL 模板不包含数据值。
- **关闭**:不涉及。
- **取消**:无独立的取消句柄——粒度是「提交前 / 提交后」,不是执行中途。
- **配置目标**:提交前复核四件事:改动计划里的表必须属于这条会话的**所属连接**、驱动类型与 database/schema 必须与当前配置一致、连接不得只读、以及结构指纹是否仍然匹配。任一不符即中止,避免把过期界面上的改动写到已经变化的表上。表所属连接是从会话反查出来的(见 [2.1](#21-连接-id-约定)),该反查依赖 owner 映射在空闲回收后仍然保留。

**数据导出**

- **创建 / 复用**:只消费调用方的 `dbSessionId`,不创建;导出前先过 SQL 安全门闸,再以流式回调逐批输出。
- **SQL 格式**:查询筛选和 SQL INSERT 导出使用驱动显式声明的字面量方言;不支持该方言或遇到无法表示的值时返回错误,不能退回共享反斜杠替换。CSV / JSON 导出不经过 SQL 字面量格式器。
- **关闭 / 取消**:**导出没有中途取消机制**——没有取消令牌也没有取消标志位,与 Data Sync / Data Transfer 的 job 标志是两种模型。导出结果里的「已取消」只表示用户在保存对话框里放弃了选文件,也就是导出根本没开始;一旦开始流式写盘,就只能由查询本身失败而结束。
- **配置目标**:database / schema 随导出请求传递,并从会话配置补齐缺省项,不读会话当前状态。

Expand Down
27 changes: 27 additions & 0 deletions packages/data-transfer/src/job/sqlfile.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,33 @@ impl DataTransferHandler {
)
.await;
match result {
Ok(res) if res.cancelled || res.partial => {
let cancelled = res.cancelled;
tracing::warn!(
destination = destination.display().to_string(),
cancelled,
table_errors = ?res.tables.iter().filter_map(|table| table.error.as_deref()).collect::<Vec<_>>(),
"SQL file was not published because generation was incomplete"
);
Ok(StageOutcome {
stage_id: spec.stage_id.clone(),
terminal: if cancelled {
StageTerminal::Cancelled
} else {
StageTerminal::Failed
},
progress: JobProgress::default(),
commit_boundaries: Vec::new(),
execution_ids: Vec::new(),
artifact_ids: Vec::new(),
effect_outcome: EffectOutcome::RolledBack,
error_code: Some(if cancelled {
ExecutionErrorCode::Cancelled
} else {
ExecutionErrorCode::SqlError
}),
})
}
Ok(res) => {
let digest = artifact_digest(destination)?;
let boundary = commit_boundary(
Expand Down
4 changes: 4 additions & 0 deletions packages/data-transfer/src/job/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,10 @@ fn unsupported<T>() -> Result<T, DriverError> {

#[async_trait]
impl DatabaseDriver for FakeDb {
fn sql_literal_dialect(&self) -> Option<SqlLiteralDialect> {
Some(SqlLiteralDialect::Postgres)
}

async fn cancel_query(&self, _: &ConnectionHandle) -> Result<(), DriverError> {
Ok(())
}
Expand Down
18 changes: 10 additions & 8 deletions packages/data-transfer/src/sql_file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -724,13 +724,15 @@ fn insert_sql_batch(
mappings.len()
)));
}
Ok(format!(
"({})",
row.iter()
.map(|value| driver.format_sql_literal(value))
.collect::<Vec<_>>()
.join(", ")
))
let values = row
.iter()
.map(|value| {
driver
.try_format_sql_literal(value)
.map_err(|error| TransferError::unsupported(error.to_string()))
})
.collect::<Result<Vec<_>, _>>()?;
Ok(format!("({})", values.join(", ")))
})
.collect::<Result<Vec<_>, TransferError>>()?;
render_sql_file_insert(
Expand Down Expand Up @@ -1513,7 +1515,7 @@ mod tests {
)
.unwrap();
assert!(sql.contains("\"display_name\""));
assert!(sql.contains("'O''Reilly'"));
assert!(sql.contains("$datazen$O'Reilly$datazen$"));
}

#[test]
Expand Down
18 changes: 17 additions & 1 deletion packages/driver-api/src/mock_driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ use crate::{
is_schema_object_command, query_command_definition, query_stream_command_definition,
schema_catalog_command_definitions, schema_object_command_definitions,
try_execute_schema_catalog_command, validate_schema_target, CommandResult, DdlAtomicity,
DriverCommandDefinition, SchemaScope,
DriverCommandDefinition, SchemaScope, SqlLiteralDialect,
};
use crate::{
ColumnInfo, ColumnSchema, ConnectionConfig, ConnectionHandle, DatabaseDriver, DatabaseType,
Expand Down Expand Up @@ -386,6 +386,18 @@ impl MockDriver {

#[async_trait]
impl DatabaseDriver for MockDriver {
fn sql_literal_dialect(&self) -> Option<SqlLiteralDialect> {
match self.db_type.to_ascii_lowercase().as_str() {
"clickhouse" => Some(SqlLiteralDialect::ClickHouse),
"duckdb" => Some(SqlLiteralDialect::DuckDb),
"mysql" | "mariadb" | "doris" | "starrocks" => Some(SqlLiteralDialect::MySql),
"postgres" | "postgresql" => Some(SqlLiteralDialect::Postgres),
"sqlserver" | "mssql" => Some(SqlLiteralDialect::SqlServer),
"sqlite" | "turso" | "rqlite" => Some(SqlLiteralDialect::Sqlite),
_ => None,
}
}

fn driver_type(&self) -> DatabaseType {
self.db_type.clone()
}
Expand All @@ -394,6 +406,10 @@ impl DatabaseDriver for MockDriver {
self.opts.category.clone()
}

fn supports_bound_writes(&self) -> bool {
self.opts.parameterized_writes
}

fn default_host(&self) -> Option<&'static str> {
self.opts.default_host
}
Expand Down
12 changes: 12 additions & 0 deletions packages/driver-api/src/reuse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,18 @@ impl DatabaseDriver for ReuseDriver {
self.inner.format_sql_literal(value)
}

fn sql_literal_dialect(&self) -> Option<SqlLiteralDialect> {
self.inner.sql_literal_dialect()
}

fn supports_bound_writes(&self) -> bool {
self.inner.supports_bound_writes()
}

fn try_format_sql_literal(&self, value: &Option<Value>) -> Result<String, DriverError> {
self.inner.try_format_sql_literal(value)
}

fn build_update_sql(
&self,
table: &str,
Expand Down
12 changes: 7 additions & 5 deletions packages/driver-api/src/sql_dump/dump.rs
Original file line number Diff line number Diff line change
Expand Up @@ -251,12 +251,14 @@ where
let tuples: Vec<String> = result
.rows
.iter()
.map(|row| {
let vals: Vec<String> =
row.iter().map(|v| driver.format_sql_literal(v)).collect();
format!("({})", vals.join(", "))
.map(|row| -> Result<String, DriverError> {
let vals: Vec<String> = row
.iter()
.map(|value| driver.try_format_sql_literal(value))
.collect::<Result<_, _>>()?;
Ok(format!("({})", vals.join(", ")))
})
.collect();
.collect::<Result<_, _>>()?;
append_batched_inserts(
out,
&rel,
Expand Down
42 changes: 42 additions & 0 deletions packages/driver-api/src/traits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,10 +161,26 @@ pub trait DatabaseDriver: Send + Sync {
DdlAtomicity::Unknown
}

/// Legacy best-effort SQL literal formatter. New code should use
/// [`Self::try_format_sql_literal`], which requires an explicit dialect.
fn format_sql_literal(&self, value: &Option<Value>) -> String {
sql_text::format_sql_literal(value)
}

/// The literal grammar this driver can safely render for SQL artifacts.
/// Drivers with session-dependent or unsupported syntax must leave this
/// unset; their export/filter callers then fail closed.
fn sql_literal_dialect(&self) -> Option<SqlLiteralDialect> {
None
}

fn try_format_sql_literal(&self, value: &Option<Value>) -> Result<String, DriverError> {
let dialect = self.sql_literal_dialect().ok_or_else(|| {
DriverError::Unsupported("this driver has no declared SQL literal formatter".into())
})?;
sql_text::format_sql_literal_for_dialect(value, dialect)
}

fn build_update_sql(
&self,
table: &str,
Expand All @@ -179,6 +195,25 @@ pub trait DatabaseDriver: Send + Sync {
sql_text::build_delete_sql(self, table, pk_columns)
}

/// Build an UPDATE with driver placeholders and values kept out of SQL.
fn build_update_statement(
&self,
table: &str,
set_columns: &[(&str, Option<Value>)],
pk_columns: &[(&str, Option<Value>)],
) -> Result<BoundSqlStatement, DriverError> {
sql_text::build_update_statement(self, table, set_columns, pk_columns)
}

/// Build a DELETE with driver placeholders and values kept out of SQL.
fn build_delete_statement(
&self,
table: &str,
pk_columns: &[(&str, Option<Value>)],
) -> Result<BoundSqlStatement, DriverError> {
sql_text::build_delete_statement(self, table, pk_columns)
}

/// The host this driver dials when the connection config leaves `host` unset.
///
/// `connect` resolves that default internally, so a config that omits
Expand Down Expand Up @@ -367,6 +402,13 @@ pub trait DatabaseDriver: Send + Sync {
params: &[Value],
) -> Result<QueryResult, DriverError>;

/// Whether this driver implements both placeholder generation and bound
/// DML execution. Callers use this to reject writes before opening a
/// transaction.
fn supports_bound_writes(&self) -> bool {
false
}

/// Render a parameter for this dialect. Unsupported drivers must fail before writes.
/// `data_type` comes from the inspected target column metadata.
fn parameter_placeholder(
Expand Down
Loading
Loading