1#[path = "../../onestore/tests/support/ops.rs"]
2mod ops;
3use notebook::{Error, Replica};
4use onestore::{
5 ExGuid, RevisionIndex, Store,
6 document::{Document, Kind},
7 op::{Edit, Op, PageOp, SectionOp},
8 page::{Outline, PageObject},
9};
10use std::{
11 collections::BTreeSet,
12 fs,
13 io::ErrorKind,
14 sync::Barrier,
15 time::{Duration, Instant},
16};
17
18#[path = "support/model_ops.rs"]
19mod model_ops;
20#[path = "support/server.rs"]
21mod server;
22#[path = "../../onestore/tests/support/sweep.rs"]
23mod sweep;
24use server::snapshot;
25
26fn 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.
48fn 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]
65fn 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]
81fn 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]
136fn 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]
199fn 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]
262fn 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]
335fn 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]
390fn 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]
458fn 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]
570fn 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}