authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-08 09:26:20-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-08 09:26:20-07:00
log0fdda4e484c31a29b80ec9becba2d5730477032e
tree3b056e9017f60dfe605dddd8d64654a77fd7c700
parentc0d4cbf0d0de5d6a8f0db4455c2165bb952802af
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

test: verify offline paragraph boundaries with twelve native and Rust clients

Extend the durable workload through split, dependent right-text editing and join. Model complete character partitions and paragraph graphs, retired identities, original uncertain receipts and every reader timing interval. Verify real cache trace output independently and retain compatibility with older captures. The 12-client outage run passes 1,536 document publications, 256 Rust appends, 256 native edits and 29,922 reader snapshots. Fresh OneNote reopening matches 261 paragraphs and preserves 655 active identities. Receipt ledgers, server-visible lock overlap and owned VM teardown pass. Add native split/join overlap regressions; 136 Python tests, affected Rust tests and Clippy pass. Assisted-by: gpt-6-astra

5 files changed, 596 insertions(+), 98 deletions(-)

crates/onestore-offline/examples/smb_offline_client.rs+160-30
......@@ -2,8 +2,8 @@
22mod concurrent;
33
44use onestore::{
5 CommitError, CommitState, ExGuid, Insertion, PreparedEdit, RevisionIndex, Store, TextAttribute,
6 document::Document,
5 CommitError, CommitState, ExGuid, Insertion, ParagraphJoin, ParagraphSplit, PreparedEdit,
6 RevisionIndex, Store, TextAttribute, document::Document,
77};
88use onestore_offline::{EditStatus, Error, Remote, Replica, SmbRemote};
99use onestore_smb::{Client, Credentials};
......@@ -41,7 +41,7 @@ impl DocumentView {
4141 }
4242}
4343
44const DOCUMENT_OPERATIONS: [&str; 3] = ["insert", "format", "text"];
44const DOCUMENT_OPERATIONS: [&str; 6] = ["insert", "format", "text", "split", "right_text", "join"];
4545
4646fn now() -> u128 {
4747 SystemTime::now()
......@@ -56,14 +56,14 @@ enum Pause {
5656 FormatReply(PathBuf),
5757}
5858
59struct Traced {
60 remote: SmbRemote,
59struct Traced<R> {
60 remote: R,
6161 before: Option<(String, Option<DocumentView>)>,
6262 pause: Option<Pause>,
6363 documents: bool,
6464}
6565
66impl Remote for Traced {
66impl<R: Remote> Remote for Traced<R> {
6767 fn read(&mut self) -> io::Result<Vec<u8>> {
6868 let started = now();
6969 let bytes = self.remote.read()?;
......@@ -329,7 +329,13 @@ fn queue_document(
329329 } else {
330330 parent.unwrap()
331331 };
332 let range = 1..u32::try_from(text.encode_utf16().count())? - 2;
332 let end = u32::try_from(text.encode_utf16().count())?;
333 let split = ParagraphSplit::new(insertion.text_object(), end - 2, "Offline document writer")?;
334 let join = ParagraphJoin::new(
335 insertion.text_object(),
336 split.text_object(),
337 "Offline document writer",
338 )?;
333339 let attributes = [
334340 TextAttribute::Bold(true),
335341 TextAttribute::FontSize(18.0 + (operation % 9) as f32),
......@@ -338,11 +344,21 @@ fn queue_document(
338344 let mut ids = [0; DOCUMENT_OPERATIONS.len()];
339345 for (step, id) in ids.iter_mut().enumerate() {
340346 let kind = DOCUMENT_OPERATIONS[step];
341 let range = if step == 2 {
342 let end = u32::try_from(text.encode_utf16().count())?;
343 end - 3..end
347 let range = match kind {
348 "text" => end - 3..end,
349 "split" => end - 2..end - 2,
350 "right_text" => 0..2,
351 _ => 1..end - 2,
352 };
353 let target = if kind == "right_text" {
354 split.text_object()
344355 } else {
345 range.clone()
356 insertion.text_object()
357 };
358 let replacement = match kind {
359 "text" => Some(" e\u{301}🐈"),
360 "right_text" => Some("B🦋"),
361 _ => None,
346362 };
347363 loop {
348364 if Instant::now() >= deadline {
......@@ -350,31 +366,28 @@ fn queue_document(
350366 }
351367 let source = cache.snapshot()?;
352368 let started = now();
353 let result = if step == 0 {
354 cache.insert(&source, space, &insertion)
355 } else if step == 1 {
356 cache.format(
357 &source,
358 space,
359 insertion.text_object(),
360 range.clone(),
361 &attributes,
362 )
363 } else {
364 cache.edit_text(
365 &source,
366 space,
367 insertion.text_object(),
368 range.clone(),
369 " e\u{301}🐈",
370 )
369 let result = match kind {
370 "insert" => cache.insert(&source, space, &insertion),
371 "format" => cache.format(&source, space, target, range.clone(), &attributes),
372 "text" | "right_text" => {
373 cache.edit_text(&source, space, target, range.clone(), replacement.unwrap())
374 }
375 "split" => cache.split(&source, space, &split),
376 "join" => cache.join(&source, space, &join),
377 _ => unreachable!(),
371378 };
372379 match result {
373380 Ok(Some(acknowledged)) => {
374381 *id = acknowledged;
375382 println!(
376383 "{}",
377 json!({"event":"local_document_commit","id":acknowledged,"operation":operation,"kind":kind,"space":space.to_string(),"object":insertion.text_object().to_string(),"text":text,"insertion":if step==0 {Some(&insertion)} else {None},"range":[range.start,range.end],"attributes":attributes,"replacement":if step==2 {Some(" e\u{301}🐈")} else {None},"started_us":started,"finished_us":now()})
384 json!({"event":"local_document_commit","id":acknowledged,"operation":operation,"kind":kind,
385 "space":space.to_string(),"object":target.to_string(),"document":insertion.text_object().to_string(),
386 "text":text,"insertion":if kind=="insert" {Some(&insertion)} else {None},
387 "split":if kind=="split" {Some(&split)} else {None},
388 "joined":if kind=="join" {Some(join.texts().map(|id| id.to_string()))} else {None},
389 "range":[range.start,range.end],"attributes":attributes,"replacement":replacement,
390 "started_us":started,"finished_us":now()})
378391 );
379392 break;
380393 }
......@@ -629,10 +642,127 @@ fn append_review_requires_a_unique_ordered_history_and_an_absent_new_token() {
629642 );
630643}
631644
645#[cfg(test)]
646#[path = "../../onestore/tests/support/disk.rs"]
647mod disk;
648
632649#[cfg(test)]
633650mod tests {
634651 use super::*;
635652
653 #[test]
654 fn document_workload_retains_dependencies_and_receipts_across_reopen() {
655 struct Server {
656 disk: disk::Disk,
657 lost_reply: bool,
658 }
659 impl Remote for Server {
660 fn read(&mut self) -> io::Result<Vec<u8>> {
661 Ok(self.disk.visible.clone())
662 }
663 fn publish(&mut self, edit: &PreparedEdit<'_>) -> Result<(), CommitError> {
664 edit.commit(&mut self.disk)?;
665 if std::mem::take(&mut self.lost_reply) {
666 Err(CommitError {
667 state: CommitState::Unknown,
668 error: io::Error::from(io::ErrorKind::ConnectionAborted),
669 })
670 } else {
671 Ok(())
672 }
673 }
674 fn confirm(&mut self, snapshot: &[u8]) -> Result<(), CommitError> {
675 onestore::confirm_snapshot(&mut self.disk, snapshot)
676 }
677 }
678 let source =
679 onestore::create_section("workload.one", "Concurrent edits:", "Author").unwrap();
680 let directory = tempfile::tempdir().unwrap();
681 let path = directory.path().join("cache.sqlite");
682 let mut cache = Replica::create(&path, &source).unwrap();
683 let mut remote = Traced {
684 remote: Server {
685 disk: disk::Disk {
686 visible: source.clone(),
687 durable: source.clone(),
688 operation: 0,
689 fail_at: None,
690 write_limit: 97,
691 random: 1,
692 },
693 lost_reply: false,
694 },
695 before: None,
696 pause: None,
697 documents: true,
698 };
699 println!(
700 "{}",
701 json!({"event":"ready", "actor":"w0", "offline":true,
702 "document_operations":true, "document_graph":true, "document_kinds":DOCUMENT_OPERATIONS})
703 );
704 let deadline = Instant::now() + Duration::from_secs(30);
705 let mut parent = None;
706 let mut ids = Vec::new();
707 for operation in 0..2 {
708 let (outline, added) =
709 queue_document(&cache, "w0", operation, parent, deadline).unwrap();
710 parent = Some(outline);
711 ids.extend(added);
712 drop(cache);
713 cache = Replica::open(&path).unwrap();
714 }
715 let local = cache.snapshot().unwrap();
716 assert_eq!(cache.pending().unwrap().len(), 12);
717 for (step, id) in ids.iter().enumerate() {
718 remote.remote.lost_reply = matches!(step % 6, 3 | 5);
719 let result = cache.sync_once(&mut remote);
720 let state = if result.is_err() {
721 assert!(matches!(
722 result,
723 Err(Error::Remote(CommitError {
724 state: CommitState::Unknown,
725 ..
726 }))
727 ));
728 drop(cache);
729 cache = Replica::open(&path).unwrap();
730 cache.sync_once(&mut remote).unwrap().unwrap()
731 } else {
732 result.unwrap().unwrap()
733 };
734 assert_eq!(state.0, *id);
735 let EditStatus::Published { revision } = state.1 else {
736 panic!("{state:?}")
737 };
738 println!(
739 "{}",
740 json!({"event":"document_receipt", "id":id, "revision":revision.to_string(), "at_us":now()})
741 );
742 drop(cache);
743 cache = Replica::open(&path).unwrap();
744 assert_eq!(
745 cache.status(*id).unwrap(),
746 Some(EditStatus::Published { revision })
747 );
748 if step + 1 < ids.len() {
749 assert_eq!(cache.snapshot().unwrap(), local);
750 }
751 }
752 assert!(cache.pending().unwrap().is_empty());
753 for id in ids {
754 let Some(EditStatus::Published { revision }) = cache.status(id).unwrap() else {
755 panic!()
756 };
757 println!(
758 "{}",
759 json!({"event":"reopened_document_receipt", "id":id, "revision":revision.to_string()})
760 );
761 }
762 remote.read().unwrap();
763 println!("{}", json!({"event":"done"}));
764 }
765
636766 #[test]
637767 fn publication_differences_retain_removed_text_and_graph_identities() {
638768 let before = DocumentView {
crates/onestore-offline/tests/sync.rs+194
......@@ -26,6 +26,200 @@ mod paragraph {
2626 (source, sid, left, parent)
2727 }
2828
29 #[test]
30 fn native_splits_retain_independent_local_identities_and_changed_boundaries() {
31 let source = include_bytes!("../../../corpus/paragraph-edit/before/notebook/synthetic.one");
32 let native = include_bytes!("../../../corpus/paragraph-edit/split/notebook/synthetic.one");
33 let manifest: serde_json::Value =
34 serde_json::from_str(include_str!("../../../corpus/paragraph-edit/manifest.json"))
35 .unwrap();
36 let store = Store::parse(source).unwrap();
37 let index = RevisionIndex::parse(&store).unwrap();
38 let document = Document::parse(&index).unwrap();
39 let mut accepted = 0;
40 let mut conflicts = 0;
41 for case in manifest["cases"].as_array().unwrap() {
42 if case["case"] == "Split before hyperlink" {
43 continue;
44 }
45 let left: ExGuid = serde_json::from_value(case["original_text"].clone()).unwrap();
46 let native_right: ExGuid = serde_json::from_value(case["new_text"].clone()).unwrap();
47 let (sid, _) = document
48 .pages()
49 .unwrap()
50 .into_iter()
51 .find(|(sid, _)| {
52 let space = &document.spaces[sid];
53 space.revisions[&space.contexts[&ExGuid::default()]]
54 .nodes
55 .contains_key(&left)
56 })
57 .unwrap();
58 let split = ParagraphSplit::new(
59 left,
60 u32::try_from(case["offset_utf16"].as_u64().unwrap()).unwrap(),
61 "Offline author",
62 )
63 .unwrap();
64 let directory = tempfile::tempdir().unwrap();
65 let path = directory.path().join("native-split.sqlite");
66 let mut cache = Replica::create(&path, source).unwrap();
67 let id = cache.split(source, sid, &split).unwrap().unwrap();
68 let current = cache.snapshot().unwrap();
69 let dependent = cache
70 .edit_text(
71 &current,
72 sid,
73 split.text_object(),
74 0..0,
75 "Local dependent edit: ",
76 )
77 .unwrap()
78 .unwrap();
79 let local = cache.snapshot().unwrap();
80 let pending = cache.pending().unwrap();
81 let mut server = Server::new(native);
82 let (published, status) = cache.sync_once(&mut server).unwrap().unwrap();
83 assert_eq!(published, id);
84 assert_eq!(cache.snapshot().unwrap(), local);
85 if matches!(case["case"].as_str().unwrap(), "Split end" | "Split empty") {
86 accepted += 1;
87 assert!(
88 matches!(status, EditStatus::Published { .. }),
89 "{}: {status:?}",
90 case["case"]
91 );
92 drop(cache);
93 cache = Replica::open(&path).unwrap();
94 assert_eq!(cache.sync_once(&mut server).unwrap().unwrap().0, dependent);
95 assert_eq!(server.publications, 2);
96 let store = Store::parse(&server.visible).unwrap();
97 let index = RevisionIndex::parse(&store).unwrap();
98 let document = Document::parse(&index).unwrap();
99 let space = &document.spaces[&sid];
100 let view = &space.revisions[&space.contexts[&ExGuid::default()]];
101 let original: ExGuid =
102 serde_json::from_value(case["original_paragraph"].clone()).unwrap();
103 let native_paragraph: ExGuid =
104 serde_json::from_value(case["new_paragraph"].clone()).unwrap();
105 let parent = view
106 .nodes
107 .values()
108 .find(|node| node.children.contains(&original))
109 .unwrap();
110 let position = parent
111 .children
112 .iter()
113 .position(|id| *id == original)
114 .unwrap();
115 assert_eq!(
116 &parent.children[position..position + 3],
117 &[original, split.object(), native_paragraph]
118 );
119 assert!(
120 matches!(&view.nodes[&native_right].kind, Kind::RichText {text,..} if text.is_empty())
121 );
122 assert!(
123 matches!(&view.nodes[&split.text_object()].kind, Kind::RichText {text,..} if text == "Local dependent edit: ")
124 );
125 } else {
126 conflicts += 1;
127 assert!(
128 matches!(
129 status,
130 EditStatus::Conflict(
131 ConflictKind::TextChanged | ConflictKind::StructureChanged
132 )
133 ),
134 "{}: {status:?}",
135 case["case"]
136 );
137 assert_eq!(server.publications, 0);
138 assert_eq!(cache.pending().unwrap(), pending);
139 drop(cache);
140 cache = Replica::open(&path).unwrap();
141 assert_eq!(cache.snapshot().unwrap(), local);
142 assert_eq!(cache.pending().unwrap(), pending);
143 assert_eq!(cache.status(id).unwrap(), Some(status));
144 }
145 }
146 assert_eq!((accepted, conflicts), (2, 10));
147 }
148
149 #[test]
150 fn native_joins_do_not_acknowledge_or_discard_an_independent_local_branch() {
151 let source = include_bytes!(
152 "../../../corpus/paragraph-edit/join-tags/before/notebook/synthetic.one"
153 );
154 let native = include_bytes!(
155 "../../../corpus/paragraph-edit/join-tags/joined/notebook/synthetic.one"
156 );
157 let store = Store::parse(source).unwrap();
158 let index = RevisionIndex::parse(&store).unwrap();
159 let document = Document::parse(&index).unwrap();
160 let mut cases = 0;
161 for (sid, page) in document.pages().unwrap() {
162 let space = &document.spaces[&sid];
163 let view = &space.revisions[&space.contexts[&ExGuid::default()]];
164 let Kind::Metadata {
165 title: Some(name), ..
166 } = &view.nodes[&view.roots[&2]].kind
167 else {
168 continue;
169 };
170 if !name.starts_with("Join ") {
171 continue;
172 }
173 cases += 1;
174 let outline = view.nodes[&page]
175 .children
176 .iter()
177 .find(|id| matches!(view.nodes[id].kind, Kind::Outline { .. }))
178 .unwrap();
179 let children = &view.nodes[outline].children;
180 let left = view.nodes[&children[0]].content[0];
181 let right = view.nodes[&children[1]].content[0];
182 let Kind::RichText {
183 text: left_text, ..
184 } = &view.nodes[&left].kind
185 else {
186 panic!()
187 };
188 let survivor = if left_text.is_empty() { right } else { left };
189 let directory = tempfile::tempdir().unwrap();
190 let path = directory.path().join("native-join.sqlite");
191 let mut cache = Replica::create(&path, source).unwrap();
192 let join = ParagraphJoin::new(left, right, "Offline author").unwrap();
193 let id = cache.join(source, sid, &join).unwrap().unwrap();
194 let current = cache.snapshot().unwrap();
195 cache
196 .edit_text(&current, sid, survivor, 0..0, "Local dependent edit: ")
197 .unwrap()
198 .unwrap();
199 let local = cache.snapshot().unwrap();
200 let pending = cache.pending().unwrap();
201 let mut server = Server::new(native);
202 let (conflicted, status) = cache.sync_once(&mut server).unwrap().unwrap();
203 assert_eq!(conflicted, id);
204 assert!(
205 matches!(
206 status,
207 EditStatus::Conflict(ConflictKind::TargetUnavailable)
208 ),
209 "{name}: {status:?}"
210 );
211 assert_eq!(server.publications, 0);
212 assert_eq!(cache.snapshot().unwrap(), local);
213 assert_eq!(cache.pending().unwrap(), pending);
214 drop(cache);
215 cache = Replica::open(&path).unwrap();
216 assert_eq!(cache.snapshot().unwrap(), local);
217 assert_eq!(cache.pending().unwrap(), pending);
218 assert_eq!(cache.status(id).unwrap(), Some(status));
219 }
220 assert_eq!(cases, 3);
221 }
222
29223 #[test]
30224 fn split_join_dependencies_reopen_and_rebase_with_remote_text_and_styles() {
31225 let (source, sid, left, parent) = fixture();
tools/offline_document_history.py+105-53
......@@ -17,7 +17,8 @@ def characters(observed):
1717
1818def operation_kinds(events):
1919 kinds = tuple(events[0].get('document_kinds', ('insert', 'format')))
20 assert kinds in (('insert', 'format'), ('insert', 'format', 'text')), 'Unknown document workload'
20 assert kinds in (('insert', 'format'), ('insert', 'format', 'text'),
21 ('insert', 'format', 'text', 'split', 'right_text', 'join')), 'Unknown document workload'
2122 return kinds
2223
2324
......@@ -39,15 +40,18 @@ def document_history(logs, operations):
3940 assert [(r['id'], r['revision']) for r in reopened] == [(r['id'], r['revision']) for r in receipts], 'Document receipts changed across reopen'
4041 linked = {}
4142 for intent, receipt in zip(edits, receipts, strict=True):
43 changed_targets = {intent['object']}
44 if intent['kind'] == 'split': changed_targets.add(identity(intent['split'], 2))
45 if intent['kind'] == 'join': changed_targets = set(intent['joined'])
4246 attempts = [row for row in events if row['event'] == 'remote_attempt' and row['revision'] == receipt['revision']
43 and set(row.get('document_changes') or {}) == {intent['object']}]
47 and set(row.get('document_changes') or {}) == changed_targets]
4448 if not attempts and intent['kind'] == 'format':
4549 attempts = [row for row in events if row['event'] == 'remote_attempt' and row['state'] == 'Unknown'
46 and set(row.get('document_changes') or {}) == {intent['object']}]
50 and set(row.get('document_changes') or {}) == changed_targets]
4751 assert len(attempts) == 1, 'Document receipt lacks one publication attempt'
4852 attempt, = attempts
4953 assert attempt['state'] in ('Committed', 'Unknown'), 'Receipt identifies an unpublished document operation'
50 assert intent['object'] in attempt['documents'] and attempt['document_changes'] == {intent['object']: attempt['documents'][intent['object']]}, 'Document publication changed another target'
54 assert intent['object'] in attempt['documents'] and attempt['document_changes'] == {target: attempt['documents'].get(target) for target in changed_targets}, 'Document publication changed another target'
5155 successful = [row for row in events if row['event'] == 'remote_attempt'
5256 and row['state'] in ('Committed', 'Unknown')
5357 and row.get('document_changes') == attempt['document_changes']]
......@@ -71,6 +75,11 @@ def document_history(logs, operations):
7175 else:
7276 assert receipt['revision'] == attempt['revision']
7377 linked[intent['id']] = {**attempt, 'acknowledged_us': receipt['at_us'], 'receipt_revision': receipt['revision']}
78 assert {(row['revision'], row['started_us']) for row in events
79 if row['event'] == 'remote_attempt' and row['state'] in ('Committed', 'Unknown')
80 and row.get('document_changes')} == {
81 (attempt['revision'], attempt['started_us']) for attempt in linked.values()
82 }, 'A document publication lacks its recorded intent and receipt'
7483 for at in range(0, len(edits), len(kinds)):
7584 inserted, formatted, *replaced = edits[at:at + len(kinds)]
7685 number = inserted['operation']
......@@ -102,7 +111,7 @@ def document_history(logs, operations):
102111 states = {'insert': {'characters': old, 'attempt': created},
103112 'format': {'characters': new, 'attempt': changed}}
104113 if replaced:
105 replacement, = replaced
114 replacement, *boundaries = replaced
106115 end = len(text.encode('utf-16-le')) // 2
107116 assert replacement['object'] == target and replacement['space'] == inserted['space'], 'Text edit addresses another object'
108117 assert replacement['text'] == text and replacement['range'] == [end-3, end], 'Cross-run text range differs from workload'
......@@ -112,65 +121,102 @@ def document_history(logs, operations):
112121 assert changed['finished_us'] <= attempt['started_us'], 'Text replacement preceded its formatting'
113122 assert characters(attempt['documents'][target]) == final, 'Text publication differs from its local intent'
114123 states['text'] = {'characters': final, 'attempt': attempt}
124 paragraph = identity(insertion, 3 if 'Outline' in insertion['placement'] else 1)
125 for state in states.values():
126 state['parts'] = [(paragraph, target, 0, len(state['characters']))]
127 if replaced and boundaries:
128 split, right_edit, join = boundaries
129 assert events[0].get('document_graph') is True, 'Boundary workload requires structural observations'
130 assert all(row.get('document') == target and row['space'] == inserted['space']
131 for row in edits[at:at + len(kinds)]), 'Boundary operation lost its owning document'
132 intent = split['split']
133 right, right_paragraph = identity(intent, 2), identity(intent, 1)
134 assert split['object'] == intent['text'] == target and intent['author'] == insertion['author'], 'Split addresses another text or author'
135 assert intent['offset'] == end-2 and split['range'] == [end-2, end-2], 'Split boundary differs from workload'
136 boundary = len(text)-1
137 assert right_edit['object'] == right and right_edit['range'] == [0, 2] and right_edit['replacement'] == 'B🦋', 'Dependent right edit differs from workload'
138 assert join['object'] == target and join['joined'] == [target, right], 'Join does not retain its original targets'
139 updated = final[:boundary] + [(char, *final[boundary][1:]) for char in 'B🦋'] + final[boundary+2:]
140 for row, value, parts in [
141 (split, final, [(paragraph, target, 0, boundary), (right_paragraph, right, boundary, len(final))]),
142 (right_edit, updated, [(paragraph, target, 0, boundary), (right_paragraph, right, boundary, len(updated))]),
143 (join, updated, [(paragraph, target, 0, len(updated))]),
144 ]:
145 attempt = linked[row['id']]
146 assert list(states.values())[-1]['attempt']['finished_us'] <= attempt['started_us'], 'Boundary publication preceded its dependency'
147 assert set(attempt['documents']) & {target, right} == {oid for _, oid, _, _ in parts}, 'Boundary publication omitted or resurrected a text object'
148 for _, oid, start, stop in parts:
149 assert characters(attempt['documents'][oid]) == value[start:stop], 'Boundary publication differs from its local intent'
150 states[row['kind']] = {'characters': value, 'parts': parts, 'attempt': attempt}
151 allocated = {oid for state in states.values() for paragraph, text, _, _ in state['parts'] for oid in (paragraph, text)}
152 existing = {oid for document in documents.values() for state in document['states'].values()
153 for paragraph, text, _, _ in state['parts'] for oid in (paragraph, text)}
154 assert not allocated & existing, 'Two document operations share allocated identities'
115155 documents[target] = {'insertion': insertion, 'space': inserted['space'], 'states': states}
116156 assert documents, 'No document operations were recorded'
117 if any(events[0].get('document_graph') for events in logs.values()):
118 for events in logs.values():
157 structural = any(events[0].get('document_graph') for events in logs.values())
158 allowed_texts = {oid for document in documents.values() for state in document['states'].values()
159 for _, oid, _, _ in state['parts']}
160 for actor, events in logs.items():
161 if structural:
119162 assert events[0].get('document_graph') is True, 'Client omitted structural observations'
120 before = None
121 for row in events:
122 if row['event'] not in ('read', 'document_read', 'remote_attempt'): continue
123 observed = row.get('documents')
124 assert isinstance(observed, dict), 'Snapshot omitted document text'
125 graph = row.get('document_graph')
126 assert isinstance(graph, dict), 'Snapshot omitted document graph'
127 expected = {}
128 for target, document in documents.items():
129 if target not in observed: continue
163 before, previous = None, {}
164 reads = [row for row in events if row['event'] in ('read', 'document_read') and row.get('documents') is not None]
165 assert reads, f'{actor} did not observe document snapshots'
166 assert any(row['documents'] for row in reads), f'{actor} never observed a created document object'
167 for row in events:
168 if row['event'] not in ('read', 'document_read', 'remote_attempt'): continue
169 is_read = row['event'] != 'remote_attempt'
170 observed = row.get('documents')
171 assert isinstance(observed, dict), 'Snapshot omitted document text'
172 assert set(previous) <= set(observed) <= allowed_texts, 'Reader lost an object or observed an unrecorded insertion'
173 graph = row.get('document_graph')
174 if structural: assert isinstance(graph, dict), 'Snapshot omitted document graph'
175 expected_graph = {}
176 for target, document in documents.items():
177 states = list(document['states'].values())
178 known_texts = {oid for state in states for _, oid, _, _ in state['parts']}
179 known_paragraphs = {oid for state in states for oid, _, _, _ in state['parts']}
180 if is_read and row['started_us'] > states[0]['attempt']['acknowledged_us']:
181 assert target in observed, 'Reader missed an acknowledged insertion'
182 if target not in observed:
183 assert not known_texts & observed.keys(), 'Snapshot omitted the original paragraph text'
184 continue
185 actual = {oid: characters(observed[oid]) for oid in known_texts & observed.keys()}
186 matches = [i for i, state in enumerate(states)
187 if actual == {oid: state['characters'][start:stop] for _, oid, start, stop in state['parts']}
188 and (not structural or known_paragraphs & graph.keys() == {oid for oid, _, _, _ in state['parts']})]
189 assert matches, 'Reader observed partial or invented document content'
190 if is_read:
191 matches = [i for i in matches if i >= previous.get(target, 0)]
192 assert matches, 'Reader reverted document content'
193 matches = [i for i in matches if row['finished_us'] >= states[i]['attempt']['started_us']]
194 assert matches, 'Reader observed future document content'
195 acknowledged = max((i for i, state in enumerate(states) if row['started_us'] > state['attempt']['acknowledged_us']), default=0)
196 matches = [i for i in matches if i >= acknowledged]
197 assert matches, 'Reader missed acknowledged document content'
198 previous[target] = min(matches)
199 if structural:
130200 insertion = document['insertion']
131201 outline = insertion['parent']
132 paragraph = identity(insertion, 1)
133202 if 'Outline' in insertion['placement']:
134 outline = paragraph
135 paragraph = identity(insertion, 3)
136 expected[outline] = {'parent': insertion['parent'], 'children': [], 'content': [],
137 'child_level': 1, 'position': insertion['placement']['Outline']}
138 assert outline in expected, 'Snapshot omitted the inserted outline'
139 expected[outline]['children'].append(paragraph)
140 expected[paragraph] = {'parent': outline, 'children': [], 'content': [target],
141 'child_level': 1, 'position': None}
142 assert graph == expected, 'Snapshot contains partial, reordered or invented paragraph structure'
143 if row['event'] == 'remote_attempt':
203 outline = identity(insertion, 1)
204 expected_graph[outline] = {'parent': insertion['parent'], 'children': [], 'content': [],
205 'child_level': 1, 'position': insertion['placement']['Outline']}
206 assert outline in expected_graph, 'Snapshot omitted the inserted outline'
207 for paragraph, oid, _, _ in states[min(matches)]['parts']:
208 expected_graph[outline]['children'].append(paragraph)
209 expected_graph[paragraph] = {'parent': outline, 'children': [], 'content': [oid],
210 'child_level': 1, 'position': None}
211 if structural:
212 assert graph == expected_graph, 'Snapshot contains partial, reordered or invented paragraph structure'
213 if not is_read:
144214 assert before is not None, 'Publication omitted its observed source'
145215 for image, delta in [('documents', 'document_changes'), ('document_graph', 'document_graph_changes')]:
146216 expected_delta = {key: row[image].get(key) for key in before[image].keys() | row[image].keys()
147217 if before[image].get(key) != row[image].get(key)}
148218 assert row.get(delta) == expected_delta, 'Publication diff omitted or invented a changed object'
149 else:
150 before = row
151 for actor, events in logs.items():
152 previous = {}
153 reads = [row for row in events if row['event'] in ('read', 'document_read') and row.get('documents') is not None]
154 assert reads, f'{actor} did not observe document snapshots'
155 assert any(row['documents'] for row in reads), f'{actor} never observed a created document object'
156 for read in reads:
157 observed = read['documents']
158 assert set(previous) <= set(observed) <= set(documents), 'Reader lost an object or observed an unrecorded insertion'
159 for target, document in documents.items():
160 states = list(document['states'].values())
161 if read['started_us'] > states[0]['attempt']['acknowledged_us']:
162 assert target in observed, 'Reader missed an acknowledged insertion'
163 if target not in observed: continue
164 actual = characters(observed[target])
165 matches = [i for i, state in enumerate(states) if actual == state['characters']]
166 assert len(matches) == 1, 'Reader observed partial or invented document content'
167 current, = matches
168 assert current >= previous.get(target, 0), 'Reader reverted document content'
169 assert read['finished_us'] >= states[current]['attempt']['started_us'], 'Reader observed future document content'
170 for i, state in enumerate(states):
171 if read['started_us'] > state['attempt']['acknowledged_us']:
172 assert current >= i, 'Reader missed acknowledged document content'
173 previous[target] = current
219 if is_read: before = row
174220 return documents
175221
176222
......@@ -184,10 +230,15 @@ def verify_model(model, documents):
184230 insertion = document['insertion']
185231 if 'Paragraph' in insertion['placement']:
186232 children[insertion['parent']].append(identity(insertion, 1))
233 retired = {oid for document in documents.values() for state in document['states'].values()
234 for paragraph, text, _, _ in state['parts'] for oid in (paragraph, text)} - {
235 oid for document in documents.values() for paragraph, text, _, _ in list(document['states'].values())[-1]['parts']
236 for oid in (paragraph, text)}
187237 found = set()
188238 for sid, _, revision, page in ordered_pages(model):
189239 nodes = revision['nodes']
190240 for target, node in walk(revision, page):
241 assert target not in retired, 'Final model resurrected a retired paragraph or text'
191242 if target not in documents: continue
192243 assert target not in found, 'Inserted text is reachable twice'
193244 found.add(target)
......@@ -198,6 +249,7 @@ def verify_model(model, documents):
198249 object_id = identity(insertion, 1)
199250 assert object_id in nodes[insertion['parent']]['children'], 'Insertion lost its parent'
200251 paragraph = identity(insertion, 3) if 'Outline' in insertion['placement'] else object_id
252 assert not nodes[paragraph]['children'], 'Final paragraph gained an unrecorded child'
201253 assert nodes[paragraph]['content'] == [target], 'Inserted paragraph content changed'
202254 if 'Outline' in insertion['placement']:
203255 position = insertion['placement']['Outline']
tools/test_offline_document_history.py+136-14
......@@ -1,4 +1,7 @@
11import copy
2import json
3from pathlib import Path
4import subprocess
25import unittest
36import uuid
47from unittest.mock import patch
......@@ -6,8 +9,11 @@ from unittest.mock import patch
69from offline_document_history import document_history, identity, verify_model, verify_native
710
811
9def history(text_edits=False, graph=False):
12def history(text_edits=False, graph=False, boundaries=False):
13 if boundaries: text_edits, graph = True, True
1014 logs = {'w0': [{'event': 'ready', 'document_operations': True, 'document_kinds': ['insert', 'format', 'text'] if text_edits else ['insert', 'format']}]}
15 if boundaries: logs['w0'][0]['document_kinds'] += ['split', 'right_text', 'join']
16 kinds = logs['w0'][0]['document_kinds']
1117 events = logs['w0']
1218 if graph: events[0]['document_graph'] = True
1319 structure = {}
......@@ -19,9 +25,13 @@ def history(text_edits=False, graph=False):
1925 'placement': {'Outline': {'x': 144, 'y': 144}} if operation == 0 else {'Paragraph': {'before': None}}}
2026 previous = insertion
2127 target = identity(insertion, 2)
22 for step, kind in enumerate(('insert', 'format', 'text') if text_edits else ('insert', 'format')):
23 timestamp = 10 + operation*30 + step*10
24 local_id = operation*3 + step + 2
28 end = len(insertion['text'].encode('utf-16-le'))//2
29 split = {'guid': list(uuid.UUID(int=100+operation).bytes_le), 'text': target,
30 'offset': end-2, 'author': insertion['author'], 'created': 1}
31 right, right_paragraph = identity(split, 2), identity(split, 1)
32 for step, kind in enumerate(kinds):
33 timestamp = 10 + operation*max(3, len(kinds))*10 + step*10
34 local_id = operation*max(3, len(kinds)) + step + 2
2535 event = {'event': 'local_document_commit', 'id': local_id, 'operation': operation, 'kind': kind,
2636 'space': 'space', 'object': target, 'text': insertion['text'], 'insertion': insertion if step == 0 else None,
2737 'range': [1, len(insertion['text'].encode('utf-16-le'))//2-2],
......@@ -30,7 +40,12 @@ def history(text_edits=False, graph=False):
3040 if kind == 'text':
3141 end = len(insertion['text'].encode('utf-16-le'))//2
3242 event.update(range=[end-3, end], replacement=' e\u0301🐈')
43 if boundaries: event['document'] = target
44 if kind == 'split': event.update(split=split, range=[end-2, end-2])
45 if kind == 'right_text': event.update(object=right, range=[0, 2], replacement='B🦋')
46 if kind == 'join': event['joined'] = [target, right]
3347 events.append(event)
48 prior_observed = copy.deepcopy(observed)
3449 prior_structure = copy.deepcopy(structure)
3550 if graph:
3651 events.append({'event': 'read', 'started_us': timestamp-1, 'finished_us': timestamp,
......@@ -44,26 +59,45 @@ def history(text_edits=False, graph=False):
4459 structure[parent]['children'].append(paragraph)
4560 structure[paragraph] = {'parent': parent, 'children': [], 'content': [target],
4661 'child_level': 1, 'position': None}
47 value = insertion['text'].replace(' 🦀', ' e\u0301🐈') if kind == 'text' else insertion['text']
48 runs = []
49 for index, char in enumerate(value):
50 selected = index > 0 and (kind == 'text' or (kind == 'format' and index < len(value)-1))
51 runs.append({'text': char, 'bold': selected, 'size': 18+operation if selected else 11,
52 'color': 0x563412 if selected else 0xff000000})
53 observed[target] = {'text': value, 'runs': runs}
62 if kind == 'split':
63 boundary = len(insertion['text'])-1
64 old = observed[target]
65 left_runs, right_runs = old['runs'][:boundary], old['runs'][boundary:]
66 observed[target] = {'text': old['text'][:boundary], 'runs': left_runs}
67 observed[right] = {'text': old['text'][boundary:], 'runs': right_runs}
68 structure[parent]['children'].insert(structure[parent]['children'].index(paragraph)+1, right_paragraph)
69 structure[right_paragraph] = {**copy.deepcopy(structure[paragraph]), 'content': [right]}
70 elif kind == 'right_text':
71 runs = [{**observed[right]['runs'][0], 'text': c} for c in 'B🦋'] + observed[right]['runs'][2:]
72 observed[right] = {'text': 'B🦋🐈', 'runs': runs}
73 elif kind == 'join':
74 observed[target]['text'] += observed[right]['text']
75 observed[target]['runs'] += observed[right]['runs']
76 del observed[right]
77 structure[parent]['children'].remove(right_paragraph)
78 del structure[right_paragraph]
79 else:
80 value = insertion['text'].replace(' 🦀', ' e\u0301🐈') if kind == 'text' else insertion['text']
81 runs = []
82 for index, char in enumerate(value):
83 selected = index > 0 and (kind == 'text' or (kind == 'format' and index < len(value)-1))
84 runs.append({'text': char, 'bold': selected, 'size': 18+operation if selected else 11,
85 'color': 0x563412 if selected else 0xff000000})
86 observed[target] = {'text': value, 'runs': runs}
5487 revision = f'revision-{local_id}'
5588 events.append({'event': 'remote_attempt', 'revision': revision, 'state': 'Committed',
56 'document_changes': {target: copy.deepcopy(observed[target])},
89 'document_changes': {key: copy.deepcopy(observed.get(key)) for key in observed.keys() | prior_observed.keys()
90 if observed.get(key) != prior_observed.get(key)},
5791 'documents': copy.deepcopy(observed), 'started_us': timestamp, 'finished_us': timestamp+1})
5892 if graph:
5993 events[-1].update(document_graph=copy.deepcopy(structure), document_graph_changes={
60 key: copy.deepcopy(value) for key, value in structure.items() if prior_structure.get(key) != value})
94 key: copy.deepcopy(structure.get(key)) for key in structure.keys() | prior_structure.keys() if prior_structure.get(key) != structure.get(key)})
6195 events.append({'event': 'document_receipt', 'id': local_id, 'revision': revision, 'at_us': timestamp+2})
6296 if text_edits:
6397 events.append({'event': 'read', 'started_us': timestamp+3, 'finished_us': timestamp+4, 'documents': copy.deepcopy(observed), **({'document_graph': copy.deepcopy(structure)} if graph else {})})
6498 events.extend({'event': 'reopened_document_receipt', 'id': row['id'], 'revision': row['revision']}
6599 for row in list(events) if row['event'] == 'document_receipt')
66 read = {'event': 'read', 'started_us': 100, 'finished_us': 101, 'documents': observed, **({'document_graph': structure} if graph else {})}
100 read = {'event': 'read', 'started_us': max(100, 20*len(kinds)+10), 'finished_us': max(101, 20*len(kinds)+11), 'documents': observed, **({'document_graph': structure} if graph else {})}
67101 events.extend([read, {'event': 'done'}])
68102 logs['r0'] = [{'event': 'ready', **({'document_graph': True} if graph else {})}, copy.deepcopy(read), {'event': 'done'}]
69103 return logs
......@@ -106,6 +140,94 @@ class DocumentHistoryTests(unittest.TestCase):
106140 with self.subTest(mutation=mutation), self.assertRaises(AssertionError):
107141 document_history(logs, 2)
108142
143 def test_production_cache_workload_matches_the_independent_history_model(self):
144 result = subprocess.run(['cargo', 'test', '--locked', '-p', 'onestore-offline', '--features', 'smb',
145 '--example', 'smb_offline_client',
146 'tests::document_workload_retains_dependencies_and_receipts_across_reopen',
147 '--', '--exact', '--nocapture'],
148 cwd=Path(__file__).resolve().parent.parent, capture_output=True, text=True, timeout=120)
149 self.assertEqual(result.returncode, 0, result.stdout + result.stderr)
150 events = [json.loads(line) for line in result.stdout.splitlines() if line.startswith('{')]
151 documents = document_history({'w0': events}, 2)
152 self.assertEqual(sum(len(document['states']) for document in documents.values()), 12)
153 self.assertEqual(sum(row['event'] == 'remote_attempt' and row['state'] == 'Unknown' for row in events), 4)
154 self.assertEqual(sum(row['event'] == 'remote_confirm' and row['state'] == 'Committed' for row in events), 4)
155
156 def test_boundary_histories_retain_every_intermediate_graph_and_retired_identity(self):
157 baseline = history(boundaries=True)
158 documents = document_history(baseline, 2)
159 for document in documents.values():
160 self.assertEqual(list(document['states']), ['insert', 'format', 'text', 'split', 'right_text', 'join'])
161 self.assertEqual(''.join(c for c, *_ in document['states']['join']['characters']),
162 document['insertion']['text'].replace('🦀', 'B🦋🐈'))
163 for mutation in ('missing-suffix', 'half-join', 'resurrected', 'reorder', 'missing-removal',
164 'wrong-split', 'wrong-dependent-target', 'wrong-join', 'stale', 'future', 'duplicate-attempt'):
165 logs = copy.deepcopy(baseline)
166 split = next(row for row in logs['w0'] if row.get('kind') == 'split')
167 right, right_paragraph = identity(split['split'], 2), identity(split['split'], 1)
168 split_read = next(row for row in logs['w0'] if row['event'] == 'read' and row['started_us'] == 43)
169 joined_read = next(row for row in logs['w0'] if row['event'] == 'read' and row['started_us'] == 63)
170 if mutation == 'missing-suffix': del split_read['documents'][right]
171 elif mutation == 'half-join':
172 joined_read['documents'][right] = copy.deepcopy(split_read['documents'][right])
173 elif mutation == 'resurrected':
174 joined_read['document_graph'][right_paragraph] = copy.deepcopy(split_read['document_graph'][right_paragraph])
175 elif mutation == 'reorder':
176 next(node for node in split_read['document_graph'].values() if node['position'])['children'].reverse()
177 elif mutation == 'missing-removal':
178 join_attempt = next(row for row in logs['w0'] if row['event'] == 'remote_attempt' and row['started_us'] == 60)
179 del join_attempt['document_changes'][right]
180 elif mutation == 'wrong-split': split['split']['offset'] -= 1
181 elif mutation == 'wrong-dependent-target':
182 next(row for row in logs['w0'] if row.get('kind') == 'right_text')['object'] = split['object']
183 elif mutation == 'wrong-join':
184 next(row for row in logs['w0'] if row.get('kind') == 'join')['joined'].reverse()
185 elif mutation == 'stale':
186 joined_read.update(documents=copy.deepcopy(split_read['documents']), document_graph=copy.deepcopy(split_read['document_graph']))
187 elif mutation == 'future': split_read.update(started_us=0, finished_us=1)
188 else:
189 attempt = copy.deepcopy(next(row for row in logs['w0'] if row['event'] == 'remote_attempt' and row['started_us'] == 40))
190 attempt.update(state='Unknown', revision='duplicate')
191 logs['w0'].insert(1, attempt)
192 with self.subTest(mutation=mutation), self.assertRaises(AssertionError):
193 document_history(logs, 2)
194
195 def test_no_reader_phase_can_omit_an_acknowledged_text_or_paragraph(self):
196 baseline = history(boundaries=True)
197 for index, row in enumerate(baseline['w0']):
198 if row['event'] != 'read': continue
199 for field in ('documents', 'document_graph'):
200 for target in row[field]:
201 changed = copy.deepcopy(baseline)
202 del changed['w0'][index][field][target]
203 with self.subTest(read=index, field=field, target=target), self.assertRaises(AssertionError):
204 document_history(changed, 2)
205 extra = copy.deepcopy(next(row for row in baseline['w0'] if row['event'] == 'remote_attempt'))
206 extra.update(revision='unreceipted', started_us=1000, finished_us=1001,
207 document_changes={'unrecorded': {'text': 'extra'}})
208 baseline['w0'].insert(-1, extra)
209 with self.assertRaisesRegex(AssertionError, 'lacks its recorded intent and receipt'):
210 document_history(baseline, 2)
211
212 def test_unknown_boundary_receipts_require_the_original_revision(self):
213 for kind in ('split', 'join'):
214 logs = history(boundaries=True)
215 events = logs['w0']
216 intent = next(row for row in events if row.get('kind') == kind)
217 receipt = next(row for row in events if row['event'] == 'document_receipt' and row['id'] == intent['id'])
218 attempt = next(row for row in events if row['event'] == 'remote_attempt' and row['revision'] == receipt['revision'])
219 attempt['state'] = 'Unknown'
220 confirmation = {'event': 'remote_confirm', 'state': 'Committed',
221 'started_us': attempt['finished_us'], 'finished_us': receipt['at_us'],
222 'revisions': {'space': [receipt['revision']]}}
223 events.insert(events.index(receipt), confirmation)
224 document_history(logs, 2)
225 confirmation['revisions']['space'] = ['later-revision']
226 receipt['revision'] = 'later-revision'
227 next(row for row in events if row['event'] == 'reopened_document_receipt' and row['id'] == intent['id'])['revision'] = 'later-revision'
228 with self.subTest(kind=kind), self.assertRaises(AssertionError):
229 document_history(logs, 2)
230
109231 def test_cross_run_text_requires_complete_ordered_states_and_exact_receipts(self):
110232 logs = history(True)
111233 documents = document_history(logs, 2)
tools/verify_offline.py+1-1
......@@ -32,7 +32,7 @@ def verify(output, max_gap=120):
3232 edits = {event['id']: event for event in events if event['event'] in ('local_commit', 'local_document_commit')}
3333 receipts = {event['id']: event for event in events if event['event'] in ('remote_receipt', 'document_receipt')}
3434 attempts = {event['revision']: event for event in events if event['event'] == 'remote_attempt'}
35 publication_starts = {id: documents[edit['object']]['states'][edit['kind']]['attempt']['started_us'] if edit['event'] == 'local_document_commit'
35 publication_starts = {id: documents[edit.get('document', edit['object'])]['states'][edit['kind']]['attempt']['started_us'] if edit['event'] == 'local_document_commit'
3636 else attempts[receipts[id]['revision']]['started_us'] for id, edit in edits.items()}
3737 queues[actor] = max(sum(edit['finished_us'] <= at < publication_starts[id] for id, edit in edits.items()) for at in [edit['finished_us'] for edit in edits.values()])
3838 assert queues[actor] >= 2, 'Writer did not establish a durable local queue before publication'