1//! Replica protocol: durable receipts, retired attempts, conflict pages and contention.
2
3use notebook::{EditStatus, Error, Remote, Replica};
4use onestore::{CommitError, CommitState, ExGuid, Transaction};
5use std::{io, ops::Range};
6
7#[path = "support/server.rs"]
8mod server;
9use server::*;
10#[path = "support/model_ops.rs"]
11mod model_ops;
12#[path = "../../onestore/tests/support/sweep.rs"]
13mod sweep;
14
15/// Saves a replacement of `range` in the page holding `text`.
16fn save(cache: &Replica, text: ExGuid, range: Range<u32>, replacement: &str) -> Option<u64> {
17 model_ops::save(cache, text, |page| {
18 model_ops::replace_text(page, text, range, replacement)
19 })
20 .unwrap()
21}
22
23#[test]
24fn asynchronous_io_keeps_the_attempt_and_recovers_a_lost_reply() {
25 use std::{
26 future::Future,
27 pin::Pin,
28 task::{Context, Poll, Waker},
29 };
30
31 struct Deferred {
32 server: Server,
33 ready: bool,
34 }
35 impl Remote for Deferred {
36 fn pending(&mut self) -> Option<Pin<Box<dyn Future<Output = ()> + '_>>> {
37 let mut yielded = false;
38 Some(Box::pin(std::future::poll_fn(move |context| {
39 if yielded {
40 self.ready = true;
41 Poll::Ready(())
42 } else {
43 yielded = true;
44 context.waker().wake_by_ref();
45 Poll::Pending
46 }
47 })))
48 }
49 fn read(&mut self) -> io::Result<Vec<u8>> {
50 if !std::mem::take(&mut self.ready) {
51 return Err(io::ErrorKind::WouldBlock.into());
52 }
53 self.server.read()
54 }
55 fn stamp(&mut self) -> io::Result<onestore::Stamp> {
56 if !std::mem::take(&mut self.ready) {
57 return Err(io::ErrorKind::WouldBlock.into());
58 }
59 self.server.stamp()
60 }
61 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
62 if !std::mem::take(&mut self.ready) {
63 return Err(CommitError {
64 state: CommitState::NotCommitted,
65 error: io::ErrorKind::WouldBlock.into(),
66 });
67 }
68 self.server.publish(transaction)
69 }
70 fn confirm(&mut self, base: &onestore::Stamp) -> Result<(), CommitError> {
71 if !std::mem::take(&mut self.ready) {
72 return Err(CommitError {
73 state: CommitState::NotCommitted,
74 error: io::ErrorKind::WouldBlock.into(),
75 });
76 }
77 self.server.confirm(base)
78 }
79 }
80 fn run<T>(work: impl Future<Output = T>) -> T {
81 let mut work = std::pin::pin!(work);
82 let mut polls = 0;
83 loop {
84 polls += 1;
85 assert!(polls < 100, "An async operation stopped making progress");
86 if let Poll::Ready(result) = work.as_mut().poll(&mut Context::from_waker(Waker::noop()))
87 {
88 assert!(polls > 1);
89 return result;
90 }
91 }
92 }
93 let directory = tempfile::tempdir().unwrap();
94 let path = directory.path().join("cache.sqlite");
95 let source = onestore::create_section("async.one", "Original", "Fixture").unwrap();
96 let (_, object, _) = text(&source);
97 let cache = Replica::create(&path, &source).unwrap();
98 let id = save(&cache, object, 0..0, "Browser ").unwrap();
99 let mut remote = Deferred {
100 server: Server::new(&source),
101 ready: false,
102 };
103 remote.server.fault = Fault::UnknownAfter;
104 assert!(matches!(
105 run(cache.sync_once_async(&mut remote)),
106 Err(Error::Remote(CommitError {
107 state: CommitState::Unknown,
108 ..
109 }))
110 ));
111 assert_eq!(remote.server.publications, 1);
112 drop(cache);
113 let cache = Replica::open(&path).unwrap();
114 assert!(
115 matches!(run(cache.sync_once_async(&mut remote)).unwrap().edit, Some((published, EditStatus::Published { .. })) if published == id)
116 );
117 assert_eq!(remote.server.publications, 1);
118 assert_eq!(text(&remote.server.visible).2, "Browser Original");
119 assert_eq!(remote.server.visible, remote.server.durable);
120 assert!(cache.pending().unwrap().is_empty());
121}
122
123#[test]
124fn an_unchanged_stamp_publishes_and_settles_without_reading_the_remote() {
125 struct Counted {
126 server: Server,
127 reads: usize,
128 }
129 impl Remote for Counted {
130 fn read(&mut self) -> io::Result<Vec<u8>> {
131 self.reads += 1;
132 self.server.read()
133 }
134 fn stamp(&mut self) -> io::Result<onestore::Stamp> {
135 self.server.stamp()
136 }
137 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
138 self.server.publish(transaction)
139 }
140 fn confirm(&mut self, base: &onestore::Stamp) -> Result<(), CommitError> {
141 self.server.confirm(base)
142 }
143 }
144 let directory = tempfile::tempdir().unwrap();
145 let source = onestore::create_section("stamp.one", "Original", "Fixture").unwrap();
146 let (sid, object, _) = text(&source);
147 let cache = Replica::create(directory.path().join("cache.sqlite"), &source).unwrap();
148 let id = save(&cache, object, 0..0, "Local ").unwrap();
149 let mut remote = Counted {
150 server: Server::new(&source),
151 reads: 0,
152 };
153 assert!(matches!(
154 cache.sync_once(&mut remote).unwrap().edit,
155 Some((published, EditStatus::Published { .. })) if published == id
156 ));
157 assert_eq!(cache.sync_once(&mut remote).unwrap().edit, None);
158 assert_eq!(remote.reads, 0);
159 // Another writer's commit moves the header, so the next step reads the file.
160 let native = typed(&remote.server.visible, sid, object, 0..0, "Native ");
161 remote.server.visible.clone_from(&native);
162 assert_eq!(cache.sync_once(&mut remote).unwrap().edit, None);
163 assert_eq!(remote.reads, 1);
164 assert_eq!(snapshot(&cache), native);
165}
166
167#[test]
168fn recovery_archive_preserves_the_queue_uncertainty_and_receipts_without_becoming_a_writer() {
169 use notebook::{Recovery, RecoverySummary};
170
171 let directory = tempfile::tempdir().unwrap();
172 let path = directory.path().join("live.sqlite");
173 let archive_path = directory.path().join("recovery.sqlite");
174 let source = onestore::create_section("recovery.one", "Original", "Fixture").unwrap();
175 let (_, object, _) = text(&source);
176 let cache = Replica::create(&path, &source).unwrap();
177 let published = save(&cache, object, 0..0, "Published ").unwrap();
178 let mut server = Server::new(&source);
179 let (_, receipt) = cache.sync_once(&mut server).unwrap().edit.unwrap();
180 let attempted = save(&cache, object, 0..0, "Uncertain ").unwrap();
181 server.fault = Fault::UnknownBefore;
182 assert!(matches!(
183 cache.sync_once(&mut server),
184 Err(Error::Remote(CommitError {
185 state: CommitState::Unknown,
186 ..
187 }))
188 ));
189 save(&cache, object, 0..0, "Queued ").unwrap();
190 let created = onestore::PageCreation::new(None, Some("Recovery 🦀"), "Fixture").unwrap();
191 section_op(&cache, onestore::op::SectionOp::Create(created));
192 save(&cache, object, 0..0, "Last ").unwrap();
193 let working = snapshot(&cache);
194 let remote = remote_snapshot(&cache);
195 let pending = cache.pending().unwrap();
196 let uncertain = cache.status(attempted).unwrap();
197 let source_file = std::fs::read(&path).unwrap();
198 let summary = RecoverySummary {
199 queued_edits: 4,
200 uncertain_edits: 1,
201 published_receipts: 1,
202 base_bytes: remote.len() as u64,
203 remote_bytes: 0,
204 cached_assets: 0,
205 cached_asset_bytes: 0,
206 };
207 assert_eq!(cache.recovery_summary().unwrap(), summary);
208 cache.export_recovery(&archive_path).unwrap();
209 assert_eq!(std::fs::read(&path).unwrap(), source_file);
210 let archive_bytes = std::fs::read(&archive_path).unwrap();
211 assert!(Replica::open(&archive_path).is_err());
212 assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes);
213 assert!(Recovery::open(&path).is_err());
214 let archive = Recovery::open(&archive_path).unwrap();
215 assert_eq!(archive.summary().unwrap(), summary);
216 assert_eq!(pages(&archive.snapshot().unwrap()), pages(&working));
217 assert_eq!(archive.remote_snapshot().unwrap(), remote);
218 assert_eq!(archive.pending().unwrap(), pending);
219 assert_eq!(archive.status(published).unwrap(), Some(receipt.clone()));
220 assert_eq!(archive.status(attempted).unwrap(), uncertain);
221 for edit in &pending {
222 assert_eq!(
223 archive.status(edit.id).unwrap(),
224 cache.status(edit.id).unwrap()
225 );
226 }
227 let EditStatus::Published { revision } = receipt.clone() else {
228 panic!()
229 };
230 assert_eq!(archive.receipts().unwrap(), [(published, revision)].into());
231 assert_eq!(
232 archive.status(u64::MAX - 1).unwrap_err().to_string(),
233 cache.status(u64::MAX - 1).unwrap_err().to_string()
234 );
235 for existing in [&path, &archive_path] {
236 assert!(
237 matches!(cache.export_recovery(existing), Err(Error::Io(error)) if error.kind() == io::ErrorKind::AlreadyExists)
238 );
239 }
240 assert_eq!(std::fs::read(&path).unwrap(), source_file);
241 assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes);
242 save(&cache, object, 0..0, "Later ").unwrap();
243 assert_eq!(pages(&archive.snapshot().unwrap()), pages(&working));
244 assert_eq!(archive.pending().unwrap(), pending);
245 drop(archive);
246 assert_eq!(
247 Recovery::open(&archive_path).unwrap().summary().unwrap(),
248 summary
249 );
250 assert_eq!(std::fs::read(&archive_path).unwrap(), archive_bytes);
251 #[cfg(unix)]
252 {
253 use std::os::unix::fs::PermissionsExt;
254 assert_eq!(
255 std::fs::metadata(&archive_path)
256 .unwrap()
257 .permissions()
258 .mode()
259 & 0o777,
260 0o600
261 );
262 }
263 let remaining: Vec<_> = std::fs::read_dir(directory.path())
264 .unwrap()
265 .map(|entry| entry.unwrap().file_name())
266 .collect();
267 assert!(
268 remaining
269 .iter()
270 .all(|name| !name.to_string_lossy().starts_with(".onestore-recovery-")),
271 "Temporary recovery files remain: {remaining:?}"
272 );
273}
274
275#[test]
276fn recovery_archive_retains_the_queue_and_rejects_foreign_or_future_archives() {
277 use notebook::Recovery;
278
279 let directory = tempfile::tempdir().unwrap();
280 let source = onestore::create_section("recovery.one", "Original", "Fixture").unwrap();
281 let (_, object, _) = text(&source);
282 let cache = Replica::create(directory.path().join("live.sqlite"), &source).unwrap();
283 let id = save(&cache, object, 0..8, "Local").unwrap();
284 let path = directory.path().join("queued.sqlite");
285 cache.export_recovery(&path).unwrap();
286 let archive = Recovery::open(&path).unwrap();
287 assert_eq!(text(&archive.snapshot().unwrap()).2, "Local");
288 assert_eq!(archive.remote_snapshot().unwrap(), source);
289 assert_eq!(archive.status(id).unwrap(), Some(EditStatus::Pending));
290 assert_eq!(archive.summary().unwrap().queued_edits, 1);
291 assert_eq!(archive.summary().unwrap().uncertain_edits, 0);
292 drop(archive);
293 for sql in [
294 "PRAGMA user_version=99",
295 "PRAGMA user_version=4; PRAGMA application_id=0",
296 ] {
297 let connection = rusqlite::Connection::open(&path).unwrap();
298 connection.execute_batch(sql).unwrap();
299 drop(connection);
300 let before = std::fs::read(&path).unwrap();
301 assert!(Recovery::open(&path).is_err());
302 assert_eq!(std::fs::read(&path).unwrap(), before);
303 }
304 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
305 assert_eq!(text(&snapshot(&cache)).2, "Local");
306}
307
308#[test]
309fn disjoint_remote_changes_merge_and_persist_the_remote_receipt() {
310 let dir = tempfile::tempdir().unwrap();
311 let path = dir.path().join("cache.sqlite");
312 let source = onestore::create_section("sync.one", "ab🦀cd", "Fixture").unwrap();
313 let (sid, oid, _) = text(&source);
314 let cache = Replica::create(&path, &source).unwrap();
315 let id = save(&cache, oid, 2..4, "🐈").unwrap();
316 let remote = typed(&source, sid, oid, 0..6, "Xab🦀cYd");
317 let mut server = Server::new(&remote);
318 let outcome = cache.sync_once(&mut server).unwrap().edit.unwrap();
319 assert_eq!(outcome.0, id);
320 assert!(matches!(outcome.1, EditStatus::Published { .. }));
321 assert_eq!(text(&server.durable).2, "Xab🐈cYd");
322 assert_eq!(snapshot(&cache), server.durable);
323 assert!(cache.pending().unwrap().is_empty());
324 drop(cache);
325 let cache = Replica::open(&path).unwrap();
326 assert_eq!(cache.status(id).unwrap(), Some(outcome.1.clone()));
327 assert_eq!(cache.sync_once(&mut server).unwrap().edit, None);
328 assert_eq!(server.publications, 1);
329 assert_eq!(cache.status(id + 1).unwrap(), None);
330 let next = save(&cache, oid, 0..0, "Later ").unwrap();
331 assert!(next > id);
332}
333
334/// As OneNote 2010 does, the remote's version stays the page and the local one becomes a
335/// conflict page under it, published with the rest of the queue.
336#[test]
337fn overlapping_changes_keep_both_versions_and_survive_restart() {
338 let dir = tempfile::tempdir().unwrap();
339 let path = dir.path().join("cache.sqlite");
340 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
341 let (sid, oid, _) = text(&source);
342 let cache = Replica::create(&path, &source).unwrap();
343 let id = save(&cache, oid, 1..2, "L").unwrap();
344 let remote = typed(&source, sid, oid, 1..2, "R");
345 let mut server = Server::new(&remote);
346 assert!(matches!(
347 cache.sync_once(&mut server).unwrap().edit,
348 Some((newest, EditStatus::Published { .. })) if newest > id
349 ));
350 assert_eq!(server.publications, 1);
351 assert_eq!(content(&server.durable, oid), "aRc");
352 assert_eq!(
353 conflicts(&server.durable),
354 [(
355 sid,
356 vec![(model_ops::AUTHOR.to_owned(), vec!["aLc".to_owned()])]
357 )]
358 );
359 assert_eq!(snapshot(&cache), server.durable);
360 assert!(cache.pending().unwrap().is_empty());
361 drop(cache);
362 let cache = Replica::open(&path).unwrap();
363 assert!(matches!(
364 cache.status(id).unwrap(),
365 Some(EditStatus::Published { .. })
366 ));
367 assert_eq!(cache.conflicts().unwrap()[0].0, sid);
368 assert_eq!(cache.sync_once(&mut server).unwrap().edit, None);
369}
370
371#[test]
372fn lost_replies_and_process_termination_never_blindly_replay_an_attempt() {
373 for fault in [
374 Fault::UnknownBefore,
375 Fault::UnknownAfter,
376 Fault::PanicBefore,
377 Fault::PanicAfter,
378 ] {
379 let dir = tempfile::tempdir().unwrap();
380 let path = dir.path().join("cache.sqlite");
381 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
382 let (_, oid, _) = text(&source);
383 let cache = Replica::create(&path, &source).unwrap();
384 let id = save(&cache, oid, 0..0, "Once ").unwrap();
385 let mut server = Server::new(&source);
386 server.fault = fault;
387 let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
388 cache.sync_once(&mut server)
389 }));
390 let attempted = cache.status(id).unwrap().unwrap();
391 assert!(matches!(attempted, EditStatus::AwaitingConfirmation { .. }));
392 drop(cache);
393 let cache = Replica::open(&path).unwrap();
394 let result = cache.sync_once(&mut server).unwrap().edit.unwrap();
395 if matches!(fault, Fault::UnknownAfter | Fault::PanicAfter) {
396 assert!(matches!(result.1, EditStatus::Published { .. }));
397 assert_eq!(text(&server.durable).2, "Once abc");
398 assert_eq!(server.confirmations, 1);
399 assert!(cache.pending().unwrap().is_empty());
400 } else {
401 assert_eq!(result.1, attempted);
402 assert_eq!(cache.pending().unwrap().len(), 1);
403 assert_eq!(server.confirmations, 0);
404 assert_eq!(server.durable, source);
405 }
406 assert_eq!(server.publications, 1);
407 }
408}
409
410#[test]
411fn failed_confirmation_does_not_promote_visible_bytes_to_a_receipt() {
412 let dir = tempfile::tempdir().unwrap();
413 let path = dir.path().join("cache.sqlite");
414 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
415 let (_, oid, _) = text(&source);
416 let cache = Replica::create(&path, &source).unwrap();
417 let id = save(&cache, oid, 0..0, "Once ").unwrap();
418 let mut server = Server::new(&source);
419 server.fault = Fault::UnknownAfter;
420 assert!(cache.sync_once(&mut server).is_err());
421 server.fault = Fault::Confirm;
422 assert!(cache.sync_once(&mut server).is_err());
423 assert_eq!(server.durable, source);
424 assert!(matches!(
425 cache.status(id).unwrap(),
426 Some(EditStatus::AwaitingConfirmation { .. })
427 ));
428 assert!(matches!(
429 cache.sync_once(&mut server).unwrap().edit,
430 Some((_, EditStatus::Published { .. }))
431 ));
432 assert_eq!(server.publications, 1);
433 assert_eq!(server.confirmations, 2);
434 assert_eq!(server.visible, server.durable);
435}
436
437#[test]
438fn confirmation_cleanup_failure_still_records_a_durable_receipt() {
439 let dir = tempfile::tempdir().unwrap();
440 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
441 let (_, oid, _) = text(&source);
442 let cache = Replica::create(dir.path().join("cache.sqlite"), &source).unwrap();
443 let id = save(&cache, oid, 0..0, "Once ").unwrap();
444 let mut server = Server::new(&source);
445 server.fault = Fault::UnknownAfter;
446 assert!(cache.sync_once(&mut server).is_err());
447 server.fault = Fault::ConfirmCommitted;
448 assert!(
449 matches!(cache.sync_once(&mut server), Err(Error::Remote(error)) if error.state == CommitState::Committed)
450 );
451 assert!(matches!(
452 cache.status(id).unwrap(),
453 Some(EditStatus::Published { .. })
454 ));
455 assert!(cache.pending().unwrap().is_empty());
456 assert_eq!(server.visible, server.durable);
457 assert_eq!(server.publications, 1);
458 assert_eq!(server.confirmations, 1);
459}
460
461#[test]
462fn proven_unpublished_attempts_retry_and_committed_cleanup_errors_keep_receipts() {
463 for fault in [Fault::Before, Fault::Committed] {
464 let dir = tempfile::tempdir().unwrap();
465 let path = dir.path().join("cache.sqlite");
466 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
467 let (_, oid, _) = text(&source);
468 let cache = Replica::create(&path, &source).unwrap();
469 let id = save(&cache, oid, 0..0, "Once ").unwrap();
470 let mut server = Server::new(&source);
471 server.fault = fault;
472 assert!(matches!(
473 cache.sync_once(&mut server),
474 Err(Error::Remote(_))
475 ));
476 if matches!(fault, Fault::Before) {
477 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
478 assert!(matches!(
479 cache.sync_once(&mut server).unwrap().edit,
480 Some((_, EditStatus::Published { .. }))
481 ));
482 assert_eq!(server.publications, 2);
483 } else {
484 assert!(matches!(
485 cache.status(id).unwrap(),
486 Some(EditStatus::Published { .. })
487 ));
488 assert_eq!(cache.sync_once(&mut server).unwrap().edit, None);
489 assert_eq!(server.publications, 1);
490 }
491 assert_eq!(text(&server.durable).2, "Once abc");
492 }
493}
494
495#[test]
496fn database_failures_before_and_after_publication_preserve_recovery_state() {
497 for (table, event) in [
498 ("attempt", "UPDATE OF attempted ON batches"),
499 ("receipts", "INSERT ON receipts"),
500 ] {
501 let dir = tempfile::tempdir().unwrap();
502 let path = dir.path().join("cache.sqlite");
503 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
504 let (_, oid, _) = text(&source);
505 let cache = Replica::create(&path, &source).unwrap();
506 let id = save(&cache, oid, 0..0, "Once ").unwrap();
507 drop(cache);
508 let db = rusqlite::Connection::open(&path).unwrap();
509 db.execute_batch(&format!("CREATE TRIGGER interrupted BEFORE {event} BEGIN SELECT RAISE(ABORT,'Injected cache failure'); END;")).unwrap();
510 drop(db);
511 let cache = Replica::open(&path).unwrap();
512 let mut server = Server::new(&source);
513 assert!(matches!(
514 cache.sync_once(&mut server),
515 Err(Error::Database(_))
516 ));
517 assert_eq!(server.publications, usize::from(table == "receipts"));
518 assert_eq!(cache.pending().unwrap().len(), 1);
519 drop(cache);
520 let db = rusqlite::Connection::open(&path).unwrap();
521 db.execute_batch("DROP TRIGGER interrupted").unwrap();
522 drop(db);
523 let cache = Replica::open(&path).unwrap();
524 assert!(matches!(
525 cache.sync_once(&mut server).unwrap().edit,
526 Some((_, EditStatus::Published { .. }))
527 ));
528 assert!(matches!(
529 cache.status(id).unwrap(),
530 Some(EditStatus::Published { .. })
531 ));
532 assert_eq!(server.publications, 1);
533 assert_eq!(server.confirmations, usize::from(table == "receipts"));
534 assert_eq!(text(&server.durable).2, "Once abc");
535 }
536}
537
538#[test]
539fn twelve_local_editors_progress_during_remote_reads_publication_and_confirmation() {
540 struct Paused {
541 server: Server,
542 phase: &'static str,
543 entered: std::sync::mpsc::Sender<()>,
544 resume: std::sync::mpsc::Receiver<()>,
545 }
546 impl Paused {
547 fn wait(&self, phase: &str) {
548 if self.phase == phase {
549 self.entered.send(()).unwrap();
550 self.resume
551 .recv_timeout(std::time::Duration::from_secs(120))
552 .unwrap();
553 }
554 }
555 }
556 impl Remote for Paused {
557 fn read(&mut self) -> io::Result<Vec<u8>> {
558 self.server.read()
559 }
560 fn stamp(&mut self) -> io::Result<onestore::Stamp> {
561 self.wait("read");
562 self.server.stamp()
563 }
564 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
565 self.wait("publish");
566 self.server.publish(transaction)
567 }
568 fn confirm(&mut self, base: &onestore::Stamp) -> Result<(), CommitError> {
569 self.wait("confirm");
570 self.server.confirm(base)
571 }
572 }
573 for phase in ["read", "publish", "confirm"] {
574 let dir = tempfile::tempdir().unwrap();
575 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
576 let (sid, oid, _) = text(&source);
577 let cache = Replica::create(dir.path().join("cache.sqlite"), &source).unwrap();
578 let first = save(&cache, oid, 0..0, "First ").unwrap();
579 let mut server = Server::new(&source);
580 if phase == "confirm" {
581 server.fault = Fault::UnknownAfter;
582 assert!(cache.sync_once(&mut server).is_err());
583 }
584 let (entered_tx, entered_rx) = std::sync::mpsc::channel();
585 let (resume_tx, resume_rx) = std::sync::mpsc::channel();
586 let mut paused = Paused {
587 server,
588 phase,
589 entered: entered_tx,
590 resume: resume_rx,
591 };
592 let mut server = std::thread::scope(|scope| {
593 let running = scope.spawn(|| {
594 let result = cache.sync_once(&mut paused);
595 assert!(matches!(
596 result,
597 Ok(notebook::Synced {
598 edit: Some((_, EditStatus::Published { .. })),
599 ..
600 })
601 ));
602 assert!(matches!(
603 cache.status(first).unwrap(),
604 Some(EditStatus::Published { .. })
605 ));
606 paused.server
607 });
608 entered_rx
609 .recv_timeout(std::time::Duration::from_secs(120))
610 .unwrap();
611 assert!(
612 matches!(cache.sync_once(&mut Server::new(&source)),Err(Error::Io(error)) if error.kind()==io::ErrorKind::WouldBlock)
613 );
614 let started = std::time::Instant::now();
615 let handles: Vec<_> = (0..12)
616 .map(|writer| {
617 let cache = &cache;
618 scope.spawn(move || {
619 // Ops, unlike models, commute with the other writers' edits.
620 cache
621 .apply(
622 model_ops::AUTHOR,
623 onestore::op::Edit {
624 at: model_ops::now(),
625 ops: vec![onestore::op::Op::Page {
626 space: sid,
627 op: onestore::op::PageOp::Text {
628 text: oid,
629 range: 0..0,
630 with: format!("[{writer}] "),
631 },
632 }],
633 },
634 )
635 .unwrap()
636 })
637 })
638 .collect();
639 let ids: std::collections::BTreeSet<_> = handles
640 .into_iter()
641 .map(|handle| handle.join().unwrap())
642 .collect();
643 assert_eq!(ids.len(), 12, "every local edit is its own");
644 // Well inside the paused remote's wait, which edits queued behind it would outlast.
645 assert!(
646 started.elapsed() < std::time::Duration::from_secs(60),
647 "Local edits waited for remote {phase}"
648 );
649 resume_tx.send(()).unwrap();
650 running.join().unwrap()
651 });
652 // Edits arriving while a batch publishes form the next batch.
653 assert_eq!(
654 cache.pending().unwrap().len(),
655 if phase == "read" { 0 } else { 12 }
656 );
657 let expected = text(&snapshot(&cache)).2;
658 for writer in 0..12 {
659 assert_eq!(expected.matches(&format!("[{writer}] ")).count(), 1);
660 }
661 while !cache.pending().unwrap().is_empty() {
662 assert!(matches!(
663 cache.sync_once(&mut server).unwrap().edit,
664 Some((_, EditStatus::Published { .. }))
665 ));
666 }
667 assert_eq!(server.publications, if phase == "read" { 1 } else { 2 });
668 assert_eq!(text(&server.durable).2, expected);
669 assert_eq!(snapshot(&cache), server.durable);
670 }
671}
672
673#[test]
674fn unrelated_remote_files_never_replace_a_local_cache() {
675 let dir = tempfile::tempdir().unwrap();
676 let source = onestore::create_section("sync.one", "abc", "Fixture").unwrap();
677 let (_, oid, _) = text(&source);
678 let cache = Replica::create(dir.path().join("cache.sqlite"), &source).unwrap();
679 let other = onestore::create_section("other.one", "abc", "Fixture").unwrap();
680 let mut server = Server::new(&other);
681 for pending in [false, true] {
682 if pending {
683 save(&cache, oid, 0..0, "Local ").unwrap();
684 }
685 let before = pages(&snapshot(&cache));
686 assert!(
687 matches!(cache.sync_once(&mut server),Err(Error::Io(error)) if error.kind()==io::ErrorKind::InvalidInput)
688 );
689 assert_eq!(pages(&snapshot(&cache)), before);
690 assert_eq!(remote_snapshot(&cache), source);
691 assert_eq!(server.publications, 0);
692 }
693}
694
695/// The text object's content in an image.
696fn content(bytes: &[u8], text: ExGuid) -> String {
697 let (_, page) = model_ops::locate(bytes, text);
698 model_ops::paragraph_with(&page, text)
699 .unwrap()
700 .text()
701 .unwrap()
702 .text
703 .text()
704 .to_owned()
705}
706
707/// A conflict never holds the queue: later edits of the page apply to the remote's version.
708#[test]
709fn edits_after_a_conflict_apply_to_the_remote_version_and_survive_reopen() {
710 let dir = tempfile::tempdir().unwrap();
711 let path = dir.path().join("cache.sqlite");
712 let source = onestore::create_section("resolve.one", "abc", "Fixture").unwrap();
713 let (sid, oid, _) = text(&source);
714 let cache = Replica::create(&path, &source).unwrap();
715 save(&cache, oid, 1..2, "L").unwrap();
716 let remote = typed(&source, sid, oid, 0..3, "aRcZ");
717 let mut server = Server::new(&remote);
718 assert!(matches!(
719 cache.sync_once(&mut server).unwrap().edit,
720 Some((_, EditStatus::Published { .. }))
721 ));
722 assert_eq!(content(&snapshot(&cache), oid), "aRcZ");
723 let mut expected = "aRcZ".to_owned();
724 for n in 0..12 {
725 save(&cache, oid, 0..0, &format!("[{n}] ")).unwrap();
726 expected.insert_str(0, &format!("[{n}] "));
727 }
728 let pending = cache.pending().unwrap();
729 drop(cache);
730 let cache = Replica::open(&path).unwrap();
731 assert_eq!(cache.pending().unwrap(), pending);
732 assert!(matches!(
733 cache.sync_once(&mut server).unwrap().edit,
734 Some((id, EditStatus::Published { .. })) if id == pending.last().unwrap().id
735 ));
736 assert_eq!(server.publications, 2);
737 assert_eq!(content(&server.durable, oid), expected);
738 assert_eq!(
739 conflicts(&server.durable),
740 [(
741 sid,
742 vec![(model_ops::AUTHOR.to_owned(), vec!["aLc".to_owned()])]
743 )]
744 );
745 assert!(cache.pending().unwrap().is_empty());
746}
747
748/// A rebase that cannot be stored leaves the queue, and every page, as it was.
749#[test]
750fn a_failed_rebase_leaves_the_queue_as_it_was() {
751 let dir = tempfile::tempdir().unwrap();
752 let path = dir.path().join("cache.sqlite");
753 let source = onestore::create_section("resolve.one", "abc", "Fixture").unwrap();
754 let (sid, oid, _) = text(&source);
755 let cache = Replica::create(&path, &source).unwrap();
756 let id = save(&cache, oid, 1..2, "L").unwrap();
757 let remote = typed(&source, sid, oid, 0..3, "XaRc");
758 let mut server = Server::new(&remote);
759 let local = pages(&snapshot(&cache));
760 let pending = cache.pending().unwrap();
761 drop(cache);
762 let connection = rusqlite::Connection::open(&path).unwrap();
763 connection.execute_batch("CREATE TRIGGER fail_clear BEFORE DELETE ON batches BEGIN SELECT RAISE(ABORT, 'test rebase failure'); END;").unwrap();
764 drop(connection);
765 let cache = Replica::open(&path).unwrap();
766 assert!(matches!(
767 cache.sync_once(&mut server),
768 Err(Error::Database(_))
769 ));
770 drop(cache);
771 let cache = Replica::open(&path).unwrap();
772 assert_eq!(cache.pending().unwrap(), pending);
773 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
774 assert_eq!(pages(&snapshot(&cache)), local);
775 assert_eq!(server.publications, 0);
776}
777
778#[test]
779fn seeded_conflicts_keep_the_remote_text_and_later_edits_publish_once() {
780 let dir = tempfile::tempdir().unwrap();
781 let original: Vec<_> = "abcdefghij🦀klmnop".chars().collect();
782 let cases = sweep::seeds(0..64, 16);
783 let mut seed = 911 + cases.start;
784 for case in cases {
785 seed ^= seed << 13;
786 seed ^= seed >> 7;
787 seed ^= seed << 17;
788 let at = seed as usize % original.len();
789 let prefix = "[".repeat((seed >> 8) as usize % 5);
790 let suffix = "]".repeat((seed >> 16) as usize % 5);
791 let start: u32 = original[..at].iter().map(|ch| ch.len_utf16() as u32).sum();
792 let end = start + original[at].len_utf16() as u32;
793 let source = onestore::create_section(
794 "resolve.one",
795 &original.iter().collect::<String>(),
796 "Fixture",
797 )
798 .unwrap();
799 let (sid, oid, _) = text(&source);
800 let cache = Replica::create(dir.path().join(format!("{case}.sqlite")), &source).unwrap();
801 save(&cache, oid, start..end, "λ🦊μ").unwrap();
802 save(&cache, oid, start + 1..start + 3, "🐕").unwrap();
803 let local = content(&snapshot(&cache), oid);
804 let mut remote_text = original.clone();
805 remote_text.splice(at..at + 1, "Ω🐈π".chars());
806 let remote_text = prefix.clone() + &remote_text.iter().collect::<String>() + &suffix;
807 let remote = typed(
808 &source,
809 sid,
810 oid,
811 0..original.iter().map(|ch| ch.len_utf16() as u32).sum(),
812 &remote_text,
813 );
814 let mut server = Server::new(&remote);
815 assert!(
816 matches!(
817 cache.sync_once(&mut server).unwrap().edit,
818 Some((_, EditStatus::Published { .. }))
819 ),
820 "case {case}, seed {seed}"
821 );
822 assert_eq!(
823 content(&server.durable, oid),
824 remote_text,
825 "case {case}, seed {seed}"
826 );
827 assert_eq!(
828 conflicts(&server.durable),
829 [(sid, vec![(model_ops::AUTHOR.to_owned(), vec![local])])],
830 "case {case}, seed {seed}"
831 );
832 let at_remote = start + prefix.len() as u32;
833 let later = save(&cache, oid, at_remote..at_remote + 4, "λ🐕μ").unwrap();
834 assert!(
835 matches!(cache.sync_once(&mut server).unwrap().edit, Some((actual, EditStatus::Published { .. })) if actual == later),
836 "case {case}, seed {seed}"
837 );
838 let mut expected = original.clone();
839 expected.splice(at..at + 1, "λ🐕μ".chars());
840 let expected = prefix + &expected.iter().collect::<String>() + &suffix;
841 assert_eq!(
842 content(&server.durable, oid),
843 expected,
844 "case {case}, seed {seed}"
845 );
846 assert_eq!(server.publications, 2);
847 assert!(cache.pending().unwrap().is_empty());
848 }
849}
850
851#[test]
852fn remote_changes_after_a_conflict_are_never_overwritten() {
853 struct ChangedAfterStamp {
854 server: Server,
855 change: Option<Vec<u8>>,
856 }
857 impl Remote for ChangedAfterStamp {
858 fn read(&mut self) -> io::Result<Vec<u8>> {
859 self.server.read()
860 }
861 fn stamp(&mut self) -> io::Result<onestore::Stamp> {
862 let stamp = self.server.stamp()?;
863 if let Some(changed) = self.change.take() {
864 self.server.visible = changed.clone();
865 self.server.durable = changed;
866 }
867 Ok(stamp)
868 }
869 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
870 self.server.publish(transaction)
871 }
872 fn confirm(&mut self, base: &onestore::Stamp) -> Result<(), CommitError> {
873 self.server.confirm(base)
874 }
875 }
876 for after_read in [false, true] {
877 let dir = tempfile::tempdir().unwrap();
878 let source = onestore::create_section("resolve.one", "abc", "Fixture").unwrap();
879 let (sid, oid, _) = text(&source);
880 let cache = Replica::create(dir.path().join("cache.sqlite"), &source).unwrap();
881 save(&cache, oid, 1..2, "L").unwrap();
882 let remote = typed(&source, sid, oid, 1..2, "R");
883 let mut remote = ChangedAfterStamp {
884 server: Server::new(&remote),
885 change: None,
886 };
887 assert!(matches!(
888 cache.sync_once(&mut remote).unwrap().edit,
889 Some((_, EditStatus::Published { .. }))
890 ));
891 // The local version again, typed over the remote's.
892 let id = save(&cache, oid, 1..2, "L").unwrap();
893 assert_eq!(content(&snapshot(&cache), oid), "aLc");
894 let changed = typed(&remote.server.visible, sid, oid, 1..2, "Q");
895 if after_read {
896 remote.change = Some(changed.clone());
897 assert!(
898 matches!(cache.sync_once(&mut remote), Err(Error::Remote(error)) if error.state == CommitState::NotCommitted)
899 );
900 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
901 } else {
902 remote.server.visible = changed.clone();
903 remote.server.durable = changed.clone();
904 }
905 assert!(matches!(
906 cache.sync_once(&mut remote).unwrap().edit,
907 Some((_, EditStatus::Published { .. }))
908 ));
909 assert_eq!(content(&remote.server.durable, oid), "aQc");
910 let versions = [(model_ops::AUTHOR.to_owned(), vec!["aLc".to_owned()])];
911 assert_eq!(
912 conflicts(&remote.server.durable),
913 [(sid, [versions.clone(), versions].concat())]
914 );
915 assert!(cache.pending().unwrap().is_empty());
916 }
917}
918
919#[test]
920fn remote_restore_retains_historical_receipts_without_replaying_them() {
921 let directory = tempfile::tempdir().unwrap();
922 let path = directory.path().join("cache.sqlite");
923 let source = onestore::create_section("restore.one", "Original", "Fixture").unwrap();
924 let (_, object, _) = text(&source);
925 let cache = Replica::create(&path, &source).unwrap();
926 let mut server = Server::new(&source);
927 let published = save(&cache, object, 0..0, "Published ").unwrap();
928 let (_, receipt) = cache.sync_once(&mut server).unwrap().edit.unwrap();
929 assert!(matches!(receipt, EditStatus::Published { .. }));
930 let published_image = snapshot(&cache);
931 cache
932 .export_recovery(directory.path().join("published.sqlite"))
933 .unwrap();
934 drop(cache);
935
936 server.visible.clone_from(&source);
937 server.durable.clone_from(&source);
938 let cache = Replica::open(&path).unwrap();
939 assert_eq!(cache.sync_once(&mut server).unwrap().edit, None);
940 assert_eq!(snapshot(&cache), source);
941 assert_eq!(remote_snapshot(&cache), source);
942 assert_eq!(cache.status(published).unwrap(), Some(receipt.clone()));
943 assert_eq!(server.publications, 1);
944 assert!(cache.pending().unwrap().is_empty());
945 let archive = notebook::Recovery::open(directory.path().join("published.sqlite")).unwrap();
946 assert_eq!(archive.snapshot().unwrap(), published_image);
947 assert_eq!(archive.status(published).unwrap(), Some(receipt.clone()));
948 assert_eq!(text(&archive.snapshot().unwrap()).2, "Published Original");
949 drop(cache);
950 let cache = Replica::open(path).unwrap();
951 assert_eq!(text(&snapshot(&cache)).2, "Original");
952 assert_eq!(cache.status(published).unwrap(), Some(receipt.clone()));
953}
954
955#[test]
956fn remote_restore_rebases_unsent_work_but_never_replays_an_uncertain_attempt() {
957 for fault in [Fault::None, Fault::UnknownBefore, Fault::UnknownAfter] {
958 let uncertain = !matches!(fault, Fault::None);
959 let directory = tempfile::tempdir().unwrap();
960 let path = directory.path().join("cache.sqlite");
961 let source = onestore::create_section("restore.one", "Original", "Fixture").unwrap();
962 let (_, object, _) = text(&source);
963 let cache = Replica::create(&path, &source).unwrap();
964 let mut server = Server::new(&source);
965 let published = save(&cache, object, 0..0, "Published ").unwrap();
966 let (_, receipt) = cache.sync_once(&mut server).unwrap().edit.unwrap();
967 let end = u32::try_from(text(&snapshot(&cache)).2.encode_utf16().count()).unwrap();
968 let queued = save(&cache, object, end..end, " Local").unwrap();
969 if uncertain {
970 server.fault = fault;
971 assert!(cache.sync_once(&mut server).is_err());
972 }
973 let prior_status = cache.status(queued).unwrap();
974 let local = snapshot(&cache);
975 let pending = cache.pending().unwrap();
976 cache
977 .export_recovery(directory.path().join("before-restore.sqlite"))
978 .unwrap();
979 drop(cache);
980 server.visible.clone_from(&source);
981 server.durable.clone_from(&source);
982 let attempts = server.publications;
983 let cache = Replica::open(&path).unwrap();
984 let (_, status) = cache.sync_once(&mut server).unwrap().edit.unwrap();
985 assert_eq!(cache.status(published).unwrap(), Some(receipt.clone()));
986 if uncertain {
987 assert!(matches!(status, EditStatus::AwaitingConfirmation { .. }));
988 assert_eq!(Some(status.clone()), prior_status);
989 assert_eq!(server.publications, attempts);
990 assert_eq!(server.visible, source);
991 assert_eq!(pages(&snapshot(&cache)), pages(&local));
992 assert_eq!(cache.pending().unwrap(), pending);
993 } else {
994 assert!(matches!(status, EditStatus::Published { .. }));
995 assert_eq!(server.publications, attempts + 1);
996 assert_eq!(text(&server.visible).2, "Original Local");
997 assert!(cache.pending().unwrap().is_empty());
998 }
999 let archive =
1000 notebook::Recovery::open(directory.path().join("before-restore.sqlite")).unwrap();
1001 assert_eq!(pages(&archive.snapshot().unwrap()), pages(&local));
1002 assert_eq!(archive.pending().unwrap(), pending);
1003 assert_eq!(archive.status(queued).unwrap(), prior_status);
1004 drop(cache);
1005 let cache = Replica::open(&path).unwrap();
1006 assert_eq!(cache.status(queued).unwrap(), Some(status.clone()));
1007 assert_eq!(cache.status(published).unwrap(), Some(receipt.clone()));
1008 }
1009}
1010
1011/// A remote restored to an older image conflicts with local work on what it undid: the
1012/// restored page stays and the local version becomes its conflict page.
1013#[test]
1014fn a_remote_restore_keeps_local_work_on_a_conflict_page() {
1015 let directory = tempfile::tempdir().unwrap();
1016 let path = directory.path().join("cache.sqlite");
1017 let source = onestore::create_section("restore.one", "Original", "Fixture").unwrap();
1018 let (space, object, _) = text(&source);
1019 let cache = Replica::create(&path, &source).unwrap();
1020 let mut server = Server::new(&source);
1021 let published = save(&cache, object, 0..0, "Published ").unwrap();
1022 let (_, receipt) = cache.sync_once(&mut server).unwrap().edit.unwrap();
1023 let head = save(&cache, object, 0..9, "Revised").unwrap();
1024 server.visible.clone_from(&source);
1025 server.durable.clone_from(&source);
1026 assert!(matches!(
1027 cache.sync_once(&mut server).unwrap().edit,
1028 Some((_, EditStatus::Published { .. }))
1029 ));
1030 assert_eq!(content(&server.durable, object), "Original");
1031 assert_eq!(
1032 conflicts(&server.durable),
1033 [(
1034 space,
1035 vec![(
1036 model_ops::AUTHOR.to_owned(),
1037 vec!["Revised Original".to_owned()]
1038 )]
1039 )]
1040 );
1041 drop(cache);
1042 let cache = Replica::open(&path).unwrap();
1043 assert!(matches!(
1044 cache.status(head).unwrap(),
1045 Some(EditStatus::Published { .. })
1046 ));
1047 assert_eq!(cache.status(published).unwrap(), Some(receipt));
1048 assert!(cache.pending().unwrap().is_empty());
1049}
1050
1051#[test]
1052fn cache_open_validates_the_tail_before_head_only_synchronization() {
1053 for damage in ["edit='{}'", "edit=replace(edit, 'Create', 'Delete')"] {
1054 let directory = tempfile::tempdir().unwrap();
1055 let path = directory.path().join("cache.sqlite");
1056 let source = onestore::create_section("tail.one", "Original", "Fixture").unwrap();
1057 let (_, object, _) = text(&source);
1058 let cache = Replica::create(&path, &source).unwrap();
1059 save(&cache, object, 0..0, "Head ").unwrap();
1060 let page = onestore::PageCreation::new(None, Some("Tail"), "Fixture").unwrap();
1061 let tail = section_op(&cache, onestore::op::SectionOp::Create(page));
1062 assert_eq!(cache.pending().unwrap().len(), 2);
1063 drop(cache);
1064 let database = rusqlite::Connection::open(&path).unwrap();
1065 database
1066 .execute(
1067 &format!("UPDATE edits SET {damage} WHERE id=?1"),
1068 [i64::try_from(tail).unwrap()],
1069 )
1070 .unwrap();
1071 drop(database);
1072 let damaged = std::fs::read(&path).unwrap();
1073 assert!(
1074 matches!(Replica::open(&path), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData)
1075 );
1076 assert_eq!(std::fs::read(&path).unwrap(), damaged);
1077 }
1078}