| 1 | #[path = "../../onestore/tests/support/ops.rs"] |
| 2 | mod ops; |
| 3 | use notebook::{Error, Replica}; |
| 4 | use onestore::{ |
| 5 | ExGuid, RevisionIndex, Store, |
| 6 | document::{Document, Kind}, |
| 7 | op::{Edit, Op, PageOp, SectionOp}, |
| 8 | page::{Outline, PageObject}, |
| 9 | }; |
| 10 | use std::{ |
| 11 | collections::BTreeSet, |
| 12 | fs, |
| 13 | io::ErrorKind, |
| 14 | sync::Barrier, |
| 15 | time::{Duration, Instant}, |
| 16 | }; |
| 17 | |
| 18 | #[path = "support/model_ops.rs"] |
| 19 | mod model_ops; |
| 20 | #[path = "support/server.rs"] |
| 21 | mod server; |
| 22 | #[path = "../../onestore/tests/support/sweep.rs"] |
| 23 | mod sweep; |
| 24 | use server::snapshot; |
| 25 | |
| 26 | fn target(source: &[u8]) -> (ExGuid, ExGuid, String) { |
| 27 | let store = Store::parse(source).unwrap(); |
| 28 | assert!(store.checksum_mismatches.is_empty()); |
| 29 | let index = RevisionIndex::parse(&store).unwrap(); |
| 30 | index.validate_current().unwrap(); |
| 31 | let doc = Document::parse(&index).unwrap(); |
| 32 | doc.spaces |
| 33 | .iter() |
| 34 | .find_map(|(sid, space)| { |
| 35 | let revision = &space.revisions[&space.contexts[&ExGuid::default()]]; |
| 36 | revision |
| 37 | .nodes |
| 38 | .iter() |
| 39 | .find_map(|(oid, node)| match &node.kind { |
| 40 | Kind::RichText { text, .. } => Some((*sid, *oid, text.clone())), |
| 41 | _ => None, |
| 42 | }) |
| 43 | }) |
| 44 | .unwrap() |
| 45 | } |
| 46 | |
| 47 | /// An edit typing `with` at the start of a text. |
| 48 | fn typed(space: ExGuid, text: ExGuid, with: &str) -> Edit { |
| 49 | Edit { |
| 50 | at: model_ops::now(), |
| 51 | ops: vec![Op::Page { |
| 52 | space, |
| 53 | op: PageOp::Text { |
| 54 | text, |
| 55 | range: 0..0, |
| 56 | with: with.into(), |
| 57 | }, |
| 58 | }], |
| 59 | } |
| 60 | } |
| 61 | |
| 62 | /// A cache another thread creates while this one reads its section, as the background |
| 63 | /// makes a section's offline copy while the section opens, is opened instead of refused. |
| 64 | #[test] |
| 65 | fn a_cache_created_meanwhile_opens() { |
| 66 | let source = onestore::create_section("Raced.one", "Raced", "Fixture").unwrap(); |
| 67 | let directory = tempfile::tempdir().unwrap(); |
| 68 | let path = directory.path().join("raced.sqlite"); |
| 69 | let replica = Replica::open_or_create(&path, None, || { |
| 70 | drop(Replica::create(&path, &source)?); |
| 71 | Ok(source.clone()) |
| 72 | }) |
| 73 | .unwrap(); |
| 74 | assert_eq!(replica.pages().unwrap()[0].1, "Raced"); |
| 75 | drop(replica); |
| 76 | let reopened = Replica::open_or_create(&path, None, || unreachable!()).unwrap(); |
| 77 | assert_eq!(reopened.pages().unwrap()[0].1, "Raced"); |
| 78 | } |
| 79 | |
| 80 | #[test] |
| 81 | fn cache_reopen_preserves_the_base_and_queued_edits() { |
| 82 | let dir = tempfile::tempdir().unwrap(); |
| 83 | let path = dir.path().join("section.sqlite"); |
| 84 | let source = onestore::create_section("section.one", "café 🦀", "Fixture").unwrap(); |
| 85 | let replica = Replica::create(&path, &source).unwrap(); |
| 86 | assert_eq!(snapshot(&replica), source); |
| 87 | assert!(replica.pending().unwrap().is_empty()); |
| 88 | let (sid, oid, _) = target(&source); |
| 89 | let unchanged = |page: &mut onestore::page::Page| { |
| 90 | model_ops::replace_text(page, oid, 0..0, ""); |
| 91 | model_ops::replace_text(page, oid, 0..4, "café"); |
| 92 | }; |
| 93 | assert_eq!(model_ops::save(&replica, oid, unchanged).unwrap(), None); |
| 94 | let first = model_ops::save(&replica, oid, |page| { |
| 95 | model_ops::replace_text(page, oid, 5..7, "🐈 日本語") |
| 96 | }) |
| 97 | .unwrap() |
| 98 | .unwrap(); |
| 99 | assert_eq!(target(&snapshot(&replica)).2, "café 🐈 日本語"); |
| 100 | let edits = replica.pending().unwrap(); |
| 101 | assert_eq!(edits.len(), 1); |
| 102 | assert_eq!( |
| 103 | (edits[0].id, edits[0].author.as_str()), |
| 104 | (first, model_ops::AUTHOR) |
| 105 | ); |
| 106 | assert!(matches!( |
| 107 | &edits[0].edit.ops[..], |
| 108 | [Op::Page { space, op: PageOp::Text { text, range, with } }] |
| 109 | if (*space, *text, range.clone(), with.as_str()) == (sid, oid, 5..7, "🐈 日本語") |
| 110 | )); |
| 111 | let beside = |suffix: &str| { |
| 112 | let mut file = path.as_os_str().to_owned(); |
| 113 | file.push(suffix); |
| 114 | std::path::PathBuf::from(file) |
| 115 | }; |
| 116 | // Commits go to the write-ahead log; its index stays in memory. |
| 117 | assert!(beside("-wal").exists()); |
| 118 | assert!(!beside("-shm").exists()); |
| 119 | drop(replica); |
| 120 | assert!( |
| 121 | !beside("-wal").exists(), |
| 122 | "closing checkpoints and removes the log" |
| 123 | ); |
| 124 | let replica = Replica::open(&path).unwrap(); |
| 125 | assert_eq!(target(&snapshot(&replica)).2, "café 🐈 日本語"); |
| 126 | assert_eq!(replica.pending().unwrap(), edits); |
| 127 | let second = replica |
| 128 | .apply(model_ops::AUTHOR, typed(sid, oid, "Recovered ")) |
| 129 | .unwrap(); |
| 130 | assert!(second > first); |
| 131 | assert_eq!(target(&snapshot(&replica)).2, "Recovered café 🐈 日本語"); |
| 132 | assert_eq!(replica.pending().unwrap().len(), 2); |
| 133 | } |
| 134 | |
| 135 | #[test] |
| 136 | fn refused_edits_and_failed_writes_leave_the_queue_as_it_was() { |
| 137 | let dir = tempfile::tempdir().unwrap(); |
| 138 | let path = dir.path().join("section.sqlite"); |
| 139 | let source = onestore::create_section("section.one", "café 🦀", "Fixture").unwrap(); |
| 140 | let replica = Replica::create(&path, &source).unwrap(); |
| 141 | let (sid, oid, _) = target(&source); |
| 142 | let empty = Outline { |
| 143 | id: onestore::page::text::new_id().unwrap(), |
| 144 | title: false, |
| 145 | min_width: None, |
| 146 | layout: onestore::document::Layout { |
| 147 | x: Some(72.0), |
| 148 | y: Some(400.0), |
| 149 | ..Default::default() |
| 150 | }, |
| 151 | indents: Vec::new(), |
| 152 | paragraphs: Vec::new(), |
| 153 | unsupported: Vec::new(), |
| 154 | }; |
| 155 | let refused = Edit { |
| 156 | at: model_ops::now(), |
| 157 | ops: vec![ |
| 158 | Op::Page { |
| 159 | space: sid, |
| 160 | op: PageOp::Text { |
| 161 | text: oid, |
| 162 | range: 0..0, |
| 163 | with: "Undone ".into(), |
| 164 | }, |
| 165 | }, |
| 166 | Op::Page { |
| 167 | space: sid, |
| 168 | op: PageOp::Add { |
| 169 | object: PageObject::Outline(empty), |
| 170 | before: None, |
| 171 | }, |
| 172 | }, |
| 173 | ], |
| 174 | }; |
| 175 | assert!(matches!( |
| 176 | replica.apply("Author", refused), |
| 177 | Err(Error::Rejected(_)) |
| 178 | )); |
| 179 | assert!(matches!( |
| 180 | replica.apply("Author", typed(ExGuid::default(), oid, "Elsewhere ")), |
| 181 | Err(Error::Rejected(_)) |
| 182 | )); |
| 183 | assert_eq!(snapshot(&replica), source); |
| 184 | assert!(replica.pending().unwrap().is_empty()); |
| 185 | drop(replica); |
| 186 | let connection = rusqlite::Connection::open(&path).unwrap(); |
| 187 | connection.execute_batch("CREATE TRIGGER fail_edit BEFORE INSERT ON edits BEGIN SELECT RAISE(ABORT, 'Injected queue failure'); END;").unwrap(); |
| 188 | drop(connection); |
| 189 | let replica = Replica::open(&path).unwrap(); |
| 190 | assert!(replica.apply("Author", typed(sid, oid, "lost? ")).is_err()); |
| 191 | assert_eq!(snapshot(&replica), source); |
| 192 | drop(replica); |
| 193 | let replica = Replica::open(&path).unwrap(); |
| 194 | assert_eq!(snapshot(&replica), source); |
| 195 | assert!(replica.pending().unwrap().is_empty()); |
| 196 | } |
| 197 | |
| 198 | #[test] |
| 199 | fn concurrent_recovery_exports_capture_one_complete_acknowledged_queue() { |
| 200 | use notebook::Recovery; |
| 201 | |
| 202 | let directory = tempfile::tempdir().unwrap(); |
| 203 | let source = onestore::create_section("recovery.one", "base", "Fixture").unwrap(); |
| 204 | let (space, object, _) = target(&source); |
| 205 | let replica = Replica::create(directory.path().join("live.sqlite"), &source).unwrap(); |
| 206 | let start = Barrier::new(4); |
| 207 | std::thread::scope(|scope| { |
| 208 | for writer in 0..3 { |
| 209 | let (replica, start) = (&replica, &start); |
| 210 | scope.spawn(move || { |
| 211 | start.wait(); |
| 212 | for edit in 0..20 { |
| 213 | replica |
| 214 | .apply( |
| 215 | "Fixture", |
| 216 | typed(space, object, &format!("[{writer}:{edit}] ")), |
| 217 | ) |
| 218 | .unwrap(); |
| 219 | } |
| 220 | }); |
| 221 | } |
| 222 | start.wait(); |
| 223 | for n in 0..12 { |
| 224 | let path = directory.path().join(format!("recovery-{n}.sqlite")); |
| 225 | replica.export_recovery(&path).unwrap(); |
| 226 | let archive = Recovery::open(path).unwrap(); |
| 227 | let pending = archive.pending().unwrap(); |
| 228 | // Each edit types at the start, so the text is the queue read backwards. |
| 229 | let expected: String = pending |
| 230 | .iter() |
| 231 | .rev() |
| 232 | .map(|edit| match &edit.edit.ops[..] { |
| 233 | [ |
| 234 | Op::Page { |
| 235 | op: PageOp::Text { with, .. }, |
| 236 | .. |
| 237 | }, |
| 238 | ] => with.as_str(), |
| 239 | other => panic!("{other:?}"), |
| 240 | }) |
| 241 | .chain(["base"]) |
| 242 | .collect(); |
| 243 | assert_eq!(target(&archive.snapshot().unwrap()).2, expected); |
| 244 | assert_eq!(archive.remote_snapshot().unwrap(), source); |
| 245 | assert_eq!( |
| 246 | archive.summary().unwrap().queued_edits, |
| 247 | pending.len() as u64 |
| 248 | ); |
| 249 | assert!(archive.receipts().unwrap().is_empty()); |
| 250 | } |
| 251 | }); |
| 252 | assert_eq!(replica.pending().unwrap().len(), 60); |
| 253 | let content = target(&snapshot(&replica)).2; |
| 254 | for writer in 0..3 { |
| 255 | for edit in 0..20 { |
| 256 | assert_eq!(content.matches(&format!("[{writer}:{edit}] ")).count(), 1); |
| 257 | } |
| 258 | } |
| 259 | } |
| 260 | |
| 261 | #[test] |
| 262 | fn ownership_and_foreign_file_rejection_preserve_existing_data() { |
| 263 | let dir = tempfile::tempdir().unwrap(); |
| 264 | let path = dir.path().join("section.sqlite"); |
| 265 | assert!(Replica::open(&path).is_err()); |
| 266 | assert!(!path.exists()); |
| 267 | assert!(Replica::create(&path, b"invalid").is_err()); |
| 268 | assert!(!path.exists()); |
| 269 | let source = onestore::create_section("section.one", "Owned", "Fixture").unwrap(); |
| 270 | let replica = Replica::create(&path, &source).unwrap(); |
| 271 | for _ in 0..3 { |
| 272 | assert!( |
| 273 | matches!(Replica::open(&path), Err(Error::Database(error)) if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy)) |
| 274 | ); |
| 275 | assert!( |
| 276 | matches!(Replica::create(&path, &source), Err(Error::Io(error)) if error.kind() == ErrorKind::AlreadyExists) |
| 277 | ); |
| 278 | assert_eq!(snapshot(&replica), source); |
| 279 | } |
| 280 | drop(replica); |
| 281 | let current: u32 = rusqlite::Connection::open(&path) |
| 282 | .unwrap() |
| 283 | .pragma_query_value(None, "user_version", |row| row.get(0)) |
| 284 | .unwrap(); |
| 285 | for sql in [ |
| 286 | "PRAGMA application_id=0".to_owned(), |
| 287 | // Schemas 14 and 15 convert; older caches do not. |
| 288 | format!( |
| 289 | "PRAGMA application_id=1330529615; PRAGMA user_version={}", |
| 290 | current - 3 |
| 291 | ), |
| 292 | format!("PRAGMA user_version={}", current + 1), |
| 293 | ] { |
| 294 | let connection = rusqlite::Connection::open(&path).unwrap(); |
| 295 | connection.execute_batch(&sql).unwrap(); |
| 296 | drop(connection); |
| 297 | let before = fs::read(&path).unwrap(); |
| 298 | let Err(Error::Io(error)) = Replica::open(&path) else { |
| 299 | panic!("{sql}") |
| 300 | }; |
| 301 | assert_eq!(error.kind(), ErrorKind::InvalidData, "{sql}"); |
| 302 | if sql.contains("user_version") { |
| 303 | let written = sql.rsplit('=').next().unwrap(); |
| 304 | let message = error.to_string(); |
| 305 | assert!( |
| 306 | message.contains(&format!("version {written} ")) |
| 307 | && message.contains(&format!("version {current}")), |
| 308 | "{message}" |
| 309 | ); |
| 310 | } |
| 311 | assert!( |
| 312 | matches!(notebook::Recovery::open(&path), Err(Error::Io(error)) if error.kind() == ErrorKind::InvalidData), |
| 313 | "{sql}" |
| 314 | ); |
| 315 | assert_eq!(fs::read(&path).unwrap(), before); |
| 316 | } |
| 317 | let foreign = dir.path().join("foreign.sqlite"); |
| 318 | let connection = rusqlite::Connection::open(&foreign).unwrap(); |
| 319 | connection |
| 320 | .execute_batch( |
| 321 | "CREATE TABLE unrelated (value TEXT); INSERT INTO unrelated VALUES ('preserve');", |
| 322 | ) |
| 323 | .unwrap(); |
| 324 | drop(connection); |
| 325 | let before = fs::read(&foreign).unwrap(); |
| 326 | assert!(Replica::open(&foreign).is_err()); |
| 327 | assert_eq!(fs::read(&foreign).unwrap(), before); |
| 328 | let incomplete = dir.path().join("incomplete.sqlite"); |
| 329 | fs::write(&incomplete, []).unwrap(); |
| 330 | assert!(Replica::open(&incomplete).is_err()); |
| 331 | assert_eq!(fs::read(&incomplete).unwrap(), b""); |
| 332 | } |
| 333 | |
| 334 | #[test] |
| 335 | fn twelve_local_editors_queue_every_edit_once_in_order() { |
| 336 | let dir = tempfile::tempdir().unwrap(); |
| 337 | let path = dir.path().join("section.sqlite"); |
| 338 | let source = onestore::create_section("section.one", "Shared café 🦀", "Fixture").unwrap(); |
| 339 | let replica = Replica::create(&path, &source).unwrap(); |
| 340 | let (sid, oid, _) = target(&source); |
| 341 | let barrier = Barrier::new(12); |
| 342 | let outcomes = std::thread::scope(|scope| { |
| 343 | let handles: Vec<_> = (0..12) |
| 344 | .map(|writer| { |
| 345 | let (replica, barrier) = (&replica, &barrier); |
| 346 | scope.spawn(move || { |
| 347 | barrier.wait(); |
| 348 | (0..20) |
| 349 | .map(|edit| { |
| 350 | replica |
| 351 | .apply("Fixture", typed(sid, oid, &format!("[{writer}-{edit}] "))) |
| 352 | .unwrap() |
| 353 | }) |
| 354 | .collect::<Vec<_>>() |
| 355 | }) |
| 356 | }) |
| 357 | .collect(); |
| 358 | handles |
| 359 | .into_iter() |
| 360 | .map(|handle| handle.join().unwrap()) |
| 361 | .collect::<Vec<_>>() |
| 362 | }); |
| 363 | for ids in &outcomes { |
| 364 | assert!( |
| 365 | ids.windows(2).all(|pair| pair[0] < pair[1]), |
| 366 | "a writer's edits keep its order" |
| 367 | ); |
| 368 | } |
| 369 | let ids: BTreeSet<_> = outcomes.into_iter().flatten().collect(); |
| 370 | assert_eq!(ids.len(), 240); |
| 371 | let content = target(&snapshot(&replica)).2; |
| 372 | assert!(content.ends_with("Shared café 🦀")); |
| 373 | for writer in 0..12 { |
| 374 | for edit in 0..20 { |
| 375 | assert_eq!(content.matches(&format!("[{writer}-{edit}] ")).count(), 1); |
| 376 | } |
| 377 | } |
| 378 | let pending = replica.pending().unwrap(); |
| 379 | assert_eq!( |
| 380 | pending.iter().map(|edit| edit.id).collect::<BTreeSet<_>>(), |
| 381 | ids |
| 382 | ); |
| 383 | drop(replica); |
| 384 | let reopened = Replica::open(&path).unwrap(); |
| 385 | assert_eq!(target(&snapshot(&reopened)).2, content); |
| 386 | assert_eq!(reopened.pending().unwrap(), pending); |
| 387 | } |
| 388 | |
| 389 | #[test] |
| 390 | fn seeded_unicode_edits_and_restarts_match_an_independent_text_model() { |
| 391 | let dir = tempfile::tempdir().unwrap(); |
| 392 | let path = dir.path().join("model.sqlite"); |
| 393 | let mut text = "ab🚀ab🦀 é repeated repeated".to_owned(); |
| 394 | let source = onestore::create_section("model.one", &text, "Fixture").unwrap(); |
| 395 | let mut replica = Replica::create(&path, &source).unwrap(); |
| 396 | let (_, object, _) = target(&source); |
| 397 | let steps = sweep::seeds(0..512, 64); |
| 398 | let mut random = 911 + steps.start; |
| 399 | let mut next = || { |
| 400 | random ^= random << 13; |
| 401 | random ^= random >> 7; |
| 402 | random ^= random << 17; |
| 403 | random |
| 404 | }; |
| 405 | let mut acknowledged = Vec::new(); |
| 406 | for step in steps.clone() { |
| 407 | let boundaries: Vec<_> = text |
| 408 | .char_indices() |
| 409 | .map(|(at, _)| at) |
| 410 | .chain([text.len()]) |
| 411 | .collect(); |
| 412 | let first = boundaries[next() as usize % boundaries.len()]; |
| 413 | let last = boundaries[next() as usize % boundaries.len()]; |
| 414 | let bytes = first.min(last)..first.max(last); |
| 415 | let range = u32::try_from(text[..bytes.start].encode_utf16().count()).unwrap() |
| 416 | ..u32::try_from(text[..bytes.end].encode_utf16().count()).unwrap(); |
| 417 | let replacement = ["", "🐈", "日本語", "repeated", "é", "ab🦀ab"][next() as usize % 6]; |
| 418 | let mut expected = text.clone(); |
| 419 | expected.replace_range(bytes, replacement); |
| 420 | let acknowledgement = model_ops::save(&replica, object, |page| { |
| 421 | model_ops::replace_text(page, object, range.clone(), replacement) |
| 422 | }) |
| 423 | .unwrap(); |
| 424 | match acknowledgement { |
| 425 | Some(id) => { |
| 426 | assert_ne!(text, expected); |
| 427 | assert!(acknowledged.last().is_none_or(|last| *last < id)); |
| 428 | acknowledged.push(id); |
| 429 | } |
| 430 | None => assert_eq!(text, expected), |
| 431 | } |
| 432 | text = expected; |
| 433 | assert_eq!(target(&snapshot(&replica)).2, text, "step {step}"); |
| 434 | if step % 37 == 0 { |
| 435 | let pending = replica.pending().unwrap(); |
| 436 | drop(replica); |
| 437 | replica = Replica::open(&path).unwrap(); |
| 438 | assert_eq!(target(&snapshot(&replica)).2, text); |
| 439 | assert_eq!(replica.pending().unwrap(), pending); |
| 440 | } |
| 441 | } |
| 442 | assert!( |
| 443 | acknowledged.len() as u64 * 2 > steps.end - steps.start, |
| 444 | "Most random edits change the text" |
| 445 | ); |
| 446 | assert_eq!( |
| 447 | replica |
| 448 | .pending() |
| 449 | .unwrap() |
| 450 | .iter() |
| 451 | .map(|edit| edit.id) |
| 452 | .collect::<Vec<_>>(), |
| 453 | acknowledged |
| 454 | ); |
| 455 | } |
| 456 | |
| 457 | #[test] |
| 458 | fn twelve_local_clients_preserve_inserted_identities_and_dependent_edits() { |
| 459 | let directory = tempfile::tempdir().unwrap(); |
| 460 | let path = directory.path().join("parallel.sqlite"); |
| 461 | let source = onestore::create_section("parallel.one", "Original", "Author").unwrap(); |
| 462 | let (sid, anchor, _) = target(&source); |
| 463 | let cache = Replica::create(&path, &source).unwrap(); |
| 464 | let barrier = Barrier::new(12); |
| 465 | let deadline = Instant::now() + Duration::from_secs(90); |
| 466 | let italic = |format: &mut onestore::document::Format| format.italic = Some(true); |
| 467 | let all = std::thread::scope(|scope| { |
| 468 | let handles: Vec<_> = (0..12) |
| 469 | .map(|client| { |
| 470 | let (cache, barrier) = (&cache, &barrier); |
| 471 | scope.spawn(move || { |
| 472 | let mut ids = Vec::new(); |
| 473 | let mut objects = Vec::new(); |
| 474 | barrier.wait(); |
| 475 | for sequence in 0..4 { |
| 476 | assert!(Instant::now() < deadline, "Client {client} stopped"); |
| 477 | let content = if sequence == 0 { |
| 478 | format!("Client {client}") |
| 479 | } else { |
| 480 | format!("Paragraph {client}:{sequence}") |
| 481 | }; |
| 482 | let mut inserted = None; |
| 483 | let id = model_ops::save(cache, anchor, |page| { |
| 484 | let text = if sequence == 0 { |
| 485 | model_ops::insert_outline( |
| 486 | page, |
| 487 | 72.0, |
| 488 | 144.0 + client as f32 * 72.0, |
| 489 | &content, |
| 490 | ) |
| 491 | .2 |
| 492 | } else { |
| 493 | model_ops::insert_after(page, *objects.last().unwrap(), &content).1 |
| 494 | }; |
| 495 | let end = content.encode_utf16().count() as u32; |
| 496 | model_ops::restyle(page, text, 0..end, |format| format.italic = None); |
| 497 | model_ops::restyle(page, text, 0..6, italic); |
| 498 | inserted = Some(text); |
| 499 | }) |
| 500 | .unwrap() |
| 501 | .unwrap(); |
| 502 | ids.push(id); |
| 503 | objects.push(inserted.unwrap()); |
| 504 | if sequence == 0 { |
| 505 | continue; |
| 506 | } |
| 507 | let text = *objects.last().unwrap(); |
| 508 | ids.push( |
| 509 | cache |
| 510 | .apply(model_ops::AUTHOR, typed(sid, text, "Edited ")) |
| 511 | .unwrap(), |
| 512 | ); |
| 513 | } |
| 514 | (ids, objects) |
| 515 | }) |
| 516 | }) |
| 517 | .collect(); |
| 518 | handles |
| 519 | .into_iter() |
| 520 | .map(|h| h.join().unwrap()) |
| 521 | .collect::<Vec<_>>() |
| 522 | }); |
| 523 | let ids: BTreeSet<_> = all |
| 524 | .iter() |
| 525 | .flat_map(|(ids, _)| ids.iter().copied()) |
| 526 | .collect(); |
| 527 | assert_eq!(ids.len(), 84); |
| 528 | let image = snapshot(&cache); |
| 529 | drop(cache); |
| 530 | let cache = Replica::open(&path).unwrap(); |
| 531 | assert_eq!(server::pages(&snapshot(&cache)), server::pages(&image)); |
| 532 | assert_eq!( |
| 533 | cache |
| 534 | .pending() |
| 535 | .unwrap() |
| 536 | .iter() |
| 537 | .map(|e| e.id) |
| 538 | .collect::<BTreeSet<_>>(), |
| 539 | ids |
| 540 | ); |
| 541 | let store = Store::parse(&image).unwrap(); |
| 542 | let index = RevisionIndex::parse(&store).unwrap(); |
| 543 | let doc = Document::parse(&index).unwrap(); |
| 544 | let s = &doc.spaces[&sid]; |
| 545 | let v = &s.revisions[&s.contexts[&ExGuid::default()]]; |
| 546 | for (client, (_, objects)) in all.iter().enumerate() { |
| 547 | for (sequence, id) in objects.iter().enumerate() { |
| 548 | let wanted = if sequence == 0 { |
| 549 | format!("Client {client}") |
| 550 | } else { |
| 551 | format!("Edited Paragraph {client}:{sequence}") |
| 552 | }; |
| 553 | assert!(matches!(&v.nodes[id].kind,Kind::RichText{text,..} if *text==wanted)); |
| 554 | let runs = v.text_runs(*id).unwrap(); |
| 555 | assert_eq!(runs[0].format.italic, Some(true)); |
| 556 | assert_eq!( |
| 557 | runs[0].text, |
| 558 | if sequence == 0 { |
| 559 | "Client" |
| 560 | } else { |
| 561 | "Edited Paragr" |
| 562 | } |
| 563 | ); |
| 564 | assert!(runs[1..].iter().all(|run| run.format.italic != Some(true))); |
| 565 | } |
| 566 | } |
| 567 | } |
| 568 | |
| 569 | #[test] |
| 570 | fn unrecognized_persisted_ops_are_rejected_without_dropping_fields() { |
| 571 | for kind in ["Text", "Create", "Pages", "Delete"] { |
| 572 | let directory = tempfile::tempdir().unwrap(); |
| 573 | let path = directory.path().join("unknown.sqlite"); |
| 574 | let first = onestore::create_section("unknown.one", "Original", "Author").unwrap(); |
| 575 | let second = onestore::PageCreation::new(None, Some("Second"), "Author").unwrap(); |
| 576 | let source = ops::section_op(&first, SectionOp::Create(second.clone())) |
| 577 | .unwrap() |
| 578 | .as_bytes() |
| 579 | .to_vec(); |
| 580 | let (space, oid, _) = target(&source); |
| 581 | let store = Store::parse(&source).unwrap(); |
| 582 | let index = RevisionIndex::parse(&store).unwrap(); |
| 583 | let sid = Document::parse(&index).unwrap().pages().unwrap()[1].0; |
| 584 | let cache = Replica::create(&path, &source).unwrap(); |
| 585 | let section = |op| { |
| 586 | server::section_op(&cache, op); |
| 587 | }; |
| 588 | match kind { |
| 589 | "Text" => { |
| 590 | cache.apply("Author", typed(space, oid, "New ")).unwrap(); |
| 591 | } |
| 592 | "Create" => section(SectionOp::Create( |
| 593 | onestore::PageCreation::new(None, Some("Created"), "Author").unwrap(), |
| 594 | )), |
| 595 | "Pages" => section(SectionOp::Pages(vec![ |
| 596 | onestore::PageEdit::set_level(sid, 2).unwrap(), |
| 597 | ])), |
| 598 | _ => section(SectionOp::Delete(vec![sid])), |
| 599 | } |
| 600 | drop(cache); |
| 601 | let encoded: String = rusqlite::Connection::open(&path) |
| 602 | .unwrap() |
| 603 | .query_row("SELECT edit FROM edits", [], |r| r.get(0)) |
| 604 | .unwrap(); |
| 605 | for target in ["", "/ops/0", "/ops/0/Page/op/Text", "/ops/0/Section"] { |
| 606 | let mut value: serde_json::Value = serde_json::from_str(&encoded).unwrap(); |
| 607 | let Some(object) = value.pointer_mut(target).and_then(|v| v.as_object_mut()) else { |
| 608 | continue; |
| 609 | }; |
| 610 | object.insert("future_option".into(), true.into()); |
| 611 | rusqlite::Connection::open(&path) |
| 612 | .unwrap() |
| 613 | .execute("UPDATE edits SET edit=?1", [value.to_string()]) |
| 614 | .unwrap(); |
| 615 | let before = fs::read(&path).unwrap(); |
| 616 | assert!( |
| 617 | matches!(Replica::open(&path),Err(Error::Io(error))if error.kind()==ErrorKind::InvalidData), |
| 618 | "{kind} {target}" |
| 619 | ); |
| 620 | assert_eq!(fs::read(&path).unwrap(), before); |
| 621 | } |
| 622 | } |
| 623 | } |