From 7f7d757f154d1e83acbe204e459b53403375538e Mon Sep 17 00:00:00 2001 From: Yuyi-Oak <1722157266@qq.com> Date: Wed, 16 Sep 2026 21:02:22 +0800 Subject: [PATCH] =?UTF-8?q?fix(sqlite):=20=E6=94=B6=E5=8F=A3=E9=95=BF?= =?UTF-8?q?=E6=9C=9F=E7=8A=B6=E6=80=81=E6=95=B0=E6=8D=AE=E5=BA=93=E7=89=88?= =?UTF-8?q?=E6=9C=AC=E5=8C=96=E8=BF=81=E7=A7=BB=E5=A5=91=E7=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CURRENT_STATUS.md | 8 + TODO.md | 30 +- .../architecture/official-resource-backend.md | 2 +- docs/reference/rpc-backend-api.md | 12 +- docs/reports/CURRENT_GAPS.md | 19 +- infrastructure/src/glossary.rs | 826 ++++++++++++--- infrastructure/src/lib.rs | 1 + infrastructure/src/sqlite_migration.rs | 377 +++++++ infrastructure/src/translation_memory.rs | 635 +++++++++-- infrastructure/src/translation_tasks.rs | 998 +++++++++++++++--- 10 files changed, 2487 insertions(+), 421 deletions(-) create mode 100644 infrastructure/src/sqlite_migration.rs diff --git a/CURRENT_STATUS.md b/CURRENT_STATUS.md index 380be6d..d77e879 100644 --- a/CURRENT_STATUS.md +++ b/CURRENT_STATUS.md @@ -46,6 +46,14 @@ TextUnit scope、source history 和 approved review。worker、TM 复用、人 按当前 QA 精确校验 identity。`translation.glossary.*` 已通过 `bat.sock` 暴露,Go `bat-api` 仅做鉴权 typed forwarding。 +三个长期 SQLite owner 现在统一使用只读 schema preflight、精确 component +fingerprint 和 `BEGIN IMMEDIATE` writer transaction:Translation Tasks 从 V1 +按显式 `v1 -> v2` step 迁移,当前版本为 V2;Translation Memory 和 Glossary +当前版本均为 V1。future、未知或版本与结构不一致的数据库在任何 schema/data +mutation 前 fail closed;migration 失败会 rollback,已知无 +`schema_migrations` 表的历史 fingerprint 可安全补建版本表后重试。正式 schema +路径不再使用 `ensure_column` 隐式补列。 + `localized.status` 现将 generic manifest schema/contract 与已发布 artifact integrity 分开报告;current、state 和 identity 存在但文件被截断或手工修改时返回 `localized.degraded`,只读检查不会自动回滚、删除或修复。双 release 的 diff --git a/TODO.md b/TODO.md index 620838f..d6d81b1 100644 --- a/TODO.md +++ b/TODO.md @@ -236,7 +236,7 @@ T02 → T13 # T03 — SQLite Schema Migration 体系化 **类型:** P2 -**状态:** Ready +**状态:** Done ## 问题 @@ -311,12 +311,23 @@ T03 → T15 T03 → T16 ``` +## 完成记录 + +已为 Translation Tasks、Translation Memory 和 Glossary 收口统一的 SQLite +migration contract:现有数据库先以只读方式读取 schema、component version 和结构 +fingerprint,future/unknown/mismatch 在任何写入前 fail closed;确认可迁移后在 +`BEGIN IMMEDIATE` 事务内执行显式版本步骤,重读并校验结构后写入版本并提交。Translation +Tasks 按真实历史从 V1 迁移到 V2,Translation Memory 和 Glossary 保持真实的 V1 +contract,没有凭空引入版本;已覆盖缺失 schema bookkeeping、业务数据保留、失败回滚 +重试、并发打开和当前版本 reopen no-op。生产路径不再依赖 `ensure_column` 隐式升级, +相关状态、架构、RPC 参考和缺口文档已同步。 + --- # T04 — Translation Memory Trusted 唯一性与 Supersede 治理 **类型:** P2 -**状态:** Blocked by T03 +**状态:** Ready ## 问题 @@ -1047,16 +1058,7 @@ Stable Resource / Release / AssetBundle Platform ## In Progress -建议当前只放: - -```text -T01 — bat-api 发布完整性与分发健康契约 -T08 — CI Gate 与质量门禁整理 -``` - -T01 是当前唯一 P1。 - -T08 可以与其并行,不涉及核心业务状态机。 +当前无。 --- @@ -1064,7 +1066,7 @@ T08 可以与其并行,不涉及核心业务状态机。 ```text T02 — ResourceRepository 查询契约统一 -T03 — SQLite Schema Migration 体系化 +T04 — Translation Memory Trusted 唯一性与 Supersede 治理 T05 — bat.sock 本地 IPC 安全与资源边界 T06 — AssetBundle / ZIP Parser Resource Budget T07 — Durable Atomic State Write @@ -1078,8 +1080,6 @@ T11 可以现在就开始采集,但 T01 后的数据才作为新的正式 dist ## Blocked / Planned ```text -T04 ← T03 - T09 ← T06 T10 ← T09 diff --git a/docs/architecture/official-resource-backend.md b/docs/architecture/official-resource-backend.md index ba98297..d75562d 100644 --- a/docs/architecture/official-resource-backend.md +++ b/docs/architecture/official-resource-backend.md @@ -179,7 +179,7 @@ `resource.index` RPC / CLI 只读查询现有 SQLite 索引;索引不存在时返回 `available=false`,不会因为查询创建空库。发布后的 TextUnit 队列还会在当前 release 根目录写入 `translation-tasks.sqlite`,由版本化 `schema_migrations` -管理 queued/running/failed/completed/skipped、provider run、lease、失败分类、 +管理 V2 queued/running/failed/completed/skipped、provider run、lease、失败分类、 重试计划和 TextUnit 级译文结果。跨 release 的 Translation Memory V1 独立存储在 `/translation-memory.sqlite`,记录 raw source/hash、完整 context、candidate/ trusted 和 release/TextUnit/provider/run provenance;`translation.tasks` 优先查询这份状态库, diff --git a/docs/reference/rpc-backend-api.md b/docs/reference/rpc-backend-api.md index 7b7e96a..7783c92 100644 --- a/docs/reference/rpc-backend-api.md +++ b/docs/reference/rpc-backend-api.md @@ -165,18 +165,24 @@ count,任何一页不一致都会丢弃整个候选快照。 - `translation-tasks.sqlite`:当前 release 的可变 worker 状态库,记录 queued / running / failed / completed / skipped、attempt count、provider run ID、provider、TextUnit 级译文结果、lease、失败分类、可重试标记和 - next attempt;schema 由 `schema_migrations` 版本表管理。 + next attempt;当前 schema version 为 V2,schema 由 `schema_migrations` 版本表 + 管理,并通过只读 fingerprint、`BEGIN IMMEDIATE` 和显式 V1 → V2 migration + 保证 future/未知 schema fail closed。 - `translation-handoff.json`:当前 release 的版本化 job/unit/provider run 交接 快照;worker 更新后的实时状态仍以 `translation-tasks.sqlite` 为准。 - `translation-memory.sqlite`:跨 release 的项目级 Translation Memory,不位于 `versions/`,也不与 `translation-tasks.sqlite` 共用;记录 raw source/hash、完整 TextUnit context、candidate/trusted、translation 和 release/TextUnit/provider/run - provenance。默认路径为 `/translation-memory.sqlite`,可由 + provenance;当前 schema version 为 V1,打开时先进行只读 fingerprint preflight, + 再在 writer transaction 内补齐 `schema_migrations` 版本记录。默认路径为 + `/translation-memory.sqlite`,可由 `BAT_TRANSLATION_MEMORY_PATH`、`[translation.worker].translation_memory_path` 或 CLI 覆盖。 - `glossary.sqlite`:跨 release 的项目级 Glossary,不位于 `versions/`,也不与 `translation-tasks.sqlite` 或 TM 共用;记录 term、alias、推荐/允许译法、scope、 - priority、review 状态、source provenance 和完整 source/review history。默认路径为 + priority、review 状态、source provenance 和完整 source/review history;当前 + schema version 为 V1,正式 schema evolution 不再通过 `ensure_column` 隐式修复。 + 默认路径为 `/glossary.sqlite`,可由 `BAT_GLOSSARY_PATH`、`[translation.worker].glossary_path` 或 CLI 覆盖。 diff --git a/docs/reports/CURRENT_GAPS.md b/docs/reports/CURRENT_GAPS.md index 1227144..9bb5314 100644 --- a/docs/reports/CURRENT_GAPS.md +++ b/docs/reports/CURRENT_GAPS.md @@ -154,11 +154,26 @@ official/localized distribution 均被阻断,普通查询不会自动重建。 ### G-012:Translation Memory V1 已实现,扩展能力仍缺失 -Rust `bat` 已提供独立项目级 SQLite TM,记录 raw source/hash、完整 context、release/TextUnit/provider/run provenance,区分 candidate/trusted,只有显式 confirm 才能建立 trusted 记录;worker 只自动复用 trusted 的 raw source + 完整 context exact match,并在复用前执行已批准 Glossary 的确定性 QA。Go `bat-api` 已提供鉴权的 summary/query 只读接口和 confirm 转发,但 Go 不持有 TM 状态。仍缺少模糊匹配和更丰富的导入导出历史能力。 +Rust `bat` 已提供独立项目级 SQLite TM,当前 schema version 为 V1;schema 打开遵守 +只读 preflight、fingerprint、transaction rollback 和 future/unknown fail-closed +契约。它记录 raw source/hash、完整 context、release/TextUnit/provider/run provenance, +区分 candidate/trusted,只有显式 confirm 才能建立 trusted 记录;worker 只自动复用 +trusted 的 raw source + 完整 context exact match,并在复用前执行已批准 Glossary 的 +确定性 QA。Go `bat-api` 已提供鉴权的 summary/query 只读接口和 confirm 转发,但 Go +不持有 TM 状态。仍缺少模糊匹配和更丰富的导入导出历史能力。 ### G-013:Glossary V1 已实现,协作视图仍缺失 -Rust `bat` 已提供独立项目级 `glossary.sqlite`:term/alias/recommended/allowed/category/priority、全局与 TextUnit scope、source history、approved review、冲突诊断、provider-neutral constraints 和确定性 QA 均由 Rust 持有。trusted TM 复用会先经过 Glossary QA;provider、TM、人工 task/workbench 结果都记录 QA,blocking deviation 必须显式提交与当前 QA 精确绑定的 `qa_identity` 及 reviewer/reason/provenance。localized publish 会把发布时重算的 QA 写入 manifest。`translation.glossary.*` 已通过 `bat.sock` 暴露,Go 仅提供鉴权后的 typed forwarding。剩余缺口是完整 Web 术语协作视图和更丰富的导入/搜索能力。 +Rust `bat` 已提供独立项目级 `glossary.sqlite`,当前 schema version 为 V1;打开 +遵守只读 fingerprint preflight、writer transaction、rollback 和 future/unknown +fail-closed 契约,不再使用隐式 `ensure_column` 修复结构。term/alias/recommended/ +allowed/category/priority、全局与 TextUnit scope、source history、approved review、 +冲突诊断、provider-neutral constraints 和确定性 QA 均由 Rust 持有。trusted TM 复用会 +先经过 Glossary QA;provider、TM、人工 task/workbench 结果都记录 QA,blocking +deviation 必须显式提交与当前 QA 精确绑定的 `qa_identity` 及 reviewer/reason/provenance。 +localized publish 会把发布时重算的 QA 写入 manifest。`translation.glossary.*` 已通过 +`bat.sock` 暴露,Go 仅提供鉴权后的 typed forwarding。剩余缺口是完整 Web 术语协作视图 +和更丰富的导入/搜索能力。 ### G-014:完整 Provider 扩展体系未实现 diff --git a/infrastructure/src/glossary.rs b/infrastructure/src/glossary.rs index 978dc50..d553e2a 100644 --- a/infrastructure/src/glossary.rs +++ b/infrastructure/src/glossary.rs @@ -3,6 +3,9 @@ use crate::path_security::{ ensure_safe_directory_path, lexical_absolute, set_file_mode, STATE_FILE_MODE, }; +use crate::sqlite_migration::{ + self, ExpectedColumn, ExpectedIndex, ExpectedTable, SqliteSchemaSnapshot, +}; use async_trait::async_trait; use bat_core::domain::{ evaluate_glossary, validate_glossary_draft, GlossaryEvaluation, GlossaryHistoryRecord, @@ -12,7 +15,7 @@ use bat_core::domain::{ use bat_core::repositories::GlossaryRepository; use bat_core::{Error, Result}; use serde::de::DeserializeOwned; -use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions}; +use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode}; use sqlx::{Row, SqlitePool}; use std::path::{Path, PathBuf}; use std::str::FromStr; @@ -76,15 +79,20 @@ impl SqliteGlossaryRepository { absolute.display() ))); } + if metadata.len() > 0 { + let snapshot = + sqlite_migration::read_only_preflight(&absolute, GLOSSARY_SCHEMA_COMPONENT) + .await + .map_err(db_error)?; + classify_glossary_schema(&snapshot)?; + } } let options = SqliteConnectOptions::from_str(&format!("sqlite://{}", absolute.display())) .map_err(|error| Error::Other(error.into()))? .create_if_missing(create_if_missing) .journal_mode(SqliteJournalMode::Wal) .busy_timeout(Duration::from_secs(30)); - let pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with(options) + let pool = sqlite_migration::connect_writable_pool(options) .await .map_err(db_error)?; if create_if_missing { @@ -97,119 +105,86 @@ impl SqliteGlossaryRepository { } async fn init_schema(&self) -> Result<()> { - sqlx::query( - "CREATE TABLE IF NOT EXISTS schema_migrations ( - component TEXT PRIMARY KEY NOT NULL, - version INTEGER NOT NULL CHECK(version >= 1) - )", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - "CREATE TABLE IF NOT EXISTS glossary_terms ( - term_id TEXT PRIMARY KEY NOT NULL, - source_term TEXT NOT NULL, - aliases_json TEXT NOT NULL, - recommended_translation TEXT NOT NULL, - allowed_translations_json TEXT NOT NULL, - source_language TEXT, - target_language TEXT, - category TEXT, - priority INTEGER NOT NULL, - scope_json TEXT NOT NULL, - review_status TEXT NOT NULL, - source_kind TEXT NOT NULL, - source_ref TEXT, - source_author TEXT, - source_note TEXT, - source_observed_unix_seconds INTEGER NOT NULL, - created_unix_seconds INTEGER NOT NULL, - updated_unix_seconds INTEGER NOT NULL, - CHECK(length(term_id) > 0), - CHECK(length(source_term) > 0), - CHECK(length(recommended_translation) > 0), - CHECK(review_status IN ('draft', 'approved', 'deprecated', 'rejected')), - CHECK(source_kind IN ('manual', 'imported')) - )", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - ensure_column( - &self.pool, - "glossary_terms", - "source_observed_unix_seconds", - "INTEGER NOT NULL DEFAULT 1", - ) - .await?; - sqlx::query( - "CREATE TABLE IF NOT EXISTS glossary_term_history ( - history_id TEXT PRIMARY KEY NOT NULL, - term_id TEXT NOT NULL, - action TEXT NOT NULL, - reviewer TEXT, - reason TEXT, - source_json TEXT NOT NULL, - review_status TEXT NOT NULL, - snapshot_json TEXT NOT NULL, - observed_unix_seconds INTEGER NOT NULL, - FOREIGN KEY(term_id) REFERENCES glossary_terms(term_id) - )", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - "CREATE TABLE IF NOT EXISTS glossary_term_deletions ( - deletion_id TEXT PRIMARY KEY NOT NULL, - term_id TEXT NOT NULL, - reviewer TEXT NOT NULL, - reason TEXT NOT NULL, - source_json TEXT NOT NULL, - snapshot_json TEXT NOT NULL, - history_json TEXT NOT NULL, - observed_unix_seconds INTEGER NOT NULL - )", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - "CREATE INDEX IF NOT EXISTS idx_glossary_status - ON glossary_terms(review_status, priority DESC, term_id)", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - "CREATE INDEX IF NOT EXISTS idx_glossary_source_term - ON glossary_terms(source_term)", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - let current: Option = - sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") - .bind(GLOSSARY_SCHEMA_COMPONENT) - .fetch_optional(&self.pool) + self.init_schema_with_failure(None).await + } + + #[cfg(test)] + async fn init_schema_with_test_failure(&self, fail_after_step: usize) -> Result<()> { + self.init_schema_with_failure(Some(fail_after_step)).await + } + + async fn init_schema_with_failure(&self, fail_after_step: Option) -> Result<()> { + let mut transaction = sqlite_migration::begin_immediate(&self.pool) + .await + .map_err(db_error)?; + let result = self + .migrate_in_transaction(&mut transaction, fail_after_step) + .await; + match result { + Ok(()) => transaction.commit().await.map_err(db_error), + Err(error) => { + let _ = transaction.rollback().await; + Err(error) + } + } + } + + async fn migrate_in_transaction( + &self, + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + fail_after_step: Option, + ) -> Result<()> { + let snapshot = + sqlite_migration::snapshot_connection(transaction.as_mut(), GLOSSARY_SCHEMA_COMPONENT) .await .map_err(db_error)?; - if current.is_some_and(|version| version > i64::from(GLOSSARY_SCHEMA_VERSION)) { - return Err(Error::InvalidArgument(format!( - "不支持的 Glossary schema 版本:{}", - current.unwrap_or_default() - ))); + match classify_glossary_schema(&snapshot)? { + GlossarySchemaState::Empty => { + create_glossary_schema(transaction, fail_after_step).await?; + sqlite_migration::write_component_version( + transaction, + GLOSSARY_SCHEMA_COMPONENT, + GLOSSARY_SCHEMA_VERSION, + ) + .await + .map_err(db_error)?; + } + GlossarySchemaState::Version(version) if version == GLOSSARY_SCHEMA_VERSION => { + if snapshot.component_version.is_none() { + if !snapshot + .tables + .contains_key(sqlite_migration::SCHEMA_MIGRATIONS_TABLE) + { + sqlite_migration::create_schema_migrations_table(transaction) + .await + .map_err(db_error)?; + } + sqlite_migration::write_component_version( + transaction, + GLOSSARY_SCHEMA_COMPONENT, + GLOSSARY_SCHEMA_VERSION, + ) + .await + .map_err(db_error)?; + } + } + GlossarySchemaState::Version(version) => { + return Err(invalid_glossary_schema(format!( + "内部不支持的迁移起点 {version}" + ))); + } + } + let final_snapshot = + sqlite_migration::snapshot_connection(transaction.as_mut(), GLOSSARY_SCHEMA_COMPONENT) + .await + .map_err(db_error)?; + if final_snapshot.component_version != Some(i64::from(GLOSSARY_SCHEMA_VERSION)) + || !matches_glossary_fingerprint(&final_snapshot) + { + return Err(invalid_glossary_schema( + "migration 结果与当前 schema fingerprint 不一致".to_string(), + )); } - sqlx::query( - "INSERT INTO schema_migrations(component, version) VALUES (?1, ?2) - ON CONFLICT(component) DO UPDATE SET version = excluded.version", - ) - .bind(GLOSSARY_SCHEMA_COMPONENT) - .bind(i64::from(GLOSSARY_SCHEMA_VERSION)) - .execute(&self.pool) - .await - .map_err(db_error)?; Ok(()) } @@ -743,6 +718,450 @@ fn row_to_history(row: sqlx::sqlite::SqliteRow) -> Result }) } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum GlossarySchemaState { + Empty, + Version(u32), +} + +fn classify_glossary_schema(snapshot: &SqliteSchemaSnapshot) -> Result { + if snapshot.is_empty() { + return Ok(GlossarySchemaState::Empty); + } + if let Some(observed) = snapshot.component_version { + if observed > i64::from(GLOSSARY_SCHEMA_VERSION) { + return Err(invalid_glossary_schema(format!( + "不支持的 Glossary schema 版本:observed={observed}, supported={GLOSSARY_SCHEMA_VERSION}" + ))); + } + if observed < 1 { + return Err(invalid_glossary_schema(format!( + "schema version 无效:observed={observed}, supported=1..={GLOSSARY_SCHEMA_VERSION}" + ))); + } + } + if (matches_glossary_fingerprint(snapshot) + || (snapshot.component_version.is_none() + && matches_glossary_component_fingerprint(snapshot))) + && snapshot + .component_version + .is_none_or(|observed| observed == i64::from(GLOSSARY_SCHEMA_VERSION)) + { + return Ok(GlossarySchemaState::Version(GLOSSARY_SCHEMA_VERSION)); + } + match snapshot.component_version { + Some(observed) => Err(invalid_glossary_schema(format!( + "schema version 与实际结构不一致:observed={observed}, supported={GLOSSARY_SCHEMA_VERSION}" + ))), + None => Err(invalid_glossary_schema( + "未识别的 legacy schema,拒绝静默修复".to_string(), + )), + } +} + +fn invalid_glossary_schema(detail: String) -> Error { + Error::InvalidArgument(format!("Glossary schema 无效:{detail}")) +} + +fn matches_glossary_fingerprint(snapshot: &SqliteSchemaSnapshot) -> bool { + let tables = [ + sqlite_migration::schema_migrations_table(), + ExpectedTable { + name: "glossary_terms", + columns: &GLOSSARY_TERM_COLUMNS, + }, + ExpectedTable { + name: "glossary_term_history", + columns: &GLOSSARY_HISTORY_COLUMNS, + }, + ExpectedTable { + name: "glossary_term_deletions", + columns: &GLOSSARY_DELETION_COLUMNS, + }, + ]; + let indexes = [ + ExpectedIndex { + table: "glossary_terms", + name: "idx_glossary_status", + columns: &["review_status", "priority", "term_id"], + }, + ExpectedIndex { + table: "glossary_terms", + name: "idx_glossary_source_term", + columns: &["source_term"], + }, + ]; + matches_glossary_fingerprint_with_tables(snapshot, &tables, &indexes) +} + +fn matches_glossary_component_fingerprint(snapshot: &SqliteSchemaSnapshot) -> bool { + let tables = [ + ExpectedTable { + name: "glossary_terms", + columns: &GLOSSARY_TERM_COLUMNS, + }, + ExpectedTable { + name: "glossary_term_history", + columns: &GLOSSARY_HISTORY_COLUMNS, + }, + ExpectedTable { + name: "glossary_term_deletions", + columns: &GLOSSARY_DELETION_COLUMNS, + }, + ]; + let indexes = [ + ExpectedIndex { + table: "glossary_terms", + name: "idx_glossary_status", + columns: &["review_status", "priority", "term_id"], + }, + ExpectedIndex { + table: "glossary_terms", + name: "idx_glossary_source_term", + columns: &["source_term"], + }, + ]; + matches_glossary_fingerprint_with_tables(snapshot, &tables, &indexes) +} + +fn matches_glossary_fingerprint_with_tables( + snapshot: &SqliteSchemaSnapshot, + tables: &[ExpectedTable<'_>], + indexes: &[ExpectedIndex<'_>], +) -> bool { + sqlite_migration::matches_fingerprint(snapshot, tables, indexes) +} + +const GLOSSARY_TERM_COLUMNS: [ExpectedColumn<'static>; 18] = [ + ExpectedColumn { + name: "term_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: true, + }, + ExpectedColumn { + name: "source_term", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "aliases_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "recommended_translation", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "allowed_translations_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_language", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "target_language", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "category", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "priority", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "scope_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "review_status", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_kind", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_ref", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_author", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_note", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_observed_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "created_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "updated_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, +]; + +const GLOSSARY_HISTORY_COLUMNS: [ExpectedColumn<'static>; 9] = [ + ExpectedColumn { + name: "history_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: true, + }, + ExpectedColumn { + name: "term_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "action", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "reviewer", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "reason", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "review_status", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "snapshot_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "observed_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, +]; + +const GLOSSARY_DELETION_COLUMNS: [ExpectedColumn<'static>; 8] = [ + ExpectedColumn { + name: "deletion_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: true, + }, + ExpectedColumn { + name: "term_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "reviewer", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "reason", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "snapshot_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "history_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "observed_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, +]; + +async fn create_glossary_schema( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + fail_after_step: Option, +) -> Result<()> { + let steps = [ + "CREATE TABLE schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + "CREATE TABLE glossary_terms ( + term_id TEXT PRIMARY KEY NOT NULL, + source_term TEXT NOT NULL, + aliases_json TEXT NOT NULL, + recommended_translation TEXT NOT NULL, + allowed_translations_json TEXT NOT NULL, + source_language TEXT, + target_language TEXT, + category TEXT, + priority INTEGER NOT NULL, + scope_json TEXT NOT NULL, + review_status TEXT NOT NULL, + source_kind TEXT NOT NULL, + source_ref TEXT, + source_author TEXT, + source_note TEXT, + source_observed_unix_seconds INTEGER NOT NULL, + created_unix_seconds INTEGER NOT NULL, + updated_unix_seconds INTEGER NOT NULL, + CHECK(length(term_id) > 0), + CHECK(length(source_term) > 0), + CHECK(length(recommended_translation) > 0), + CHECK(review_status IN ('draft', 'approved', 'deprecated', 'rejected')), + CHECK(source_kind IN ('manual', 'imported')) + )", + "CREATE TABLE glossary_term_history ( + history_id TEXT PRIMARY KEY NOT NULL, + term_id TEXT NOT NULL, + action TEXT NOT NULL, + reviewer TEXT, + reason TEXT, + source_json TEXT NOT NULL, + review_status TEXT NOT NULL, + snapshot_json TEXT NOT NULL, + observed_unix_seconds INTEGER NOT NULL, + FOREIGN KEY(term_id) REFERENCES glossary_terms(term_id) + )", + "CREATE TABLE glossary_term_deletions ( + deletion_id TEXT PRIMARY KEY NOT NULL, + term_id TEXT NOT NULL, + reviewer TEXT NOT NULL, + reason TEXT NOT NULL, + source_json TEXT NOT NULL, + snapshot_json TEXT NOT NULL, + history_json TEXT NOT NULL, + observed_unix_seconds INTEGER NOT NULL + )", + "CREATE INDEX idx_glossary_status + ON glossary_terms(review_status, priority DESC, term_id)", + "CREATE INDEX idx_glossary_source_term + ON glossary_terms(source_term)", + ]; + for (index, statement) in steps.iter().enumerate() { + sqlx::query(statement) + .execute(&mut **transaction) + .await + .map_err(db_error)?; + if fail_after_step == Some(index + 1) { + return Err(Error::Other(anyhow::anyhow!( + "Glossary migration failed after step {}", + index + 1 + ))); + } + } + Ok(()) +} + fn parse_json(value: String) -> Result { serde_json::from_str(&value).map_err(|error| Error::Serialization(error.to_string())) } @@ -769,35 +1188,10 @@ fn db_error(error: sqlx::Error) -> Error { Error::Other(error.into()) } -async fn ensure_column( - pool: &SqlitePool, - table: &str, - column: &str, - definition: &str, -) -> Result<()> { - let columns = sqlx::query(&format!("PRAGMA table_info({table})")) - .fetch_all(pool) - .await - .map_err(db_error)?; - let exists = columns.iter().any(|row| { - row.try_get::("name") - .map(|name| name == column) - .unwrap_or(false) - }); - if !exists { - sqlx::query(&format!( - "ALTER TABLE {table} ADD COLUMN {column} {definition}" - )) - .execute(pool) - .await - .map_err(db_error)?; - } - Ok(()) -} - #[cfg(test)] mod tests { use super::*; + use sqlx::sqlite::SqlitePoolOptions; use std::collections::BTreeMap; fn draft(status: GlossaryReviewStatus) -> GlossaryTermDraft { @@ -825,6 +1219,158 @@ mod tests { } } + async fn raw_repository( + path: &std::path::Path, + create_if_missing: bool, + ) -> SqliteGlossaryRepository { + let options = SqliteConnectOptions::from_str(&format!("sqlite://{}", path.display())) + .unwrap() + .create_if_missing(create_if_missing); + let pool = SqlitePoolOptions::new() + .max_connections(1) + .connect_with(options) + .await + .unwrap(); + SqliteGlossaryRepository { pool } + } + + #[tokio::test] + async fn sqlite_glossary_preserves_data_across_reopen_and_missing_row() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(GLOSSARY_REPOSITORY_FILE); + let repository = SqliteGlossaryRepository::new(&path).await.unwrap(); + repository + .add(draft(GlossaryReviewStatus::Draft)) + .await + .unwrap(); + repository + .review( + "term-sensei", + GlossaryReviewStatus::Approved, + "reviewer", + Some("accepted".to_string()), + ) + .await + .unwrap(); + sqlx::query("DROP TABLE schema_migrations") + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + + let reopened = SqliteGlossaryRepository::open(&path).await.unwrap(); + let term = reopened.find("term-sensei").await.unwrap(); + assert_eq!(term.definition.recommended_translation, "老师"); + assert_eq!(term.history.len(), 2); + assert_eq!(term.history[1].action, "approved"); + assert_eq!( + reopened.summary().await.unwrap().schema_version, + GLOSSARY_SCHEMA_VERSION + ); + } + + #[tokio::test] + async fn sqlite_glossary_future_schema_is_read_only_failure() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(GLOSSARY_REPOSITORY_FILE); + let repository = SqliteGlossaryRepository::new(&path).await.unwrap(); + sqlx::query("UPDATE schema_migrations SET version = ?2 WHERE component = ?1") + .bind(GLOSSARY_SCHEMA_COMPONENT) + .bind(i64::from(GLOSSARY_SCHEMA_VERSION) + 1) + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + let before = std::fs::read(&path).unwrap(); + let wal_path = std::path::PathBuf::from(format!("{}-wal", path.display())); + let shm_path = std::path::PathBuf::from(format!("{}-shm", path.display())); + let wal_before = std::fs::read(&wal_path).ok(); + let shm_before = std::fs::read(&shm_path).ok(); + let error = SqliteGlossaryRepository::new(&path).await.unwrap_err(); + assert!(error.to_string().contains("不支持的 Glossary schema")); + assert_eq!(std::fs::read(&path).unwrap(), before); + assert_eq!(std::fs::read(&wal_path).ok(), wal_before); + assert_eq!(std::fs::read(&shm_path).ok(), shm_before); + } + + #[tokio::test] + async fn sqlite_glossary_unknown_schema_fails_closed() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(GLOSSARY_REPOSITORY_FILE); + let repository = raw_repository(&path, true).await; + sqlx::query( + "CREATE TABLE schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + ) + .execute(&repository.pool) + .await + .unwrap(); + sqlx::query("INSERT INTO schema_migrations(component, version) VALUES (?1, 1)") + .bind(GLOSSARY_SCHEMA_COMPONENT) + .execute(&repository.pool) + .await + .unwrap(); + sqlx::query("CREATE TABLE glossary_terms (term_id TEXT PRIMARY KEY)") + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + let before = std::fs::read(&path).unwrap(); + let error = SqliteGlossaryRepository::open(&path).await.unwrap_err(); + assert!(error + .to_string() + .contains("schema version 与实际结构不一致")); + assert_eq!(std::fs::read(&path).unwrap(), before); + } + + #[tokio::test] + async fn sqlite_glossary_failed_new_schema_rolls_back_and_retries() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(GLOSSARY_REPOSITORY_FILE); + let repository = raw_repository(&path, true).await; + assert!(repository.init_schema_with_test_failure(3).await.is_err()); + let table_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM sqlite_master WHERE type = 'table'") + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(table_count, 0); + repository.init_schema().await.unwrap(); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(GLOSSARY_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, 1); + } + + #[tokio::test] + async fn sqlite_glossary_concurrent_new_open_has_one_current_schema() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(GLOSSARY_REPOSITORY_FILE); + let (left, right) = tokio::join!( + SqliteGlossaryRepository::new(&path), + SqliteGlossaryRepository::new(&path) + ); + assert!(left.is_ok(), "left open failed: {left:?}"); + assert!(right.is_ok(), "right open failed: {right:?}"); + let repository = left.unwrap(); + drop(right); + assert_eq!( + repository.summary().await.unwrap().schema_version, + GLOSSARY_SCHEMA_VERSION + ); + repository.pool.close().await; + let reopened = SqliteGlossaryRepository::open(&path).await.unwrap(); + assert_eq!( + reopened.summary().await.unwrap().schema_version, + GLOSSARY_SCHEMA_VERSION + ); + } + #[tokio::test] async fn sqlite_glossary_preserves_history_and_only_approved_terms_match() { let temp = tempfile::TempDir::new().unwrap(); diff --git a/infrastructure/src/lib.rs b/infrastructure/src/lib.rs index 0e92176..4755292 100644 --- a/infrastructure/src/lib.rs +++ b/infrastructure/src/lib.rs @@ -31,6 +31,7 @@ pub mod path_security; pub mod release_flow; pub mod release_ops; pub mod resources; +mod sqlite_migration; pub mod translation_memory; pub mod translation_tasks; pub mod translation_worker; diff --git a/infrastructure/src/sqlite_migration.rs b/infrastructure/src/sqlite_migration.rs new file mode 100644 index 0000000..51e9f2d --- /dev/null +++ b/infrastructure/src/sqlite_migration.rs @@ -0,0 +1,377 @@ +//! Shared, deliberately small SQLite schema-migration primitives. +//! +//! Component owners still define their own schema fingerprints and migration +//! steps. This module only owns the read-only preflight, schema snapshot, and +//! writer-lock mechanics shared by the long-lived SQLite stores. + +use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; +use sqlx::{Row, SqliteConnection, SqlitePool}; +use std::collections::{BTreeMap, BTreeSet}; +use std::path::Path; +use std::str::FromStr; +use std::time::Duration; + +pub const SCHEMA_MIGRATIONS_TABLE: &str = "schema_migrations"; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SqliteColumn { + pub data_type: String, + pub not_null: bool, + pub default_value: Option, + pub primary_key: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SqliteSchemaSnapshot { + /// Non-internal SQLite objects, including tables, indexes, views, and + /// triggers. Internal `sqlite_autoindex_*` objects are omitted. + pub objects: BTreeSet<(String, String)>, + pub tables: BTreeMap>, + pub indexes: BTreeMap>>, + pub component_version: Option, +} + +impl SqliteSchemaSnapshot { + pub fn is_empty(&self) -> bool { + self.objects.is_empty() + } +} + +#[derive(Debug, Clone, Copy)] +pub struct ExpectedColumn<'a> { + pub name: &'a str, + pub data_type: &'a str, + pub not_null: bool, + pub default_value: Option<&'a str>, + pub primary_key: bool, +} + +#[derive(Debug, Clone, Copy)] +pub struct ExpectedTable<'a> { + pub name: &'a str, + pub columns: &'a [ExpectedColumn<'a>], +} + +#[derive(Debug, Clone, Copy)] +pub struct ExpectedIndex<'a> { + pub table: &'a str, + pub name: &'a str, + pub columns: &'a [&'a str], +} + +/// Opens an existing database with SQLite's read-only flag and snapshots its +/// schema before a writable connection can perform any mutation. +pub async fn read_only_preflight( + path: &Path, + component: &str, +) -> Result { + let has_wal_sidecar = + sidecar_path(path, "-wal").exists() || sidecar_path(path, "-shm").exists(); + let options = SqliteConnectOptions::from_str(&format!("sqlite://{}", path.display()))? + .read_only(true) + .create_if_missing(false) + // A cleanly closed WAL database has all committed pages in the main + // file. Immutable read-only mode prevents SQLite from creating a new + // `-shm` sidecar during future-schema rejection. Live WAL sidecars + // must remain visible to the preflight reader. + .immutable(!has_wal_sidecar) + .busy_timeout(Duration::from_secs(30)); + let pool = SqlitePoolOptions::new() + .max_connections(1) + .connect_with(options) + .await?; + let snapshot = { + let mut connection = pool.acquire().await?; + snapshot_connection(&mut connection, component).await + }; + pool.close().await; + snapshot +} + +/// Connects a writable single-connection pool, retrying the SQLite-specific +/// exclusive lock needed when a connection switches an existing database to +/// WAL mode. SQLite's busy timeout cannot wait for that PRAGMA, so the retry +/// belongs around connection establishment rather than only around writes. +pub async fn connect_writable_pool( + options: SqliteConnectOptions, +) -> Result { + const MAX_ATTEMPTS: usize = 32; + + for attempt in 0..=MAX_ATTEMPTS { + match SqlitePoolOptions::new() + .max_connections(1) + .connect_with(options.clone()) + .await + { + Ok(pool) => return Ok(pool), + Err(error) if attempt < MAX_ATTEMPTS && is_sqlite_lock_error(&error) => { + let delay_millis = (25 * (attempt as u64 + 1)).min(250); + tokio::time::sleep(Duration::from_millis(delay_millis)).await; + } + Err(error) => return Err(error), + } + } + + unreachable!("SQLite connection retry loop always returns") +} + +fn is_sqlite_lock_error(error: &sqlx::Error) -> bool { + error.to_string().contains("database is locked") +} + +/// Snapshots the schema using an already-open connection. The caller may use +/// this both for read-only preflight and inside the migration transaction. +pub async fn snapshot_connection( + connection: &mut SqliteConnection, + component: &str, +) -> Result { + let object_rows = sqlx::query( + "SELECT type, name FROM sqlite_master + WHERE name NOT LIKE 'sqlite_%' + ORDER BY type, name", + ) + .fetch_all(&mut *connection) + .await?; + + let mut objects = BTreeSet::new(); + let mut table_names = BTreeSet::new(); + for row in object_rows { + let object_type: String = row.try_get("type")?; + let name: String = row.try_get("name")?; + if object_type == "table" { + table_names.insert(name.clone()); + } + objects.insert((object_type, name)); + } + + let mut tables = BTreeMap::new(); + let mut indexes = BTreeMap::new(); + for table in table_names { + let quoted_table = quote_identifier(&table); + let column_rows = sqlx::query(&format!("PRAGMA table_info({quoted_table})")) + .fetch_all(&mut *connection) + .await?; + let mut columns = BTreeMap::new(); + for row in column_rows { + let name: String = row.try_get("name")?; + let data_type: String = row.try_get("type")?; + let not_null: i64 = row.try_get("notnull")?; + let default_value: Option = row.try_get("dflt_value")?; + let primary_key: i64 = row.try_get("pk")?; + columns.insert( + name, + SqliteColumn { + data_type, + not_null: not_null != 0, + default_value, + primary_key: primary_key != 0, + }, + ); + } + tables.insert(table.clone(), columns); + + let index_rows = sqlx::query(&format!("PRAGMA index_list({quoted_table})")) + .fetch_all(&mut *connection) + .await?; + let mut table_indexes = BTreeMap::new(); + for row in index_rows { + let index_name: String = row.try_get("name")?; + if index_name.starts_with("sqlite_autoindex_") { + continue; + } + let quoted_index = quote_identifier(&index_name); + let index_columns = sqlx::query(&format!("PRAGMA index_info({quoted_index})")) + .fetch_all(&mut *connection) + .await?; + let mut columns = Vec::new(); + for index_column in index_columns { + let sequence: i64 = index_column.try_get("seqno")?; + let name: Option = index_column.try_get("name")?; + if sequence < 0 { + continue; + } + let name = name.ok_or_else(|| { + sqlx::Error::Protocol(format!( + "SQLite index {index_name} has an unnamed column" + )) + })?; + columns.push((sequence, name)); + } + columns.sort_by_key(|(sequence, _)| *sequence); + table_indexes.insert( + index_name, + columns + .into_iter() + .map(|(_, name)| name) + .collect::>(), + ); + } + if !table_indexes.is_empty() { + indexes.insert(table, table_indexes); + } + } + + let component_version = if tables.contains_key(SCHEMA_MIGRATIONS_TABLE) { + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(component) + .fetch_optional(&mut *connection) + .await? + } else { + None + }; + + Ok(SqliteSchemaSnapshot { + objects, + tables, + indexes, + component_version, + }) +} + +/// Starts a real SQLite writer transaction. `BEGIN IMMEDIATE` serializes DDL +/// migration writers instead of allowing two preflight results to race. +pub async fn begin_immediate( + pool: &SqlitePool, +) -> Result, sqlx::Error> { + pool.begin_with("BEGIN IMMEDIATE").await +} + +pub async fn write_component_version( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + component: &str, + version: u32, +) -> Result<(), sqlx::Error> { + sqlx::query( + "INSERT INTO schema_migrations(component, version) VALUES (?1, ?2) + ON CONFLICT(component) DO UPDATE SET version = excluded.version", + ) + .bind(component) + .bind(i64::from(version)) + .execute(&mut **transaction) + .await?; + Ok(()) +} + +pub async fn create_schema_migrations_table( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, +) -> Result<(), sqlx::Error> { + sqlx::query( + "CREATE TABLE IF NOT EXISTS schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + ) + .execute(&mut **transaction) + .await?; + Ok(()) +} + +pub fn matches_fingerprint( + snapshot: &SqliteSchemaSnapshot, + expected_tables: &[ExpectedTable<'_>], + expected_indexes: &[ExpectedIndex<'_>], +) -> bool { + let expected_table_names = expected_tables + .iter() + .map(|table| table.name) + .collect::>(); + let actual_table_names: BTreeSet<&str> = snapshot.tables.keys().map(String::as_str).collect(); + if actual_table_names != expected_table_names { + return false; + } + + let expected_object_names = expected_tables + .iter() + .map(|table| ("table", table.name)) + .chain(expected_indexes.iter().map(|index| ("index", index.name))) + .collect::>(); + let actual_object_names = snapshot + .objects + .iter() + .map(|(object_type, name)| (object_type.as_str(), name.as_str())) + .collect::>(); + if actual_object_names != expected_object_names { + return false; + } + + for table in expected_tables { + let Some(actual_columns) = snapshot.tables.get(table.name) else { + return false; + }; + if actual_columns.len() != table.columns.len() { + return false; + } + for expected in table.columns { + let Some(actual) = actual_columns.get(expected.name) else { + return false; + }; + if actual.data_type.to_ascii_uppercase() != expected.data_type + || actual.not_null != expected.not_null + || actual.primary_key != expected.primary_key + || normalize_default(actual.default_value.as_deref()) + != normalize_default(expected.default_value) + { + return false; + } + } + } + + let expected_indexes = expected_indexes + .iter() + .map(|index| { + ( + index.table.to_string(), + index.name.to_string(), + index + .columns + .iter() + .map(|column| (*column).to_string()) + .collect::>(), + ) + }) + .collect::>(); + let actual_indexes = snapshot + .indexes + .iter() + .flat_map(|(table, indexes)| { + indexes + .iter() + .map(|(name, columns)| (table.clone(), name.clone(), columns.clone())) + }) + .collect::>(); + actual_indexes == expected_indexes +} + +pub fn schema_migrations_table() -> ExpectedTable<'static> { + ExpectedTable { + name: SCHEMA_MIGRATIONS_TABLE, + columns: &[ + ExpectedColumn { + name: "component", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: true, + }, + ExpectedColumn { + name: "version", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ], + } +} + +fn normalize_default(value: Option<&str>) -> Option { + value.map(|value| value.trim().to_ascii_lowercase()) +} + +fn quote_identifier(value: &str) -> String { + format!("\"{}\"", value.replace('"', "\"\"")) +} + +fn sidecar_path(path: &Path, suffix: &str) -> std::path::PathBuf { + std::path::PathBuf::from(format!("{}{}", path.display(), suffix)) +} diff --git a/infrastructure/src/translation_memory.rs b/infrastructure/src/translation_memory.rs index 56101fe..67f6ae1 100644 --- a/infrastructure/src/translation_memory.rs +++ b/infrastructure/src/translation_memory.rs @@ -1,6 +1,9 @@ //! 跨 official release 的 Translation Memory SQLite 仓储。 use crate::path_security::{set_file_mode, STATE_FILE_MODE}; +use crate::sqlite_migration::{ + self, ExpectedColumn, ExpectedIndex, ExpectedTable, SqliteSchemaSnapshot, +}; use async_trait::async_trait; use bat_core::domain::{ TranslationMemoryContext, TranslationMemoryDraft, TranslationMemoryEntry, @@ -10,7 +13,7 @@ use bat_core::domain::{ use bat_core::repositories::TranslationMemoryRepository; use bat_core::{Error, Result}; use serde::de::DeserializeOwned; -use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions}; +use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode}; use sqlx::{Row, SqlitePool}; use std::fs; use std::path::{Path, PathBuf}; @@ -85,14 +88,24 @@ impl SqliteTranslationMemoryRepository { return Err(Error::NotFound(absolute.display().to_string())); } + if let Ok(metadata) = fs::symlink_metadata(&absolute) { + if metadata.len() > 0 { + let snapshot = sqlite_migration::read_only_preflight( + &absolute, + TRANSLATION_MEMORY_SCHEMA_COMPONENT, + ) + .await + .map_err(db_error)?; + classify_translation_memory_schema(&snapshot)?; + } + } + let options = SqliteConnectOptions::from_str(&format!("sqlite://{}", absolute.display())) .map_err(|error| Error::Other(error.into()))? .create_if_missing(create_if_missing) .journal_mode(SqliteJournalMode::Wal) .busy_timeout(Duration::from_secs(30)); - let pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with(options) + let pool = sqlite_migration::connect_writable_pool(options) .await .map_err(db_error)?; set_file_mode(&absolute, STATE_FILE_MODE, "Translation Memory 数据库") @@ -103,98 +116,93 @@ impl SqliteTranslationMemoryRepository { } async fn init_schema(&self) -> Result<()> { - sqlx::query( - r#" - CREATE TABLE IF NOT EXISTS schema_migrations ( - component TEXT PRIMARY KEY NOT NULL, - version INTEGER NOT NULL CHECK(version >= 1) - ) - "#, - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - r#" - CREATE TABLE IF NOT EXISTS translation_memory ( - record_id TEXT PRIMARY KEY NOT NULL, - source_text TEXT NOT NULL, - source_hash TEXT NOT NULL, - normalized_source_text TEXT NOT NULL, - source_context_json TEXT NOT NULL, - source_context_hash TEXT NOT NULL, - translated_text TEXT NOT NULL, - translation_source_kind TEXT NOT NULL, - trust_status TEXT NOT NULL, - official_release_id TEXT NOT NULL, - source_trace_json TEXT NOT NULL, - provider TEXT, - provider_run_id TEXT, - created_unix_seconds INTEGER NOT NULL, - updated_unix_seconds INTEGER NOT NULL, - trusted_unix_seconds INTEGER, - trusted_by TEXT, - trusted_reason TEXT, - supersedes_record_id TEXT, - superseded_by_record_id TEXT, - CHECK (length(source_text) > 0), - CHECK (length(source_hash) > 0), - CHECK (length(source_context_hash) > 0), - CHECK (length(official_release_id) > 0), - CHECK (translation_source_kind IN ('provider', 'manual', 'imported')), - CHECK (trust_status IN ('candidate', 'trusted', 'superseded', 'rejected')) - ) - "#, - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - "CREATE INDEX IF NOT EXISTS idx_translation_memory_source_hash \ - ON translation_memory(source_hash)", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - "CREATE INDEX IF NOT EXISTS idx_translation_memory_normalized_source \ - ON translation_memory(normalized_source_text)", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - sqlx::query( - "CREATE INDEX IF NOT EXISTS idx_translation_memory_context \ - ON translation_memory(source_hash, source_context_hash)", - ) - .execute(&self.pool) - .await - .map_err(db_error)?; + self.init_schema_with_failure(None).await + } - let current: Option = - sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") - .bind(TRANSLATION_MEMORY_SCHEMA_COMPONENT) - .fetch_optional(&self.pool) + #[cfg(test)] + async fn init_schema_with_test_failure(&self, fail_after_step: usize) -> Result<()> { + self.init_schema_with_failure(Some(fail_after_step)).await + } + + async fn init_schema_with_failure(&self, fail_after_step: Option) -> Result<()> { + let mut transaction = sqlite_migration::begin_immediate(&self.pool) + .await + .map_err(db_error)?; + let result = self + .migrate_in_transaction(&mut transaction, fail_after_step) + .await; + match result { + Ok(()) => transaction.commit().await.map_err(db_error), + Err(error) => { + let _ = transaction.rollback().await; + Err(error) + } + } + } + + async fn migrate_in_transaction( + &self, + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + fail_after_step: Option, + ) -> Result<()> { + let snapshot = sqlite_migration::snapshot_connection( + transaction.as_mut(), + TRANSLATION_MEMORY_SCHEMA_COMPONENT, + ) + .await + .map_err(db_error)?; + match classify_translation_memory_schema(&snapshot)? { + TranslationMemorySchemaState::Empty => { + create_translation_memory_schema(transaction, fail_after_step).await?; + sqlite_migration::write_component_version( + transaction, + TRANSLATION_MEMORY_SCHEMA_COMPONENT, + TRANSLATION_MEMORY_SCHEMA_VERSION, + ) .await .map_err(db_error)?; - if current.is_some_and(|version| version > i64::from(TRANSLATION_MEMORY_SCHEMA_VERSION)) { - return Err(Error::InvalidArgument(format!( - "不支持的 Translation Memory schema 版本:{}", - current.unwrap_or_default() - ))); + } + TranslationMemorySchemaState::Version(version) + if version == TRANSLATION_MEMORY_SCHEMA_VERSION => + { + if snapshot.component_version.is_none() { + if !snapshot + .tables + .contains_key(sqlite_migration::SCHEMA_MIGRATIONS_TABLE) + { + sqlite_migration::create_schema_migrations_table(transaction) + .await + .map_err(db_error)?; + } + sqlite_migration::write_component_version( + transaction, + TRANSLATION_MEMORY_SCHEMA_COMPONENT, + TRANSLATION_MEMORY_SCHEMA_VERSION, + ) + .await + .map_err(db_error)?; + } + } + TranslationMemorySchemaState::Version(version) => { + return Err(invalid_translation_memory_schema(format!( + "内部不支持的迁移起点 {version}" + ))); + } } - sqlx::query( - r#" - INSERT INTO schema_migrations(component, version) - VALUES (?1, ?2) - ON CONFLICT(component) DO UPDATE SET version = excluded.version - "#, + + let final_snapshot = sqlite_migration::snapshot_connection( + transaction.as_mut(), + TRANSLATION_MEMORY_SCHEMA_COMPONENT, ) - .bind(TRANSLATION_MEMORY_SCHEMA_COMPONENT) - .bind(i64::from(TRANSLATION_MEMORY_SCHEMA_VERSION)) - .execute(&self.pool) .await .map_err(db_error)?; + if final_snapshot.component_version != Some(i64::from(TRANSLATION_MEMORY_SCHEMA_VERSION)) + || !matches_translation_memory_fingerprint(&final_snapshot) + { + return Err(invalid_translation_memory_schema( + "migration 结果与当前 schema fingerprint 不一致".to_string(), + )); + } Ok(()) } @@ -219,6 +227,297 @@ impl SqliteTranslationMemoryRepository { } } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum TranslationMemorySchemaState { + Empty, + Version(u32), +} + +fn classify_translation_memory_schema( + snapshot: &SqliteSchemaSnapshot, +) -> Result { + if snapshot.is_empty() { + return Ok(TranslationMemorySchemaState::Empty); + } + if let Some(observed) = snapshot.component_version { + if observed > i64::from(TRANSLATION_MEMORY_SCHEMA_VERSION) { + return Err(invalid_translation_memory_schema(format!( + "不支持的 Translation Memory schema 版本:observed={observed}, supported={TRANSLATION_MEMORY_SCHEMA_VERSION}" + ))); + } + if observed < 1 { + return Err(invalid_translation_memory_schema(format!( + "schema version 无效:observed={observed}, supported=1..={TRANSLATION_MEMORY_SCHEMA_VERSION}" + ))); + } + } + if (matches_translation_memory_fingerprint(snapshot) + || (snapshot.component_version.is_none() + && matches_translation_memory_component_fingerprint(snapshot))) + && snapshot + .component_version + .is_none_or(|observed| observed == i64::from(TRANSLATION_MEMORY_SCHEMA_VERSION)) + { + return Ok(TranslationMemorySchemaState::Version( + TRANSLATION_MEMORY_SCHEMA_VERSION, + )); + } + match snapshot.component_version { + Some(observed) => Err(invalid_translation_memory_schema(format!( + "schema version 与实际结构不一致:observed={observed}, supported={TRANSLATION_MEMORY_SCHEMA_VERSION}" + ))), + None => Err(invalid_translation_memory_schema( + "未识别的 legacy schema,拒绝静默修复".to_string(), + )), + } +} + +fn invalid_translation_memory_schema(detail: String) -> Error { + Error::InvalidArgument(format!("Translation Memory schema 无效:{detail}")) +} + +fn matches_translation_memory_fingerprint(snapshot: &SqliteSchemaSnapshot) -> bool { + matches_translation_memory_fingerprint_with_migrations(snapshot, true) +} + +fn matches_translation_memory_component_fingerprint(snapshot: &SqliteSchemaSnapshot) -> bool { + matches_translation_memory_fingerprint_with_migrations(snapshot, false) +} + +fn matches_translation_memory_fingerprint_with_migrations( + snapshot: &SqliteSchemaSnapshot, + include_migrations: bool, +) -> bool { + let tables = [ExpectedTable { + name: "translation_memory", + columns: &TRANSLATION_MEMORY_COLUMNS, + }]; + let indexes = [ + ExpectedIndex { + table: "translation_memory", + name: "idx_translation_memory_source_hash", + columns: &["source_hash"], + }, + ExpectedIndex { + table: "translation_memory", + name: "idx_translation_memory_normalized_source", + columns: &["normalized_source_text"], + }, + ExpectedIndex { + table: "translation_memory", + name: "idx_translation_memory_context", + columns: &["source_hash", "source_context_hash"], + }, + ]; + if include_migrations { + let tables = [sqlite_migration::schema_migrations_table(), tables[0]]; + return sqlite_migration::matches_fingerprint(snapshot, &tables, &indexes); + } + sqlite_migration::matches_fingerprint(snapshot, &tables, &indexes) +} + +const TRANSLATION_MEMORY_COLUMNS: [ExpectedColumn<'static>; 20] = [ + ExpectedColumn { + name: "record_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: true, + }, + ExpectedColumn { + name: "source_text", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_hash", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "normalized_source_text", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_context_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_context_hash", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "translated_text", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "translation_source_kind", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "trust_status", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "official_release_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "source_trace_json", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "provider", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "provider_run_id", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "created_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "updated_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "trusted_unix_seconds", + data_type: "INTEGER", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "trusted_by", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "trusted_reason", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "supersedes_record_id", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "superseded_by_record_id", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, +]; + +async fn create_translation_memory_schema( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + fail_after_step: Option, +) -> Result<()> { + let steps = [ + "CREATE TABLE schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + "CREATE TABLE translation_memory ( + record_id TEXT PRIMARY KEY NOT NULL, + source_text TEXT NOT NULL, + source_hash TEXT NOT NULL, + normalized_source_text TEXT NOT NULL, + source_context_json TEXT NOT NULL, + source_context_hash TEXT NOT NULL, + translated_text TEXT NOT NULL, + translation_source_kind TEXT NOT NULL, + trust_status TEXT NOT NULL, + official_release_id TEXT NOT NULL, + source_trace_json TEXT NOT NULL, + provider TEXT, + provider_run_id TEXT, + created_unix_seconds INTEGER NOT NULL, + updated_unix_seconds INTEGER NOT NULL, + trusted_unix_seconds INTEGER, + trusted_by TEXT, + trusted_reason TEXT, + supersedes_record_id TEXT, + superseded_by_record_id TEXT, + CHECK (length(source_text) > 0), + CHECK (length(source_hash) > 0), + CHECK (length(source_context_hash) > 0), + CHECK (length(official_release_id) > 0), + CHECK (translation_source_kind IN ('provider', 'manual', 'imported')), + CHECK (trust_status IN ('candidate', 'trusted', 'superseded', 'rejected')) + )", + "CREATE INDEX idx_translation_memory_source_hash + ON translation_memory(source_hash)", + "CREATE INDEX idx_translation_memory_normalized_source + ON translation_memory(normalized_source_text)", + "CREATE INDEX idx_translation_memory_context + ON translation_memory(source_hash, source_context_hash)", + ]; + for (index, statement) in steps.iter().enumerate() { + sqlx::query(statement) + .execute(&mut **transaction) + .await + .map_err(db_error)?; + if fail_after_step == Some(index + 1) { + return Err(Error::Other(anyhow::anyhow!( + "Translation Memory migration failed after step {}", + index + 1 + ))); + } + } + Ok(()) +} + #[async_trait] impl TranslationMemoryRepository for SqliteTranslationMemoryRepository { async fn upsert_candidate( @@ -721,6 +1020,7 @@ fn ensure_safe_tm_parent(parent: &Path) -> Result<()> { mod tests { use super::*; use bat_core::domain::TranslationMemorySourceTrace; + use sqlx::sqlite::SqlitePoolOptions; fn draft(release: &str, source: &str, translated: &str) -> TranslationMemoryDraft { let source_trace = TranslationMemorySourceTrace { @@ -762,6 +1062,163 @@ mod tests { } } + async fn raw_repository( + path: &std::path::Path, + create_if_missing: bool, + ) -> SqliteTranslationMemoryRepository { + let options = SqliteConnectOptions::from_str(&format!("sqlite://{}", path.display())) + .unwrap() + .create_if_missing(create_if_missing); + let pool = SqlitePoolOptions::new() + .max_connections(1) + .connect_with(options) + .await + .unwrap(); + SqliteTranslationMemoryRepository { pool } + } + + #[tokio::test] + async fn sqlite_translation_memory_preserves_data_across_reopen_and_missing_row() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_MEMORY_REPOSITORY_FILE); + let repository = SqliteTranslationMemoryRepository::new(&path).await.unwrap(); + let entry = repository + .upsert_candidate(draft("release-1", "Hello", "你好")) + .await + .unwrap(); + let trusted = repository + .confirm(&entry.record_id, "reviewer", Some("accepted".to_string())) + .await + .unwrap(); + sqlx::query("DROP TABLE schema_migrations") + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + + let reopened = SqliteTranslationMemoryRepository::open(&path) + .await + .unwrap(); + let restored = reopened.find(&trusted.record_id).await.unwrap(); + assert_eq!(restored.translated_text, "你好"); + assert_eq!(restored.trust_status, TranslationMemoryTrustStatus::Trusted); + assert_eq!(restored.source_trace.official_release_id, "release-1"); + assert_eq!( + reopened.summary().await.unwrap().schema_version, + TRANSLATION_MEMORY_SCHEMA_VERSION + ); + } + + #[tokio::test] + async fn sqlite_translation_memory_future_schema_is_read_only_failure() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_MEMORY_REPOSITORY_FILE); + let repository = SqliteTranslationMemoryRepository::new(&path).await.unwrap(); + sqlx::query("UPDATE schema_migrations SET version = ?2 WHERE component = ?1") + .bind(TRANSLATION_MEMORY_SCHEMA_COMPONENT) + .bind(i64::from(TRANSLATION_MEMORY_SCHEMA_VERSION) + 1) + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + let before = std::fs::read(&path).unwrap(); + let wal_path = std::path::PathBuf::from(format!("{}-wal", path.display())); + let shm_path = std::path::PathBuf::from(format!("{}-shm", path.display())); + let wal_before = std::fs::read(&wal_path).ok(); + let shm_before = std::fs::read(&shm_path).ok(); + let error = SqliteTranslationMemoryRepository::new(&path) + .await + .unwrap_err(); + assert!(error + .to_string() + .contains("不支持的 Translation Memory schema")); + assert_eq!(std::fs::read(&path).unwrap(), before); + assert_eq!(std::fs::read(&wal_path).ok(), wal_before); + assert_eq!(std::fs::read(&shm_path).ok(), shm_before); + } + + #[tokio::test] + async fn sqlite_translation_memory_unknown_schema_fails_closed() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_MEMORY_REPOSITORY_FILE); + let repository = raw_repository(&path, true).await; + sqlx::query( + "CREATE TABLE schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + ) + .execute(&repository.pool) + .await + .unwrap(); + sqlx::query("INSERT INTO schema_migrations(component, version) VALUES (?1, 1)") + .bind(TRANSLATION_MEMORY_SCHEMA_COMPONENT) + .execute(&repository.pool) + .await + .unwrap(); + sqlx::query("CREATE TABLE translation_memory (record_id TEXT PRIMARY KEY)") + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + let before = std::fs::read(&path).unwrap(); + let error = SqliteTranslationMemoryRepository::open(&path) + .await + .unwrap_err(); + assert!(error + .to_string() + .contains("schema version 与实际结构不一致")); + assert_eq!(std::fs::read(&path).unwrap(), before); + } + + #[tokio::test] + async fn sqlite_translation_memory_failed_new_schema_rolls_back_and_retries() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_MEMORY_REPOSITORY_FILE); + let repository = raw_repository(&path, true).await; + assert!(repository.init_schema_with_test_failure(2).await.is_err()); + let table_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM sqlite_master WHERE type = 'table'") + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(table_count, 0); + repository.init_schema().await.unwrap(); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_MEMORY_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, 1); + } + + #[tokio::test] + async fn sqlite_translation_memory_concurrent_new_open_has_one_current_schema() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_MEMORY_REPOSITORY_FILE); + let (left, right) = tokio::join!( + SqliteTranslationMemoryRepository::new(&path), + SqliteTranslationMemoryRepository::new(&path) + ); + assert!(left.is_ok(), "left open failed: {left:?}"); + assert!(right.is_ok(), "right open failed: {right:?}"); + let repository = left.unwrap(); + drop(right); + assert_eq!( + repository.summary().await.unwrap().schema_version, + TRANSLATION_MEMORY_SCHEMA_VERSION + ); + repository.pool.close().await; + let reopened = SqliteTranslationMemoryRepository::open(&path) + .await + .unwrap(); + assert_eq!( + reopened.summary().await.unwrap().schema_version, + TRANSLATION_MEMORY_SCHEMA_VERSION + ); + } + #[tokio::test] async fn initializes_schema_and_reuses_trusted_entry_across_releases() { let temp = tempfile::TempDir::new().unwrap(); diff --git a/infrastructure/src/translation_tasks.rs b/infrastructure/src/translation_tasks.rs index d535d35..b40e110 100644 --- a/infrastructure/src/translation_tasks.rs +++ b/infrastructure/src/translation_tasks.rs @@ -12,9 +12,10 @@ use crate::path_security::{ ensure_path_within_root, ensure_safe_file_target, read_file_no_symlink, write_file_atomic, STATE_FILE_MODE, }; +use crate::sqlite_migration::{self, ExpectedColumn, ExpectedTable, SqliteSchemaSnapshot}; use bat_core::Result; use serde::{Deserialize, Serialize}; -use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions}; +use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode}; use sqlx::{QueryBuilder, Sqlite, SqlitePool}; use std::collections::{BTreeMap, BTreeSet}; use std::path::Path; @@ -588,14 +589,31 @@ impl SqliteTranslationTaskRepository { } } + let existing = std::fs::symlink_metadata(path).ok(); + if let Some(metadata) = existing.as_ref() { + if metadata.file_type().is_symlink() || !metadata.is_file() { + return Err(bat_core::Error::InvalidArgument(format!( + "翻译任务数据库必须是普通文件:{}", + path.display() + ))); + } + if metadata.len() > 0 { + let snapshot = + sqlite_migration::read_only_preflight(path, TRANSLATION_TASK_SCHEMA_COMPONENT) + .await + .map_err(db_error)?; + classify_task_schema(&snapshot)?; + } + } else if !create_if_missing { + return Err(bat_core::Error::NotFound(path.display().to_string())); + } + let options = SqliteConnectOptions::from_str(&format!("sqlite://{}", path.display())) .map_err(|error| bat_core::Error::Other(error.into()))? .create_if_missing(create_if_missing) .journal_mode(SqliteJournalMode::Wal) .busy_timeout(Duration::from_secs(30)); - let pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with(options) + let pool = sqlite_migration::connect_writable_pool(options) .await .map_err(|error| bat_core::Error::Other(error.into()))?; let repository = Self { pool }; @@ -609,163 +627,109 @@ impl SqliteTranslationTaskRepository { } async fn init_schema(&self) -> Result<()> { - sqlx::query( - r#" - CREATE TABLE IF NOT EXISTS schema_migrations ( - component TEXT PRIMARY KEY NOT NULL, - version INTEGER NOT NULL CHECK(version >= 1) - ) - "#, + self.init_schema_with_failure(None).await + } + + #[cfg(test)] + async fn init_schema_with_test_failure(&self, fail_after_step: usize) -> Result<()> { + self.init_schema_with_failure(Some(fail_after_step)).await + } + + async fn init_schema_with_failure(&self, fail_after_step: Option) -> Result<()> { + let mut transaction = sqlite_migration::begin_immediate(&self.pool) + .await + .map_err(db_error)?; + let result = self + .migrate_in_transaction(&mut transaction, fail_after_step) + .await; + match result { + Ok(()) => transaction.commit().await.map_err(db_error), + Err(error) => { + let _ = transaction.rollback().await; + Err(error) + } + } + } + + async fn migrate_in_transaction( + &self, + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + fail_after_step: Option, + ) -> Result<()> { + let snapshot = sqlite_migration::snapshot_connection( + transaction.as_mut(), + TRANSLATION_TASK_SCHEMA_COMPONENT, ) - .execute(&self.pool) .await .map_err(db_error)?; - - sqlx::query( - r#" - CREATE TABLE IF NOT EXISTS translation_tasks ( - task_id TEXT PRIMARY KEY NOT NULL, - official_release_id TEXT NOT NULL, - destination TEXT NOT NULL, - archive_entry TEXT, - queue_status TEXT NOT NULL, - queue_reason TEXT, - parse_status TEXT, - text_unit_formats_json TEXT NOT NULL DEFAULT '[]', - task_json TEXT NOT NULL DEFAULT '{}', - worker_status TEXT NOT NULL DEFAULT 'queued', - failure_reason TEXT, - attempt_count INTEGER NOT NULL DEFAULT 0 CHECK(attempt_count >= 0), - created_unix_seconds INTEGER NOT NULL, - updated_unix_seconds INTEGER NOT NULL, - completed_unix_seconds INTEGER, - provider_run_id TEXT, - translation_results_json TEXT NOT NULL DEFAULT '[]', - provider TEXT, - lease_owner TEXT, - lease_expires_unix_seconds INTEGER, - failure_class TEXT, - failure_retryable INTEGER NOT NULL DEFAULT 0, - next_attempt_unix_seconds INTEGER - ) - "#, - ) - .execute(&self.pool) - .await - .map_err(db_error)?; - - // These defaults keep old experimental databases readable while the - // schema version table records the migration boundary explicitly. - ensure_column(&self.pool, "translation_tasks", "queue_reason", "TEXT").await?; - ensure_column(&self.pool, "translation_tasks", "parse_status", "TEXT").await?; - ensure_column( - &self.pool, - "translation_tasks", - "text_unit_formats_json", - "TEXT NOT NULL DEFAULT '[]'", - ) - .await?; - ensure_column( - &self.pool, - "translation_tasks", - "task_json", - "TEXT NOT NULL DEFAULT '{}'", - ) - .await?; - ensure_column( - &self.pool, - "translation_tasks", - "worker_status", - "TEXT NOT NULL DEFAULT 'queued'", - ) - .await?; - ensure_column(&self.pool, "translation_tasks", "failure_reason", "TEXT").await?; - ensure_column( - &self.pool, - "translation_tasks", - "attempt_count", - "INTEGER NOT NULL DEFAULT 0", - ) - .await?; - ensure_column( - &self.pool, - "translation_tasks", - "created_unix_seconds", - "INTEGER NOT NULL DEFAULT 0", - ) - .await?; - ensure_column( - &self.pool, - "translation_tasks", - "updated_unix_seconds", - "INTEGER NOT NULL DEFAULT 0", - ) - .await?; - ensure_column( - &self.pool, - "translation_tasks", - "completed_unix_seconds", - "INTEGER", - ) - .await?; - ensure_column(&self.pool, "translation_tasks", "provider_run_id", "TEXT").await?; - ensure_column( - &self.pool, - "translation_tasks", - "translation_results_json", - "TEXT NOT NULL DEFAULT '[]'", - ) - .await?; - ensure_column(&self.pool, "translation_tasks", "provider", "TEXT").await?; - ensure_column(&self.pool, "translation_tasks", "lease_owner", "TEXT").await?; - ensure_column( - &self.pool, - "translation_tasks", - "lease_expires_unix_seconds", - "INTEGER", - ) - .await?; - ensure_column(&self.pool, "translation_tasks", "failure_class", "TEXT").await?; - ensure_column( - &self.pool, - "translation_tasks", - "failure_retryable", - "INTEGER NOT NULL DEFAULT 0", - ) - .await?; - ensure_column( - &self.pool, - "translation_tasks", - "next_attempt_unix_seconds", - "INTEGER", - ) - .await?; - - let current: Option = - sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") - .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) - .fetch_optional(&self.pool) + match classify_task_schema(&snapshot)? { + TaskSchemaState::Empty => { + create_translation_task_schema(transaction).await?; + sqlite_migration::write_component_version( + transaction, + TRANSLATION_TASK_SCHEMA_COMPONENT, + TRANSLATION_TASK_SCHEMA_VERSION, + ) .await .map_err(db_error)?; - if current.is_some_and(|version| version > i64::from(TRANSLATION_TASK_SCHEMA_VERSION)) { - return Err(bat_core::Error::InvalidArgument(format!( - "不支持的翻译任务 schema 版本:{}", - current.unwrap_or_default() - ))); + } + TaskSchemaState::Version(version) if version == TRANSLATION_TASK_SCHEMA_VERSION => { + if snapshot.component_version.is_none() { + if !snapshot + .tables + .contains_key(sqlite_migration::SCHEMA_MIGRATIONS_TABLE) + { + sqlite_migration::create_schema_migrations_table(transaction) + .await + .map_err(db_error)?; + } + sqlite_migration::write_component_version( + transaction, + TRANSLATION_TASK_SCHEMA_COMPONENT, + TRANSLATION_TASK_SCHEMA_VERSION, + ) + .await + .map_err(db_error)?; + } + } + TaskSchemaState::Version(1) => { + if !snapshot + .tables + .contains_key(sqlite_migration::SCHEMA_MIGRATIONS_TABLE) + { + sqlite_migration::create_schema_migrations_table(transaction) + .await + .map_err(db_error)?; + } + migrate_translation_tasks_v1_to_v2(transaction, fail_after_step).await?; + sqlite_migration::write_component_version( + transaction, + TRANSLATION_TASK_SCHEMA_COMPONENT, + TRANSLATION_TASK_SCHEMA_VERSION, + ) + .await + .map_err(db_error)?; + } + TaskSchemaState::Version(version) => { + return Err(invalid_task_schema(format!( + "内部不支持的迁移起点 {version}" + ))); + } } - sqlx::query( - r#" - INSERT INTO schema_migrations(component, version) - VALUES (?1, ?2) - ON CONFLICT(component) DO UPDATE SET version = excluded.version - "#, + + let final_snapshot = sqlite_migration::snapshot_connection( + transaction.as_mut(), + TRANSLATION_TASK_SCHEMA_COMPONENT, ) - .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) - .bind(i64::from(TRANSLATION_TASK_SCHEMA_VERSION)) - .execute(&self.pool) .await .map_err(db_error)?; - + if final_snapshot.component_version != Some(i64::from(TRANSLATION_TASK_SCHEMA_VERSION)) + || !matches_task_fingerprint(&final_snapshot, TRANSLATION_TASK_SCHEMA_VERSION) + { + return Err(invalid_task_schema( + "migration 结果与当前 schema fingerprint 不一致".to_string(), + )); + } Ok(()) } @@ -1579,28 +1543,432 @@ fn parse_status_label(status: crate::official_parse::OfficialParseStatus) -> &'s } } -async fn ensure_column( - pool: &SqlitePool, - table: &str, - column: &str, - column_type: &str, +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum TaskSchemaState { + Empty, + Version(u32), +} + +fn classify_task_schema(snapshot: &SqliteSchemaSnapshot) -> Result { + if snapshot.is_empty() { + return Ok(TaskSchemaState::Empty); + } + if let Some(observed) = snapshot.component_version { + if observed > i64::from(TRANSLATION_TASK_SCHEMA_VERSION) { + return Err(invalid_task_schema(format!( + "不支持的翻译任务 schema 版本:observed={observed}, supported={TRANSLATION_TASK_SCHEMA_VERSION}" + ))); + } + if observed < 1 { + return Err(invalid_task_schema(format!( + "schema version 无效:observed={observed}, supported=1..={TRANSLATION_TASK_SCHEMA_VERSION}" + ))); + } + } + for version in 1..=TRANSLATION_TASK_SCHEMA_VERSION { + if (matches_task_fingerprint(snapshot, version) + || (snapshot.component_version.is_none() + && matches_task_component_fingerprint(snapshot, version))) + && snapshot + .component_version + .is_none_or(|observed| observed == i64::from(version)) + { + return Ok(TaskSchemaState::Version(version)); + } + } + match snapshot.component_version { + Some(observed) => Err(invalid_task_schema(format!( + "schema version 与实际结构不一致:observed={observed}, supported={TRANSLATION_TASK_SCHEMA_VERSION}" + ))), + None => Err(invalid_task_schema( + "未识别的 legacy schema,拒绝静默修复".to_string(), + )), + } +} + +fn invalid_task_schema(detail: String) -> bat_core::Error { + bat_core::Error::InvalidArgument(format!("Translation Tasks schema 无效:{detail}")) +} + +fn matches_task_fingerprint(snapshot: &SqliteSchemaSnapshot, version: u32) -> bool { + let tables = task_tables(version); + sqlite_migration::matches_fingerprint(snapshot, &tables, &[]) +} + +fn matches_task_component_fingerprint(snapshot: &SqliteSchemaSnapshot, version: u32) -> bool { + let tables = [ExpectedTable { + name: "translation_tasks", + columns: if version == 1 { + &TASK_V1_COLUMNS + } else { + &TASK_V2_COLUMNS + }, + }]; + sqlite_migration::matches_fingerprint(snapshot, &tables, &[]) +} + +fn task_tables(version: u32) -> [ExpectedTable<'static>; 2] { + [ + sqlite_migration::schema_migrations_table(), + ExpectedTable { + name: "translation_tasks", + columns: if version == 1 { + &TASK_V1_COLUMNS + } else { + &TASK_V2_COLUMNS + }, + }, + ] +} + +const TASK_V1_COLUMNS: [ExpectedColumn<'static>; 16] = [ + ExpectedColumn { + name: "task_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: true, + }, + ExpectedColumn { + name: "official_release_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "destination", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "archive_entry", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "queue_status", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "queue_reason", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "parse_status", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "text_unit_formats_json", + data_type: "TEXT", + not_null: true, + default_value: Some("'[]'"), + primary_key: false, + }, + ExpectedColumn { + name: "task_json", + data_type: "TEXT", + not_null: true, + default_value: Some("'{}'"), + primary_key: false, + }, + ExpectedColumn { + name: "worker_status", + data_type: "TEXT", + not_null: true, + default_value: Some("'queued'"), + primary_key: false, + }, + ExpectedColumn { + name: "failure_reason", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "attempt_count", + data_type: "INTEGER", + not_null: true, + default_value: Some("0"), + primary_key: false, + }, + ExpectedColumn { + name: "created_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "updated_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "completed_unix_seconds", + data_type: "INTEGER", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "provider_run_id", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, +]; + +const TASK_V2_COLUMNS: [ExpectedColumn<'static>; 23] = [ + ExpectedColumn { + name: "task_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: true, + }, + ExpectedColumn { + name: "official_release_id", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "destination", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "archive_entry", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "queue_status", + data_type: "TEXT", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "queue_reason", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "parse_status", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "text_unit_formats_json", + data_type: "TEXT", + not_null: true, + default_value: Some("'[]'"), + primary_key: false, + }, + ExpectedColumn { + name: "task_json", + data_type: "TEXT", + not_null: true, + default_value: Some("'{}'"), + primary_key: false, + }, + ExpectedColumn { + name: "worker_status", + data_type: "TEXT", + not_null: true, + default_value: Some("'queued'"), + primary_key: false, + }, + ExpectedColumn { + name: "failure_reason", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "attempt_count", + data_type: "INTEGER", + not_null: true, + default_value: Some("0"), + primary_key: false, + }, + ExpectedColumn { + name: "created_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "updated_unix_seconds", + data_type: "INTEGER", + not_null: true, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "completed_unix_seconds", + data_type: "INTEGER", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "provider_run_id", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "translation_results_json", + data_type: "TEXT", + not_null: true, + default_value: Some("'[]'"), + primary_key: false, + }, + ExpectedColumn { + name: "provider", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "lease_owner", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "lease_expires_unix_seconds", + data_type: "INTEGER", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "failure_class", + data_type: "TEXT", + not_null: false, + default_value: None, + primary_key: false, + }, + ExpectedColumn { + name: "failure_retryable", + data_type: "INTEGER", + not_null: true, + default_value: Some("0"), + primary_key: false, + }, + ExpectedColumn { + name: "next_attempt_unix_seconds", + data_type: "INTEGER", + not_null: false, + default_value: None, + primary_key: false, + }, +]; + +async fn create_translation_task_schema( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, ) -> Result<()> { - let exists: i64 = - sqlx::query_scalar("SELECT COUNT(*) FROM pragma_table_info(?1) WHERE name = ?2") - .bind(table) - .bind(column) - .fetch_one(pool) + sqlx::query( + "CREATE TABLE schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + ) + .execute(&mut **transaction) + .await + .map_err(db_error)?; + sqlx::query( + "CREATE TABLE translation_tasks ( + task_id TEXT PRIMARY KEY NOT NULL, + official_release_id TEXT NOT NULL, + destination TEXT NOT NULL, + archive_entry TEXT, + queue_status TEXT NOT NULL, + queue_reason TEXT, + parse_status TEXT, + text_unit_formats_json TEXT NOT NULL DEFAULT '[]', + task_json TEXT NOT NULL DEFAULT '{}', + worker_status TEXT NOT NULL DEFAULT 'queued', + failure_reason TEXT, + attempt_count INTEGER NOT NULL DEFAULT 0 CHECK(attempt_count >= 0), + created_unix_seconds INTEGER NOT NULL, + updated_unix_seconds INTEGER NOT NULL, + completed_unix_seconds INTEGER, + provider_run_id TEXT, + translation_results_json TEXT NOT NULL DEFAULT '[]', + provider TEXT, + lease_owner TEXT, + lease_expires_unix_seconds INTEGER, + failure_class TEXT, + failure_retryable INTEGER NOT NULL DEFAULT 0, + next_attempt_unix_seconds INTEGER + )", + ) + .execute(&mut **transaction) + .await + .map_err(db_error)?; + Ok(()) +} + +async fn migrate_translation_tasks_v1_to_v2( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + fail_after_step: Option, +) -> Result<()> { + let steps = [ + "ALTER TABLE translation_tasks ADD COLUMN translation_results_json TEXT NOT NULL DEFAULT '[]'", + "ALTER TABLE translation_tasks ADD COLUMN provider TEXT", + "ALTER TABLE translation_tasks ADD COLUMN lease_owner TEXT", + "ALTER TABLE translation_tasks ADD COLUMN lease_expires_unix_seconds INTEGER", + "ALTER TABLE translation_tasks ADD COLUMN failure_class TEXT", + "ALTER TABLE translation_tasks ADD COLUMN failure_retryable INTEGER NOT NULL DEFAULT 0", + "ALTER TABLE translation_tasks ADD COLUMN next_attempt_unix_seconds INTEGER", + ]; + for (index, statement) in steps.iter().enumerate() { + sqlx::query(statement) + .execute(&mut **transaction) .await .map_err(db_error)?; - if exists == 0 { - let mut query = QueryBuilder::::new("ALTER TABLE "); - query - .push(table) - .push(" ADD COLUMN ") - .push(column) - .push(" "); - query.push(column_type); - query.build().execute(pool).await.map_err(db_error)?; + if fail_after_step == Some(index + 1) { + return Err(bat_core::Error::Other(anyhow::anyhow!( + "Translation Tasks migration failed after step {}", + index + 1 + ))); + } } Ok(()) } @@ -1641,6 +2009,8 @@ mod tests { OfficialTextUnitTaskStatus, OfficialTextUnitTaskSummary, OFFICIAL_TEXTUNIT_TASK_QUEUE_VERSION, }; + use sqlx::sqlite::SqlitePoolOptions; + use sqlx::Row; fn task( task_id: &str, @@ -1708,6 +2078,292 @@ mod tests { } } + async fn raw_repository( + path: &std::path::Path, + create_if_missing: bool, + ) -> SqliteTranslationTaskRepository { + let options = SqliteConnectOptions::from_str(&format!("sqlite://{}", path.display())) + .unwrap() + .create_if_missing(create_if_missing); + let pool = SqlitePoolOptions::new() + .max_connections(1) + .connect_with(options) + .await + .unwrap(); + SqliteTranslationTaskRepository { pool } + } + + async fn create_v1_database(path: &std::path::Path, with_version: bool) { + let repository = raw_repository(path, true).await; + sqlx::query( + "CREATE TABLE schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + ) + .execute(&repository.pool) + .await + .unwrap(); + sqlx::query( + "CREATE TABLE translation_tasks ( + task_id TEXT PRIMARY KEY NOT NULL, + official_release_id TEXT NOT NULL, + destination TEXT NOT NULL, + archive_entry TEXT, + queue_status TEXT NOT NULL, + queue_reason TEXT, + parse_status TEXT, + text_unit_formats_json TEXT NOT NULL DEFAULT '[]', + task_json TEXT NOT NULL DEFAULT '{}', + worker_status TEXT NOT NULL DEFAULT 'queued', + failure_reason TEXT, + attempt_count INTEGER NOT NULL DEFAULT 0 CHECK(attempt_count >= 0), + created_unix_seconds INTEGER NOT NULL, + updated_unix_seconds INTEGER NOT NULL, + completed_unix_seconds INTEGER, + provider_run_id TEXT + )", + ) + .execute(&repository.pool) + .await + .unwrap(); + if with_version { + sqlx::query("INSERT INTO schema_migrations(component, version) VALUES (?1, ?2)") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .bind(1_i64) + .execute(&repository.pool) + .await + .unwrap(); + } + sqlx::query( + "INSERT INTO translation_tasks ( + task_id, official_release_id, destination, queue_status, parse_status, + text_unit_formats_json, task_json, worker_status, attempt_count, + created_unix_seconds, updated_unix_seconds, provider_run_id + ) VALUES ('legacy-task', 'release-1', 'Bundles/legacy.bundle', 'queued', + 'parsed', '[\"plain\"]', '{}', 'failed', 3, 10, 20, 'run-legacy')", + ) + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + } + + #[tokio::test] + async fn sqlite_translation_tasks_migrates_v1_and_preserves_business_state() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_TASK_REPOSITORY_FILE); + create_v1_database(&path, true).await; + + let repository = SqliteTranslationTaskRepository::open(&path).await.unwrap(); + let row = sqlx::query( + "SELECT worker_status, attempt_count, provider_run_id, + translation_results_json, failure_retryable + FROM translation_tasks WHERE task_id = 'legacy-task'", + ) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(row.try_get::("worker_status").unwrap(), "failed"); + assert_eq!(row.try_get::("attempt_count").unwrap(), 3); + assert_eq!( + row.try_get::, _>("provider_run_id") + .unwrap() + .as_deref(), + Some("run-legacy") + ); + assert_eq!( + row.try_get::("translation_results_json") + .unwrap(), + "[]" + ); + assert_eq!(row.try_get::("failure_retryable").unwrap(), 0); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, i64::from(TRANSLATION_TASK_SCHEMA_VERSION)); + + let snapshot = sqlite_migration::snapshot_connection( + &mut repository.pool.acquire().await.unwrap(), + TRANSLATION_TASK_SCHEMA_COMPONENT, + ) + .await + .unwrap(); + assert!(matches_task_fingerprint( + &snapshot, + TRANSLATION_TASK_SCHEMA_VERSION + )); + } + + #[tokio::test] + async fn sqlite_translation_tasks_recognizes_legacy_without_version_row() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_TASK_REPOSITORY_FILE); + create_v1_database(&path, false).await; + let repository = SqliteTranslationTaskRepository::open(&path).await.unwrap(); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, 2); + } + + #[tokio::test] + async fn sqlite_translation_tasks_recognizes_legacy_without_schema_table() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_TASK_REPOSITORY_FILE); + create_v1_database(&path, false).await; + let repository = raw_repository(&path, false).await; + sqlx::query("DROP TABLE schema_migrations") + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + + let repository = SqliteTranslationTaskRepository::open(&path).await.unwrap(); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, 2); + } + + #[tokio::test] + async fn sqlite_translation_tasks_future_schema_is_read_only_failure() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_TASK_REPOSITORY_FILE); + let repository = SqliteTranslationTaskRepository::new(&path).await.unwrap(); + sqlx::query("UPDATE schema_migrations SET version = ?2 WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .bind(i64::from(TRANSLATION_TASK_SCHEMA_VERSION) + 1) + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + let before = std::fs::read(&path).unwrap(); + let wal_path = std::path::PathBuf::from(format!("{}-wal", path.display())); + let shm_path = std::path::PathBuf::from(format!("{}-shm", path.display())); + let wal_before = std::fs::read(&wal_path).ok(); + let shm_before = std::fs::read(&shm_path).ok(); + let error = SqliteTranslationTaskRepository::new(&path) + .await + .unwrap_err(); + assert!(error.to_string().contains("不支持的翻译任务 schema")); + assert_eq!(std::fs::read(&path).unwrap(), before); + assert_eq!(std::fs::read(&wal_path).ok(), wal_before); + assert_eq!(std::fs::read(&shm_path).ok(), shm_before); + } + + #[tokio::test] + async fn sqlite_translation_tasks_unknown_schema_fails_closed() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_TASK_REPOSITORY_FILE); + let repository = raw_repository(&path, true).await; + sqlx::query( + "CREATE TABLE schema_migrations ( + component TEXT PRIMARY KEY NOT NULL, + version INTEGER NOT NULL CHECK(version >= 1) + )", + ) + .execute(&repository.pool) + .await + .unwrap(); + sqlx::query("INSERT INTO schema_migrations(component, version) VALUES (?1, 1)") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .execute(&repository.pool) + .await + .unwrap(); + sqlx::query("CREATE TABLE translation_tasks (task_id TEXT PRIMARY KEY)") + .execute(&repository.pool) + .await + .unwrap(); + repository.pool.close().await; + let before = std::fs::read(&path).unwrap(); + let error = SqliteTranslationTaskRepository::open(&path) + .await + .unwrap_err(); + assert!(error + .to_string() + .contains("schema version 与实际结构不一致")); + assert_eq!(std::fs::read(&path).unwrap(), before); + } + + #[tokio::test] + async fn sqlite_translation_tasks_failed_migration_rolls_back_and_retries() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_TASK_REPOSITORY_FILE); + create_v1_database(&path, true).await; + let repository = raw_repository(&path, false).await; + assert!(repository.init_schema_with_test_failure(1).await.is_err()); + let columns: Vec = sqlx::query_scalar( + "SELECT name FROM pragma_table_info('translation_tasks') ORDER BY cid", + ) + .fetch_all(&repository.pool) + .await + .unwrap(); + assert!(!columns.iter().any(|column| column == "provider")); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, 1); + let attempt_count: i64 = sqlx::query_scalar( + "SELECT attempt_count FROM translation_tasks WHERE task_id = 'legacy-task'", + ) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(attempt_count, 3); + + repository.init_schema().await.unwrap(); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, 2); + } + + #[tokio::test] + async fn sqlite_translation_tasks_concurrent_new_open_has_one_current_schema() { + let temp = tempfile::TempDir::new().unwrap(); + let path = temp.path().join(TRANSLATION_TASK_REPOSITORY_FILE); + let (left, right) = tokio::join!( + SqliteTranslationTaskRepository::new(&path), + SqliteTranslationTaskRepository::new(&path) + ); + assert!(left.is_ok(), "left open failed: {left:?}"); + assert!(right.is_ok(), "right open failed: {right:?}"); + let repository = left.unwrap(); + drop(right); + let version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .fetch_one(&repository.pool) + .await + .unwrap(); + assert_eq!(version, 2); + repository.pool.close().await; + let reopened = SqliteTranslationTaskRepository::open(&path).await.unwrap(); + let reopened_version: i64 = + sqlx::query_scalar("SELECT version FROM schema_migrations WHERE component = ?1") + .bind(TRANSLATION_TASK_SCHEMA_COMPONENT) + .fetch_one(&reopened.pool) + .await + .unwrap(); + assert_eq!(reopened_version, 2); + } + #[tokio::test] async fn sqlite_translation_tasks_sync_and_preserve_worker_state() { let temp = tempfile::TempDir::new().unwrap();