| author | |
| committer | |
| log | e076b4012533d329f3c9c7aca05c8288d84afb4b |
| tree | 7dfda8504ff0867f96816a816c691a68c8a83da2 |
| parent | b253974a925fe86b3f2ee31e969477fe2ae43fd3 |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
Capture images, typed intents, attempts, conflicts, receipts and ID allocation in a consistent SQLite backup. Publish to a new path and identify archives separately so they cannot become a second active replica. Add content-free summaries and read-only inspection.
Validated with 54 offline tests, four doctests, concurrent export histories, Clippy and iOS device/simulator linking.
Assisted-by: gpt-6-astra7 files changed, 541 insertions(+), 84 deletions(-)
crates/onestore-offline/Cargo.toml+1-3| ... | ... | @@ -10,12 +10,10 @@ smb = ["dep:onestore-smb"] |
| 10 | 10 | [dependencies] |
| 11 | 11 | onestore = { path = "../onestore" } |
| 12 | 12 | onestore-smb = { path = "../onestore-smb", optional = true } |
| 13 | rusqlite = { version = "=0.40.2", features = ["bundled"] } | |
| 13 | rusqlite = { version = "=0.40.2", features = ["bundled", "backup"] } | |
| 14 | 14 | thiserror = "2" |
| 15 | 15 | serde = { version = "1", features = ["derive"] } |
| 16 | 16 | serde_json = "1" |
| 17 | ||
| 18 | [dev-dependencies] | |
| 19 | 17 | tempfile = "3" |
| 20 | 18 | |
| 21 | 19 | [[example]] |
crates/onestore-offline/README.md+32| ... | ... | @@ -189,3 +189,35 @@ The SMB-enabled `smb_offline_client` example is an owned-lab workload for |
| 189 | 189 | acknowledgements, remote publication attempts, persisted receipts and cache reopen |
| 190 | 190 | checks. Its append-specific conflict review policy lives in the test client; |
| 191 | 191 | the library continues to preserve conflicts requiring an explicit decision. |
| 192 | ||
| 193 | ## Recovery archives | |
| 194 | ||
| 195 | `export_recovery(new_path)` captures both complete notebook images, the typed | |
| 196 | queue, uncertain attempts, conflicts, receipts and the edit-ID sequence in one | |
| 197 | SQLite snapshot. It refuses existing destinations and leaves the live queue | |
| 198 | unchanged. Export to a local directory from a background thread: copying holds | |
| 199 | the cache mutex while capturing the database. Failure after the final rename can | |
| 200 | leave a complete archive at the requested path; it never acknowledges a remote | |
| 201 | edit. Archives contain notebook content and use a separate database identity, so | |
| 202 | `Replica::open` rejects them as writable caches. | |
| 203 | ||
| 204 | ```no_run | |
| 205 | use onestore_offline::{Recovery, Replica}; | |
| 206 | # fn example(cache: &Replica) -> Result<(), Box<dyn std::error::Error>> { | |
| 207 | cache.export_recovery("review.sqlite")?; | |
| 208 | let review = Recovery::open("review.sqlite")?; | |
| 209 | let counts = review.summary()?; | |
| 210 | let local = review.snapshot()?; | |
| 211 | let remote = review.remote_snapshot()?; | |
| 212 | let pending = review.pending()?; | |
| 213 | let receipts = review.receipts()?; | |
| 214 | # Ok(()) | |
| 215 | # } | |
| 216 | ``` | |
| 217 | ||
| 218 | `Recovery` provides read-only inspection and no synchronization or restore method. | |
| 219 | `status(id)` preserves the same state interpretation as the live replica. | |
| 220 | `recovery_summary()` on the live replica and `summary()` on an archive return | |
| 221 | counts and byte sizes without notebook text, paths, authors or credentials. | |
| 222 | Opening an archive validates its schema and images without migration. A recovery | |
| 223 | archive is evidence for a reviewed recovery decision, not a second active queue. |
crates/onestore-offline/src/lib.rs+54-46| ... | ... | @@ -10,8 +10,10 @@ use std::{fs::OpenOptions, io, ops::Range, path::Path, sync::Mutex, time::Durati |
| 10 | 10 | |
| 11 | 11 | mod formatting; |
| 12 | 12 | mod rebase; |
| 13 | mod recovery; | |
| 13 | 14 | mod schema; |
| 14 | 15 | pub use formatting::FormatEdit; |
| 16 | pub use recovery::{Recovery, RecoverySummary}; | |
| 15 | 17 | mod sync; |
| 16 | 18 | pub use sync::{ConflictKind, EditStatus, Remote}; |
| 17 | 19 | mod worker; |
| ... | ... | @@ -97,31 +99,7 @@ impl Replica { |
| 97 | 99 | } |
| 98 | 100 | |
| 99 | 101 | fn connect(path: &Path, source: Option<&[u8]>) -> Result<Self> { |
| 100 | let mut connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; | |
| 101 | connection.busy_timeout(Duration::ZERO)?; | |
| 102 | connection.execute_batch( | |
| 103 | "PRAGMA locking_mode=EXCLUSIVE; PRAGMA synchronous=EXTRA; PRAGMA fullfsync=ON; PRAGMA foreign_keys=ON;", | |
| 104 | )?; | |
| 105 | for (name, expected) in [("locking_mode", "exclusive"), ("journal_mode", "delete")] { | |
| 106 | let actual: String = connection.pragma_query_value(None, name, |row| row.get(0))?; | |
| 107 | if actual != expected { | |
| 108 | return Err(io::Error::new( | |
| 109 | io::ErrorKind::InvalidData, | |
| 110 | "Unsupported cache locking or journal mode", | |
| 111 | ) | |
| 112 | .into()); | |
| 113 | } | |
| 114 | } | |
| 115 | for (name, expected) in [("synchronous", 3), ("fullfsync", 1), ("foreign_keys", 1)] { | |
| 116 | let actual: i64 = connection.pragma_query_value(None, name, |row| row.get(0))?; | |
| 117 | if actual != expected { | |
| 118 | return Err(io::Error::new( | |
| 119 | io::ErrorKind::Unsupported, | |
| 120 | "Required cache synchronization is unavailable", | |
| 121 | ) | |
| 122 | .into()); | |
| 123 | } | |
| 124 | } | |
| 102 | let mut connection = cache_connection(path)?; | |
| 125 | 103 | let transaction = connection.transaction_with_behavior(TransactionBehavior::Exclusive)?; |
| 126 | 104 | let application: u32 = |
| 127 | 105 | transaction.pragma_query_value(None, "application_id", |row| row.get(0))?; |
| ... | ... | @@ -159,27 +137,7 @@ impl Replica { |
| 159 | 137 | ) |
| 160 | 138 | .into()); |
| 161 | 139 | } |
| 162 | let integrity: String = | |
| 163 | transaction.query_row("PRAGMA quick_check", [], |row| row.get(0))?; | |
| 164 | if integrity != "ok" { | |
| 165 | return Err(io::Error::new( | |
| 166 | io::ErrorKind::InvalidData, | |
| 167 | "Cache integrity check failed", | |
| 168 | ) | |
| 169 | .into()); | |
| 170 | } | |
| 171 | let (base, working): (Vec<u8>, Vec<u8>) = transaction.query_row( | |
| 172 | "SELECT base, working FROM replica WHERE id=1", | |
| 173 | [], | |
| 174 | |row| Ok((row.get(0)?, row.get(1)?)), | |
| 175 | )?; | |
| 176 | if validate(&base)? != validate(&working)? { | |
| 177 | return Err(io::Error::new( | |
| 178 | io::ErrorKind::InvalidData, | |
| 179 | "Cache images belong to different documents", | |
| 180 | ) | |
| 181 | .into()); | |
| 182 | } | |
| 140 | validate_images(&transaction)?; | |
| 183 | 141 | if version < SCHEMA_VERSION { |
| 184 | 142 | schema::migrate(&transaction, version)?; |
| 185 | 143 | transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?; |
| ... | ... | @@ -360,3 +318,53 @@ fn pending(connection: &Connection) -> Result<Vec<PendingEdit>> { |
| 360 | 318 | } |
| 361 | 319 | Ok(edits) |
| 362 | 320 | } |
| 321 | ||
| 322 | fn cache_connection(path: &Path) -> Result<Connection> { | |
| 323 | let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; | |
| 324 | connection.busy_timeout(Duration::ZERO)?; | |
| 325 | connection.execute_batch( | |
| 326 | "PRAGMA locking_mode=EXCLUSIVE; PRAGMA synchronous=EXTRA; PRAGMA fullfsync=ON; PRAGMA foreign_keys=ON;", | |
| 327 | )?; | |
| 328 | for (name, expected) in [("locking_mode", "exclusive"), ("journal_mode", "delete")] { | |
| 329 | let actual: String = connection.pragma_query_value(None, name, |row| row.get(0))?; | |
| 330 | if actual != expected { | |
| 331 | return Err(io::Error::new( | |
| 332 | io::ErrorKind::InvalidData, | |
| 333 | "Unsupported cache locking or journal mode", | |
| 334 | ) | |
| 335 | .into()); | |
| 336 | } | |
| 337 | } | |
| 338 | for (name, expected) in [("synchronous", 3), ("fullfsync", 1), ("foreign_keys", 1)] { | |
| 339 | let actual: i64 = connection.pragma_query_value(None, name, |row| row.get(0))?; | |
| 340 | if actual != expected { | |
| 341 | return Err(io::Error::new( | |
| 342 | io::ErrorKind::Unsupported, | |
| 343 | "Required cache synchronization is unavailable", | |
| 344 | ) | |
| 345 | .into()); | |
| 346 | } | |
| 347 | } | |
| 348 | Ok(connection) | |
| 349 | } | |
| 350 | ||
| 351 | fn validate_images(connection: &Connection) -> Result<()> { | |
| 352 | let integrity: String = connection.query_row("PRAGMA quick_check", [], |row| row.get(0))?; | |
| 353 | if integrity != "ok" { | |
| 354 | return Err( | |
| 355 | io::Error::new(io::ErrorKind::InvalidData, "Cache integrity check failed").into(), | |
| 356 | ); | |
| 357 | } | |
| 358 | let (base, working): (Vec<u8>, Vec<u8>) = | |
| 359 | connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { | |
| 360 | Ok((row.get(0)?, row.get(1)?)) | |
| 361 | })?; | |
| 362 | if validate(&base)? != validate(&working)? { | |
| 363 | return Err(io::Error::new( | |
| 364 | io::ErrorKind::InvalidData, | |
| 365 | "Cache images belong to different documents", | |
| 366 | ) | |
| 367 | .into()); | |
| 368 | } | |
| 369 | Ok(()) | |
| 370 | } |
crates/onestore-offline/src/recovery.rs created+171| ... | ... | @@ -0,0 +1,171 @@ |
| 1 | use super::*; | |
| 2 | use rusqlite::backup::{Backup, StepResult}; | |
| 3 | use std::{collections::BTreeMap, fs::File}; | |
| 4 | ||
| 5 | const RECOVERY_ID: u32 = 0x4f4e4552; | |
| 6 | ||
| 7 | /// Counts and image sizes without notebook text, paths, authors or credentials. | |
| 8 | #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] | |
| 9 | pub struct RecoverySummary { | |
| 10 | pub queued_edits: u64, | |
| 11 | pub conflicts: u64, | |
| 12 | pub uncertain_edits: u64, | |
| 13 | pub published_receipts: u64, | |
| 14 | pub working_bytes: u64, | |
| 15 | pub remote_bytes: u64, | |
| 16 | } | |
| 17 | ||
| 18 | /// Read-only recovery evidence; it cannot publish or acknowledge an edit. | |
| 19 | pub struct Recovery { | |
| 20 | connection: Connection, | |
| 21 | } | |
| 22 | ||
| 23 | impl Recovery { | |
| 24 | /// Opens an exported archive without migration or conversion into a writable replica. | |
| 25 | pub fn open(path: impl AsRef<Path>) -> Result<Self> { | |
| 26 | let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY)?; | |
| 27 | connection.busy_timeout(Duration::ZERO)?; | |
| 28 | connection.execute_batch("BEGIN")?; | |
| 29 | let application: u32 = | |
| 30 | connection.pragma_query_value(None, "application_id", |row| row.get(0))?; | |
| 31 | let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; | |
| 32 | if application != RECOVERY_ID || version != SCHEMA_VERSION { | |
| 33 | return Err(io::Error::new( | |
| 34 | io::ErrorKind::InvalidData, | |
| 35 | "Unrecognized recovery archive or unsupported schema version", | |
| 36 | ) | |
| 37 | .into()); | |
| 38 | } | |
| 39 | validate_images(&connection)?; | |
| 40 | for edit in pending(&connection)? { | |
| 41 | sync::status(&connection, edit.id)?; | |
| 42 | } | |
| 43 | receipts(&connection)?; | |
| 44 | Ok(Self { connection }) | |
| 45 | } | |
| 46 | ||
| 47 | pub fn summary(&self) -> Result<RecoverySummary> { | |
| 48 | summary(&self.connection) | |
| 49 | } | |
| 50 | ||
| 51 | pub fn snapshot(&self) -> Result<Vec<u8>> { | |
| 52 | Ok(self | |
| 53 | .connection | |
| 54 | .query_row("SELECT working FROM replica WHERE id=1", [], |row| { | |
| 55 | row.get(0) | |
| 56 | })?) | |
| 57 | } | |
| 58 | ||
| 59 | pub fn remote_snapshot(&self) -> Result<Vec<u8>> { | |
| 60 | Ok(self | |
| 61 | .connection | |
| 62 | .query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) | |
| 63 | } | |
| 64 | ||
| 65 | pub fn pending(&self) -> Result<Vec<PendingEdit>> { | |
| 66 | pending(&self.connection) | |
| 67 | } | |
| 68 | ||
| 69 | pub fn status(&self, id: u64) -> Result<Option<EditStatus>> { | |
| 70 | sync::status(&self.connection, id) | |
| 71 | } | |
| 72 | ||
| 73 | pub fn receipts(&self) -> Result<BTreeMap<u64, ExGuid>> { | |
| 74 | receipts(&self.connection) | |
| 75 | } | |
| 76 | } | |
| 77 | ||
| 78 | impl Replica { | |
| 79 | /// Summarizes durable state without exposing notebook content. | |
| 80 | pub fn recovery_summary(&self) -> Result<RecoverySummary> { | |
| 81 | let connection = self | |
| 82 | .connection | |
| 83 | .lock() | |
| 84 | .map_err(|_| io::Error::other("Cache owner panicked"))?; | |
| 85 | summary(&connection) | |
| 86 | } | |
| 87 | ||
| 88 | /// Exports a consistent archive to a new local path without changing the live queue. | |
| 89 | /// Archives contain notebook content and cannot be opened as writable replicas. | |
| 90 | /// An error after publication can leave a complete archive at the destination. | |
| 91 | pub fn export_recovery(&self, path: impl AsRef<Path>) -> Result<()> { | |
| 92 | let path = path.as_ref(); | |
| 93 | if path.file_name().is_none() { | |
| 94 | return Err(io::Error::new( | |
| 95 | io::ErrorKind::InvalidInput, | |
| 96 | "Provide a new recovery archive filename", | |
| 97 | ) | |
| 98 | .into()); | |
| 99 | } | |
| 100 | let parent = path | |
| 101 | .parent() | |
| 102 | .filter(|parent| !parent.as_os_str().is_empty()) | |
| 103 | .unwrap_or_else(|| Path::new(".")); | |
| 104 | let temporary = tempfile::Builder::new() | |
| 105 | .prefix(".onestore-recovery-") | |
| 106 | .tempfile_in(parent)?; | |
| 107 | let mut destination = cache_connection(temporary.path())?; | |
| 108 | { | |
| 109 | let source = self | |
| 110 | .connection | |
| 111 | .lock() | |
| 112 | .map_err(|_| io::Error::other("Cache owner panicked"))?; | |
| 113 | let backup = Backup::new(&source, &mut destination)?; | |
| 114 | if backup.step(-1)? != StepResult::Done { | |
| 115 | return Err(io::Error::new( | |
| 116 | io::ErrorKind::WouldBlock, | |
| 117 | "Recovery export could not acquire the database snapshot", | |
| 118 | ) | |
| 119 | .into()); | |
| 120 | } | |
| 121 | } | |
| 122 | destination.pragma_update(None, "application_id", RECOVERY_ID)?; | |
| 123 | destination.close().map_err(|(_, error)| error)?; | |
| 124 | temporary.as_file().sync_all()?; | |
| 125 | temporary | |
| 126 | .persist_noclobber(path) | |
| 127 | .map_err(|error| error.error)?; | |
| 128 | File::open(parent)?.sync_all()?; | |
| 129 | Ok(()) | |
| 130 | } | |
| 131 | } | |
| 132 | ||
| 133 | fn summary(connection: &Connection) -> Result<RecoverySummary> { | |
| 134 | Ok(connection.query_row( | |
| 135 | "SELECT (SELECT count(*) FROM edits), (SELECT count(*) FROM conflicts), | |
| 136 | (SELECT count(*) FROM attempt), (SELECT count(*) FROM receipts), | |
| 137 | length(working), length(base) FROM replica WHERE id=1", | |
| 138 | [], | |
| 139 | |row| { | |
| 140 | Ok(RecoverySummary { | |
| 141 | queued_edits: unsigned(row, 0)?, | |
| 142 | conflicts: unsigned(row, 1)?, | |
| 143 | uncertain_edits: unsigned(row, 2)?, | |
| 144 | published_receipts: unsigned(row, 3)?, | |
| 145 | working_bytes: unsigned(row, 4)?, | |
| 146 | remote_bytes: unsigned(row, 5)?, | |
| 147 | }) | |
| 148 | }, | |
| 149 | )?) | |
| 150 | } | |
| 151 | ||
| 152 | fn receipts(connection: &Connection) -> Result<BTreeMap<u64, ExGuid>> { | |
| 153 | let mut statement = | |
| 154 | connection.prepare("SELECT edit_id, revision FROM receipts ORDER BY edit_id")?; | |
| 155 | let mut rows = statement.query([])?; | |
| 156 | let mut receipts = BTreeMap::new(); | |
| 157 | while let Some(row) = rows.next()? { | |
| 158 | receipts.insert(unsigned(row, 0)?, row.get::<_, String>(1)?.parse()?); | |
| 159 | } | |
| 160 | Ok(receipts) | |
| 161 | } | |
| 162 | ||
| 163 | fn unsigned(row: &rusqlite::Row<'_>, column: usize) -> rusqlite::Result<u64> { | |
| 164 | u64::try_from(row.get::<_, i64>(column)?).map_err(|error| { | |
| 165 | rusqlite::Error::FromSqlConversionFailure( | |
| 166 | column, | |
| 167 | rusqlite::types::Type::Integer, | |
| 168 | Box::new(error), | |
| 169 | ) | |
| 170 | }) | |
| 171 | } |
crates/onestore-offline/src/sync.rs+39-35| ... | ... | @@ -36,45 +36,11 @@ pub enum EditStatus { |
| 36 | 36 | impl Replica { |
| 37 | 37 | /// Returns a durable receipt or the persisted state of a locally acknowledged edit. |
| 38 | 38 | pub fn status(&self, id: u64) -> Result<Option<EditStatus>> { |
| 39 | let id = i64::try_from(id).map_err(io::Error::other)?; | |
| 40 | 39 | let connection = self |
| 41 | 40 | .connection |
| 42 | 41 | .lock() |
| 43 | 42 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 44 | if let Some(revision) = connection | |
| 45 | .query_row( | |
| 46 | "SELECT revision FROM receipts WHERE edit_id=?1", | |
| 47 | [id], | |
| 48 | |row| row.get::<_, String>(0), | |
| 49 | ) | |
| 50 | .optional()? | |
| 51 | { | |
| 52 | return Ok(Some(EditStatus::Published { | |
| 53 | revision: revision.parse()?, | |
| 54 | })); | |
| 55 | } | |
| 56 | let record: Option<(Option<String>, Option<i64>)> = connection.query_row( | |
| 57 | "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()?; | |
| 58 | Ok(match record { | |
| 59 | None => None, | |
| 60 | Some((Some(revision), _)) => Some(EditStatus::AwaitingConfirmation { | |
| 61 | revision: revision.parse()?, | |
| 62 | }), | |
| 63 | Some((None, Some(kind))) => Some(EditStatus::Conflict(match kind { | |
| 64 | 0 => ConflictKind::TextChanged, | |
| 65 | 1 => ConflictKind::TargetUnavailable, | |
| 66 | 2 => ConflictKind::UnsupportedEdit, | |
| 67 | 3 => ConflictKind::FormattingChanged, | |
| 68 | _ => { | |
| 69 | return Err(io::Error::new( | |
| 70 | io::ErrorKind::InvalidData, | |
| 71 | "Unknown cached conflict kind", | |
| 72 | ) | |
| 73 | .into()); | |
| 74 | } | |
| 75 | })), | |
| 76 | Some((None, None)) => Some(EditStatus::Pending), | |
| 77 | }) | |
| 43 | status(&connection, id) | |
| 78 | 44 | } |
| 79 | 45 | |
| 80 | 46 | /// The last observed remote image, retained alongside the complete local working image. |
| ... | ... | @@ -456,3 +422,41 @@ impl Replica { |
| 456 | 422 | Ok(()) |
| 457 | 423 | } |
| 458 | 424 | } |
| 425 | ||
| 426 | pub(crate) fn status(connection: &Connection, id: u64) -> Result<Option<EditStatus>> { | |
| 427 | let id = i64::try_from(id).map_err(io::Error::other)?; | |
| 428 | if let Some(revision) = connection | |
| 429 | .query_row( | |
| 430 | "SELECT revision FROM receipts WHERE edit_id=?1", | |
| 431 | [id], | |
| 432 | |row| row.get::<_, String>(0), | |
| 433 | ) | |
| 434 | .optional()? | |
| 435 | { | |
| 436 | return Ok(Some(EditStatus::Published { | |
| 437 | revision: revision.parse()?, | |
| 438 | })); | |
| 439 | } | |
| 440 | let record: Option<(Option<String>, Option<i64>)> = connection.query_row( | |
| 441 | "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()?; | |
| 442 | Ok(match record { | |
| 443 | None => None, | |
| 444 | Some((Some(revision), _)) => Some(EditStatus::AwaitingConfirmation { | |
| 445 | revision: revision.parse()?, | |
| 446 | }), | |
| 447 | Some((None, Some(kind))) => Some(EditStatus::Conflict(match kind { | |
| 448 | 0 => ConflictKind::TextChanged, | |
| 449 | 1 => ConflictKind::TargetUnavailable, | |
| 450 | 2 => ConflictKind::UnsupportedEdit, | |
| 451 | 3 => ConflictKind::FormattingChanged, | |
| 452 | _ => { | |
| 453 | return Err(io::Error::new( | |
| 454 | io::ErrorKind::InvalidData, | |
| 455 | "Unknown cached conflict kind", | |
| 456 | ) | |
| 457 | .into()); | |
| 458 | } | |
| 459 | })), | |
| 460 | Some((None, None)) => Some(EditStatus::Pending), | |
| 461 | }) | |
| 462 | } |
crates/onestore-offline/tests/cache.rs+61| ... | ... | @@ -112,6 +112,67 @@ fn failed_edits_preserve_both_intent_queue_and_working_image() { |
| 112 | 112 | assert!(replica.pending().unwrap().is_empty()); |
| 113 | 113 | } |
| 114 | 114 | |
| 115 | #[test] | |
| 116 | fn concurrent_recovery_exports_capture_one_complete_acknowledged_queue() { | |
| 117 | use onestore_offline::{Operation, Recovery}; | |
| 118 | ||
| 119 | let directory = tempfile::tempdir().unwrap(); | |
| 120 | let source = onestore::create_section("recovery.one", "base", "Fixture").unwrap(); | |
| 121 | let (space, object, _) = target(&source); | |
| 122 | let replica = Replica::create(directory.path().join("live.sqlite"), &source).unwrap(); | |
| 123 | let start = Barrier::new(4); | |
| 124 | std::thread::scope(|scope| { | |
| 125 | for writer in 0..3 { | |
| 126 | let (replica, start) = (&replica, &start); | |
| 127 | scope.spawn(move || { | |
| 128 | start.wait(); | |
| 129 | for edit in 0..20 { | |
| 130 | loop { | |
| 131 | let source = replica.snapshot().unwrap(); | |
| 132 | match replica.edit_text( | |
| 133 | &source, | |
| 134 | space, | |
| 135 | object, | |
| 136 | 0..0, | |
| 137 | &format!("[{writer}:{edit}] "), | |
| 138 | ) { | |
| 139 | Ok(Some(_)) => break, | |
| 140 | Err(Error::Io(error)) if error.kind() == ErrorKind::ResourceBusy => { | |
| 141 | continue; | |
| 142 | } | |
| 143 | other => panic!("Unexpected local edit: {other:?}"), | |
| 144 | } | |
| 145 | } | |
| 146 | } | |
| 147 | }); | |
| 148 | } | |
| 149 | start.wait(); | |
| 150 | for n in 0..12 { | |
| 151 | let path = directory.path().join(format!("recovery-{n}.sqlite")); | |
| 152 | replica.export_recovery(&path).unwrap(); | |
| 153 | let archive = Recovery::open(path).unwrap(); | |
| 154 | let pending = archive.pending().unwrap(); | |
| 155 | let mut expected = String::new(); | |
| 156 | for edit in pending.iter().rev() { | |
| 157 | let Operation::Text(edit) = &edit.operation else { | |
| 158 | panic!() | |
| 159 | }; | |
| 160 | assert_eq!(edit.range, 0..0); | |
| 161 | expected.push_str(&edit.replacement); | |
| 162 | } | |
| 163 | expected.push_str("base"); | |
| 164 | assert_eq!(target(&archive.snapshot().unwrap()).2, expected); | |
| 165 | assert_eq!(archive.remote_snapshot().unwrap(), source); | |
| 166 | assert_eq!( | |
| 167 | archive.summary().unwrap().queued_edits, | |
| 168 | pending.len() as u64 | |
| 169 | ); | |
| 170 | assert!(archive.receipts().unwrap().is_empty()); | |
| 171 | } | |
| 172 | }); | |
| 173 | assert_eq!(replica.pending().unwrap().len(), 60); | |
| 174 | } | |
| 175 | ||
| 115 | 176 | #[test] |
| 116 | 177 | fn ownership_and_foreign_file_rejection_preserve_existing_data() { |
| 117 | 178 | let dir = tempfile::tempdir().unwrap(); |
crates/onestore-offline/tests/sync.rs+183| ... | ... | @@ -128,6 +128,189 @@ fn text(source: &[u8]) -> (ExGuid, ExGuid, String) { |
| 128 | 128 | .unwrap() |
| 129 | 129 | } |
| 130 | 130 | |
| 131 | #[test] | |
| 132 | fn recovery_archive_preserves_typed_queue_uncertainty_and_receipts_without_becoming_a_writer() { | |
| 133 | use onestore::{Insertion, TextAttribute}; | |
| 134 | use onestore_offline::{Recovery, RecoverySummary}; | |
| 135 | ||
| 136 | let directory = tempfile::tempdir().unwrap(); | |
| 137 | let path = directory.path().join("live.sqlite"); | |
| 138 | let archive_path = directory.path().join("recovery.sqlite"); | |
| 139 | let source = onestore::create_section("recovery.one", "Original", "Fixture").unwrap(); | |
| 140 | let (space, object, _) = text(&source); | |
| 141 | let cache = Replica::create(&path, &source).unwrap(); | |
| 142 | let published = cache | |
| 143 | .edit_text(&source, space, object, 0..0, "Published ") | |
| 144 | .unwrap() | |
| 145 | .unwrap(); | |
| 146 | let mut server = Server::new(&source); | |
| 147 | let (_, receipt) = cache.sync_once(&mut server).unwrap().unwrap(); | |
| 148 | let source = cache.snapshot().unwrap(); | |
| 149 | let store = Store::parse(&source).unwrap(); | |
| 150 | let index = RevisionIndex::parse(&store).unwrap(); | |
| 151 | let document = Document::parse(&index).unwrap(); | |
| 152 | let (page_space, page) = document.pages().unwrap()[0]; | |
| 153 | let outline = Insertion::outline(page, 100.0, 200.0, "Outline", "Fixture").unwrap(); | |
| 154 | let inserted = cache | |
| 155 | .insert(&source, page_space, &outline) | |
| 156 | .unwrap() | |
| 157 | .unwrap(); | |
| 158 | let paragraph = Insertion::paragraph(outline.object(), None, "Recovery 🦀", "Fixture").unwrap(); | |
| 159 | cache | |
| 160 | .insert(&cache.snapshot().unwrap(), page_space, &paragraph) | |
| 161 | .unwrap(); | |
| 162 | cache | |
| 163 | .format( | |
| 164 | &cache.snapshot().unwrap(), | |
| 165 | page_space, | |
| 166 | paragraph.text_object(), | |
| 167 | 0..3, | |
| 168 | &[TextAttribute::Bold(true)], | |
| 169 | ) | |
| 170 | .unwrap(); | |
| 171 | cache | |
| 172 | .edit_text( | |
| 173 | &cache.snapshot().unwrap(), | |
| 174 | page_space, | |
| 175 | paragraph.text_object(), | |
| 176 | 0..0, | |
| 177 | "Pending ", | |
| 178 | ) | |
| 179 | .unwrap(); | |
| 180 | server.fault = Fault::UnknownBefore; | |
| 181 | assert!(matches!( | |
| 182 | cache.sync_once(&mut server), | |
| 183 | Err(Error::Remote(CommitError { | |
| 184 | state: CommitState::Unknown, | |
| 185 | .. | |
| 186 | })) | |
| 187 | )); | |
| 188 | let working = cache.snapshot().unwrap(); | |
| 189 | let remote = cache.remote_snapshot().unwrap(); | |
| 190 | let pending = cache.pending().unwrap(); | |
| 191 | let uncertain = cache.status(inserted).unwrap(); | |
| 192 | let source_file = std::fs::read(&path).unwrap(); | |
| 193 | let summary = RecoverySummary { | |
| 194 | queued_edits: 4, | |
| 195 | conflicts: 0, | |
| 196 | uncertain_edits: 1, | |
| 197 | published_receipts: 1, | |
| 198 | working_bytes: working.len() as u64, | |
| 199 | remote_bytes: remote.len() as u64, | |
| 200 | }; | |
| 201 | assert_eq!(cache.recovery_summary().unwrap(), summary); | |
| 202 | cache.export_recovery(&archive_path).unwrap(); | |
| 203 | assert_eq!(std::fs::read(&path).unwrap(), source_file); | |
| 204 | let archive_bytes = std::fs::read(&archive_path).unwrap(); | |
| 205 | assert!(Replica::open(&archive_path).is_err()); | |
| 206 | assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes); | |
| 207 | assert!(Recovery::open(&path).is_err()); | |
| 208 | let archive = Recovery::open(&archive_path).unwrap(); | |
| 209 | assert_eq!(archive.summary().unwrap(), summary); | |
| 210 | assert_eq!(archive.snapshot().unwrap(), working); | |
| 211 | assert_eq!(archive.remote_snapshot().unwrap(), remote); | |
| 212 | assert_eq!(archive.pending().unwrap(), pending); | |
| 213 | assert_eq!(archive.status(published).unwrap(), Some(receipt)); | |
| 214 | assert_eq!(archive.status(inserted).unwrap(), uncertain); | |
| 215 | for edit in &pending { | |
| 216 | assert_eq!( | |
| 217 | archive.status(edit.id).unwrap(), | |
| 218 | cache.status(edit.id).unwrap() | |
| 219 | ); | |
| 220 | } | |
| 221 | let EditStatus::Published { revision } = receipt else { | |
| 222 | panic!() | |
| 223 | }; | |
| 224 | assert_eq!(archive.receipts().unwrap(), [(published, revision)].into()); | |
| 225 | assert_eq!( | |
| 226 | archive.status(u64::MAX - 1).unwrap_err().to_string(), | |
| 227 | cache.status(u64::MAX - 1).unwrap_err().to_string() | |
| 228 | ); | |
| 229 | for existing in [&path, &archive_path] { | |
| 230 | assert!( | |
| 231 | matches!(cache.export_recovery(existing), Err(Error::Io(error)) if error.kind() == io::ErrorKind::AlreadyExists) | |
| 232 | ); | |
| 233 | } | |
| 234 | assert_eq!(std::fs::read(&path).unwrap(), source_file); | |
| 235 | assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes); | |
| 236 | cache | |
| 237 | .edit_text(&working, space, object, 0..0, "Later ") | |
| 238 | .unwrap(); | |
| 239 | assert_eq!(archive.snapshot().unwrap(), working); | |
| 240 | assert_eq!(archive.pending().unwrap(), pending); | |
| 241 | drop(archive); | |
| 242 | assert_eq!( | |
| 243 | Recovery::open(&archive_path).unwrap().summary().unwrap(), | |
| 244 | summary | |
| 245 | ); | |
| 246 | assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes); | |
| 247 | #[cfg(unix)] | |
| 248 | { | |
| 249 | use std::os::unix::fs::PermissionsExt; | |
| 250 | assert_eq!( | |
| 251 | std::fs::metadata(&archive_path) | |
| 252 | .unwrap() | |
| 253 | .permissions() | |
| 254 | .mode() | |
| 255 | & 0o777, | |
| 256 | 0o600 | |
| 257 | ); | |
| 258 | } | |
| 259 | let remaining: Vec<_> = std::fs::read_dir(directory.path()) | |
| 260 | .unwrap() | |
| 261 | .map(|entry| entry.unwrap().file_name()) | |
| 262 | .collect(); | |
| 263 | assert!( | |
| 264 | remaining | |
| 265 | .iter() | |
| 266 | .all(|name| !name.to_string_lossy().starts_with(".onestore-recovery-")), | |
| 267 | "Temporary recovery files remain: {remaining:?}" | |
| 268 | ); | |
| 269 | } | |
| 270 | ||
| 271 | #[test] | |
| 272 | fn recovery_archive_retains_conflict_images_and_rejects_foreign_or_future_archives() { | |
| 273 | use onestore_offline::Recovery; | |
| 274 | ||
| 275 | let directory = tempfile::tempdir().unwrap(); | |
| 276 | let source = onestore::create_section("recovery.one", "Original", "Fixture").unwrap(); | |
| 277 | let (space, object, _) = text(&source); | |
| 278 | let cache = Replica::create(directory.path().join("live.sqlite"), &source).unwrap(); | |
| 279 | let id = cache | |
| 280 | .edit_text(&source, space, object, 0..8, "Local") | |
| 281 | .unwrap() | |
| 282 | .unwrap(); | |
| 283 | let changed = onestore::replace_text(&source, space, object, 0..8, "Remote").unwrap(); | |
| 284 | let mut server = Server::new(&changed); | |
| 285 | let outcome = cache.sync_once(&mut server).unwrap().unwrap(); | |
| 286 | assert_eq!( | |
| 287 | outcome, | |
| 288 | (id, EditStatus::Conflict(ConflictKind::TextChanged)) | |
| 289 | ); | |
| 290 | let path = directory.path().join("conflict.sqlite"); | |
| 291 | cache.export_recovery(&path).unwrap(); | |
| 292 | let archive = Recovery::open(&path).unwrap(); | |
| 293 | assert_eq!(text(&archive.snapshot().unwrap()).2, "Local"); | |
| 294 | assert_eq!(archive.remote_snapshot().unwrap(), changed); | |
| 295 | assert_eq!(archive.status(id).unwrap(), Some(outcome.1)); | |
| 296 | assert_eq!(archive.summary().unwrap().conflicts, 1); | |
| 297 | assert_eq!(archive.summary().unwrap().uncertain_edits, 0); | |
| 298 | drop(archive); | |
| 299 | for sql in [ | |
| 300 | "PRAGMA user_version=99", | |
| 301 | "PRAGMA user_version=4; PRAGMA application_id=0", | |
| 302 | ] { | |
| 303 | let connection = rusqlite::Connection::open(&path).unwrap(); | |
| 304 | connection.execute_batch(sql).unwrap(); | |
| 305 | drop(connection); | |
| 306 | let before = std::fs::read(&path).unwrap(); | |
| 307 | assert!(Recovery::open(&path).is_err()); | |
| 308 | assert_eq!(std::fs::read(&path).unwrap(), before); | |
| 309 | } | |
| 310 | assert_eq!(cache.status(id).unwrap(), Some(outcome.1)); | |
| 311 | assert_eq!(text(&cache.snapshot().unwrap()).2, "Local"); | |
| 312 | } | |
| 313 | ||
| 131 | 314 | #[test] |
| 132 | 315 | fn rebases_multiple_disjoint_remote_changes_and_persists_the_remote_receipt() { |
| 133 | 316 | let dir = tempfile::tempdir().unwrap(); |