From c60be3fccb0753710124e11e3059f636f9c227bd Mon Sep 17 00:00:00 2001 From: ghbvf <104540935+ghbvf@users.noreply.github.com> Date: Sat, 20 Jun 2026 14:09:59 +0800 Subject: [PATCH 1/2] =?UTF-8?q?feat:=20SQLite=20=E7=BB=9F=E4=B8=80?= =?UTF-8?q?=E6=8C=81=E4=B9=85=E5=8C=96=20+=20=E9=A1=B9=E7=9B=AE/PR=20?= =?UTF-8?q?=E5=90=88=E5=B9=B6=E5=AF=BC=E8=88=AA=EF=BC=88#67=20#70=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #70(后端):引入 rusqlite(bundled)作为统一本地存储。 - 新增横切 db.rs:Database 句柄(Mutex)、PRAGMA user_version 迁移 runner、 全部 DDL(config_blob / tracked_pr / dispatch_key|event / review_session|history_item)、 with_conn/with_tx;setup 内 manage 为 tauri::State。 - config/pr-registry/pr-ledger 的 load/save 由 tauri-plugin-store 改 SQLite(保留 migrate_value/validate 与 WRITE_LOCK 语义、serde golden)。 - 新增 review/history_store:会话元数据 + 历史内容持久化;pump emit 后 best-effort 落库(按 item_id 合并 delta),状态转移镜像到 review_session。 - 新命令 get_session_history / get_pr_sessions(注册于 lib.rs)。 - 一次性 legacy JSON→SQLite 导入(meta 守卫,旧文件保留可回退)。 - webhook PR 经现有 ingest 落 tracked_pr,gh 不可用时仍可存可查。 #67(前端):合并 ProjectSwitcher + PrList 为单列 ProjectNav(项目 H2 一级分组、 活动项目展开其 PR),消除两条分离边栏;会话面板保持独立。 - useReviewStore.focus 改为从持久历史水合(打开历史会话看之前内容)。 - ReviewSessions 改为按选中 PR 展示其会话(durable getPrSessions + 实时叠加)。 - api/types 增 getSessionHistory / getPrSessions。 Closes #67 Closes #70 Co-Authored-By: Claude Opus 4.8 (1M context) --- src-tauri/Cargo.lock | 94 ++++++++ src-tauri/Cargo.toml | 7 + src-tauri/src/config/service.rs | 128 ++++++++-- src-tauri/src/db.rs | 274 +++++++++++++++++++++ src-tauri/src/lib.rs | 82 +++++++ src-tauri/src/pr/ledger.rs | 327 ++++++++++++++++++-------- src-tauri/src/pr/registry.rs | 261 +++++++++++++++----- src-tauri/src/review/commands.rs | 33 ++- src-tauri/src/review/history_store.rs | 247 +++++++++++++++++++ src-tauri/src/review/mod.rs | 1 + src-tauri/src/review/session.rs | 94 +++++++- src/App.vue | 19 +- src/pr/PrList.vue | 200 ---------------- src/pr/ProjectNav.vue | 301 ++++++++++++++++++++++++ src/pr/ProjectSwitcher.vue | 124 ---------- src/review/ReviewSessions.vue | 71 ++++-- src/review/api.ts | 19 +- src/review/useReviewStore.test.ts | 40 ++++ src/review/useReviewStore.ts | 29 ++- 19 files changed, 1814 insertions(+), 537 deletions(-) create mode 100644 src-tauri/src/db.rs create mode 100644 src-tauri/src/review/history_store.rs delete mode 100644 src/pr/PrList.vue create mode 100644 src/pr/ProjectNav.vue delete mode 100644 src/pr/ProjectSwitcher.vue diff --git a/src-tauri/Cargo.lock b/src-tauri/Cargo.lock index c672ab2..695f533 100644 --- a/src-tauri/Cargo.lock +++ b/src-tauri/Cargo.lock @@ -8,6 +8,18 @@ version = "2.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "320119579fcad9c21884f5c4861d16174d0e06250625266f50fe6898340abefa" +[[package]] +name = "ahash" +version = "0.8.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" +dependencies = [ + "cfg-if", + "once_cell", + "version_check", + "zerocopy", +] + [[package]] name = "aho-corasick" version = "1.1.4" @@ -980,6 +992,18 @@ dependencies = [ "pin-project-lite", ] +[[package]] +name = "fallible-iterator" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649" + +[[package]] +name = "fallible-streaming-iterator" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a" + [[package]] name = "fastrand" version = "2.4.1" @@ -1475,6 +1499,15 @@ version = "0.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" +dependencies = [ + "ahash", +] + [[package]] name = "hashbrown" version = "0.15.5" @@ -1490,6 +1523,15 @@ version = "0.17.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a" +[[package]] +name = "hashlink" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ba4ff7128dee98c7dc9794b6a411377e1404dba1c97deb8d1a55297bd25d8af" +dependencies = [ + "hashbrown 0.14.5", +] + [[package]] name = "heck" version = "0.4.1" @@ -2009,6 +2051,17 @@ dependencies = [ "libc", ] +[[package]] +name = "libsqlite3-sys" +version = "0.30.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e99fb7a497b1e3339bc746195567ed8d3e24945ecd636e3619d20b9de9e9149" +dependencies = [ + "cc", + "pkg-config", + "vcpkg", +] + [[package]] name = "linux-raw-sys" version = "0.12.1" @@ -2647,6 +2700,7 @@ dependencies = [ "futures", "hex", "hmac", + "rusqlite", "serde", "serde_json", "sha2", @@ -2858,6 +2912,20 @@ dependencies = [ "web-sys", ] +[[package]] +name = "rusqlite" +version = "0.32.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7753b721174eb8ff87a9a0e799e2d7bc3749323e773db92e0984debb00019d6e" +dependencies = [ + "bitflags 2.13.0", + "fallible-iterator", + "fallible-streaming-iterator", + "hashlink", + "libsqlite3-sys", + "smallvec", +] + [[package]] name = "rustc-hash" version = "2.1.2" @@ -4198,6 +4266,12 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "vcpkg" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" + [[package]] name = "version-compare" version = "0.2.1" @@ -5111,6 +5185,26 @@ dependencies = [ "zvariant", ] +[[package]] +name = "zerocopy" +version = "0.8.52" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ce1022995ff5ff5d841ad7d994facc23098cd40152f2c1d11cd607c6f530653f" +dependencies = [ + "zerocopy-derive", +] + +[[package]] +name = "zerocopy-derive" +version = "0.8.52" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ae7f38b72ec2a254e2b87ef277cf2cd4fb97cbebf944faa6f33354da0867930" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.117", +] + [[package]] name = "zerofrom" version = "0.1.8" diff --git a/src-tauri/Cargo.toml b/src-tauri/Cargo.toml index 57ae522..7a447cf 100644 --- a/src-tauri/Cargo.toml +++ b/src-tauri/Cargo.toml @@ -42,6 +42,13 @@ hex = "0.4" # review would require `'static` owned handles; `join_all` over borrowing futures # keeps the shared-borrow design. Only `futures` (no extra runtime) is pulled. futures = "0.3" +# Unified local SQLite store (#70): the single persistence backend for config / +# tracked PRs / dispatch ledger / review sessions + history / webhook PRs. `bundled` +# compiles a pinned libsqlite3 from source (no system sqlite dependency), and the +# synchronous `rusqlite` API fits the existing "sync store behind a std Mutex" pattern +# (registry/ledger) — the pump's async path calls these sub-millisecond sync writes +# directly without holding the connection lock across an `.await`. +rusqlite = { version = "0.32", features = ["bundled"] } # Pin: brotli 8.0.3 (pulled transitively by tauri-codegen/tauri-utils) requires # alloc-no-stdlib "^2.0", but alloc-stdlib/brotli-decompressor accept ">=2.0.4,<4" diff --git a/src-tauri/src/config/service.rs b/src-tauri/src/config/service.rs index 8caeb71..139f822 100644 --- a/src-tauri/src/config/service.rs +++ b/src-tauri/src/config/service.rs @@ -1,13 +1,17 @@ //! Config slice logic. //! -//! Backend-owned persistence via `tauri-plugin-store`'s Rust `StoreExt`. The -//! frontend calls the `get_config` / `set_config` commands (not the store plugin -//! directly), so all reads/writes funnel through here. +//! Backend-owned persistence in the unified SQLite store (#70): the whole [`AppConfig`] +//! lives as one camelCase-JSON blob in the single-row `config_blob` table (swapping only +//! the storage backend — the [`migrate_value`] / `validate` shape logic is unchanged). +//! The frontend calls the `get_config` / `set_config` commands (not SQLite directly), so +//! all reads/writes funnel through here. +use rusqlite::OptionalExtension; use serde_json::{json, Map, Value}; -use tauri_plugin_store::StoreExt; +use tauri::Manager; use super::model::AppConfig; +use crate::db::{map_err, Database}; use crate::error::{AppError, AppResult}; /// Re-export the project domain type THROUGH the config public service surface (#35, @@ -18,10 +22,6 @@ use crate::error::{AppError, AppResult}; /// use `Project` through this same re-export. pub use super::model::Project; -/// Store file holding the persisted config. -const STORE_FILE: &str = "config.json"; -/// Key under which the [`AppConfig`] value lives in the store. -const CONFIG_KEY: &str = "appConfig"; /// `id`/`name` assigned to the single project lifted out of a legacy flat config by /// [`migrate_value`] (#35). One source so the migration and its tests agree on the /// id the active-project pointer (`activeProjectId`) is also set to. @@ -120,16 +120,30 @@ fn migrate_value(raw: Value) -> Value { /// [`migrate_value`] first so a legacy flat single-project config (#35) upgrades to /// the multi-project shape before deserialization. pub fn load(app: &tauri::AppHandle) -> AppResult { - let store = app - .store(STORE_FILE) - .map_err(|e| AppError::new(format!("打开配置存储失败: {e}")))?; + load_db(app.state::().inner()) +} - match store.get(CONFIG_KEY) { +/// SQLite-level load (no Tauri app) — reads the `config_blob` row, runs [`migrate_value`] +/// on it, then deserializes. Split from [`load`] so the blob path + legacy migration are +/// testable against an in-memory [`Database`]. +pub(crate) fn load_db(db: &Database) -> AppResult { + let raw: Option = db.with_conn(|conn| { + conn.query_row("SELECT json FROM config_blob WHERE id = 1", [], |r| { + r.get::<_, String>(0) + }) + .optional() + })?; + + match raw { None => Ok(AppConfig::default()), - // Surface (don't silently discard) a corrupt/incompatible persisted - // config so the user can fix it rather than lose their settings. - Some(value) => serde_json::from_value(migrate_value(value)) - .map_err(|e| AppError::new(format!("解析持久化配置失败: {e}"))), + // Surface (don't silently discard) a corrupt/incompatible persisted config so the + // user can fix it rather than lose their settings. + Some(json) => { + let value: Value = serde_json::from_str(&json) + .map_err(|e| AppError::new(format!("解析持久化配置失败: {e}")))?; + serde_json::from_value(migrate_value(value)) + .map_err(|e| AppError::new(format!("解析持久化配置失败: {e}"))) + } } } @@ -201,16 +215,35 @@ pub fn save(app: &tauri::AppHandle, config: AppConfig) -> /// and the lenient [`set_active_project`] both funnel through here so the /// store-write plumbing lives in one place. fn persist(app: &tauri::AppHandle, config: &AppConfig) -> AppResult<()> { - let store = app - .store(STORE_FILE) - .map_err(|e| AppError::new(format!("打开配置存储失败: {e}")))?; - - let value = serde_json::to_value(config).map_err(|e| AppError::new(e.to_string()))?; - // tauri-plugin-store 2.x: `Store::set` is infallible and returns `()`. - store.set(CONFIG_KEY, value); - store - .save() - .map_err(|e| AppError::new(format!("写入配置存储失败: {e}")))?; + persist_db(app.state::().inner(), config) +} + +/// SQLite-level write of the whole [`AppConfig`] as the single `config_blob` row (#70). +/// Split from [`persist`] so the blob round-trip is testable against an in-memory +/// [`Database`]. Stores the camelCase JSON; [`load_db`] runs `migrate_value` on read +/// (identity for an already-new shape). +pub(crate) fn persist_db(db: &Database, config: &AppConfig) -> AppResult<()> { + let json = serde_json::to_string(config).map_err(|e| AppError::new(e.to_string()))?; + db.with_conn(|conn| { + conn.execute( + "INSERT OR REPLACE INTO config_blob (id, json) VALUES (1, ?1)", + [&json], + ) + .map(|_| ()) + }) +} + +/// One-time legacy import (#70) of the old `config.json` `appConfig` value into the +/// `config_blob` row. Stores the RAW legacy value as-is (which may be the pre-#35 flat +/// shape) — [`load_db`]'s [`migrate_value`] lifts it on the next read, so this needs no +/// shape knowledge. Runs inside the composition root's import transaction. +pub fn import_legacy_config(tx: &rusqlite::Transaction, value: &Value) -> AppResult<()> { + let json = serde_json::to_string(value).map_err(|e| AppError::new(e.to_string()))?; + tx.execute( + "INSERT OR REPLACE INTO config_blob (id, json) VALUES (1, ?1)", + [&json], + ) + .map_err(map_err)?; Ok(()) } @@ -246,6 +279,49 @@ pub fn set_active_project(app: &tauri::AppHandle, id: &str mod tests { use super::*; + // SQLite blob round-trip (#70): an empty DB loads the default; a persisted config + // reads back equal. The blob path is the storage swap — `migrate_value`/`validate` + // (tested below) are unchanged. + #[test] + fn sqlite_blob_empty_loads_default_and_round_trips() { + let db = Database::open_in_memory().expect("open db"); + // No row yet → default (first launch, empty projects). + let first = load_db(&db).expect("load empty"); + assert!(first.projects.is_empty()); + assert_eq!(first.active_project_id, ""); + + // Persist a non-default value (a webhook port) and read it back through the blob. + let config = AppConfig { + webhook_port: 9123, + ..Default::default() + }; + persist_db(&db, &config).expect("persist"); + let back = load_db(&db).expect("load"); + assert_eq!(back.webhook_port, 9123); + } + + // One-time legacy import (#70) of the highest-risk case: a pre-#35 FLAT `config.json` + // value imported into `config_blob` must, on the next `load_db`, surface as the + // migrated multi-project shape (the import stores it raw; `migrate_value` lifts it on + // read). This is the migration existing users depend on to keep their settings. + #[test] + fn legacy_flat_config_import_migrates_on_load() { + let db = Database::open_in_memory().expect("open db"); + let legacy = json!({ + "repo": "octocat/hello", + "repoRoot": "/tmp/hello", + "autoReview": true + }); + db.with_tx(|tx| import_legacy_config(tx, &legacy)) + .expect("import"); + + let config = load_db(&db).expect("load"); + assert_eq!(config.projects.len(), 1, "flat shape lifted to one project"); + assert_eq!(config.projects[0].repo, "octocat/hello"); + assert_eq!(config.active_project_id, MIGRATED_PROJECT_ID); + assert!(config.projects[0].auto_review); + } + #[test] fn migrate_old_flat_produces_single_default_project() { let raw = json!({ diff --git a/src-tauri/src/db.rs b/src-tauri/src/db.rs new file mode 100644 index 0000000..490cb4c --- /dev/null +++ b/src-tauri/src/db.rs @@ -0,0 +1,274 @@ +//! Unified SQLite persistence backend (#70). +//! +//! Horizontal infra (like [`crate::error`] / [`crate::state`]): owns the single +//! SQLite connection, the schema migration runner (`PRAGMA user_version`), and the +//! generic [`Database::with_conn`] / [`Database::with_tx`] accessors. It is the +//! schema's composition root — **ALL DDL lives here**, in ordered migration steps; +//! slices never `CREATE TABLE`. Each slice owns its own tables' *queries* in its own +//! `*_store` module, reaching the connection through the [`Database`] handle (a +//! horizontal dependency, exactly like [`crate::error::AppResult`] — NOT a cross-slice +//! import). Adding a slice table means editing this file's migration steps: the single +//! intended choke point for schema evolution, mirroring how [`crate::lib`]'s +//! `generate_handler!` is the command-registration choke point. +//! +//! **Why a `tauri::State`, not an [`crate::state::AppState`] field:** `app_data_dir()` +//! only resolves inside `setup`, while `AppState` is `.manage()`d at builder time +//! (keeping `AppState: Default`). So the composition root opens the DB in `setup` and +//! `app.manage(Database::open(..)?)`. Slices/commands already carry `app: &AppHandle`, +//! so they reach it via `app.state::()`, mirroring `app.state::()`. +//! A command running before that manage would panic on `app.state::()` +//! (fail-fast) — but `setup` completes before any command is served, so it never does. + +use std::sync::Mutex; + +use rusqlite::{Connection, OptionalExtension, Transaction}; +use tauri::Manager; + +use crate::error::{AppError, AppResult}; + +/// Current schema version. Bump + add an `apply_vN` step for every schema change; the +/// migration runner replays only the steps newer than the DB's `user_version`. +const SCHEMA_VERSION: i64 = 1; + +/// `meta` guard key marking the one-time legacy JSON → SQLite import done (#70). Kept +/// SEPARATE from `user_version` so the import runs exactly once even across future +/// schema bumps (a schema migration must not re-trigger the data import). +const META_LEGACY_IMPORTED: &str = "legacyImported"; + +/// The single SQLite connection behind a `std::sync::Mutex` — the same "sync store +/// behind a Mutex" shape the JSON stores used, so all the existing write-lock reasoning +/// (registry/ledger) carries over unchanged. `Connection: Send` ⇒ `Mutex: +/// Send + Sync`, so this is a valid `tauri::State`. +pub struct Database { + conn: Mutex, +} + +impl Database { + /// Opens (creating if absent) the app's `prmonitor.db` under `app_data_dir`, sets + /// pragmas, and runs schema migrations. The one-time legacy JSON import is driven + /// SEPARATELY by the composition root (it must read the old `tauri-plugin-store` + /// files via `app`), gated by [`Database::legacy_imported`]. + pub fn open(app: &tauri::AppHandle) -> AppResult { + let dir = app + .path() + .app_data_dir() + .map_err(|e| AppError::new(format!("解析应用数据目录失败: {e}")))?; + std::fs::create_dir_all(&dir) + .map_err(|e| AppError::new(format!("创建应用数据目录失败: {e}")))?; + let path = dir.join("prmonitor.db"); + let conn = + Connection::open(&path).map_err(|e| AppError::new(format!("打开 SQLite 失败: {e}")))?; + Self::from_conn(conn) + } + + /// In-memory database — schema migrated, no legacy import. Used by store round-trip + /// tests across slices (each slice's `*_store` tests open one of these directly). + pub fn open_in_memory() -> AppResult { + let conn = Connection::open_in_memory() + .map_err(|e| AppError::new(format!("打开内存 SQLite 失败: {e}")))?; + Self::from_conn(conn) + } + + fn from_conn(conn: Connection) -> AppResult { + // `execute_batch` (not `pragma_update`) for journal_mode: it returns the + // resulting mode as a row, which `pragma_update`'s `execute` would reject. WAL + // silently stays "memory" for in-memory DBs — harmless. + conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;") + .map_err(map_err)?; + run_migrations(&conn)?; + Ok(Self { + conn: Mutex::new(conn), + }) + } + + /// Runs `f` with a shared `&Connection` (a read or a single-statement write). + /// Serializes all DB access process-wide via the connection mutex. The closure is + /// synchronous; callers in async paths (the pump) never hold the guard across an + /// `.await` (the closure body has none). + pub fn with_conn(&self, f: impl FnOnce(&Connection) -> rusqlite::Result) -> AppResult { + let conn = self.conn.lock().expect("db mutex poisoned"); + f(&conn).map_err(map_err) + } + + /// Runs `f` inside a transaction (atomic multi-statement read-modify-write), + /// committing on `Ok` and rolling back on `Err`. The closure returns [`AppResult`] + /// so it can interleave slice query helpers (which already map to [`AppError`]). + pub fn with_tx(&self, f: impl FnOnce(&Transaction) -> AppResult) -> AppResult { + let mut conn = self.conn.lock().expect("db mutex poisoned"); + let tx = conn.transaction().map_err(map_err)?; + let out = f(&tx)?; + tx.commit().map_err(map_err)?; + Ok(out) + } + + /// Whether the one-time legacy JSON import has already run (#70 guard). The + /// composition root checks this before reading the old JSON stores. + pub fn legacy_imported(&self) -> AppResult { + self.with_conn(|conn| meta_get(conn, META_LEGACY_IMPORTED).map(|v| v.is_some())) + } +} + +/// Maps a rusqlite error into the app error funnel ([`AppError`]). +pub(crate) fn map_err(e: rusqlite::Error) -> AppError { + AppError::new(format!("数据库错误: {e}")) +} + +/// Reads a `meta` value by key (rusqlite-level so it composes inside `with_conn`/`with_tx`). +fn meta_get(conn: &Connection, key: &str) -> rusqlite::Result> { + conn.query_row("SELECT value FROM meta WHERE key = ?1", [key], |row| { + row.get::<_, String>(0) + }) + .optional() +} + +/// Marks the one-time legacy import done (#70). MUST be called inside the SAME +/// transaction as the imported-row inserts so a crash mid-import rolls back the guard +/// too and re-runs cleanly. +pub fn mark_legacy_imported(tx: &Transaction) -> AppResult<()> { + tx.execute( + "INSERT OR REPLACE INTO meta (key, value) VALUES (?1, '1')", + [META_LEGACY_IMPORTED], + ) + .map_err(map_err)?; + Ok(()) +} + +/// Replays schema steps newer than the DB's `user_version`, then stamps the current +/// version. Each future schema change = a new `apply_vN` + a `< N` gate here. +fn run_migrations(conn: &Connection) -> AppResult<()> { + let version: i64 = conn + .pragma_query_value(None, "user_version", |row| row.get(0)) + .map_err(map_err)?; + if version < 1 { + apply_v1(conn)?; + } + conn.pragma_update(None, "user_version", SCHEMA_VERSION) + .map_err(map_err)?; + Ok(()) +} + +fn apply_v1(conn: &Connection) -> AppResult<()> { + conn.execute_batch(SCHEMA_V1).map_err(map_err)?; + Ok(()) +} + +/// v1 schema — the unified store (#70). `review_session` precedes `review_history_item` +/// (the FK target must exist first under `foreign_keys=ON`). Per-project partitioning +/// is a real `project_id TEXT` column (replacing the JSON stores' `prefix:{pid}` keys). +const SCHEMA_V1: &str = r#" +CREATE TABLE IF NOT EXISTS meta ( + key TEXT PRIMARY KEY, + value TEXT +); + +CREATE TABLE IF NOT EXISTS config_blob ( + id INTEGER PRIMARY KEY CHECK (id = 1), + json TEXT NOT NULL +); + +CREATE TABLE IF NOT EXISTS tracked_pr ( + project_id TEXT NOT NULL, + number INTEGER NOT NULL, + title TEXT NOT NULL, + labels_json TEXT NOT NULL, + url TEXT NOT NULL, + kind TEXT NOT NULL, + skip_reason TEXT, + first_seen_epoch INTEGER NOT NULL, + last_seen_epoch INTEGER NOT NULL, + archived INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (project_id, number) +); +CREATE INDEX IF NOT EXISTS idx_tracked_pr_project ON tracked_pr(project_id, number DESC); + +CREATE TABLE IF NOT EXISTS dispatch_key ( + project_id TEXT NOT NULL, + key TEXT NOT NULL, + PRIMARY KEY (project_id, key) +); + +CREATE TABLE IF NOT EXISTS dispatch_event ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + project_id TEXT NOT NULL, + pr INTEGER NOT NULL, + kind TEXT NOT NULL, + head_sha TEXT NOT NULL, + key TEXT NOT NULL, + dispatched_at_epoch INTEGER NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_dispatch_event_lookup ON dispatch_event(project_id, pr, kind); + +CREATE TABLE IF NOT EXISTS review_session ( + thread_id TEXT PRIMARY KEY, + project_id TEXT NOT NULL, + pr_number INTEGER NOT NULL, + turn_id TEXT NOT NULL DEFAULT '', + kind TEXT NOT NULL, + status TEXT NOT NULL, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_review_session_pr ON review_session(project_id, pr_number, created_at); + +CREATE TABLE IF NOT EXISTS review_history_item ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + thread_id TEXT NOT NULL, + item_id TEXT NOT NULL, + kind TEXT NOT NULL, + text TEXT NOT NULL, + UNIQUE (thread_id, item_id), + FOREIGN KEY (thread_id) REFERENCES review_session(thread_id) +); +CREATE INDEX IF NOT EXISTS idx_history_thread ON review_history_item(thread_id, id); +"#; + +#[cfg(test)] +mod tests { + use super::*; + + /// v1 migration lock (Medium per ai-robust.md): a missing/renamed table here means + /// a slice store's first query fails at runtime, not at compile time — so pin the + /// table set + the stamped `user_version`. + #[test] + fn migrations_create_all_tables_and_stamp_version() { + let db = Database::open_in_memory().expect("open"); + db.with_conn(|conn| { + let version: i64 = conn.pragma_query_value(None, "user_version", |r| r.get(0))?; + assert_eq!(version, SCHEMA_VERSION, "user_version stamped"); + + let mut stmt = + conn.prepare("SELECT name FROM sqlite_master WHERE type='table' ORDER BY name")?; + let names: Vec = stmt + .query_map([], |r| r.get::<_, String>(0))? + .collect::>()?; + for expected in [ + "config_blob", + "dispatch_event", + "dispatch_key", + "meta", + "review_history_item", + "review_session", + "tracked_pr", + ] { + assert!(names.contains(&expected.to_string()), "missing {expected}"); + } + Ok(()) + }) + .expect("query"); + } + + /// The legacy-import guard flips exactly once and is observable through the public + /// `legacy_imported()` the composition root gates on. + #[test] + fn legacy_import_guard_flips_once() { + let db = Database::open_in_memory().expect("open"); + assert!(!db.legacy_imported().expect("read guard"), "starts unset"); + + db.with_tx(mark_legacy_imported).expect("mark"); + assert!(db.legacy_imported().expect("read guard"), "set after mark"); + + // Idempotent: marking again (INSERT OR REPLACE) keeps it true with no dup row. + db.with_tx(mark_legacy_imported).expect("re-mark"); + assert!(db.legacy_imported().expect("read guard")); + } +} diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index 1f58e98..043a176 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -19,6 +19,7 @@ // `pr::source::PrSource`, `review::engine::ReviewEngine`) count as reachable API // in this skeleton rather than tripping `dead_code` before their first use. pub mod config; +pub mod db; pub mod dispatch; pub mod error; pub mod events; @@ -40,6 +41,15 @@ pub fn run() { .plugin(tauri_plugin_store::Builder::new().build()) .manage(AppState::default()) .setup(|app| { + // Open + migrate the unified SQLite store and manage it as a `tauri::State` + // BEFORE anything that reads persistence (config load / poll start). It is a + // `State` rather than an `AppState` field because `app_data_dir()` only + // resolves here in `setup`, while `AppState` is `.manage()`d at builder time. + // Then run the one-time legacy JSON → SQLite import (#70) so existing users' + // config / tracked PRs / ledger carry over before the first read. + app.manage(db::Database::open(app.handle())?); + import_legacy_stores(app.handle())?; + let state = app.state::(); // Install the auto-trigger dispatcher BEFORE starting the loop, so the // immediate first tick already auto-starts dispatchable reviews. The @@ -105,6 +115,8 @@ pub fn run() { review::commands::start_review, review::commands::stop_review, review::commands::list_review_sessions, + review::commands::get_session_history, + review::commands::get_pr_sessions, config::commands::set_active_project, ]) .build(tauri::generate_context!()) @@ -123,6 +135,76 @@ pub fn run() { }); } +/// One-time legacy JSON → SQLite import (#70). Reads the pre-SQLite `tauri-plugin-store` +/// files (`config.json` / `prs.json` / `ledger.json`) and hands each value to the OWNING +/// slice's `import_legacy_*` helper (column knowledge stays in the slice), inserting all +/// rows + the done-guard in ONE transaction so a crash mid-import rolls back and re-runs +/// cleanly. A no-op once [`db::Database::legacy_imported`] is set, and on fresh installs +/// (the legacy stores are empty, so nothing imports). The old JSON files are LEFT in +/// place (recoverable / downgradeable); the guard makes them inert. +/// +/// This is composition (it spans `config` + `pr` slices), so it lives at the root, not +/// in `db` (which stays a pure horizontal owning only schema + connection). +fn import_legacy_stores(app: &tauri::AppHandle) -> error::AppResult<()> { + use tauri_plugin_store::StoreExt; + + let db = app.state::(); + if db.legacy_imported()? { + return Ok(()); + } + + // Gather the legacy values up front (reads, outside the write transaction). A + // missing store file just yields an empty store → nothing to import. + let config_value = app + .store("config.json") + .ok() + .and_then(|s| s.get("appConfig")); + + let mut tracked: Vec<(String, serde_json::Value)> = Vec::new(); + if let Ok(store) = app.store("prs.json") { + for key in store.keys() { + if let Some(pid) = key.strip_prefix("tracked:") { + if let Some(v) = store.get(&key) { + tracked.push((pid.to_string(), v)); + } + } + } + } + + let mut dispatched: Vec<(String, serde_json::Value)> = Vec::new(); + let mut events: Vec<(String, serde_json::Value)> = Vec::new(); + if let Ok(store) = app.store("ledger.json") { + for key in store.keys() { + if let Some(pid) = key.strip_prefix("dispatched:") { + if let Some(v) = store.get(&key) { + dispatched.push((pid.to_string(), v)); + } + } else if let Some(pid) = key.strip_prefix("events:") { + if let Some(v) = store.get(&key) { + events.push((pid.to_string(), v)); + } + } + } + } + + db.with_tx(|tx| { + if let Some(v) = &config_value { + config::service::import_legacy_config(tx, v)?; + } + for (pid, v) in &tracked { + pr::registry::import_legacy_tracked(tx, pid, v)?; + } + for (pid, v) in &dispatched { + pr::ledger::import_legacy_dispatched(tx, pid, v)?; + } + for (pid, v) in &events { + pr::ledger::import_legacy_events(tx, pid, v)?; + } + db::mark_legacy_imported(tx)?; + Ok(()) + }) +} + /// Build the per-cycle [`pr::scheduler::ProjectDispatcher`] both auto-trigger sources /// share — the poll scheduler and the webhook ingestor. Both drive a dispatchable /// `(project_id, candidates)` through the SAME [`run_auto_dispatch`] (the composition diff --git a/src-tauri/src/pr/ledger.rs b/src-tauri/src/pr/ledger.rs index 6bb9af5..b77e933 100644 --- a/src-tauri/src/pr/ledger.rs +++ b/src-tauri/src/pr/ledger.rs @@ -1,61 +1,41 @@ //! Dispatch de-duplication ledger + cooldown source. //! //! Port of `router.py`'s two state files (`dispatched` keys + `dispatch-events` -//! epochs), unified into one `ledger.json` persisted via `tauri-plugin-store`'s -//! `StoreExt` — the same backend-owned store pattern as the config slice. +//! epochs), persisted in the unified SQLite store (#70) — the `dispatch_key` (dedup +//! set) and `dispatch_event` (cooldown log) tables, each partitioned by a real +//! `project_id` column (replacing the old `ledger.json` `prefix:{pid}` store keys). //! //! PR3 uses the **read** path (`has_dispatched` / `last_dispatch_at`) to annotate //! the PR list with "already dispatched" / cooldown skip reasons. The **write** //! path (`record_many`) is invoked by the auto-trigger dispatcher //! ([`crate::dispatch`]) once review turns actually start; recording it here keeps -//! the dedup machinery complete. The write is batched (one persist for the whole +//! the dedup machinery complete. The write is batched (one transaction for the whole //! cycle's started candidates) so unbounded concurrent starts can't race the store. use std::collections::HashSet; use std::sync::Mutex; use serde::{Deserialize, Serialize}; -use tauri_plugin_store::StoreExt; +use tauri::Manager; -use crate::error::{AppError, AppResult}; +use crate::db::{map_err, Database}; +use crate::error::AppResult; use crate::model::Candidate; -/// Store file holding the persisted ledger. -const STORE_FILE: &str = "ledger.json"; -/// Key PREFIX holding the per-project set of dispatched dedup keys (#35). The -/// effective key is `dispatched:{project_id}` (see [`dispatched_key`]); a single -/// `ledger.json` holds every project's partition under its own key. -const DISPATCHED_KEY_PREFIX: &str = "dispatched"; -/// Key PREFIX holding the per-project dispatch-event log for cooldown (#35). The -/// effective key is `events:{project_id}` (see [`events_key`]). -const EVENTS_KEY_PREFIX: &str = "events"; - -/// Serializes EVERY load→stage→save of `ledger.json` across projects (#35). The -/// store is one file holding all projects' partitions (`dispatched:{pid}` / -/// `events:{pid}`); `Store::save` rewrites the WHOLE file, so two parallel project -/// cycles each doing a load→stage→save would interleave and one would clobber the -/// other's just-written partition (a lost dispatch record → re-review storm). A -/// process-global `Mutex<()>` (the data lives in the store, not behind the lock) -/// guards the critical section in [`Ledger::record_many`]; a module static so the -/// lock IDENTITY is fixed (a caller cannot serialize on the wrong mutex). Mirrors -/// the registry's `WRITE_LOCK` rationale. `std` (not `tokio`) `Mutex`: the guarded -/// section is fully synchronous (`tauri-plugin-store` reads/writes are sync), so no -/// `.await` is ever held across the guard. +/// Serializes EVERY load→stage→save of the dispatch ledger across projects (#35). Two +/// parallel project cycles each doing a load→stage→save would interleave and one would +/// clobber the other's just-recorded partition (a lost dispatch record → re-review +/// storm). A process-global `Mutex<()>` (the data lives in SQLite, not behind the lock) +/// guards the critical section in [`Ledger::record_many`]; a module static so the lock +/// IDENTITY is fixed (a caller cannot serialize on the wrong mutex). Mirrors the +/// registry's `WRITE_LOCK` rationale. `std` (not `tokio`) `Mutex`: the guarded section +/// is fully synchronous (the SQLite calls never `.await`), so no `.await` is held across +/// the guard. The single SQLite connection's own mutex additionally serializes +/// individual statements, but THIS lock is what makes a project's load→save one atomic +/// critical section relative to other projects' cycles (so a `load` here can't read a +/// partition another project is mid-rewrite of). static LEDGER_WRITE_LOCK: Mutex<()> = Mutex::new(()); -/// Store key for a project's dispatched dedup-key set: `dispatched:{project_id}` -/// (#35). Partitions the shared `ledger.json` so two projects' identical -/// `(number, head_sha, kind)` dedup keys never collide. -fn dispatched_key(project_id: &str) -> String { - format!("{DISPATCHED_KEY_PREFIX}:{project_id}") -} - -/// Store key for a project's dispatch-event (cooldown) log: `events:{project_id}` -/// (#35). Same partitioning rationale as [`dispatched_key`]. -fn events_key(project_id: &str) -> String { - format!("{EVENTS_KEY_PREFIX}:{project_id}") -} - /// One recorded dispatch — the cooldown source (mirrors `router.py` /// dispatch-events: `(pr, kind, dispatchedAtEpoch)`). #[derive(Debug, Clone, Serialize, Deserialize)] @@ -107,13 +87,13 @@ pub fn record_dispatched( cands: &[Candidate], ) -> AppResult<()> { // Hold the cross-project write lock across the WHOLE load→stage→save (#35): N - // parallel project cycles each rewrite the same `ledger.json` (whole-file save), - // so a load here racing another project's save would drop that project's + // parallel project cycles each replace their `dispatch_*` partition (delete + + // re-insert), so a load here racing another project's save would drop that project's // just-recorded partition. The guard makes load + persist one atomic section. // `.unwrap()` matches the registry's std-Mutex convention; the section is - // synchronous (store reads/writes are sync) so no `.await` is held across it, and - // a panic mid-section can't leave torn state (the data lives in the store, each - // key rewritten wholesale by `record_many`). Poisoning is therefore benign. + // synchronous (the SQLite calls never `.await`) so no `.await` is held across it, and + // a panic mid-section can't leave torn state (the data lives in SQLite, the partition + // rewritten wholesale by `record_many` in one transaction). Poisoning is benign. let _guard = LEDGER_WRITE_LOCK.lock().unwrap(); let mut ledger = Ledger::load(app, project_id)?; ledger.record_many(app, project_id, cands, now_epoch()) @@ -149,20 +129,42 @@ impl Ledger { /// the reservation then rejects. The write path ([`record_dispatched`]) DOES hold /// the lock across its own load→stage→save (a lost write there is unrecoverable). pub fn load(app: &tauri::AppHandle, project_id: &str) -> AppResult { - let store = app - .store(STORE_FILE) - .map_err(|e| AppError::new(format!("打开 ledger 存储失败: {e}")))?; - - let dispatched = store - .get(dispatched_key(project_id)) - .and_then(|v| serde_json::from_value::>(v).ok()) - .unwrap_or_default(); - let events = store - .get(events_key(project_id)) - .and_then(|v| serde_json::from_value::>(v).ok()) - .unwrap_or_default(); - - Ok(Self { dispatched, events }) + Self::load_db(app.state::().inner(), project_id) + } + + /// SQLite-level load (no Tauri app) — reads this project's `dispatch_key` set + + /// `dispatch_event` log. Split from [`Self::load`] so store round-trips are + /// testable against an in-memory [`Database`]. Epochs/PR numbers are stored as + /// `i64` (SQLite's only integer type) and read back as `u64`. + pub(crate) fn load_db(db: &Database, project_id: &str) -> AppResult { + db.with_conn(|conn| { + let mut dispatched = HashSet::new(); + let mut stmt = conn.prepare("SELECT key FROM dispatch_key WHERE project_id = ?1")?; + let rows = stmt.query_map([project_id], |r| r.get::<_, String>(0))?; + for k in rows { + dispatched.insert(k?); + } + + let mut events = Vec::new(); + let mut stmt = conn.prepare( + "SELECT pr, kind, head_sha, key, dispatched_at_epoch \ + FROM dispatch_event WHERE project_id = ?1 ORDER BY id", + )?; + let rows = stmt.query_map([project_id], |r| { + Ok(DispatchEvent { + pr: r.get::<_, i64>(0)? as u64, + kind: r.get(1)?, + head_sha: r.get(2)?, + key: r.get(3)?, + dispatched_at_epoch: r.get::<_, i64>(4)? as u64, + }) + })?; + for e in rows { + events.push(e?); + } + + Ok(Self { dispatched, events }) + }) } /// Whether `key` has already been dispatched. @@ -192,12 +194,13 @@ impl Ledger { /// `record` calls would produce. An empty `cands` slice still touches the store /// (a harmless no-op save) — callers gate on non-empty before calling. /// - /// **Concurrency (#35):** writes ONLY this project's keys - /// (`dispatched:{project_id}` / `events:{project_id}`), but `Store::save` rewrites - /// the whole `ledger.json`. The cross-project lost-update race that creates is - /// closed by [`record_dispatched`], which holds [`LEDGER_WRITE_LOCK`] across its - /// `load` → this `record_many`, so the load this method's `self` came from and the - /// save below are one atomic critical section relative to other projects' cycles. + /// **Concurrency (#35):** writes ONLY this project's `dispatch_key` / + /// `dispatch_event` rows, but does so by replacing the whole partition (delete + + /// re-insert `self`). The cross-project lost-update race that the read→mutate→write + /// shape creates is closed by [`record_dispatched`], which holds [`LEDGER_WRITE_LOCK`] + /// across its `load` → this `record_many`, so the load this method's `self` came from + /// and the save below are one atomic critical section relative to other projects' + /// cycles. The delete+insert runs in one transaction (atomic on its own too). pub fn record_many( &mut self, app: &tauri::AppHandle, @@ -206,22 +209,56 @@ impl Ledger { epoch: u64, ) -> AppResult<()> { self.stage_all(cands, epoch); + self.save_db(app.state::().inner(), project_id) + } - let store = app - .store(STORE_FILE) - .map_err(|e| AppError::new(format!("打开 ledger 存储失败: {e}")))?; - store.set( - dispatched_key(project_id), - serde_json::to_value(&self.dispatched).map_err(|e| AppError::new(e.to_string()))?, - ); - store.set( - events_key(project_id), - serde_json::to_value(&self.events).map_err(|e| AppError::new(e.to_string()))?, - ); - store - .save() - .map_err(|e| AppError::new(format!("写入 ledger 存储失败: {e}")))?; - Ok(()) + /// SQLite-level save (no Tauri app) — replaces this project's partition with the + /// full in-memory `self` (delete-all + insert-all, the SQLite analogue of the old + /// whole-partition `Store::set`). Split from [`Self::record_many`] so the round-trip + /// is testable against an in-memory [`Database`]. + pub(crate) fn save_db(&self, db: &Database, project_id: &str) -> AppResult<()> { + db.with_tx(|tx| { + tx.execute( + "DELETE FROM dispatch_key WHERE project_id = ?1", + [project_id], + ) + .map_err(map_err)?; + tx.execute( + "DELETE FROM dispatch_event WHERE project_id = ?1", + [project_id], + ) + .map_err(map_err)?; + { + let mut stmt = tx + .prepare("INSERT INTO dispatch_key (project_id, key) VALUES (?1, ?2)") + .map_err(map_err)?; + for k in &self.dispatched { + stmt.execute(rusqlite::params![project_id, k]) + .map_err(map_err)?; + } + } + { + let mut stmt = tx + .prepare( + "INSERT INTO dispatch_event \ + (project_id, pr, kind, head_sha, key, dispatched_at_epoch) \ + VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + ) + .map_err(map_err)?; + for e in &self.events { + stmt.execute(rusqlite::params![ + project_id, + e.pr as i64, + e.kind, + e.head_sha, + e.key, + e.dispatched_at_epoch as i64 + ]) + .map_err(map_err)?; + } + } + Ok(()) + }) } /// Stages a batch into the in-memory ledger (the dedup key set + cooldown event @@ -243,6 +280,56 @@ impl Ledger { } } +/// One-time legacy import (#70) of a project's `dispatched:{pid}` dedup-key set from the +/// old `ledger.json`. Parses the JSON array of keys and inserts `dispatch_key` rows. +/// Lenient: a corrupt value imports nothing (parity with `load`'s `unwrap_or_default`). +/// Runs inside the composition root's import transaction (see `lib::import_legacy_stores`). +pub fn import_legacy_dispatched( + tx: &rusqlite::Transaction, + project_id: &str, + value: &serde_json::Value, +) -> AppResult<()> { + let keys: HashSet = serde_json::from_value(value.clone()).unwrap_or_default(); + let mut stmt = tx + .prepare("INSERT OR IGNORE INTO dispatch_key (project_id, key) VALUES (?1, ?2)") + .map_err(map_err)?; + for k in &keys { + stmt.execute(rusqlite::params![project_id, k]) + .map_err(map_err)?; + } + Ok(()) +} + +/// One-time legacy import (#70) of a project's `events:{pid}` cooldown log from the old +/// `ledger.json`. Parses the JSON array of [`DispatchEvent`] and inserts `dispatch_event` +/// rows (order preserved by insert order → `id`). Lenient like [`import_legacy_dispatched`]. +pub fn import_legacy_events( + tx: &rusqlite::Transaction, + project_id: &str, + value: &serde_json::Value, +) -> AppResult<()> { + let events: Vec = serde_json::from_value(value.clone()).unwrap_or_default(); + let mut stmt = tx + .prepare( + "INSERT INTO dispatch_event \ + (project_id, pr, kind, head_sha, key, dispatched_at_epoch) \ + VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + ) + .map_err(map_err)?; + for e in &events { + stmt.execute(rusqlite::params![ + project_id, + e.pr as i64, + e.kind, + e.head_sha, + e.key, + e.dispatched_at_epoch as i64 + ]) + .map_err(map_err)?; + } + Ok(()) +} + #[cfg(test)] mod tests { use super::*; @@ -322,21 +409,75 @@ mod tests { assert_eq!(dispatch_key(7, "deadbeef", "check"), "7@deadbeef:check"); } - // Project store-key partitioning (#35). The persisted `ledger.json` holds every - // project's dedup set / cooldown log under a project-scoped store key; the dedup - // KEY format (`{number}@{head_sha}:{kind}`) is unchanged. This pins that two - // projects' keys differ so a same-(number, head, kind) dispatch in one project - // can't be read as already-dispatched in another (the `Store::set` slot is - // distinct), and that the suffix is the raw project id. + // SQLite store round-trip (#70, Medium carrier): staging a batch, `save_db` then + // `load_db` against an in-memory DB must round-trip the dedup set + cooldown log + // intact (a column/SQL drift surfaces here), and a different project's partition + // must read empty (the `project_id` column is the partitioning seam that replaced + // the old `tracked:{pid}` store keys). + #[test] + fn sqlite_round_trip_and_per_project_isolation() { + let db = Database::open_in_memory().expect("open db"); + let mut ledger = Ledger::default(); + ledger.stage_all(&[cand(12, "review"), cand(12, "check")], 1_700_000_000); + ledger.save_db(&db, "alpha").expect("save"); + + let back = Ledger::load_db(&db, "alpha").expect("load"); + assert!(back.has_dispatched(&dispatch_key(12, "sha", "review"))); + assert!(back.has_dispatched(&dispatch_key(12, "sha", "check"))); + assert_eq!(back.events.len(), 2); + assert_eq!(back.last_dispatch_at(12, "review"), Some(1_700_000_000)); + + // A different project's partition is empty — same (number, head, kind) is not + // visible across projects. + let other = Ledger::load_db(&db, "beta").expect("load other"); + assert!(other.dispatched.is_empty()); + assert!(other.events.is_empty()); + } + + // `save_db` replaces the whole partition (delete + re-insert `self`), so a later + // save with FEWER rows shrinks the stored set rather than leaving orphans. + #[test] + fn save_db_replaces_partition() { + let db = Database::open_in_memory().expect("open db"); + let mut full = Ledger::default(); + full.stage_all(&[cand(1, "review"), cand(2, "review")], 1_000); + full.save_db(&db, "p").expect("save full"); + + let mut fewer = Ledger::default(); + fewer.stage_all(&[cand(1, "review")], 1_000); + fewer.save_db(&db, "p").expect("save fewer"); + + let back = Ledger::load_db(&db, "p").expect("load"); + assert_eq!(back.dispatched.len(), 1); + assert!(back.has_dispatched(&dispatch_key(1, "sha", "review"))); + assert!(!back.has_dispatched(&dispatch_key(2, "sha", "review"))); + } + + // One-time legacy import (#70): the old `ledger.json` shapes (a JSON array of dedup + // keys, a JSON array of `DispatchEvent`) import into the SQLite partition and read + // back through `load_db`. Guards the migration path existing users rely on. #[test] - fn project_store_keys_are_partitioned() { - assert_eq!(dispatched_key("alpha"), "dispatched:alpha"); - assert_eq!(events_key("alpha"), "events:alpha"); - assert_ne!(dispatched_key("alpha"), dispatched_key("beta")); - assert_ne!(events_key("alpha"), events_key("beta")); - // The dedup KEY format itself is project-agnostic and unchanged — isolation - // comes from the STORE key, not from baking the project into the dedup key. - assert_eq!(dispatch_key(1, "sha", "review"), "1@sha:review"); + fn legacy_import_round_trips_through_load() { + let db = Database::open_in_memory().expect("open db"); + let dispatched_v = serde_json::json!(["12@sha:review", "13@sha:check"]); + let events_v = serde_json::to_value(vec![ + event(12, "review", 1_700_000_000), + event(13, "check", 1_700_000_100), + ]) + .expect("events serialize"); + + db.with_tx(|tx| { + import_legacy_dispatched(tx, "alpha", &dispatched_v)?; + import_legacy_events(tx, "alpha", &events_v)?; + Ok(()) + }) + .expect("import"); + + let back = Ledger::load_db(&db, "alpha").expect("load"); + assert!(back.has_dispatched("12@sha:review")); + assert!(back.has_dispatched("13@sha:check")); + assert_eq!(back.last_dispatch_at(12, "review"), Some(1_700_000_000)); + assert_eq!(back.last_dispatch_at(13, "check"), Some(1_700_000_100)); } // Ledger isolation (#35): two projects whose dedup sets are loaded from distinct diff --git a/src-tauri/src/pr/registry.rs b/src-tauri/src/pr/registry.rs index 164aacd..df2bbb2 100644 --- a/src-tauri/src/pr/registry.rs +++ b/src-tauri/src/pr/registry.rs @@ -1,13 +1,14 @@ //! Persisted PR-retention registry — the "ghost flicker" fix. //! -//! Each poll round upserts the discovered PRs into a persisted set (`prs.json` -//! via `tauri-plugin-store`'s `StoreExt`, the same backend-owned store pattern as -//! [`super::ledger`] / the config slice). The emitted list is the *retained* set, -//! not the raw per-round discovery: a transient one-round `gh` miss no longer -//! drops a row — it just flips that row's [`crate::model::PrPresence`] from -//! `Current` to `Stale` once it ages past the presence grace window. PRs are never -//! auto-evicted (users archive inactive ones); [`TrackedPrs::prune`] is only the -//! unbounded-growth backstop. +//! Each poll round upserts the discovered PRs into a persisted set (the unified SQLite +//! store's `tracked_pr` table, #70 — replacing the old `prs.json`; partitioned by a real +//! `project_id` column). The emitted list is the *retained* set, not the raw per-round +//! discovery: a transient one-round `gh` miss no longer drops a row — it just flips that +//! row's [`crate::model::PrPresence`] from `Current` to `Stale` once it ages past the +//! presence grace window. PRs are never auto-evicted (users archive inactive ones); +//! [`TrackedPrs::prune`] is only the unbounded-growth backstop. This table is ALSO the +//! durable store for webhook-received PRs (#70): the webhook ingest upserts here without +//! calling `gh`, so a webhook PR survives a `gh` outage and a restart. //! //! Slice boundary: presence is computed purely from `last_seen_epoch` vs the grace //! window — the `pr` slice stays review-agnostic and never reads `state.sessions`. @@ -15,47 +16,34 @@ use std::sync::Mutex; use serde::{Deserialize, Serialize}; -use tauri_plugin_store::StoreExt; +use tauri::Manager; use crate::config::service as config_service; -use crate::error::{AppError, AppResult}; +use crate::db::{map_err, Database}; +use crate::error::AppResult; use crate::model::{PrPresence, PullRequestView, TrackedPrView}; -/// Store file holding the persisted tracked-PR set. -const STORE_FILE: &str = "prs.json"; -/// Key PREFIX holding the per-project list of tracked PRs (#35). The effective key -/// is `tracked:{project_id}` (see [`tracked_key`]); a single `prs.json` holds every -/// project's tracked set under its own key, so two projects' PRs never mingle in one -/// list. -const TRACKED_KEY_PREFIX: &str = "tracked"; /// Unbounded-growth cap, applied PER PROJECT (#35). Beyond this, [`TrackedPrs::prune`] /// drops the least recently seen records (never the recent working set) — see its doc. const MAX_TRACKED: usize = 500; -/// Store key for a project's tracked-PR set: `tracked:{project_id}` (#35). -/// Partitions the shared `prs.json` so each project's retained list is isolated. -fn tracked_key(project_id: &str) -> String { - format!("{TRACKED_KEY_PREFIX}:{project_id}") -} - /// Serializes every read-modify-write of the persisted set (F1, PR #43). The two /// writers — the poll cycle's upsert and the `set_pr_archived` command — each do a -/// load→mutate→save of the whole `prs.json`; without a shared critical section they -/// interleave and silently lose each other's write (an archive overwritten by a poll -/// that loaded the pre-archive snapshot, or vice versa). A process-global `Mutex<()>` -/// (the data lives in the store, not behind the lock) is the gate, and +/// load→mutate→save of a project's `tracked_pr` partition; without a shared critical +/// section they interleave and silently lose each other's write (an archive overwritten +/// by a poll that loaded the pre-archive snapshot, or vice versa). A process-global +/// `Mutex<()>` (the data lives in SQLite, not behind the lock) is the gate, and /// [`mutate_tracked`] is its only acquirer. A module static — not an injected /// `AppState` field — so the lock *identity* is fixed: a caller cannot accidentally /// serialize on the wrong mutex, which closes the funnel downstream as well as up. /// `std` (not `tokio`) `Mutex`: the guarded section is fully synchronous, so no /// `.await` is ever held across the guard. /// -/// **Multi-project (#35):** the lock stays GLOBAL (not per-project) on purpose. Each -/// project's set lives under its own store key (`tracked:{project_id}`), but -/// `Store::save` rewrites the WHOLE `prs.json` — so two projects' parallel poll cycles -/// each doing a load→mutate→save would still clobber each other's just-written key. A -/// single global gate over the shared file is the correct granularity; a per-project -/// lock would reopen that cross-project lost-update race. +/// **Multi-project (#35):** the lock stays GLOBAL (not per-project) on purpose. `save` +/// replaces a project's whole partition (delete + re-insert) — the SQLite connection's +/// own mutex serializes individual statements, but only THIS lock makes a project's +/// load→mutate→save one atomic critical section. Keeping it global (rather than +/// per-project) matches the prior reasoning and keeps the lost-update analysis intact. static WRITE_LOCK: Mutex<()> = Mutex::new(()); /// One persisted PR. `first_seen_epoch` is set once on insert and preserved across @@ -88,20 +76,42 @@ impl TrackedPrs { /// this project's key (`tracked:{project_id}`), so one project's retained list never /// shows another's PRs. pub fn load(app: &tauri::AppHandle, project_id: &str) -> AppResult { - let store = app - .store(STORE_FILE) - .map_err(|e| AppError::new(format!("打开 PR 存储失败: {e}")))?; - - let prs = store - .get(tracked_key(project_id)) - .and_then(|v| serde_json::from_value::>(v).ok()) - .unwrap_or_default(); + Self::load_db(app.state::().inner(), project_id) + } - Ok(Self { prs }) + /// SQLite-level load (no Tauri app) — reads this project's `tracked_pr` rows. Split + /// from [`Self::load`] so the store round-trip is testable against an in-memory + /// [`Database`]. `labels` is stored as a JSON-text column (`labels_json`); a corrupt + /// value degrades to an empty list (parity with the old `unwrap_or_default` leniency). + /// Epochs are stored as `i64` and read back as `u64`. + pub(crate) fn load_db(db: &Database, project_id: &str) -> AppResult { + db.with_conn(|conn| { + let mut stmt = conn.prepare( + "SELECT number, title, labels_json, url, kind, skip_reason, \ + first_seen_epoch, last_seen_epoch, archived \ + FROM tracked_pr WHERE project_id = ?1 ORDER BY number DESC", + )?; + let rows = stmt.query_map([project_id], |r| { + let labels_json: String = r.get(2)?; + Ok(TrackedPr { + number: r.get::<_, i64>(0)? as u64, + title: r.get(1)?, + labels: serde_json::from_str(&labels_json).unwrap_or_default(), + url: r.get(3)?, + kind: r.get(4)?, + skip_reason: r.get(5)?, + first_seen_epoch: r.get::<_, i64>(6)? as u64, + last_seen_epoch: r.get::<_, i64>(7)? as u64, + archived: r.get::<_, i64>(8)? != 0, + }) + })?; + let prs = rows.collect::>>()?; + Ok(Self { prs }) + }) } /// Persists the tracked set. **Module-private — the F1 funnel's upstream gate.** - /// This is the only write path to `prs.json`, and it is reachable solely from + /// This is the only write path to the `tracked_pr` table, reachable solely from /// [`mutate_tracked`] (same module), which holds [`WRITE_LOCK`] across the whole /// load→mutate→save. Keeping `save` private makes a lock-free read-modify-write /// *not expressible* outside this module: a new writer has no way to call `save`, @@ -112,18 +122,45 @@ impl TrackedPrs { app: &tauri::AppHandle, project_id: &str, ) -> AppResult<()> { - let store = app - .store(STORE_FILE) - .map_err(|e| AppError::new(format!("打开 PR 存储失败: {e}")))?; - // tauri-plugin-store 2.x: `Store::set` is infallible and returns `()`. - store.set( - tracked_key(project_id), - serde_json::to_value(&self.prs).map_err(|e| AppError::new(e.to_string()))?, - ); - store - .save() - .map_err(|e| AppError::new(format!("写入 PR 存储失败: {e}")))?; - Ok(()) + self.save_db(app.state::().inner(), project_id) + } + + /// SQLite-level save (no Tauri app) — replaces this project's partition with the + /// full in-memory `self` (delete-all + insert-all, the SQLite analogue of the old + /// whole-key `Store::set`). `pub(crate)` only so round-trip tests can drive it + /// against an in-memory [`Database`]; the F1 funnel still holds because production + /// writers reach persistence only through the module-private [`Self::save`] → + /// [`mutate_tracked`]. + pub(crate) fn save_db(&self, db: &Database, project_id: &str) -> AppResult<()> { + db.with_tx(|tx| { + tx.execute("DELETE FROM tracked_pr WHERE project_id = ?1", [project_id]) + .map_err(map_err)?; + let mut stmt = tx + .prepare( + "INSERT INTO tracked_pr \ + (project_id, number, title, labels_json, url, kind, skip_reason, \ + first_seen_epoch, last_seen_epoch, archived) \ + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", + ) + .map_err(map_err)?; + for pr in &self.prs { + let labels_json = serde_json::to_string(&pr.labels).unwrap_or_else(|_| "[]".into()); + stmt.execute(rusqlite::params![ + project_id, + pr.number as i64, + pr.title, + labels_json, + pr.url, + pr.kind, + pr.skip_reason, + pr.first_seen_epoch as i64, + pr.last_seen_epoch as i64, + pr.archived as i64, + ]) + .map_err(map_err)?; + } + Ok(()) + }) } /// Upserts this round's discovered views into the tracked set. An existing PR @@ -261,6 +298,49 @@ where Ok(out) } +/// One-time legacy import (#70) of a project's `tracked:{pid}` list from the old +/// `prs.json`. Parses the JSON array of [`TrackedPr`] and inserts `tracked_pr` rows via +/// the same [`TrackedPrs::save_db`] used by production (so the column mapping has one +/// source). Lenient: a corrupt value imports an empty set (parity with `load`). +/// Runs inside the composition root's import transaction — but `save_db` opens its own +/// transaction on the shared connection, which the import's outer `with_tx` would +/// deadlock against; so it is called with a freshly-built `TrackedPrs` and inserts +/// directly here rather than nesting `save_db`. See the inline note. +pub fn import_legacy_tracked( + tx: &rusqlite::Transaction, + project_id: &str, + value: &serde_json::Value, +) -> AppResult<()> { + let prs: Vec = serde_json::from_value(value.clone()).unwrap_or_default(); + // Insert directly on the import transaction (do NOT call `save_db`, which would open + // a NESTED transaction on the same connection and fail). Mirrors `save_db`'s columns. + let mut stmt = tx + .prepare( + "INSERT OR REPLACE INTO tracked_pr \ + (project_id, number, title, labels_json, url, kind, skip_reason, \ + first_seen_epoch, last_seen_epoch, archived) \ + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", + ) + .map_err(map_err)?; + for pr in &prs { + let labels_json = serde_json::to_string(&pr.labels).unwrap_or_else(|_| "[]".into()); + stmt.execute(rusqlite::params![ + project_id, + pr.number as i64, + pr.title, + labels_json, + pr.url, + pr.kind, + pr.skip_reason, + pr.first_seen_epoch as i64, + pr.last_seen_epoch as i64, + pr.archived as i64, + ]) + .map_err(map_err)?; + } + Ok(()) +} + /// Projects the tracked set into the frontend wire rows. Each record's presence is /// `Current` when last seen within `grace_secs` of `now`, else `Stale` /// (`saturating_sub` so a backwards clock reads age 0 → `Current`). Rows are sorted @@ -352,6 +432,77 @@ mod tests { } } + // SQLite store round-trip (#70, Medium carrier): an upserted set saved + loaded + // against an in-memory DB must preserve every field — including `labels` (stored as + // the `labels_json` text column), the nullable `skip_reason`, the presence clocks, + // and the `archived` flag. A column/SQL drift surfaces here. + #[test] + fn sqlite_round_trip_preserves_all_fields() { + let db = Database::open_in_memory().expect("open db"); + let mut t = TrackedPrs::default(); + t.upsert(&[view(12, "Add feature")], 1_700_000_000); + t.set_archived(12, true); + t.save_db(&db, "alpha").expect("save"); + + let back = TrackedPrs::load_db(&db, "alpha").expect("load"); + assert_eq!(back.prs.len(), 1); + let pr = &back.prs[0]; + assert_eq!(pr.number, 12); + assert_eq!(pr.title, "Add feature"); + assert_eq!(pr.labels, vec!["review-label".to_string()]); + assert_eq!(pr.url, "https://x/12"); + assert_eq!(pr.first_seen_epoch, 1_700_000_000); + assert!(pr.archived); + + // A different project's partition is empty (the `project_id` column is the seam). + assert!(TrackedPrs::load_db(&db, "beta") + .expect("load other") + .prs + .is_empty()); + } + + // `save_db` replaces the whole partition: a row dropped from `self` (e.g. by a future + // prune) disappears from storage rather than lingering as an orphan. + #[test] + fn save_db_replaces_partition() { + let db = Database::open_in_memory().expect("open db"); + let mut two = TrackedPrs::default(); + two.upsert(&[view(1, "one"), view(2, "two")], 1_000); + two.save_db(&db, "p").expect("save two"); + + let mut one = TrackedPrs::default(); + one.upsert(&[view(1, "one")], 1_000); + one.save_db(&db, "p").expect("save one"); + + let back = TrackedPrs::load_db(&db, "p").expect("load"); + assert_eq!(back.prs.len(), 1); + assert_eq!(back.prs[0].number, 1); + } + + // One-time legacy import (#70): the old `prs.json` shape (a JSON array of `TrackedPr`) + // imports into the SQLite partition and reads back through `load_db`, preserving the + // archived flag + presence clock existing users rely on across the migration. + #[test] + fn legacy_import_round_trips_through_load() { + let db = Database::open_in_memory().expect("open db"); + let mut archived = tracked(7, 1_700_000_000); + archived.archived = true; + archived.labels = vec!["needs-review".to_string()]; + let value = serde_json::to_value(vec![tracked(9, 1_700_000_500), archived]) + .expect("legacy list serializes"); + + db.with_tx(|tx| import_legacy_tracked(tx, "alpha", &value)) + .expect("import"); + + let back = TrackedPrs::load_db(&db, "alpha").expect("load"); + assert_eq!(back.prs.len(), 2); + let pr7 = back.prs.iter().find(|p| p.number == 7).expect("pr 7"); + assert!(pr7.archived, "archived flag survives import"); + assert_eq!(pr7.labels, vec!["needs-review".to_string()]); + assert_eq!(pr7.first_seen_epoch, 0); + assert_eq!(pr7.last_seen_epoch, 1_700_000_000); + } + // Wire-shape lock for the persisted `prs.json` records (Medium carrier per // ai-robust.md). A field rename would make `TrackedPrs::load` silently drop the // records (deserialize → `unwrap_or_default()`), wiping the retained set and diff --git a/src-tauri/src/review/commands.rs b/src-tauri/src/review/commands.rs index 159ad96..f846dd7 100644 --- a/src-tauri/src/review/commands.rs +++ b/src-tauri/src/review/commands.rs @@ -1,9 +1,13 @@ //! Review slice Tauri commands. +use tauri::Manager; + use crate::config::service as config_service; +use crate::db::Database; use crate::error::{AppError, AppResult}; use crate::review::engine::{ReviewEngine, SessionId, StartReviewOutcome}; use crate::review::engines::codex::{CodexEngine, CodexStatus}; +use crate::review::history_store::{self, HistoryItem}; use crate::review::session::SessionInfo; use crate::state::AppState; @@ -141,12 +145,39 @@ pub async fn stop_review( engine.stop(&session_id).await } -/// Snapshot of all review sessions (running + finished) for the UI. +/// Snapshot of all review sessions (running + finished) for the UI — the IN-MEMORY +/// registry (live status). After a restart this is empty; [`get_pr_sessions`] reads the +/// durable table instead. #[tauri::command] pub fn list_review_sessions(state: tauri::State<'_, AppState>) -> Vec { state.sessions.list() } +/// A session's persisted history items in stream order (#70) — the message/reasoning +/// blocks produced before the user opened the session. Drives the review panel's "open a +/// history session and see prior content" (#67) by hydrating the per-session buffer. +#[tauri::command] +pub fn get_session_history( + app: tauri::AppHandle, + thread_id: String, +) -> AppResult> { + let db = app.state::(); + history_store::get_history(db.inner(), &thread_id) +} + +/// A PR's persisted sessions, newest first, from the durable `review_session` table (#70) +/// — so each PR can restore its session list after a restart (the #67 nav associates +/// sessions per PR). Distinct from [`list_review_sessions`] (in-memory live snapshot). +#[tauri::command] +pub fn get_pr_sessions( + app: tauri::AppHandle, + project_id: String, + pr_number: u64, +) -> AppResult> { + let db = app.state::(); + history_store::get_pr_sessions(db.inner(), &project_id, pr_number) +} + /// Resolve the absolute path to the pr-review skill file codex attaches to the /// turn. `repo_root` is an absolute dir and `skill_rel_path` a relative path under /// it (both config-validated), so the join is absolute and infallible. diff --git a/src-tauri/src/review/history_store.rs b/src-tauri/src/review/history_store.rs new file mode 100644 index 0000000..e14ff08 --- /dev/null +++ b/src-tauri/src/review/history_store.rs @@ -0,0 +1,247 @@ +//! Review session + history persistence (#70). +//! +//! The in-memory [`super::session::SessionRegistry`] stays the authority for dedup / +//! active-pairs / live status; this module is its DURABLE MIRROR in the unified SQLite +//! store. It persists (1) session metadata (`review_session` — so a PR's sessions +//! survive a restart and are listable per-PR) and (2) session HISTORY content +//! (`review_history_item` — the streamed message/reasoning deltas that were previously +//! emitted-and-discarded, so a history session can be reopened with its prior output). +//! +//! Slice-local: the review slice owns these tables' queries, reaching SQLite only through +//! the horizontal [`crate::db::Database`] handle (not a cross-slice import). + +use serde::Serialize; + +use super::session::{SessionInfo, SessionStatus}; +use crate::db::Database; +use crate::error::AppResult; + +/// One persisted history item — a coalesced message/reasoning block (#70). Same wire +/// shape as the frontend's `StreamItem` (`itemId` / `kind` / `text`), so +/// `get_session_history` can hydrate the review panel directly. Slice-private; mirrored +/// in `src/review/types.ts`. +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct HistoryItem { + pub item_id: String, + pub kind: String, + pub text: String, +} + +/// Wall-clock seconds for the session timestamps. Review-local (the `pr` slice has its +/// own `now_epoch`; keeping one here avoids a cross-slice import — the review slice stays +/// self-contained). Degrades to 0 on a pre-epoch clock rather than panicking. +fn now_epoch() -> u64 { + use std::time::{SystemTime, UNIX_EPOCH}; + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0) +} + +/// [`SessionStatus`] → its pinned camelCase wire string (the form stored in the `status` +/// column), via the same serde contract `list_review_sessions` uses. +fn status_wire(status: SessionStatus) -> String { + serde_json::to_value(status) + .ok() + .and_then(|v| v.as_str().map(str::to_string)) + .unwrap_or_default() +} + +/// Wire string → [`SessionStatus`] (reverse of [`status_wire`]), for projecting stored +/// rows back into [`SessionInfo`]. An unknown string degrades to `Failed` (a terminal +/// status — never resurrects a dead session as live). +fn status_from_wire(s: &str) -> SessionStatus { + serde_json::from_value(serde_json::Value::String(s.to_string())) + .unwrap_or(SessionStatus::Failed) +} + +/// Upserts a session row from its in-memory [`SessionInfo`] (#70). First insert stamps +/// `created_at`; a later call (turn/status transition) updates the mutable fields + +/// `updated_at`, preserving `created_at`. Persistence is best-effort — callers log + +/// swallow errors so a DB hiccup never breaks the live session. +pub fn upsert_session(db: &Database, info: &SessionInfo) -> AppResult<()> { + let now = now_epoch() as i64; + let status = status_wire(info.status); + db.with_conn(|conn| { + conn.execute( + "INSERT INTO review_session \ + (thread_id, project_id, pr_number, turn_id, kind, status, created_at, updated_at) \ + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?7) \ + ON CONFLICT(thread_id) DO UPDATE SET \ + project_id = excluded.project_id, \ + pr_number = excluded.pr_number, \ + turn_id = excluded.turn_id, \ + kind = excluded.kind, \ + status = excluded.status, \ + updated_at = excluded.updated_at", + rusqlite::params![ + info.thread_id, + info.project_id, + info.pr_number as i64, + info.turn_id, + info.kind, + status, + now, + ], + ) + .map(|_| ()) + }) +} + +/// Updates a session's status (#70) without the full [`SessionInfo`] — the pump's +/// terminal branch + `set_status` callsites have only the thread id. A no-op if the row +/// doesn't exist yet (the session insert always precedes any status update in practice). +pub fn set_status(db: &Database, thread_id: &str, status: SessionStatus) -> AppResult<()> { + let now = now_epoch() as i64; + let status = status_wire(status); + db.with_conn(|conn| { + conn.execute( + "UPDATE review_session SET status = ?2, updated_at = ?3 WHERE thread_id = ?1", + rusqlite::params![thread_id, status, now], + ) + .map(|_| ()) + }) +} + +/// Appends a streamed delta to a session's history (#70), COALESCING by `(thread_id, +/// item_id)`: the first delta for an item inserts a row; later deltas concatenate onto +/// its `text`. Mirrors the frontend's `appendDelta` so thousands of deltas collapse into +/// a handful of rows. `kind` is `"message"` or `"reasoning"` (constant per item id). +pub fn append_item( + db: &Database, + thread_id: &str, + item_id: &str, + kind: &str, + text: &str, +) -> AppResult<()> { + db.with_conn(|conn| { + conn.execute( + "INSERT INTO review_history_item (thread_id, item_id, kind, text) \ + VALUES (?1, ?2, ?3, ?4) \ + ON CONFLICT(thread_id, item_id) DO UPDATE SET text = text || excluded.text", + rusqlite::params![thread_id, item_id, kind, text], + ) + .map(|_| ()) + }) +} + +/// A session's stored history items in stream order (#70) — what `get_session_history` +/// returns so the UI can show content produced before the user opened the session. +pub fn get_history(db: &Database, thread_id: &str) -> AppResult> { + db.with_conn(|conn| { + let mut stmt = conn.prepare( + "SELECT item_id, kind, text FROM review_history_item \ + WHERE thread_id = ?1 ORDER BY id", + )?; + let rows = stmt.query_map([thread_id], |r| { + Ok(HistoryItem { + item_id: r.get(0)?, + kind: r.get(1)?, + text: r.get(2)?, + }) + })?; + rows.collect() + }) +} + +/// A PR's persisted sessions (#70), newest first — what `get_pr_sessions` returns so each +/// PR can restore its session list after a restart. Reads the DURABLE `review_session` +/// table (vs `list_review_sessions`'s in-memory snapshot, which is empty after restart). +pub fn get_pr_sessions( + db: &Database, + project_id: &str, + pr_number: u64, +) -> AppResult> { + db.with_conn(|conn| { + let mut stmt = conn.prepare( + "SELECT thread_id, project_id, pr_number, turn_id, kind, status \ + FROM review_session \ + WHERE project_id = ?1 AND pr_number = ?2 ORDER BY created_at DESC, thread_id", + )?; + let rows = stmt.query_map(rusqlite::params![project_id, pr_number as i64], |r| { + let status: String = r.get(5)?; + Ok(SessionInfo { + thread_id: r.get(0)?, + project_id: r.get(1)?, + pr_number: r.get::<_, i64>(2)? as u64, + turn_id: r.get(3)?, + kind: r.get(4)?, + status: status_from_wire(&status), + }) + })?; + rows.collect() + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn info(thread: &str, pr: u64, status: SessionStatus) -> SessionInfo { + SessionInfo { + project_id: "alpha".to_string(), + thread_id: thread.to_string(), + turn_id: "t1".to_string(), + pr_number: pr, + kind: "review".to_string(), + status, + } + } + + // Session metadata round-trip (#70, Medium carrier): upsert then `get_pr_sessions` + // must surface the row with its status, and a status transition must update (not + // duplicate) it. Restart durability rests on this read path. + #[test] + fn session_upsert_and_status_transition_round_trip() { + let db = Database::open_in_memory().expect("open db"); + upsert_session(&db, &info("th-1", 12, SessionStatus::Starting)).expect("insert"); + set_status(&db, "th-1", SessionStatus::Done).expect("transition"); + + let sessions = get_pr_sessions(&db, "alpha", 12).expect("list"); + assert_eq!(sessions.len(), 1, "transition updates, not duplicates"); + assert_eq!(sessions[0].thread_id, "th-1"); + assert_eq!(sessions[0].status, SessionStatus::Done); + + // Scoped per (project, PR): a different PR sees nothing. + assert!(get_pr_sessions(&db, "alpha", 99).expect("list").is_empty()); + } + + // History capture (#70): deltas for one item id COALESCE into a single concatenated + // row, distinct item ids are separate rows, and the read is in stream (`id`) order. + #[test] + fn append_item_coalesces_by_item_id_and_reads_in_order() { + let db = Database::open_in_memory().expect("open db"); + upsert_session(&db, &info("th-1", 12, SessionStatus::Running)).expect("session"); + + append_item(&db, "th-1", "i1", "reasoning", "Plan").expect("a"); + append_item(&db, "th-1", "i2", "message", "Hello ").expect("b"); + append_item(&db, "th-1", "i1", "reasoning", "ning done").expect("c"); + append_item(&db, "th-1", "i2", "message", "world").expect("d"); + + let items = get_history(&db, "th-1").expect("history"); + assert_eq!(items.len(), 2, "two item ids → two coalesced rows"); + // i1 inserted first → comes first; deltas concatenated in arrival order. + assert_eq!(items[0].item_id, "i1"); + assert_eq!(items[0].kind, "reasoning"); + assert_eq!(items[0].text, "Planning done"); + assert_eq!(items[1].item_id, "i2"); + assert_eq!(items[1].text, "Hello world"); + } + + // `HistoryItem` wire-shape lock (#70, Medium carrier): camelCase `itemId` present, + // snake_case absent — keeps the Rust↔`src/review/types.ts` (`StreamItem`) contract. + #[test] + fn history_item_wire_shape_is_camel_case() { + let item = HistoryItem { + item_id: "i1".to_string(), + kind: "message".to_string(), + text: "hi".to_string(), + }; + let v = serde_json::to_value(&item).expect("serializes"); + assert!(v.get("itemId").is_some()); + assert!(v.get("kind").is_some()); + assert!(v.get("text").is_some()); + assert!(v.get("item_id").is_none()); + } +} diff --git a/src-tauri/src/review/mod.rs b/src-tauri/src/review/mod.rs index c7826ad..bbe61a2 100644 --- a/src-tauri/src/review/mod.rs +++ b/src-tauri/src/review/mod.rs @@ -9,4 +9,5 @@ pub mod commands; pub mod engine; pub mod engines; +pub mod history_store; pub mod session; diff --git a/src-tauri/src/review/session.rs b/src-tauri/src/review/session.rs index 28c062e..c774163 100644 --- a/src-tauri/src/review/session.rs +++ b/src-tauri/src/review/session.rs @@ -12,8 +12,8 @@ use std::collections::{HashMap, HashSet}; use std::sync::{Arc, Mutex}; -use serde::Serialize; -use tauri::Emitter; +use serde::{Deserialize, Serialize}; +use tauri::{Emitter, Manager}; use tokio::sync::broadcast; use super::engines::codex::process; @@ -34,8 +34,9 @@ pub type ThreadId = String; const PR_REVIEW_SKILL: &str = "pr-review"; /// Lifecycle of one review session (the state machine). Serialized camelCase for -/// `list_review_sessions`. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +/// `list_review_sessions`; `Deserialize` so the persisted `review_session.status` wire +/// string (#70) projects back into this enum in `history_store::get_pr_sessions`. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub enum SessionStatus { /// `thread/start` / `turn/start` in flight (not yet observed on the stream). @@ -370,14 +371,18 @@ pub async fn start_review( // dispatch never slips between `thread/start` success and this insert. A later // `turn/start` failure flips it to `Failed` (visible to `list_review_sessions`, // not vanished); `turn_id` is filled once the turn starts. - registry.promote_reservation(SessionInfo { + let starting = SessionInfo { project_id: project_id.to_string(), thread_id: thread_id.clone(), turn_id: String::new(), pr_number, kind: kind.to_string(), status: SessionStatus::Starting, - }); + }; + registry.promote_reservation(starting.clone()); + // Mirror the in-memory session into the durable `review_session` table (#70) so this + // PR's session list survives a restart and its history can be reopened. Best-effort. + persist_session(app, &starting); reservation.disarm(); let prompt = review_prompt(repo, &skill_command(pr_number, kind)); @@ -409,11 +414,24 @@ pub async fn start_review( Ok(turn_id) => turn_id, Err(e) => { registry.set_status(&thread_id, SessionStatus::Failed); + persist_status(app, &thread_id, SessionStatus::Failed); return Err(e); } }; - registry.set_running(&thread_id, turn_id); + registry.set_running(&thread_id, turn_id.clone()); + // Mirror the Running transition (+ the now-known turn id) into `review_session` (#70). + persist_session( + app, + &SessionInfo { + project_id: project_id.to_string(), + thread_id: thread_id.clone(), + turn_id, + pr_number, + kind: kind.to_string(), + status: SessionStatus::Running, + }, + ); // Capture `project_id` as an owned String at spawn time so the pump stamps every // emitted `ReviewEvent` with it WITHOUT re-looking-up the session per event (#35): @@ -497,6 +515,58 @@ pub async fn stop_review( Ok(()) } +/// Best-effort mirror of an in-memory [`SessionInfo`] into the durable `review_session` +/// table (#70). Logs + swallows errors: a persistence hiccup must never break the live +/// session (the in-memory registry stays the authority for dedup / status). +fn persist_session(app: &tauri::AppHandle, info: &SessionInfo) { + let db = app.state::(); + if let Err(e) = super::history_store::upsert_session(db.inner(), info) { + eprintln!( + "review session 持久化失败({}):{}", + info.thread_id, e.message + ); + } +} + +/// Best-effort mirror of a session status transition into `review_session` (#70). Used at +/// terminal transitions in the pump where only the thread id is at hand. +fn persist_status( + app: &tauri::AppHandle, + thread_id: &str, + status: SessionStatus, +) { + let db = app.state::(); + if let Err(e) = super::history_store::set_status(db.inner(), thread_id, status) { + eprintln!( + "review session 状态持久化失败({thread_id}):{}", + e.message + ); + } +} + +/// Best-effort capture of a streamed delta into the persisted session history (#70). +/// Called AFTER `app.emit` so the live stream never waits on the DB; a non-delta event is +/// a no-op, and a persist error is logged + swallowed (the rendered stream is unaffected +/// — at worst the last delta before a crash is missing from the reopened history). +fn persist_delta( + app: &tauri::AppHandle, + thread_id: &str, + event: &ReviewEvent, +) { + let (item_id, kind, text) = match event { + ReviewEvent::MessageDelta { item_id, text, .. } => (item_id, "message", text), + ReviewEvent::ReasoningDelta { item_id, text, .. } => (item_id, "reasoning", text), + _ => return, + }; + let db = app.state::(); + if let Err(e) = super::history_store::append_item(db.inner(), thread_id, item_id, kind, text) { + eprintln!( + "review history 持久化失败({thread_id}/{item_id}):{}", + e.message + ); + } +} + /// Pump task: forward this session's notifications to the frontend as /// [`ReviewEvent`]s until the turn completes (or the connection drops). Filters by /// `thread_id` since the broadcast carries every session's stream; stamps every @@ -526,11 +596,17 @@ async fn pump( continue; }; if let ReviewEvent::TurnCompleted { status, .. } = &event { - registry.set_status(&thread_id, terminal_status(status)); + let terminal = terminal_status(status); + registry.set_status(&thread_id, terminal); let _ = app.emit(REVIEW_EVENT, &event); + persist_status(&app, &thread_id, terminal); // mirror terminal to DB (#70) break; // terminal — the turn is over. } + // Emit FIRST (streaming latency must not wait on the DB), THEN persist the + // delta to the session history (#70) best-effort — a persist error is + // logged, never breaks the live stream. let _ = app.emit(REVIEW_EVENT, &event); + persist_delta(&app, &thread_id, &event); } // The pump fell behind the shared ring and `n` notifications were // evicted. The terminal `turn/completed` may have been among them @@ -541,6 +617,7 @@ async fn pump( Err(broadcast::error::RecvError::Lagged(n)) => { eprintln!("review pump({thread_id})滞后,丢弃 {n} 条通知"); registry.set_status(&thread_id, SessionStatus::Failed); + persist_status(&app, &thread_id, SessionStatus::Failed); // mirror to DB (#70) let _ = app.emit( REVIEW_EVENT, &ReviewEvent::Error { @@ -573,6 +650,7 @@ fn fail_connection_closed( thread_id: &str, ) { registry.set_status(thread_id, SessionStatus::Failed); + persist_status(app, thread_id, SessionStatus::Failed); // mirror to DB (#70) let _ = app.emit( REVIEW_EVENT, &ReviewEvent::Error { diff --git a/src/App.vue b/src/App.vue index 0db7391..8cae330 100644 --- a/src/App.vue +++ b/src/App.vue @@ -8,8 +8,7 @@ import { useConfigStore } from "./config/useConfigStore"; import SettingsView from "./config/SettingsView.vue"; import OnboardingWizard from "./config/OnboardingWizard.vue"; import PollControls from "./pr/PollControls.vue"; -import PrList from "./pr/PrList.vue"; -import ProjectSwitcher from "./pr/ProjectSwitcher.vue"; +import ProjectNav from "./pr/ProjectNav.vue"; import WebhookPanel from "./pr/WebhookPanel.vue"; import { usePrStore } from "./pr/usePrStore"; import StatusBar from "./StatusBar.vue"; @@ -218,10 +217,13 @@ watch(activeProjectId, () => {
- -
- + @@ -331,7 +336,7 @@ body { min-height: 0; } -.sidebar { +.nav-column { width: var(--sidebar-width); border-right: 1px solid var(--color-border-strong); overflow-y: auto; diff --git a/src/pr/PrList.vue b/src/pr/PrList.vue deleted file mode 100644 index ae8db8d..0000000 --- a/src/pr/PrList.vue +++ /dev/null @@ -1,200 +0,0 @@ - - - - - diff --git a/src/pr/ProjectNav.vue b/src/pr/ProjectNav.vue new file mode 100644 index 0000000..f405d02 --- /dev/null +++ b/src/pr/ProjectNav.vue @@ -0,0 +1,301 @@ + + + + + diff --git a/src/pr/ProjectSwitcher.vue b/src/pr/ProjectSwitcher.vue deleted file mode 100644 index 2901e7a..0000000 --- a/src/pr/ProjectSwitcher.vue +++ /dev/null @@ -1,124 +0,0 @@ - - - - - diff --git a/src/review/ReviewSessions.vue b/src/review/ReviewSessions.vue index 8e162c5..f5bd580 100644 --- a/src/review/ReviewSessions.vue +++ b/src/review/ReviewSessions.vue @@ -1,23 +1,56 @@