| 1 | //! Converting schema-14 caches: whole-page edits become ops, verified page by page against |
| 2 | //! the old working image, with a recovery archive first and nothing changed on a mismatch; |
| 3 | //! and schema-15 caches, whose batches lose the review conflict. |
| 4 | |
| 5 | #[path = "../../onestore/tests/support/ops.rs"] |
| 6 | mod ops; |
| 7 | |
| 8 | use notebook::{EditStatus, Replica}; |
| 9 | use onestore::{ExGuid, PageCreation, PageEdit, op::SectionOp, page::Page}; |
| 10 | use rusqlite::{Connection, params}; |
| 11 | use serde_json::json; |
| 12 | use std::path::Path; |
| 13 | |
| 14 | #[path = "support/server.rs"] |
| 15 | mod server; |
| 16 | use server::*; |
| 17 | #[path = "support/model_ops.rs"] |
| 18 | mod model_ops; |
| 19 | |
| 20 | const SOURCE: &[u8] = include_bytes!("../../../corpus/outline-edit/before/notebook/synthetic.one"); |
| 21 | |
| 22 | /// The body pages of the fixture with their first body text. |
| 23 | fn body_pages(source: &[u8]) -> Vec<(ExGuid, ExGuid)> { |
| 24 | pages(source) |
| 25 | .into_iter() |
| 26 | .filter_map(|(space, page)| { |
| 27 | let text = page.objects.iter().find_map(|object| match object { |
| 28 | onestore::page::PageObject::Outline(outline) => { |
| 29 | outline.paragraphs.first()?.text().map(|text| text.id) |
| 30 | } |
| 31 | _ => None, |
| 32 | })?; |
| 33 | Some((space, text)) |
| 34 | }) |
| 35 | .collect() |
| 36 | } |
| 37 | |
| 38 | /// A schema-14 queue entry. |
| 39 | enum Queued { |
| 40 | Page(ExGuid, Page), |
| 41 | Create(PageCreation), |
| 42 | Pages(Vec<PageEdit>), |
| 43 | Delete(Vec<ExGuid>), |
| 44 | } |
| 45 | |
| 46 | /// Writes a schema-14 cache holding `queue` on `base`, as that schema stored it; returns |
| 47 | /// the working image. |
| 48 | fn v14(path: &Path, base: &[u8], queue: &[Queued], working: Option<&[u8]>) -> Vec<u8> { |
| 49 | let mut image = base.to_vec(); |
| 50 | let mut edits = Vec::new(); |
| 51 | for entry in queue { |
| 52 | let (space, operation, next) = match entry { |
| 53 | Queued::Page(space, after) => { |
| 54 | let before = model_ops::page_of(&image, *space); |
| 55 | let next = ops::saved(&image, *space, after).unwrap(); |
| 56 | ( |
| 57 | *space, |
| 58 | json!({"Page": {"before": before, "after": after, "author": "Author"}}), |
| 59 | next.as_slice().to_vec(), |
| 60 | ) |
| 61 | } |
| 62 | Queued::Create(creation) => ( |
| 63 | creation.space(), |
| 64 | json!({"CreatePage": creation}), |
| 65 | ops::section_op(&image, SectionOp::Create(creation.clone())) |
| 66 | .unwrap() |
| 67 | .as_bytes() |
| 68 | .to_vec(), |
| 69 | ), |
| 70 | Queued::Pages(batch) => { |
| 71 | let observed: Vec<_> = pages(&image) |
| 72 | .into_iter() |
| 73 | .map(|(space, _)| (space, 1)) |
| 74 | .collect(); |
| 75 | ( |
| 76 | root(&image), |
| 77 | json!({"Pages": {"edits": batch, "observed": observed}}), |
| 78 | ops::section_op(&image, SectionOp::Pages(batch.to_vec())) |
| 79 | .unwrap() |
| 80 | .as_bytes() |
| 81 | .to_vec(), |
| 82 | ) |
| 83 | } |
| 84 | Queued::Delete(spaces) => ( |
| 85 | root(&image), |
| 86 | json!({"DeletePages": spaces}), |
| 87 | ops::section_op(&image, SectionOp::Delete(spaces.to_vec())) |
| 88 | .unwrap() |
| 89 | .as_bytes() |
| 90 | .to_vec(), |
| 91 | ), |
| 92 | }; |
| 93 | edits.push((space, operation)); |
| 94 | image = next; |
| 95 | } |
| 96 | let working = working.map_or(image, <[u8]>::to_vec); |
| 97 | std::fs::File::create_new(path).unwrap(); |
| 98 | let connection = Connection::open(path).unwrap(); |
| 99 | connection |
| 100 | .execute_batch( |
| 101 | "PRAGMA application_id=1330529615; PRAGMA user_version=14; |
| 102 | CREATE TABLE replica (id INTEGER PRIMARY KEY CHECK(id=1), base BLOB NOT NULL, working BLOB NOT NULL) STRICT; |
| 103 | CREATE TABLE edits (id INTEGER PRIMARY KEY AUTOINCREMENT CHECK(id>0), space TEXT NOT NULL, operation TEXT NOT NULL) STRICT; |
| 104 | CREATE TABLE attempt (id INTEGER PRIMARY KEY CHECK(id=1), edit_id INTEGER NOT NULL UNIQUE REFERENCES edits(id) ON DELETE CASCADE, revisions TEXT NOT NULL) STRICT; |
| 105 | CREATE TABLE receipts (edit_id INTEGER PRIMARY KEY CHECK(edit_id>0), revision TEXT NOT NULL) STRICT; |
| 106 | CREATE TABLE archived (edit_id INTEGER PRIMARY KEY CHECK(edit_id>0), archive TEXT NOT NULL) STRICT; |
| 107 | CREATE TABLE conflicts (edit_id INTEGER PRIMARY KEY REFERENCES edits(id) ON DELETE CASCADE, kind INTEGER NOT NULL CHECK(kind BETWEEN 0 AND 3)) STRICT; |
| 108 | CREATE TABLE assets (name TEXT PRIMARY KEY NOT NULL, data BLOB NOT NULL, sha256 BLOB NOT NULL CHECK(length(sha256)=32)) STRICT;", |
| 109 | ) |
| 110 | .unwrap(); |
| 111 | // Schema 14 kept the working image as its length and the runs differing from the base. |
| 112 | let mut patch = (working.len() as u64).to_le_bytes().to_vec(); |
| 113 | patch.extend_from_slice(&0_u64.to_le_bytes()); |
| 114 | patch.extend_from_slice(&(working.len() as u64).to_le_bytes()); |
| 115 | patch.extend_from_slice(&working); |
| 116 | connection |
| 117 | .execute( |
| 118 | "INSERT INTO replica VALUES (1, ?1, ?2)", |
| 119 | params![base, patch], |
| 120 | ) |
| 121 | .unwrap(); |
| 122 | connection |
| 123 | .execute( |
| 124 | "INSERT INTO receipts VALUES (1, ?1)", |
| 125 | [ExGuid { |
| 126 | guid: [7; 16], |
| 127 | n: 1, |
| 128 | } |
| 129 | .to_string()], |
| 130 | ) |
| 131 | .unwrap(); |
| 132 | connection |
| 133 | .execute( |
| 134 | "INSERT INTO sqlite_sequence(name, seq) VALUES ('edits', 1)", |
| 135 | [], |
| 136 | ) |
| 137 | .unwrap(); |
| 138 | for (space, operation) in edits { |
| 139 | connection |
| 140 | .execute( |
| 141 | "INSERT INTO edits(space, operation) VALUES (?1, ?2)", |
| 142 | params![space.to_string(), operation.to_string()], |
| 143 | ) |
| 144 | .unwrap(); |
| 145 | } |
| 146 | working |
| 147 | } |
| 148 | |
| 149 | fn root(image: &[u8]) -> ExGuid { |
| 150 | onestore::RevisionIndex::parse(&onestore::Store::parse(image).unwrap()) |
| 151 | .unwrap() |
| 152 | .root |
| 153 | } |
| 154 | |
| 155 | fn appended(source: &[u8], space: ExGuid, text: ExGuid, suffix: &str) -> Page { |
| 156 | let mut page = model_ops::page_of(source, space); |
| 157 | let end = model_ops::paragraph_with(&page, text) |
| 158 | .unwrap() |
| 159 | .text() |
| 160 | .unwrap() |
| 161 | .text |
| 162 | .text() |
| 163 | .encode_utf16() |
| 164 | .count() as u32; |
| 165 | model_ops::replace_text(&mut page, text, end..end, suffix); |
| 166 | page |
| 167 | } |
| 168 | |
| 169 | fn archive(path: &Path) -> std::path::PathBuf { |
| 170 | let mut archive = path.as_os_str().to_owned(); |
| 171 | archive.push(".v14-recovery"); |
| 172 | archive.into() |
| 173 | } |
| 174 | |
| 175 | #[test] |
| 176 | fn a_schema_14_queue_converts_to_ops_that_reach_its_working_pages_and_publish() { |
| 177 | let body = body_pages(SOURCE); |
| 178 | let ((a, a_text), (b, b_text)) = (body[0], body[1]); |
| 179 | let directory = tempfile::tempdir().unwrap(); |
| 180 | let path = directory.path().join("cache.sqlite"); |
| 181 | let creation = PageCreation::new(None, Some("Created offline"), "Author").unwrap(); |
| 182 | let first = appended(SOURCE, a, a_text, " one"); |
| 183 | let working = v14( |
| 184 | &path, |
| 185 | SOURCE, |
| 186 | &[ |
| 187 | Queued::Page(a, first.clone()), |
| 188 | Queued::Page(b, appended(SOURCE, b, b_text, " other")), |
| 189 | Queued::Create(creation.clone()), |
| 190 | Queued::Pages(vec![PageEdit::move_to(b, Some(a), 1).unwrap()]), |
| 191 | Queued::Page( |
| 192 | a, |
| 193 | appended( |
| 194 | ops::saved(SOURCE, a, &first).unwrap().as_slice(), |
| 195 | a, |
| 196 | a_text, |
| 197 | " two", |
| 198 | ), |
| 199 | ), |
| 200 | Queued::Delete(vec![body[2].0]), |
| 201 | ], |
| 202 | None, |
| 203 | ); |
| 204 | let cache = Replica::open(&path).unwrap(); |
| 205 | assert!(archive(&path).exists()); |
| 206 | let pending = cache.pending().unwrap(); |
| 207 | assert_eq!( |
| 208 | pending.iter().map(|edit| edit.id).collect::<Vec<_>>(), |
| 209 | [2, 3, 4, 5, 6, 7] |
| 210 | ); |
| 211 | assert!(pending.iter().all(|edit| !edit.edit.ops.is_empty())); |
| 212 | assert_eq!( |
| 213 | cache.status(1).unwrap(), |
| 214 | Some(EditStatus::Published { |
| 215 | revision: ExGuid { |
| 216 | guid: [7; 16], |
| 217 | n: 1 |
| 218 | } |
| 219 | }) |
| 220 | ); |
| 221 | let local = pages(&snapshot(&cache)); |
| 222 | assert_eq!(local, pages(&working)); |
| 223 | drop(cache); |
| 224 | // A converted cache opens as it is. |
| 225 | let cache = Replica::open(&path).unwrap(); |
| 226 | assert_eq!(pages(&snapshot(&cache)), local); |
| 227 | let next = cache |
| 228 | .apply( |
| 229 | "Author", |
| 230 | onestore::op::Edit { |
| 231 | at: 133_000_000_000_000_000, |
| 232 | ops: vec![onestore::op::Op::Page { |
| 233 | space: b, |
| 234 | op: onestore::op::PageOp::Text { |
| 235 | text: b_text, |
| 236 | range: 0..0, |
| 237 | with: " three".into(), |
| 238 | }, |
| 239 | }], |
| 240 | }, |
| 241 | ) |
| 242 | .unwrap(); |
| 243 | assert!(next > 7); |
| 244 | let mut server = Server::new(SOURCE); |
| 245 | assert!(matches!( |
| 246 | cache.sync_once(&mut server).unwrap().edit, |
| 247 | Some((id, EditStatus::Published { .. })) if id == next |
| 248 | )); |
| 249 | let published = pages(&server.durable); |
| 250 | assert_eq!(published, pages(&snapshot(&cache))); |
| 251 | assert!( |
| 252 | published |
| 253 | .iter() |
| 254 | .any(|(space, _)| *space == creation.space()) |
| 255 | ); |
| 256 | assert!(published.iter().all(|(space, _)| *space != body[2].0)); |
| 257 | } |
| 258 | |
| 259 | #[test] |
| 260 | fn an_uncertain_attempt_keeps_its_state_and_a_conflict_becomes_a_conflict_page() { |
| 261 | let body = body_pages(SOURCE); |
| 262 | let (a, a_text) = body[0]; |
| 263 | let directory = tempfile::tempdir().unwrap(); |
| 264 | |
| 265 | // The oldest edit was attempted: its schema-14 evidence names the revision it published. |
| 266 | let path = directory.path().join("attempted.sqlite"); |
| 267 | let after = appended(SOURCE, a, a_text, " attempted"); |
| 268 | let published = ops::saved(SOURCE, a, &after).unwrap(); |
| 269 | let revision = |
| 270 | onestore::RevisionIndex::parse(&onestore::Store::parse(published.as_slice()).unwrap()) |
| 271 | .unwrap() |
| 272 | .active(a) |
| 273 | .unwrap(); |
| 274 | v14(&path, SOURCE, &[Queued::Page(a, after.clone())], None); |
| 275 | Connection::open(&path) |
| 276 | .unwrap() |
| 277 | .execute( |
| 278 | "INSERT INTO attempt VALUES (1, 2, ?1)", |
| 279 | [json!({a.to_string(): revision.to_string()}).to_string()], |
| 280 | ) |
| 281 | .unwrap(); |
| 282 | let cache = Replica::open(&path).unwrap(); |
| 283 | let awaiting = EditStatus::AwaitingConfirmation { revision }; |
| 284 | assert_eq!(cache.status(2).unwrap(), Some(awaiting.clone())); |
| 285 | let mut unchanged = Server::new(SOURCE); |
| 286 | assert_eq!( |
| 287 | cache.sync_once(&mut unchanged).unwrap().edit, |
| 288 | Some((2, awaiting)) |
| 289 | ); |
| 290 | assert_eq!(unchanged.publications, 0); |
| 291 | let mut landed = Server::new(published.as_slice()); |
| 292 | assert_eq!( |
| 293 | cache.sync_once(&mut landed).unwrap().edit, |
| 294 | Some((2, EditStatus::Published { revision })) |
| 295 | ); |
| 296 | assert_eq!(landed.publications, 0); |
| 297 | drop(cache); |
| 298 | |
| 299 | // The oldest edit was in conflict: the conversion keeps the remote page the conflict was |
| 300 | // recorded against and queues the local version as its conflict page. |
| 301 | let path = directory.path().join("conflicted.sqlite"); |
| 302 | let remote = ops::saved(SOURCE, a, &appended(SOURCE, a, a_text, " remote")) |
| 303 | .unwrap() |
| 304 | .as_slice() |
| 305 | .to_vec(); |
| 306 | let local = appended(SOURCE, a, a_text, " local"); |
| 307 | let working = ops::saved(SOURCE, a, &local).unwrap().as_slice().to_vec(); |
| 308 | { |
| 309 | // Schema 14 stored the remote as the base and kept the local working image. |
| 310 | v14(&path, &remote, &[], Some(&working)); |
| 311 | let connection = Connection::open(&path).unwrap(); |
| 312 | connection |
| 313 | .execute( |
| 314 | "INSERT INTO edits(space, operation) VALUES (?1, ?2)", |
| 315 | params![ |
| 316 | a.to_string(), |
| 317 | json!({"Page": {"before": model_ops::page_of(SOURCE, a), "after": local, "author": "Author"}}).to_string() |
| 318 | ], |
| 319 | ) |
| 320 | .unwrap(); |
| 321 | connection |
| 322 | .execute("INSERT INTO conflicts VALUES (2, 3)", []) |
| 323 | .unwrap(); |
| 324 | } |
| 325 | let cache = Replica::open(&path).unwrap(); |
| 326 | assert_eq!(cache.status(2).unwrap(), Some(EditStatus::Pending)); |
| 327 | assert_eq!( |
| 328 | model_ops::page_of(&snapshot(&cache), a), |
| 329 | model_ops::page_of(&remote, a) |
| 330 | ); |
| 331 | let mut server = Server::new(&remote); |
| 332 | assert!(matches!( |
| 333 | cache.sync_once(&mut server).unwrap().edit, |
| 334 | Some((2, EditStatus::Published { .. })) |
| 335 | )); |
| 336 | assert_eq!( |
| 337 | model_ops::page_of(&server.durable, a), |
| 338 | model_ops::page_of(&remote, a) |
| 339 | ); |
| 340 | assert_eq!( |
| 341 | conflicts(&server.durable), |
| 342 | [( |
| 343 | a, |
| 344 | vec![( |
| 345 | "Author".to_owned(), |
| 346 | page_texts(&model_ops::page_of(&working, a)) |
| 347 | )] |
| 348 | )] |
| 349 | ); |
| 350 | } |
| 351 | |
| 352 | #[test] |
| 353 | fn a_page_the_conversion_cannot_reproduce_leaves_the_cache_untouched() { |
| 354 | let body = body_pages(SOURCE); |
| 355 | let (a, a_text) = body[0]; |
| 356 | let directory = tempfile::tempdir().unwrap(); |
| 357 | let path = directory.path().join("cache.sqlite"); |
| 358 | // The working image holds a page the queued edit does not reach. |
| 359 | let stray = ops::saved(SOURCE, a, &appended(SOURCE, a, a_text, " stray")) |
| 360 | .unwrap() |
| 361 | .as_slice() |
| 362 | .to_vec(); |
| 363 | v14( |
| 364 | &path, |
| 365 | SOURCE, |
| 366 | &[Queued::Page(a, appended(SOURCE, a, a_text, " queued"))], |
| 367 | Some(&stray), |
| 368 | ); |
| 369 | let before = std::fs::read(&path).unwrap(); |
| 370 | let Err(notebook::Error::Io(error)) = Replica::open(&path) else { |
| 371 | panic!("the conversion must fail") |
| 372 | }; |
| 373 | assert_eq!(error.kind(), std::io::ErrorKind::InvalidData); |
| 374 | assert!(error.to_string().contains(".v14-recovery"), "{error}"); |
| 375 | assert_eq!(std::fs::read(&path).unwrap(), before); |
| 376 | assert!(archive(&path).exists()); |
| 377 | let version: u32 = Connection::open(archive(&path)) |
| 378 | .unwrap() |
| 379 | .pragma_query_value(None, "user_version", |row| row.get(0)) |
| 380 | .unwrap(); |
| 381 | assert_eq!(version, 14); |
| 382 | // A second open tries again from the same untouched cache. |
| 383 | assert!(Replica::open(&path).is_err()); |
| 384 | assert_eq!(std::fs::read(&path).unwrap(), before); |
| 385 | } |
| 386 | |
| 387 | /// A schema-15 cache upgrades in place: its queue, receipts and recovery archives stay, and |
| 388 | /// a conflict it recorded is queued work the next sync keeps as a conflict page. |
| 389 | #[test] |
| 390 | fn a_schema_15_cache_upgrades_in_place_and_its_conflict_becomes_a_conflict_page() { |
| 391 | let (space, text) = body_pages(SOURCE)[0]; |
| 392 | let directory = tempfile::tempdir().unwrap(); |
| 393 | let path = directory.path().join("cache.sqlite"); |
| 394 | let cache = Replica::create(&path, SOURCE).unwrap(); |
| 395 | let id = model_ops::save(&cache, text, |page| { |
| 396 | model_ops::replace_text(page, text, 0..0, "Local ") |
| 397 | }) |
| 398 | .unwrap() |
| 399 | .unwrap(); |
| 400 | let pending = cache.pending().unwrap(); |
| 401 | let archive = directory.path().join("before.sqlite"); |
| 402 | drop(cache); |
| 403 | { |
| 404 | // Schema 15 recorded the review conflict and its page on the batch. |
| 405 | let connection = Connection::open(&path).unwrap(); |
| 406 | connection |
| 407 | .execute_batch(&format!( |
| 408 | "PRAGMA foreign_keys=OFF; |
| 409 | BEGIN; |
| 410 | CREATE TABLE old ( |
| 411 | id INTEGER PRIMARY KEY AUTOINCREMENT CHECK(id>0), |
| 412 | sealed TEXT, |
| 413 | revisions TEXT, |
| 414 | attempted INTEGER NOT NULL DEFAULT 0 CHECK(attempted IN (0,1)), |
| 415 | conflict INTEGER CHECK(conflict BETWEEN 0 AND 3), |
| 416 | space TEXT, |
| 417 | CHECK((sealed IS NULL) = (revisions IS NULL)), |
| 418 | CHECK(attempted=0 OR sealed IS NOT NULL), |
| 419 | CHECK((conflict IS NULL) = (space IS NULL)) |
| 420 | ) STRICT; |
| 421 | INSERT INTO old SELECT id, sealed, revisions, attempted, 3, '{space}' FROM batches; |
| 422 | DROP TABLE batches; |
| 423 | ALTER TABLE old RENAME TO batches; |
| 424 | PRAGMA user_version=15; |
| 425 | COMMIT;" |
| 426 | )) |
| 427 | .unwrap(); |
| 428 | } |
| 429 | let cache = Replica::open(&path).unwrap(); |
| 430 | assert_eq!(cache.pending().unwrap(), pending); |
| 431 | assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending)); |
| 432 | cache.export_recovery(&archive).unwrap(); |
| 433 | drop(cache); |
| 434 | let connection = Connection::open(&path).unwrap(); |
| 435 | let version: u32 = connection |
| 436 | .pragma_query_value(None, "user_version", |row| row.get(0)) |
| 437 | .unwrap(); |
| 438 | let columns: Vec<String> = connection |
| 439 | .prepare("SELECT name FROM pragma_table_info('batches')") |
| 440 | .unwrap() |
| 441 | .query_map([], |row| row.get(0)) |
| 442 | .unwrap() |
| 443 | .collect::<Result<_, _>>() |
| 444 | .unwrap(); |
| 445 | assert_eq!(version, 16); |
| 446 | assert_eq!(columns, ["id", "sealed", "revisions", "attempted"]); |
| 447 | drop(connection); |
| 448 | assert_eq!( |
| 449 | notebook::Recovery::open(&archive) |
| 450 | .unwrap() |
| 451 | .pending() |
| 452 | .unwrap(), |
| 453 | pending |
| 454 | ); |
| 455 | let cache = Replica::open(&path).unwrap(); |
| 456 | let remote = typed(SOURCE, space, text, 0..0, "Remote "); |
| 457 | let mut server = Server::new(&remote); |
| 458 | assert!(matches!( |
| 459 | cache.sync_once(&mut server).unwrap().edit, |
| 460 | Some((_, EditStatus::Published { .. })) |
| 461 | )); |
| 462 | assert!(conflicted(&server.durable, space)); |
| 463 | } |