diff --git a/crates/merge-request/src/lib.rs b/crates/merge-request/src/lib.rs index 31c847d3..86662fb7 100644 --- a/crates/merge-request/src/lib.rs +++ b/crates/merge-request/src/lib.rs @@ -677,7 +677,28 @@ impl SqliteMergeRequestStore { } pub fn migrate(conn: &Connection) -> Result<()> { + conn.busy_timeout(Duration::from_secs(5)).map_err(db)?; conn.pragma_update(None, "foreign_keys", "ON").map_err(db)?; + // Acquire the writer lock before reading either the version or table layout so + // concurrent store initialization cannot act on a stale migration decision. + conn.execute_batch("BEGIN IMMEDIATE").map_err(db)?; + let result = migrate_transaction(conn); + match result { + Ok(()) => match conn.execute_batch("COMMIT") { + Ok(()) => Ok(()), + Err(error) => { + let _ = conn.execute_batch("ROLLBACK"); + Err(db(error)) + } + }, + Err(error) => { + let _ = conn.execute_batch("ROLLBACK"); + Err(error) + } + } +} + +fn migrate_transaction(conn: &Connection) -> Result<()> { conn.execute_batch("CREATE TABLE IF NOT EXISTS merge_request_schema_migrations (version INTEGER PRIMARY KEY, applied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP);").map_err(db)?; let version: i64 = conn .query_row( @@ -1014,26 +1035,41 @@ fn archive_incompatible_legacy_tables(conn: &Connection, version: i64) -> Result "merge_request_revisions", "merge_requests", ]; - conn.pragma_update(None, "foreign_keys", "OFF") - .map_err(db)?; for table in tables { if !table_exists(conn, table)? { continue; } let archive = format!("legacy_v6_{table}"); if table_exists(conn, &archive)? { - conn.pragma_update(None, "foreign_keys", "ON").map_err(db)?; - return Err(MergeRequestError::Database(format!( - "legacy archive table {archive} already exists" - ))); + // The retired non-transactional migration could archive a table, recreate + // its empty replacement, and then fail. Resume that exact state without + // ever choosing between two populated copies. + if !table_is_empty(conn, table)? { + return Err(MergeRequestError::Database(format!( + "legacy archive table {archive} already exists while {table} still contains data" + ))); + } + conn.execute_batch(&format!("DROP TABLE {table};")) + .map_err(db)?; + continue; } conn.execute_batch(&format!("ALTER TABLE {table} RENAME TO {archive};")) .map_err(db)?; } - conn.pragma_update(None, "foreign_keys", "ON").map_err(db)?; Ok(()) } +fn table_is_empty(conn: &Connection, table: &str) -> Result { + let has_row: i64 = conn + .query_row( + &format!("SELECT EXISTS(SELECT 1 FROM {table} LIMIT 1)"), + [], + |row| row.get(0), + ) + .map_err(db)?; + Ok(has_row == 0) +} + fn table_has_columns(conn: &Connection, table: &str, required: &[&str]) -> Result { if !table_exists(conn, table)? { return Ok(false); diff --git a/crates/merge-request/tests/store.rs b/crates/merge-request/tests/store.rs index 26e6743d..4cf3c9a4 100644 --- a/crates/merge-request/tests/store.rs +++ b/crates/merge-request/tests/store.rs @@ -192,6 +192,83 @@ fn rejected_v6_schema_missing_diff_digest_is_archived_before_fresh_v7() { } } +#[test] +fn interrupted_legacy_archive_with_empty_recreated_table_resumes() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("interrupted.db"); + let conn = Connection::open(&path).unwrap(); + conn.execute_batch( + "CREATE TABLE merge_request_schema_migrations(version INTEGER PRIMARY KEY,name TEXT NOT NULL,applied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP);\ + INSERT INTO merge_request_schema_migrations(version,name) VALUES(6,'rejected_merge_request_v6');\ + CREATE TABLE merge_requests(workspace_id TEXT NOT NULL,merge_request_id TEXT NOT NULL,ticket_id TEXT NOT NULL);\ + CREATE TABLE merge_request_review_findings(workspace_id TEXT NOT NULL,attempt_id TEXT NOT NULL,ordinal INTEGER NOT NULL,severity TEXT NOT NULL,code TEXT,path TEXT,line INTEGER,body TEXT NOT NULL);\ + CREATE TABLE legacy_v6_merge_request_review_findings(workspace_id TEXT NOT NULL,attempt_id TEXT NOT NULL,ordinal INTEGER NOT NULL,severity TEXT NOT NULL,code TEXT,path TEXT,line INTEGER,body TEXT NOT NULL);\ + INSERT INTO legacy_v6_merge_request_review_findings VALUES('ws-a','AT1',0,'warning',NULL,NULL,NULL,'preserved evidence');", + ) + .unwrap(); + drop(conn); + + SqliteMergeRequestStore::open(&path, "ws-a").unwrap(); + let conn = Connection::open(&path).unwrap(); + let archived_body: String = conn + .query_row( + "SELECT body FROM legacy_v6_merge_request_review_findings", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(archived_body, "preserved evidence"); + let version: i64 = conn + .query_row( + "SELECT MAX(version) FROM merge_request_schema_migrations", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(version, 8); + drop(conn); + + SqliteMergeRequestStore::open(&path, "ws-a").unwrap(); +} + +#[test] +fn conflicting_legacy_archive_rolls_back_all_table_renames() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("conflict.db"); + let conn = Connection::open(&path).unwrap(); + conn.execute_batch( + "CREATE TABLE merge_request_schema_migrations(version INTEGER PRIMARY KEY,name TEXT NOT NULL,applied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP);\ + INSERT INTO merge_request_schema_migrations(version,name) VALUES(6,'rejected_merge_request_v6');\ + CREATE TABLE merge_requests(workspace_id TEXT NOT NULL,merge_request_id TEXT NOT NULL,ticket_id TEXT NOT NULL);\ + CREATE TABLE merge_request_review_findings(body TEXT NOT NULL);\ + CREATE TABLE merge_request_reviews(body TEXT NOT NULL);\ + INSERT INTO merge_request_reviews VALUES('unarchived evidence');\ + CREATE TABLE legacy_v6_merge_request_reviews(body TEXT NOT NULL);", + ) + .unwrap(); + + let error = migrate(&conn).unwrap_err(); + assert!(error.to_string().contains( + "legacy archive table legacy_v6_merge_request_reviews already exists while merge_request_reviews still contains data" + )); + let current_findings: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='merge_request_review_findings'", + [], + |row| row.get(0), + ) + .unwrap(); + let archived_findings: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='legacy_v6_merge_request_review_findings'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(current_findings, 1); + assert_eq!(archived_findings, 0); +} + #[test] fn v7_completion_operations_are_preserved_as_legacy_assigned_coder_authority() { let dir = tempfile::tempdir().unwrap();