diff --git a/crates/onestore-offline/Cargo.toml b/crates/onestore-offline/Cargo.toml index a9859360d60b3db96106c14ce3d5c147411de05f..2f14c77276dc108b22823d8ccf42efc58fb4ae3a 100644 --- a/crates/onestore-offline/Cargo.toml +++ b/crates/onestore-offline/Cargo.toml @@ -10,12 +10,10 @@ smb = ["dep:onestore-smb"] [dependencies] onestore = { path = "../onestore" } onestore-smb = { path = "../onestore-smb", optional = true } -rusqlite = { version = "=0.40.2", features = ["bundled"] } +rusqlite = { version = "=0.40.2", features = ["bundled", "backup"] } thiserror = "2" serde = { version = "1", features = ["derive"] } serde_json = "1" - -[dev-dependencies] tempfile = "3" [[example]] diff --git a/crates/onestore-offline/README.md b/crates/onestore-offline/README.md index 4084707e6ed4de780b69a75032f94757e1dbc221..d25d4d4bbbb379f1c1ab17430153d63789a30e97 100644 --- a/crates/onestore-offline/README.md +++ b/crates/onestore-offline/README.md @@ -189,3 +189,35 @@ The SMB-enabled `smb_offline_client` example is an owned-lab workload for acknowledgements, remote publication attempts, persisted receipts and cache reopen checks. Its append-specific conflict review policy lives in the test client; the library continues to preserve conflicts requiring an explicit decision. + +## Recovery archives + +`export_recovery(new_path)` captures both complete notebook images, the typed +queue, uncertain attempts, conflicts, receipts and the edit-ID sequence in one +SQLite snapshot. It refuses existing destinations and leaves the live queue +unchanged. Export to a local directory from a background thread: copying holds +the cache mutex while capturing the database. Failure after the final rename can +leave a complete archive at the requested path; it never acknowledges a remote +edit. Archives contain notebook content and use a separate database identity, so +`Replica::open` rejects them as writable caches. + +```no_run +use onestore_offline::{Recovery, Replica}; +# fn example(cache: &Replica) -> Result<(), Box> { +cache.export_recovery("review.sqlite")?; +let review = Recovery::open("review.sqlite")?; +let counts = review.summary()?; +let local = review.snapshot()?; +let remote = review.remote_snapshot()?; +let pending = review.pending()?; +let receipts = review.receipts()?; +# Ok(()) +# } +``` + +`Recovery` provides read-only inspection and no synchronization or restore method. +`status(id)` preserves the same state interpretation as the live replica. +`recovery_summary()` on the live replica and `summary()` on an archive return +counts and byte sizes without notebook text, paths, authors or credentials. +Opening an archive validates its schema and images without migration. A recovery +archive is evidence for a reviewed recovery decision, not a second active queue. diff --git a/crates/onestore-offline/src/lib.rs b/crates/onestore-offline/src/lib.rs index f96a06aea852e09c5ff0e753a41075bddcc7d8cd..5186c411f1053e8e2646ed2eab27c8549f613157 100644 --- a/crates/onestore-offline/src/lib.rs +++ b/crates/onestore-offline/src/lib.rs @@ -10,8 +10,10 @@ use std::{fs::OpenOptions, io, ops::Range, path::Path, sync::Mutex, time::Durati mod formatting; mod rebase; +mod recovery; mod schema; pub use formatting::FormatEdit; +pub use recovery::{Recovery, RecoverySummary}; mod sync; pub use sync::{ConflictKind, EditStatus, Remote}; mod worker; @@ -97,31 +99,7 @@ impl Replica { } fn connect(path: &Path, source: Option<&[u8]>) -> Result { - let mut connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; - connection.busy_timeout(Duration::ZERO)?; - connection.execute_batch( - "PRAGMA locking_mode=EXCLUSIVE; PRAGMA synchronous=EXTRA; PRAGMA fullfsync=ON; PRAGMA foreign_keys=ON;", - )?; - for (name, expected) in [("locking_mode", "exclusive"), ("journal_mode", "delete")] { - let actual: String = connection.pragma_query_value(None, name, |row| row.get(0))?; - if actual != expected { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Unsupported cache locking or journal mode", - ) - .into()); - } - } - for (name, expected) in [("synchronous", 3), ("fullfsync", 1), ("foreign_keys", 1)] { - let actual: i64 = connection.pragma_query_value(None, name, |row| row.get(0))?; - if actual != expected { - return Err(io::Error::new( - io::ErrorKind::Unsupported, - "Required cache synchronization is unavailable", - ) - .into()); - } - } + let mut connection = cache_connection(path)?; let transaction = connection.transaction_with_behavior(TransactionBehavior::Exclusive)?; let application: u32 = transaction.pragma_query_value(None, "application_id", |row| row.get(0))?; @@ -159,27 +137,7 @@ impl Replica { ) .into()); } - let integrity: String = - transaction.query_row("PRAGMA quick_check", [], |row| row.get(0))?; - if integrity != "ok" { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Cache integrity check failed", - ) - .into()); - } - let (base, working): (Vec, Vec) = transaction.query_row( - "SELECT base, working FROM replica WHERE id=1", - [], - |row| Ok((row.get(0)?, row.get(1)?)), - )?; - if validate(&base)? != validate(&working)? { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Cache images belong to different documents", - ) - .into()); - } + validate_images(&transaction)?; if version < SCHEMA_VERSION { schema::migrate(&transaction, version)?; transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?; @@ -360,3 +318,53 @@ fn pending(connection: &Connection) -> Result> { } Ok(edits) } + +fn cache_connection(path: &Path) -> Result { + let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; + connection.busy_timeout(Duration::ZERO)?; + connection.execute_batch( + "PRAGMA locking_mode=EXCLUSIVE; PRAGMA synchronous=EXTRA; PRAGMA fullfsync=ON; PRAGMA foreign_keys=ON;", + )?; + for (name, expected) in [("locking_mode", "exclusive"), ("journal_mode", "delete")] { + let actual: String = connection.pragma_query_value(None, name, |row| row.get(0))?; + if actual != expected { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Unsupported cache locking or journal mode", + ) + .into()); + } + } + for (name, expected) in [("synchronous", 3), ("fullfsync", 1), ("foreign_keys", 1)] { + let actual: i64 = connection.pragma_query_value(None, name, |row| row.get(0))?; + if actual != expected { + return Err(io::Error::new( + io::ErrorKind::Unsupported, + "Required cache synchronization is unavailable", + ) + .into()); + } + } + Ok(connection) +} + +fn validate_images(connection: &Connection) -> Result<()> { + let integrity: String = connection.query_row("PRAGMA quick_check", [], |row| row.get(0))?; + if integrity != "ok" { + return Err( + io::Error::new(io::ErrorKind::InvalidData, "Cache integrity check failed").into(), + ); + } + let (base, working): (Vec, Vec) = + connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { + Ok((row.get(0)?, row.get(1)?)) + })?; + if validate(&base)? != validate(&working)? { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Cache images belong to different documents", + ) + .into()); + } + Ok(()) +} diff --git a/crates/onestore-offline/src/recovery.rs b/crates/onestore-offline/src/recovery.rs new file mode 100644 index 0000000000000000000000000000000000000000..4530d2d1de1aeeebc49aebeb545d2192884af82e --- /dev/null +++ b/crates/onestore-offline/src/recovery.rs @@ -0,0 +1,171 @@ +use super::*; +use rusqlite::backup::{Backup, StepResult}; +use std::{collections::BTreeMap, fs::File}; + +const RECOVERY_ID: u32 = 0x4f4e4552; + +/// Counts and image sizes without notebook text, paths, authors or credentials. +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] +pub struct RecoverySummary { + pub queued_edits: u64, + pub conflicts: u64, + pub uncertain_edits: u64, + pub published_receipts: u64, + pub working_bytes: u64, + pub remote_bytes: u64, +} + +/// Read-only recovery evidence; it cannot publish or acknowledge an edit. +pub struct Recovery { + connection: Connection, +} + +impl Recovery { + /// Opens an exported archive without migration or conversion into a writable replica. + pub fn open(path: impl AsRef) -> Result { + let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY)?; + connection.busy_timeout(Duration::ZERO)?; + connection.execute_batch("BEGIN")?; + let application: u32 = + connection.pragma_query_value(None, "application_id", |row| row.get(0))?; + let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; + if application != RECOVERY_ID || version != SCHEMA_VERSION { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Unrecognized recovery archive or unsupported schema version", + ) + .into()); + } + validate_images(&connection)?; + for edit in pending(&connection)? { + sync::status(&connection, edit.id)?; + } + receipts(&connection)?; + Ok(Self { connection }) + } + + pub fn summary(&self) -> Result { + summary(&self.connection) + } + + pub fn snapshot(&self) -> Result> { + Ok(self + .connection + .query_row("SELECT working FROM replica WHERE id=1", [], |row| { + row.get(0) + })?) + } + + pub fn remote_snapshot(&self) -> Result> { + Ok(self + .connection + .query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) + } + + pub fn pending(&self) -> Result> { + pending(&self.connection) + } + + pub fn status(&self, id: u64) -> Result> { + sync::status(&self.connection, id) + } + + pub fn receipts(&self) -> Result> { + receipts(&self.connection) + } +} + +impl Replica { + /// Summarizes durable state without exposing notebook content. + pub fn recovery_summary(&self) -> Result { + let connection = self + .connection + .lock() + .map_err(|_| io::Error::other("Cache owner panicked"))?; + summary(&connection) + } + + /// Exports a consistent archive to a new local path without changing the live queue. + /// Archives contain notebook content and cannot be opened as writable replicas. + /// An error after publication can leave a complete archive at the destination. + pub fn export_recovery(&self, path: impl AsRef) -> Result<()> { + let path = path.as_ref(); + if path.file_name().is_none() { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "Provide a new recovery archive filename", + ) + .into()); + } + let parent = path + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + .unwrap_or_else(|| Path::new(".")); + let temporary = tempfile::Builder::new() + .prefix(".onestore-recovery-") + .tempfile_in(parent)?; + let mut destination = cache_connection(temporary.path())?; + { + let source = self + .connection + .lock() + .map_err(|_| io::Error::other("Cache owner panicked"))?; + let backup = Backup::new(&source, &mut destination)?; + if backup.step(-1)? != StepResult::Done { + return Err(io::Error::new( + io::ErrorKind::WouldBlock, + "Recovery export could not acquire the database snapshot", + ) + .into()); + } + } + destination.pragma_update(None, "application_id", RECOVERY_ID)?; + destination.close().map_err(|(_, error)| error)?; + temporary.as_file().sync_all()?; + temporary + .persist_noclobber(path) + .map_err(|error| error.error)?; + File::open(parent)?.sync_all()?; + Ok(()) + } +} + +fn summary(connection: &Connection) -> Result { + Ok(connection.query_row( + "SELECT (SELECT count(*) FROM edits), (SELECT count(*) FROM conflicts), + (SELECT count(*) FROM attempt), (SELECT count(*) FROM receipts), + length(working), length(base) FROM replica WHERE id=1", + [], + |row| { + Ok(RecoverySummary { + queued_edits: unsigned(row, 0)?, + conflicts: unsigned(row, 1)?, + uncertain_edits: unsigned(row, 2)?, + published_receipts: unsigned(row, 3)?, + working_bytes: unsigned(row, 4)?, + remote_bytes: unsigned(row, 5)?, + }) + }, + )?) +} + +fn receipts(connection: &Connection) -> Result> { + let mut statement = + connection.prepare("SELECT edit_id, revision FROM receipts ORDER BY edit_id")?; + let mut rows = statement.query([])?; + let mut receipts = BTreeMap::new(); + while let Some(row) = rows.next()? { + receipts.insert(unsigned(row, 0)?, row.get::<_, String>(1)?.parse()?); + } + Ok(receipts) +} + +fn unsigned(row: &rusqlite::Row<'_>, column: usize) -> rusqlite::Result { + u64::try_from(row.get::<_, i64>(column)?).map_err(|error| { + rusqlite::Error::FromSqlConversionFailure( + column, + rusqlite::types::Type::Integer, + Box::new(error), + ) + }) +} diff --git a/crates/onestore-offline/src/sync.rs b/crates/onestore-offline/src/sync.rs index f1bfc013ad8c2a959204ff8dde9f702e487dc9d9..286fca19a5c4bb7456eafbdf1eba69cd789ef114 100644 --- a/crates/onestore-offline/src/sync.rs +++ b/crates/onestore-offline/src/sync.rs @@ -36,45 +36,11 @@ pub enum EditStatus { impl Replica { /// Returns a durable receipt or the persisted state of a locally acknowledged edit. pub fn status(&self, id: u64) -> Result> { - let id = i64::try_from(id).map_err(io::Error::other)?; let connection = self .connection .lock() .map_err(|_| io::Error::other("Cache owner panicked"))?; - if let Some(revision) = connection - .query_row( - "SELECT revision FROM receipts WHERE edit_id=?1", - [id], - |row| row.get::<_, String>(0), - ) - .optional()? - { - return Ok(Some(EditStatus::Published { - revision: revision.parse()?, - })); - } - let record: Option<(Option, Option)> = connection.query_row( - "SELECT attempt.revision, conflicts.kind FROM edits LEFT JOIN attempt ON attempt.edit_id=edits.id LEFT JOIN conflicts ON conflicts.edit_id=edits.id WHERE edits.id=?1", [id], |row| Ok((row.get(0)?,row.get(1)?))).optional()?; - Ok(match record { - None => None, - Some((Some(revision), _)) => Some(EditStatus::AwaitingConfirmation { - revision: revision.parse()?, - }), - Some((None, Some(kind))) => Some(EditStatus::Conflict(match kind { - 0 => ConflictKind::TextChanged, - 1 => ConflictKind::TargetUnavailable, - 2 => ConflictKind::UnsupportedEdit, - 3 => ConflictKind::FormattingChanged, - _ => { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Unknown cached conflict kind", - ) - .into()); - } - })), - Some((None, None)) => Some(EditStatus::Pending), - }) + status(&connection, id) } /// The last observed remote image, retained alongside the complete local working image. @@ -456,3 +422,41 @@ impl Replica { Ok(()) } } + +pub(crate) fn status(connection: &Connection, id: u64) -> Result> { + let id = i64::try_from(id).map_err(io::Error::other)?; + if let Some(revision) = connection + .query_row( + "SELECT revision FROM receipts WHERE edit_id=?1", + [id], + |row| row.get::<_, String>(0), + ) + .optional()? + { + return Ok(Some(EditStatus::Published { + revision: revision.parse()?, + })); + } + let record: Option<(Option, Option)> = connection.query_row( + "SELECT attempt.revision, conflicts.kind FROM edits LEFT JOIN attempt ON attempt.edit_id=edits.id LEFT JOIN conflicts ON conflicts.edit_id=edits.id WHERE edits.id=?1", [id], |row| Ok((row.get(0)?,row.get(1)?))).optional()?; + Ok(match record { + None => None, + Some((Some(revision), _)) => Some(EditStatus::AwaitingConfirmation { + revision: revision.parse()?, + }), + Some((None, Some(kind))) => Some(EditStatus::Conflict(match kind { + 0 => ConflictKind::TextChanged, + 1 => ConflictKind::TargetUnavailable, + 2 => ConflictKind::UnsupportedEdit, + 3 => ConflictKind::FormattingChanged, + _ => { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Unknown cached conflict kind", + ) + .into()); + } + })), + Some((None, None)) => Some(EditStatus::Pending), + }) +} diff --git a/crates/onestore-offline/tests/cache.rs b/crates/onestore-offline/tests/cache.rs index 0623378170bb8b45c642c54fb9692d840c5254cc..7e762d412e19449e22aa147f432f5ef0adde3858 100644 --- a/crates/onestore-offline/tests/cache.rs +++ b/crates/onestore-offline/tests/cache.rs @@ -112,6 +112,67 @@ fn failed_edits_preserve_both_intent_queue_and_working_image() { assert!(replica.pending().unwrap().is_empty()); } +#[test] +fn concurrent_recovery_exports_capture_one_complete_acknowledged_queue() { + use onestore_offline::{Operation, Recovery}; + + let directory = tempfile::tempdir().unwrap(); + let source = onestore::create_section("recovery.one", "base", "Fixture").unwrap(); + let (space, object, _) = target(&source); + let replica = Replica::create(directory.path().join("live.sqlite"), &source).unwrap(); + let start = Barrier::new(4); + std::thread::scope(|scope| { + for writer in 0..3 { + let (replica, start) = (&replica, &start); + scope.spawn(move || { + start.wait(); + for edit in 0..20 { + loop { + let source = replica.snapshot().unwrap(); + match replica.edit_text( + &source, + space, + object, + 0..0, + &format!("[{writer}:{edit}] "), + ) { + Ok(Some(_)) => break, + Err(Error::Io(error)) if error.kind() == ErrorKind::ResourceBusy => { + continue; + } + other => panic!("Unexpected local edit: {other:?}"), + } + } + } + }); + } + start.wait(); + for n in 0..12 { + let path = directory.path().join(format!("recovery-{n}.sqlite")); + replica.export_recovery(&path).unwrap(); + let archive = Recovery::open(path).unwrap(); + let pending = archive.pending().unwrap(); + let mut expected = String::new(); + for edit in pending.iter().rev() { + let Operation::Text(edit) = &edit.operation else { + panic!() + }; + assert_eq!(edit.range, 0..0); + expected.push_str(&edit.replacement); + } + expected.push_str("base"); + assert_eq!(target(&archive.snapshot().unwrap()).2, expected); + assert_eq!(archive.remote_snapshot().unwrap(), source); + assert_eq!( + archive.summary().unwrap().queued_edits, + pending.len() as u64 + ); + assert!(archive.receipts().unwrap().is_empty()); + } + }); + assert_eq!(replica.pending().unwrap().len(), 60); +} + #[test] fn ownership_and_foreign_file_rejection_preserve_existing_data() { let dir = tempfile::tempdir().unwrap(); diff --git a/crates/onestore-offline/tests/sync.rs b/crates/onestore-offline/tests/sync.rs index 16cc13ed761c55efb76e6b6736511abcb7fd4a69..29318e9f84a282db881901fade29ba70c6274043 100644 --- a/crates/onestore-offline/tests/sync.rs +++ b/crates/onestore-offline/tests/sync.rs @@ -128,6 +128,189 @@ fn text(source: &[u8]) -> (ExGuid, ExGuid, String) { .unwrap() } +#[test] +fn recovery_archive_preserves_typed_queue_uncertainty_and_receipts_without_becoming_a_writer() { + use onestore::{Insertion, TextAttribute}; + use onestore_offline::{Recovery, RecoverySummary}; + + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("live.sqlite"); + let archive_path = directory.path().join("recovery.sqlite"); + let source = onestore::create_section("recovery.one", "Original", "Fixture").unwrap(); + let (space, object, _) = text(&source); + let cache = Replica::create(&path, &source).unwrap(); + let published = cache + .edit_text(&source, space, object, 0..0, "Published ") + .unwrap() + .unwrap(); + let mut server = Server::new(&source); + let (_, receipt) = cache.sync_once(&mut server).unwrap().unwrap(); + let source = cache.snapshot().unwrap(); + let store = Store::parse(&source).unwrap(); + let index = RevisionIndex::parse(&store).unwrap(); + let document = Document::parse(&index).unwrap(); + let (page_space, page) = document.pages().unwrap()[0]; + let outline = Insertion::outline(page, 100.0, 200.0, "Outline", "Fixture").unwrap(); + let inserted = cache + .insert(&source, page_space, &outline) + .unwrap() + .unwrap(); + let paragraph = Insertion::paragraph(outline.object(), None, "Recovery 🦀", "Fixture").unwrap(); + cache + .insert(&cache.snapshot().unwrap(), page_space, ¶graph) + .unwrap(); + cache + .format( + &cache.snapshot().unwrap(), + page_space, + paragraph.text_object(), + 0..3, + &[TextAttribute::Bold(true)], + ) + .unwrap(); + cache + .edit_text( + &cache.snapshot().unwrap(), + page_space, + paragraph.text_object(), + 0..0, + "Pending ", + ) + .unwrap(); + server.fault = Fault::UnknownBefore; + assert!(matches!( + cache.sync_once(&mut server), + Err(Error::Remote(CommitError { + state: CommitState::Unknown, + .. + })) + )); + let working = cache.snapshot().unwrap(); + let remote = cache.remote_snapshot().unwrap(); + let pending = cache.pending().unwrap(); + let uncertain = cache.status(inserted).unwrap(); + let source_file = std::fs::read(&path).unwrap(); + let summary = RecoverySummary { + queued_edits: 4, + conflicts: 0, + uncertain_edits: 1, + published_receipts: 1, + working_bytes: working.len() as u64, + remote_bytes: remote.len() as u64, + }; + assert_eq!(cache.recovery_summary().unwrap(), summary); + cache.export_recovery(&archive_path).unwrap(); + assert_eq!(std::fs::read(&path).unwrap(), source_file); + let archive_bytes = std::fs::read(&archive_path).unwrap(); + assert!(Replica::open(&archive_path).is_err()); + assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes); + assert!(Recovery::open(&path).is_err()); + let archive = Recovery::open(&archive_path).unwrap(); + assert_eq!(archive.summary().unwrap(), summary); + assert_eq!(archive.snapshot().unwrap(), working); + assert_eq!(archive.remote_snapshot().unwrap(), remote); + assert_eq!(archive.pending().unwrap(), pending); + assert_eq!(archive.status(published).unwrap(), Some(receipt)); + assert_eq!(archive.status(inserted).unwrap(), uncertain); + for edit in &pending { + assert_eq!( + archive.status(edit.id).unwrap(), + cache.status(edit.id).unwrap() + ); + } + let EditStatus::Published { revision } = receipt else { + panic!() + }; + assert_eq!(archive.receipts().unwrap(), [(published, revision)].into()); + assert_eq!( + archive.status(u64::MAX - 1).unwrap_err().to_string(), + cache.status(u64::MAX - 1).unwrap_err().to_string() + ); + for existing in [&path, &archive_path] { + assert!( + matches!(cache.export_recovery(existing), Err(Error::Io(error)) if error.kind() == io::ErrorKind::AlreadyExists) + ); + } + assert_eq!(std::fs::read(&path).unwrap(), source_file); + assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes); + cache + .edit_text(&working, space, object, 0..0, "Later ") + .unwrap(); + assert_eq!(archive.snapshot().unwrap(), working); + assert_eq!(archive.pending().unwrap(), pending); + drop(archive); + assert_eq!( + Recovery::open(&archive_path).unwrap().summary().unwrap(), + summary + ); + assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + assert_eq!( + std::fs::metadata(&archive_path) + .unwrap() + .permissions() + .mode() + & 0o777, + 0o600 + ); + } + let remaining: Vec<_> = std::fs::read_dir(directory.path()) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .collect(); + assert!( + remaining + .iter() + .all(|name| !name.to_string_lossy().starts_with(".onestore-recovery-")), + "Temporary recovery files remain: {remaining:?}" + ); +} + +#[test] +fn recovery_archive_retains_conflict_images_and_rejects_foreign_or_future_archives() { + use onestore_offline::Recovery; + + let directory = tempfile::tempdir().unwrap(); + let source = onestore::create_section("recovery.one", "Original", "Fixture").unwrap(); + let (space, object, _) = text(&source); + let cache = Replica::create(directory.path().join("live.sqlite"), &source).unwrap(); + let id = cache + .edit_text(&source, space, object, 0..8, "Local") + .unwrap() + .unwrap(); + let changed = onestore::replace_text(&source, space, object, 0..8, "Remote").unwrap(); + let mut server = Server::new(&changed); + let outcome = cache.sync_once(&mut server).unwrap().unwrap(); + assert_eq!( + outcome, + (id, EditStatus::Conflict(ConflictKind::TextChanged)) + ); + let path = directory.path().join("conflict.sqlite"); + cache.export_recovery(&path).unwrap(); + let archive = Recovery::open(&path).unwrap(); + assert_eq!(text(&archive.snapshot().unwrap()).2, "Local"); + assert_eq!(archive.remote_snapshot().unwrap(), changed); + assert_eq!(archive.status(id).unwrap(), Some(outcome.1)); + assert_eq!(archive.summary().unwrap().conflicts, 1); + assert_eq!(archive.summary().unwrap().uncertain_edits, 0); + drop(archive); + for sql in [ + "PRAGMA user_version=99", + "PRAGMA user_version=4; PRAGMA application_id=0", + ] { + let connection = rusqlite::Connection::open(&path).unwrap(); + connection.execute_batch(sql).unwrap(); + drop(connection); + let before = std::fs::read(&path).unwrap(); + assert!(Recovery::open(&path).is_err()); + assert_eq!(std::fs::read(&path).unwrap(), before); + } + assert_eq!(cache.status(id).unwrap(), Some(outcome.1)); + assert_eq!(text(&cache.snapshot().unwrap()).2, "Local"); +} + #[test] fn rebases_multiple_disjoint_remote_changes_and_persists_the_remote_receipt() { let dir = tempfile::tempdir().unwrap();