1//! Devices sharing a section through a cloud drive that, like iCloud Drive, syncs whole files
2//! with no lock and no compare-and-swap: an upload made from a copy the server has moved past
3//! makes one of the two the file and keeps the other as a conflict version, which every device
4//! sees. Each device commits to its own copy through its replica and merges the conflict
5//! versions it sees (`Remote::versions`). Random schedules of typing, new paragraphs, new,
6//! deleted and moved pages, offline spans and late deliveries must end with every device on
7//! the same file, holding each typed string once.
8
9#[path = "../../onestore/tests/support/sweep.rs"]
10mod sweep;
11
12use notebook::{Remote, Replica, Version};
13use onestore::{
14 CommitError, CommitState, ExGuid, PageCreation, PageEdit, Section, Stamp, Transaction,
15 document::{Format, Layout},
16 op::{Edit, Op, PageOp, SectionOp, lower_page},
17 page::{
18 Outline, Page, PageObject, PageParagraph, Paragraph, ParagraphContent, TextObject,
19 text::new_id,
20 },
21};
22use std::{collections::BTreeSet, io};
23
24/// The file as the server holds it, and the versions it keeps beside it.
25struct Server {
26 current: Vec<u8>,
27 /// Names what `current` is, so a device knows whether its copy is.
28 id: u64,
29 next: u64,
30 conflicts: Vec<(u64, Vec<u8>, String)>,
31 /// Versions a device kept beside the file as files of their own.
32 kept: usize,
33}
34
35impl Server {
36 fn next(&mut self) -> u64 {
37 self.next += 1;
38 self.next
39 }
40}
41
42/// A device's copy of the file and what it has heard of the server.
43struct Copy {
44 bytes: Vec<u8>,
45 /// The server's file this copy was, before any local commit.
46 base: u64,
47 dirty: bool,
48 /// The conflict versions delivered to this device.
49 versions: Vec<(u64, Vec<u8>, String)>,
50 online: bool,
51}
52
53struct Device {
54 name: String,
55 copy: Copy,
56 replica: Replica,
57}
58
59/// A device's copy of the file as its replica publishes to it.
60struct Local<'a> {
61 copy: &'a mut Copy,
62 server: &'a mut Server,
63}
64
65fn refused() -> CommitError {
66 CommitError {
67 state: CommitState::NotCommitted,
68 error: io::Error::from(io::ErrorKind::WouldBlock),
69 }
70}
71
72impl Remote for Local<'_> {
73 fn read(&mut self) -> io::Result<Vec<u8>> {
74 Ok(self.copy.bytes.clone())
75 }
76
77 fn stamp(&mut self) -> io::Result<Stamp> {
78 Stamp::of(&self.copy.bytes).map_err(io::Error::other)
79 }
80
81 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
82 transaction
83 .apply(&mut self.copy.bytes)
84 .map_err(|_| refused())?;
85 self.copy.dirty = true;
86 Ok(())
87 }
88
89 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
90 match Stamp::of(&self.copy.bytes) {
91 Ok(stamp) if stamp == *base => Ok(()),
92 _ => Err(refused()),
93 }
94 }
95
96 fn versions(&mut self) -> io::Result<Vec<Version>> {
97 Ok(self
98 .copy
99 .versions
100 .iter()
101 .map(|(id, _, device)| Version {
102 id: id.to_string(),
103 device: Some(device.clone()),
104 })
105 .collect())
106 }
107
108 fn version(&mut self, id: &str) -> io::Result<Vec<u8>> {
109 self.copy
110 .versions
111 .iter()
112 .find(|(held, ..)| held.to_string() == id)
113 .map(|(_, bytes, _)| bytes.clone())
114 .ok_or_else(|| io::Error::from(io::ErrorKind::NotFound))
115 }
116
117 fn retire(&mut self, id: &str, keep: bool) -> io::Result<()> {
118 self.copy
119 .versions
120 .retain(|(held, ..)| held.to_string() != id);
121 let before = self.server.conflicts.len();
122 self.server
123 .conflicts
124 .retain(|(held, ..)| held.to_string() != id);
125 if keep && before != self.server.conflicts.len() {
126 self.server.kept += 1;
127 }
128 Ok(())
129 }
130}
131
132/// A small xorshift generator; schedules replay from their seed.
133struct Random(u64);
134
135impl Random {
136 fn next(&mut self) -> u64 {
137 self.0 ^= self.0 << 13;
138 self.0 ^= self.0 >> 7;
139 self.0 ^= self.0 << 17;
140 self.0
141 }
142
143 fn below(&mut self, n: usize) -> usize {
144 (self.next() % n.max(1) as u64) as usize
145 }
146}
147
148fn format() -> Format {
149 Format {
150 font: Some("Calibri".to_owned()),
151 font_size: Some(11.0),
152 language: Some(0x409),
153 ..Default::default()
154 }
155}
156
157fn paragraph(text: &str) -> PageParagraph {
158 PageParagraph {
159 id: new_id().unwrap(),
160 parent: None,
161 level: 1,
162 style: None,
163 format: Default::default(),
164 content: ParagraphContent::Text(TextObject {
165 id: new_id().unwrap(),
166 date_field: None,
167 text: Paragraph::new(text.into(), format()),
168 tags: Vec::new(),
169 }),
170 lists: Vec::new(),
171 tags: Vec::new(),
172 media: Default::default(),
173 collapsed: false,
174 }
175}
176
177/// `page` with a body outline holding one paragraph of `text`.
178fn with_body(page: &Page, text: &str) -> Page {
179 let mut page = page.clone();
180 let outline = Outline {
181 id: new_id().unwrap(),
182 title: false,
183 min_width: None,
184 layout: Layout {
185 x: Some(36.0),
186 y: Some(86.0),
187 ..Default::default()
188 },
189 indents: Vec::new(),
190 paragraphs: vec![paragraph(text)],
191 unsupported: Vec::new(),
192 };
193 page.objects.push(PageObject::Outline(outline));
194 page
195}
196
197/// A section of three pages, each with a body paragraph.
198fn fixture() -> Vec<u8> {
199 let source = onestore::create_empty_section("cloud.one", None).unwrap();
200 let arena = onestore::Arena::default();
201 let mut section = Section::open(&arena, source).unwrap();
202 for n in 0..3 {
203 let creation = PageCreation::new(None, Some(&format!("Page {n}")), "Fixture").unwrap();
204 let space = creation.space();
205 edit(
206 &mut section,
207 "Fixture",
208 vec![Op::Section(SectionOp::Create(creation))],
209 );
210 let page = section.page(space).unwrap();
211 let body = with_body(&page, &format!("Body {n} with words"));
212 let ops = lower_page(&page, &body)
213 .unwrap()
214 .into_iter()
215 .map(|op| Op::Page { space, op })
216 .collect();
217 edit(&mut section, "Fixture", ops);
218 }
219 section.seal().unwrap();
220 section.image()
221}
222
223fn edit(section: &mut Section<'_>, author: &str, ops: Vec<Op>) {
224 section
225 .apply(
226 author,
227 &Edit {
228 at: 133_000_000_000_000_000,
229 ops,
230 },
231 )
232 .unwrap();
233}
234
235/// Every body text of `page`: its identity and characters.
236fn bodies(page: &Page) -> Vec<(ExGuid, String)> {
237 page.objects
238 .iter()
239 .filter_map(|object| match object {
240 PageObject::Outline(outline) => Some(outline),
241 _ => None,
242 })
243 .flat_map(|outline| &outline.paragraphs)
244 .filter_map(|paragraph| match &paragraph.content {
245 ParagraphContent::Text(text) => Some((text.id, text.text.text().to_owned())),
246 _ => None,
247 })
248 .collect()
249}
250
251/// What a schedule typed, and where.
252#[derive(Default)]
253struct Typed {
254 /// Each typed string with the page it went into.
255 strings: Vec<(String, ExGuid)>,
256 /// Pages some device deleted, whose strings may be gone.
257 deleted: BTreeSet<ExGuid>,
258}
259
260impl Device {
261 fn remote<'a>(&'a mut self, server: &'a mut Server) -> (Local<'a>, &'a Replica) {
262 (
263 Local {
264 copy: &mut self.copy,
265 server,
266 },
267 &self.replica,
268 )
269 }
270
271 /// Commits what the replica holds and merges what it sees, as the sync thread does.
272 fn sync(&mut self, server: &mut Server) {
273 let (mut remote, replica) = self.remote(server);
274 for _ in 0..8 {
275 match replica.sync_once(&mut remote) {
276 Ok(synced)
277 if matches!(
278 synced.edit,
279 Some((_, notebook::EditStatus::Published { .. }))
280 ) => {}
281 Ok(_) => return,
282 Err(notebook::Error::Remote(error)) if error.state == CommitState::NotCommitted => {
283 }
284 Err(error) => panic!("{}: {error}", self.name),
285 }
286 }
287 }
288
289 /// Sends a changed copy to the server; one made from a file the server has moved past
290 /// becomes the file or a conflict version, either way round.
291 fn upload(&mut self, server: &mut Server, random: &mut Random) {
292 if !self.copy.online || !self.copy.dirty {
293 return;
294 }
295 self.copy.dirty = false;
296 if self.copy.base == server.id {
297 server.current = self.copy.bytes.clone();
298 server.id = server.next();
299 self.copy.base = server.id;
300 return;
301 }
302 let id = server.next();
303 if random.below(2) == 0 {
304 // This copy becomes the file; the server's goes beside it.
305 let lost = std::mem::replace(&mut server.current, self.copy.bytes.clone());
306 server.conflicts.push((id, lost, "The server".into()));
307 server.id = server.next();
308 self.copy.base = server.id;
309 } else {
310 server
311 .conflicts
312 .push((id, self.copy.bytes.clone(), self.name.clone()));
313 self.copy.bytes = server.current.clone();
314 self.copy.base = server.id;
315 }
316 }
317
318 /// Takes the server's file, or conflicts with it while this copy has changes.
319 fn download(&mut self, server: &mut Server, random: &mut Random) {
320 if !self.copy.online || self.copy.base == server.id {
321 return;
322 }
323 if self.copy.dirty {
324 return self.upload(server, random);
325 }
326 self.copy.bytes = server.current.clone();
327 self.copy.base = server.id;
328 }
329
330 /// Learns which conflict versions the server keeps.
331 fn deliver(&mut self, server: &Server) {
332 if self.copy.online {
333 self.copy.versions = server.conflicts.clone();
334 }
335 }
336
337 /// Makes a random edit, as someone using the device would.
338 fn edit(&mut self, random: &mut Random, typed: &mut Typed, serial: &mut u32) {
339 let pages = self.replica.pages().unwrap();
340 let (space, ..) = pages[random.below(pages.len())].clone();
341 let page = self.replica.page(space).unwrap();
342 *serial += 1;
343 let token = format!("⟨{}-{serial}⟩", self.name);
344 let at = crate_now();
345 let ops = match random.below(20) {
346 0..=12 => {
347 let texts = bodies(&page);
348 if texts.is_empty() {
349 return;
350 }
351 let (text, characters) = &texts[random.below(texts.len())];
352 // Anywhere but inside another typed string, so each stays whole to count.
353 let mut inside = false;
354 let mut boundaries = Vec::new();
355 for (at, character) in characters.char_indices() {
356 if !inside {
357 boundaries.push(at);
358 }
359 inside = (inside || character == '⟨') && character != '⟩';
360 }
361 boundaries.push(characters.len());
362 let byte = boundaries[random.below(boundaries.len())];
363 let offset = characters[..byte].encode_utf16().count() as u32;
364 vec![Op::Page {
365 space,
366 op: PageOp::Text {
367 text: *text,
368 range: offset..offset,
369 with: token.clone(),
370 },
371 }]
372 }
373 13..=15 => {
374 let mut after = page.clone();
375 let Some(outline) = after.objects.iter_mut().find_map(|object| match object {
376 PageObject::Outline(outline) => Some(outline),
377 _ => None,
378 }) else {
379 return;
380 };
381 let top: Vec<usize> = (0..outline.paragraphs.len())
382 .filter(|at| outline.paragraphs[*at].parent.is_none())
383 .collect();
384 let at = top[random.below(top.len())];
385 outline.paragraphs.insert(at + 1, paragraph(&token));
386 lower_page(&page, &after)
387 .unwrap()
388 .into_iter()
389 .map(|op| Op::Page { space, op })
390 .collect()
391 }
392 16 => {
393 // First or last, as a series starts at the first page.
394 let before = (random.below(2) == 0).then_some(pages[0].0);
395 let creation = PageCreation::new(before, Some(&token), &self.name).unwrap();
396 let space = creation.space();
397 return self.apply(
398 at,
399 vec![Op::Section(SectionOp::Create(creation))],
400 typed,
401 (token, space),
402 );
403 }
404 17 if pages.len() > 1 => {
405 typed.deleted.insert(space);
406 vec![Op::Section(SectionOp::Delete(vec![space]))]
407 }
408 _ => {
409 let before = pages
410 .get(random.below(pages.len() + 1))
411 .map(|page| page.0)
412 .filter(|before| *before != space);
413 let first = pages[0].0;
414 // A page going first stays at the top level.
415 let level = if before == Some(first) {
416 1
417 } else {
418 1 + random.below(2) as u32
419 };
420 vec![Op::Section(SectionOp::Pages(vec![
421 PageEdit::move_to(space, before, level).unwrap(),
422 ]))]
423 }
424 };
425 self.apply(at, ops, typed, (token, space));
426 }
427
428 /// Applies an edit, noting the string it typed once the section takes it.
429 fn apply(&self, at: u64, ops: Vec<Op>, typed: &mut Typed, token: (String, ExGuid)) {
430 let text = matches!(
431 ops.first(),
432 Some(Op::Page { .. } | Op::Section(SectionOp::Create(_)))
433 );
434 match self.replica.apply(&self.name, Edit { at, ops }) {
435 Ok(_) if text => typed.strings.push(token),
436 Ok(_) => {}
437 Err(notebook::Error::Rejected(error)) if !text => {
438 let _ = error;
439 }
440 Err(error) => panic!("{}: {error}", self.name),
441 }
442 }
443}
444
445fn crate_now() -> u64 {
446 let unix = std::time::SystemTime::now()
447 .duration_since(std::time::UNIX_EPOCH)
448 .unwrap();
449 (unix.as_secs() + 11_644_473_600) * 10_000_000
450}
451
452/// Every text on the section's pages, and on its conflict pages.
453fn texts_of(image: &[u8]) -> (String, String) {
454 let arena = onestore::Arena::default();
455 let mut section = Section::open(&arena, image.to_vec()).unwrap();
456 let mut listed = String::new();
457 for (space, title, _) in section.pages().unwrap() {
458 listed += &title;
459 listed += "\n";
460 for (_, text) in bodies(&section.page(space).unwrap()) {
461 listed += &text;
462 listed += "\n";
463 }
464 }
465 let mut conflicts = String::new();
466 for (_, pages) in section.conflicts().unwrap() {
467 for conflict in pages {
468 let page = section.page(conflict.space).unwrap();
469 conflicts += &page.title;
470 for (_, text) in bodies(&page) {
471 conflicts += &text;
472 conflicts += "\n";
473 }
474 }
475 }
476 (listed, conflicts)
477}
478
479/// Runs a schedule of `events` among `count` devices from `seed`, then lets every device
480/// come online until nothing changes, and checks where they ended.
481fn run(seed: u64, count: usize, events: usize) {
482 let source = fixture();
483 let directory = tempfile::tempdir().unwrap();
484 let mut server = Server {
485 current: source.clone(),
486 id: 0,
487 next: 0,
488 conflicts: Vec::new(),
489 kept: 0,
490 };
491 let mut devices: Vec<Device> = (0..count)
492 .map(|n| Device {
493 name: format!("D{n}"),
494 copy: Copy {
495 bytes: source.clone(),
496 base: 0,
497 dirty: false,
498 versions: Vec::new(),
499 online: true,
500 },
501 replica: Replica::create(directory.path().join(format!("{n}.sqlite")), &source)
502 .unwrap(),
503 })
504 .collect();
505 let mut random = Random(seed.wrapping_mul(0x9e37_79b9_7f4a_7c15) | 1);
506 let mut typed = Typed::default();
507 let mut serial = 0;
508 let trace = std::env::var("CLOUD_TRACE").ok();
509 for event in 0..events {
510 let index = random.below(count);
511 let device = &mut devices[index];
512 let kind = random.below(12);
513 match kind {
514 0..=3 => device.edit(&mut random, &mut typed, &mut serial),
515 4..=5 => device.sync(&mut server),
516 6..=7 => device.upload(&mut server, &mut random),
517 8 => device.download(&mut server, &mut random),
518 9 => device.deliver(&server),
519 _ => device.copy.online = !device.copy.online,
520 }
521 if let Some(token) = &trace {
522 report(
523 token,
524 &format!("{event} D{index} {kind}"),
525 &server,
526 &devices,
527 );
528 }
529 }
530 // Everyone online until a whole round changes nothing.
531 for device in &mut devices {
532 device.copy.online = true;
533 }
534 let mut rounds = 0;
535 loop {
536 rounds += 1;
537 assert!(rounds < 64, "seed {seed}: no quiescence");
538 let before = (server.id, server.conflicts.len());
539 for device in 0..devices.len() {
540 devices[device].download(&mut server, &mut random);
541 devices[device].deliver(&server);
542 devices[device].sync(&mut server);
543 devices[device].upload(&mut server, &mut random);
544 if let Some(token) = &trace {
545 report(
546 token,
547 &format!("round {rounds} D{device}"),
548 &server,
549 &devices,
550 );
551 }
552 }
553 let settled = server.conflicts.is_empty()
554 && (server.id, server.conflicts.len()) == before
555 && devices.iter().all(|device| {
556 !device.copy.dirty
557 && device.copy.base == server.id
558 && device.replica.pending().unwrap().is_empty()
559 });
560 if settled {
561 break;
562 }
563 }
564 for device in &devices {
565 assert_eq!(
566 device.copy.bytes, server.current,
567 "seed {seed}: {}",
568 device.name
569 );
570 }
571 assert_eq!(server.kept, 0, "seed {seed}: a version could not merge");
572 let (listed, conflicts) = texts_of(&server.current);
573 for (string, space) in &typed.strings {
574 // A page deleted on one device while another typed into it comes back as a copy, one
575 // for the version merged and one for the edits queued since, if both hold edits.
576 if typed.deleted.contains(space) {
577 continue;
578 }
579 let shown = listed.matches(string.as_str()).count();
580 assert!(
581 shown <= 1,
582 "seed {seed}: {string} shows {shown} times:\n{listed}"
583 );
584 assert!(
585 shown == 1 || conflicts.contains(string.as_str()),
586 "seed {seed}: {string} is gone:\n{listed}\n--\n{conflicts}"
587 );
588 }
589}
590
591#[test]
592fn devices_on_a_cloud_drive_converge_with_every_edit_once() {
593 for seed in sweep::seeds(0..200, 60) {
594 run(seed, 3, 60);
595 }
596}
597
598#[test]
599fn two_devices_merging_one_version_at_once_add_it_once() {
600 let source = fixture();
601 let directory = tempfile::tempdir().unwrap();
602 let mut server = Server {
603 current: source.clone(),
604 id: 0,
605 next: 0,
606 conflicts: Vec::new(),
607 kept: 0,
608 };
609 let mut devices: Vec<Device> = (0..3)
610 .map(|n| Device {
611 name: format!("D{n}"),
612 copy: Copy {
613 bytes: source.clone(),
614 base: 0,
615 dirty: false,
616 versions: Vec::new(),
617 online: true,
618 },
619 replica: Replica::create(directory.path().join(format!("{n}.sqlite")), &source)
620 .unwrap(),
621 })
622 .collect();
623 let space = devices[0].replica.pages().unwrap()[0].0;
624 let type_at = |device: &Device, token: &str| {
625 let page = device.replica.page(space).unwrap();
626 let (text, _) = bodies(&page)[0].clone();
627 device
628 .replica
629 .apply(
630 &device.name,
631 Edit {
632 at: crate_now(),
633 ops: vec![Op::Page {
634 space,
635 op: PageOp::Text {
636 text,
637 range: 0..0,
638 with: token.into(),
639 },
640 }],
641 },
642 )
643 .unwrap();
644 };
645 // D0 and D1 edit the same page; D1's upload loses and becomes a conflict version.
646 type_at(&devices[0], "⟨zero⟩");
647 type_at(&devices[1], "⟨one⟩");
648 let mut random = Random(1);
649 for device in &mut devices[..2] {
650 device.sync(&mut server);
651 }
652 devices[0].upload(&mut server, &mut random);
653 devices[1].upload(&mut server, &mut random);
654 assert_eq!(server.conflicts.len(), 1);
655 // D0 and D2 both see the version and merge it before either hears the other did.
656 for device in [0, 2] {
657 devices[device].download(&mut server, &mut random);
658 devices[device].deliver(&server);
659 }
660 for device in [0, 2] {
661 devices[device].sync(&mut server);
662 }
663 assert!(server.conflicts.is_empty());
664 devices[0].upload(&mut server, &mut random);
665 // D2's merge lands on the file D0's already changed: a conflict version again.
666 devices[2].upload(&mut server, &mut random);
667 assert_eq!(server.conflicts.len(), 1);
668 for _ in 0..8 {
669 for device in &mut devices {
670 device.download(&mut server, &mut random);
671 device.deliver(&server);
672 device.sync(&mut server);
673 device.upload(&mut server, &mut random);
674 }
675 }
676 assert!(server.conflicts.is_empty());
677 let (listed, conflicts) = texts_of(&server.current);
678 for token in ["⟨zero⟩", "⟨one⟩"] {
679 let shown = listed.matches(token).count() + conflicts.matches(token).count();
680 assert!(
681 (1..=2).contains(&shown) && listed.matches(token).count() <= 1,
682 "{token} shows {shown} times:\n{listed}\n--\n{conflicts}"
683 );
684 }
685 let arena = onestore::Arena::default();
686 let mut section = Section::open(&arena, server.current.clone()).unwrap();
687 let conflict_pages: usize = section
688 .conflicts()
689 .unwrap()
690 .iter()
691 .map(|(_, pages)| pages.len())
692 .sum();
693 assert!(conflict_pages <= 1, "{conflict_pages} conflict pages");
694}
695
696/// A version of a OneNote 2010 file merges into it: an edit to another paragraph replays, a
697/// page the version made comes along, and an edit both made at one place becomes a conflict
698/// page. `ONESTORE_ICLOUD_EXPORT` names a folder to write the merged notebook to, for a cold
699/// read in OneNote 2010.
700#[test]
701fn a_version_of_a_onenote_file_merges_into_it() {
702 let notebook = format!(
703 "{}/../../corpus/conflict-page/native/initial/notebook",
704 env!("CARGO_MANIFEST_DIR")
705 );
706 let source = std::fs::read(format!("{notebook}/synthetic.one")).unwrap();
707 let directory = tempfile::tempdir().unwrap();
708 let mut server = Server {
709 current: source.clone(),
710 id: 0,
711 next: 0,
712 conflicts: Vec::new(),
713 kept: 0,
714 };
715 let mut devices: Vec<Device> = ["Mac", "iPhone"]
716 .into_iter()
717 .enumerate()
718 .map(|(n, name)| Device {
719 name: name.into(),
720 copy: Copy {
721 bytes: source.clone(),
722 base: 0,
723 dirty: false,
724 versions: Vec::new(),
725 online: true,
726 },
727 replica: Replica::create(directory.path().join(format!("{n}.sqlite")), &source)
728 .unwrap(),
729 })
730 .collect();
731 let space = devices[0].replica.pages().unwrap()[0].0;
732 let texts = bodies(&devices[0].replica.page(space).unwrap());
733 let (first, second) = (texts[0].clone(), texts[1].clone());
734 let typing = |text: ExGuid, at: u32, with: &str| Op::Page {
735 space,
736 op: PageOp::Text {
737 text,
738 range: at..at,
739 with: with.into(),
740 },
741 };
742 let apply = |device: &Device, ops: Vec<Op>| {
743 device
744 .replica
745 .apply(
746 &device.name,
747 Edit {
748 at: crate_now(),
749 ops,
750 },
751 )
752 .unwrap();
753 };
754 apply(&devices[0], vec![typing(first.0, 0, "Mac: ")]);
755 let end = second.1.encode_utf16().count() as u32;
756 apply(
757 &devices[1],
758 vec![
759 typing(first.0, 0, "iPhone: "),
760 typing(second.0, end, " (typed on the iPhone)"),
761 ],
762 );
763 let creation = PageCreation::new(None, Some("From the iPhone"), "iPhone").unwrap();
764 let made = creation.space();
765 apply(&devices[1], vec![Op::Section(SectionOp::Create(creation))]);
766 let page = devices[1].replica.page(made).unwrap();
767 let body = with_body(&page, "Written offline on the iPhone.");
768 apply(
769 &devices[1],
770 lower_page(&page, &body)
771 .unwrap()
772 .into_iter()
773 .map(|op| Op::Page { space: made, op })
774 .collect(),
775 );
776 let mut random = Random(3);
777 for device in &mut devices {
778 device.sync(&mut server);
779 device.upload(&mut server, &mut random);
780 }
781 assert_eq!(server.conflicts.len(), 1);
782 let version = server.conflicts[0].clone();
783 for _ in 0..6 {
784 for device in &mut devices {
785 device.download(&mut server, &mut random);
786 device.deliver(&server);
787 device.sync(&mut server);
788 device.upload(&mut server, &mut random);
789 }
790 }
791 assert!(server.conflicts.is_empty());
792 assert!(
793 devices
794 .iter()
795 .all(|device| device.copy.bytes == server.current)
796 );
797 // The version once more, as a device that has not heard it was resolved sees it: merged
798 // already, it adds nothing.
799 devices[1].copy.versions = vec![version];
800 devices[1].sync(&mut server);
801 assert!(!devices[1].copy.dirty && devices[1].copy.versions.is_empty());
802 let (listed, conflicts) = texts_of(&server.current);
803 assert!(listed.contains("(typed on the iPhone)"), "{listed}");
804 assert!(
805 listed.contains("Written offline on the iPhone."),
806 "{listed}"
807 );
808 let arena = onestore::Arena::default();
809 let merged = Section::open(&arena, server.current.clone()).unwrap();
810 let shown: String = bodies(&merged.page(space).unwrap())
811 .into_iter()
812 .map(|(_, text)| text + "\n")
813 .collect();
814 assert_eq!(
815 shown.matches("Mac: ").count() + shown.matches("iPhone: ").count(),
816 1,
817 "{shown}"
818 );
819 assert!(
820 conflicts.contains("Mac: ") || conflicts.contains("iPhone: "),
821 "{conflicts}"
822 );
823 if let Some(output) = std::env::var_os("ONESTORE_ICLOUD_EXPORT") {
824 let output = std::path::Path::new(&output);
825 std::fs::create_dir_all(output).unwrap();
826 std::fs::write(output.join("synthetic.one"), &server.current).unwrap();
827 std::fs::copy(
828 format!("{notebook}/Open Notebook.onetoc2"),
829 output.join("Open Notebook.onetoc2"),
830 )
831 .unwrap();
832 }
833}
834
835/// Where `token` stands after `label`, for tracing a schedule (`CLOUD_TRACE`).
836fn report(token: &str, label: &str, server: &Server, devices: &[Device]) {
837 let has = |image: &[u8]| {
838 let (listed, conflicts) = texts_of(image);
839 format!(
840 "{}{}",
841 listed.matches(token).count(),
842 conflicts.matches(token).count()
843 )
844 };
845 let mut line = format!("{label}: server {}", has(&server.current));
846 for (id, bytes, _) in &server.conflicts {
847 line += &format!(" v{id}:{}", has(bytes));
848 }
849 for device in devices {
850 let working: String = device
851 .replica
852 .pages()
853 .unwrap()
854 .iter()
855 .map(|(space, title, _)| {
856 title.clone()
857 + &bodies(&device.replica.page(*space).unwrap())
858 .into_iter()
859 .map(|(_, text)| text)
860 .collect::<String>()
861 })
862 .collect();
863 line += &format!(
864 " | {} copy {} base {} dirty {} replica {} pending {}",
865 device.name,
866 has(&device.copy.bytes),
867 device.copy.base,
868 u8::from(device.copy.dirty),
869 working.matches(token).count(),
870 device.replica.pending().unwrap().len()
871 );
872 }
873 eprintln!("{line}");
874}