From 0fdda4e484c31a29b80ec9becba2d5730477032e Mon Sep 17 00:00:00 2001 From: clover caruso Date: Tue, 8 Sep 2026 09:26:20 -0700 Subject: [PATCH] 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 --- .../examples/smb_offline_client.rs | 190 ++++++++++++++--- crates/onestore-offline/tests/sync.rs | 194 ++++++++++++++++++ tools/offline_document_history.py | 158 +++++++++----- tools/test_offline_document_history.py | 150 ++++++++++++-- tools/verify_offline.py | 2 +- 5 files changed, 596 insertions(+), 98 deletions(-) diff --git a/crates/onestore-offline/examples/smb_offline_client.rs b/crates/onestore-offline/examples/smb_offline_client.rs index 54b85c47b1da226bf799d7e2fa37415d5e3371c1..3afd08aa5ccbfa29883df2225dcb87a6b9614c37 100644 --- a/crates/onestore-offline/examples/smb_offline_client.rs +++ b/crates/onestore-offline/examples/smb_offline_client.rs @@ -2,8 +2,8 @@ mod concurrent; use onestore::{ - CommitError, CommitState, ExGuid, Insertion, PreparedEdit, RevisionIndex, Store, TextAttribute, - document::Document, + CommitError, CommitState, ExGuid, Insertion, ParagraphJoin, ParagraphSplit, PreparedEdit, + RevisionIndex, Store, TextAttribute, document::Document, }; use onestore_offline::{EditStatus, Error, Remote, Replica, SmbRemote}; use onestore_smb::{Client, Credentials}; @@ -41,7 +41,7 @@ impl DocumentView { } } -const DOCUMENT_OPERATIONS: [&str; 3] = ["insert", "format", "text"]; +const DOCUMENT_OPERATIONS: [&str; 6] = ["insert", "format", "text", "split", "right_text", "join"]; fn now() -> u128 { SystemTime::now() @@ -56,14 +56,14 @@ enum Pause { FormatReply(PathBuf), } -struct Traced { - remote: SmbRemote, +struct Traced { + remote: R, before: Option<(String, Option)>, pause: Option, documents: bool, } -impl Remote for Traced { +impl Remote for Traced { fn read(&mut self) -> io::Result> { let started = now(); let bytes = self.remote.read()?; @@ -329,7 +329,13 @@ fn queue_document( } else { parent.unwrap() }; - let range = 1..u32::try_from(text.encode_utf16().count())? - 2; + let end = u32::try_from(text.encode_utf16().count())?; + let split = ParagraphSplit::new(insertion.text_object(), end - 2, "Offline document writer")?; + let join = ParagraphJoin::new( + insertion.text_object(), + split.text_object(), + "Offline document writer", + )?; let attributes = [ TextAttribute::Bold(true), TextAttribute::FontSize(18.0 + (operation % 9) as f32), @@ -338,11 +344,21 @@ fn queue_document( let mut ids = [0; DOCUMENT_OPERATIONS.len()]; for (step, id) in ids.iter_mut().enumerate() { let kind = DOCUMENT_OPERATIONS[step]; - let range = if step == 2 { - let end = u32::try_from(text.encode_utf16().count())?; - end - 3..end + let range = match kind { + "text" => end - 3..end, + "split" => end - 2..end - 2, + "right_text" => 0..2, + _ => 1..end - 2, + }; + let target = if kind == "right_text" { + split.text_object() } else { - range.clone() + insertion.text_object() + }; + let replacement = match kind { + "text" => Some(" e\u{301}🐈"), + "right_text" => Some("B🦋"), + _ => None, }; loop { if Instant::now() >= deadline { @@ -350,31 +366,28 @@ fn queue_document( } let source = cache.snapshot()?; let started = now(); - let result = if step == 0 { - cache.insert(&source, space, &insertion) - } else if step == 1 { - cache.format( - &source, - space, - insertion.text_object(), - range.clone(), - &attributes, - ) - } else { - cache.edit_text( - &source, - space, - insertion.text_object(), - range.clone(), - " e\u{301}🐈", - ) + let result = match kind { + "insert" => cache.insert(&source, space, &insertion), + "format" => cache.format(&source, space, target, range.clone(), &attributes), + "text" | "right_text" => { + cache.edit_text(&source, space, target, range.clone(), replacement.unwrap()) + } + "split" => cache.split(&source, space, &split), + "join" => cache.join(&source, space, &join), + _ => unreachable!(), }; match result { Ok(Some(acknowledged)) => { *id = acknowledged; println!( "{}", - 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()}) + json!({"event":"local_document_commit","id":acknowledged,"operation":operation,"kind":kind, + "space":space.to_string(),"object":target.to_string(),"document":insertion.text_object().to_string(), + "text":text,"insertion":if kind=="insert" {Some(&insertion)} else {None}, + "split":if kind=="split" {Some(&split)} else {None}, + "joined":if kind=="join" {Some(join.texts().map(|id| id.to_string()))} else {None}, + "range":[range.start,range.end],"attributes":attributes,"replacement":replacement, + "started_us":started,"finished_us":now()}) ); break; } @@ -629,10 +642,127 @@ fn append_review_requires_a_unique_ordered_history_and_an_absent_new_token() { ); } +#[cfg(test)] +#[path = "../../onestore/tests/support/disk.rs"] +mod disk; + #[cfg(test)] mod tests { use super::*; + #[test] + fn document_workload_retains_dependencies_and_receipts_across_reopen() { + struct Server { + disk: disk::Disk, + lost_reply: bool, + } + impl Remote for Server { + fn read(&mut self) -> io::Result> { + Ok(self.disk.visible.clone()) + } + fn publish(&mut self, edit: &PreparedEdit<'_>) -> Result<(), CommitError> { + edit.commit(&mut self.disk)?; + if std::mem::take(&mut self.lost_reply) { + Err(CommitError { + state: CommitState::Unknown, + error: io::Error::from(io::ErrorKind::ConnectionAborted), + }) + } else { + Ok(()) + } + } + fn confirm(&mut self, snapshot: &[u8]) -> Result<(), CommitError> { + onestore::confirm_snapshot(&mut self.disk, snapshot) + } + } + let source = + onestore::create_section("workload.one", "Concurrent edits:", "Author").unwrap(); + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("cache.sqlite"); + let mut cache = Replica::create(&path, &source).unwrap(); + let mut remote = Traced { + remote: Server { + disk: disk::Disk { + visible: source.clone(), + durable: source.clone(), + operation: 0, + fail_at: None, + write_limit: 97, + random: 1, + }, + lost_reply: false, + }, + before: None, + pause: None, + documents: true, + }; + println!( + "{}", + json!({"event":"ready", "actor":"w0", "offline":true, + "document_operations":true, "document_graph":true, "document_kinds":DOCUMENT_OPERATIONS}) + ); + let deadline = Instant::now() + Duration::from_secs(30); + let mut parent = None; + let mut ids = Vec::new(); + for operation in 0..2 { + let (outline, added) = + queue_document(&cache, "w0", operation, parent, deadline).unwrap(); + parent = Some(outline); + ids.extend(added); + drop(cache); + cache = Replica::open(&path).unwrap(); + } + let local = cache.snapshot().unwrap(); + assert_eq!(cache.pending().unwrap().len(), 12); + for (step, id) in ids.iter().enumerate() { + remote.remote.lost_reply = matches!(step % 6, 3 | 5); + let result = cache.sync_once(&mut remote); + let state = if result.is_err() { + assert!(matches!( + result, + Err(Error::Remote(CommitError { + state: CommitState::Unknown, + .. + })) + )); + drop(cache); + cache = Replica::open(&path).unwrap(); + cache.sync_once(&mut remote).unwrap().unwrap() + } else { + result.unwrap().unwrap() + }; + assert_eq!(state.0, *id); + let EditStatus::Published { revision } = state.1 else { + panic!("{state:?}") + }; + println!( + "{}", + json!({"event":"document_receipt", "id":id, "revision":revision.to_string(), "at_us":now()}) + ); + drop(cache); + cache = Replica::open(&path).unwrap(); + assert_eq!( + cache.status(*id).unwrap(), + Some(EditStatus::Published { revision }) + ); + if step + 1 < ids.len() { + assert_eq!(cache.snapshot().unwrap(), local); + } + } + assert!(cache.pending().unwrap().is_empty()); + for id in ids { + let Some(EditStatus::Published { revision }) = cache.status(id).unwrap() else { + panic!() + }; + println!( + "{}", + json!({"event":"reopened_document_receipt", "id":id, "revision":revision.to_string()}) + ); + } + remote.read().unwrap(); + println!("{}", json!({"event":"done"})); + } + #[test] fn publication_differences_retain_removed_text_and_graph_identities() { let before = DocumentView { diff --git a/crates/onestore-offline/tests/sync.rs b/crates/onestore-offline/tests/sync.rs index 3b991310b001555db1d28de3590002f3c01b078a..9062a51cdfdae824e39ac4c086ce85f85cf9a24f 100644 --- a/crates/onestore-offline/tests/sync.rs +++ b/crates/onestore-offline/tests/sync.rs @@ -26,6 +26,200 @@ mod paragraph { (source, sid, left, parent) } + #[test] + fn native_splits_retain_independent_local_identities_and_changed_boundaries() { + let source = include_bytes!("../../../corpus/paragraph-edit/before/notebook/synthetic.one"); + let native = include_bytes!("../../../corpus/paragraph-edit/split/notebook/synthetic.one"); + let manifest: serde_json::Value = + serde_json::from_str(include_str!("../../../corpus/paragraph-edit/manifest.json")) + .unwrap(); + let store = Store::parse(source).unwrap(); + let index = RevisionIndex::parse(&store).unwrap(); + let document = Document::parse(&index).unwrap(); + let mut accepted = 0; + let mut conflicts = 0; + for case in manifest["cases"].as_array().unwrap() { + if case["case"] == "Split before hyperlink" { + continue; + } + let left: ExGuid = serde_json::from_value(case["original_text"].clone()).unwrap(); + let native_right: ExGuid = serde_json::from_value(case["new_text"].clone()).unwrap(); + let (sid, _) = document + .pages() + .unwrap() + .into_iter() + .find(|(sid, _)| { + let space = &document.spaces[sid]; + space.revisions[&space.contexts[&ExGuid::default()]] + .nodes + .contains_key(&left) + }) + .unwrap(); + let split = ParagraphSplit::new( + left, + u32::try_from(case["offset_utf16"].as_u64().unwrap()).unwrap(), + "Offline author", + ) + .unwrap(); + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("native-split.sqlite"); + let mut cache = Replica::create(&path, source).unwrap(); + let id = cache.split(source, sid, &split).unwrap().unwrap(); + let current = cache.snapshot().unwrap(); + let dependent = cache + .edit_text( + ¤t, + sid, + split.text_object(), + 0..0, + "Local dependent edit: ", + ) + .unwrap() + .unwrap(); + let local = cache.snapshot().unwrap(); + let pending = cache.pending().unwrap(); + let mut server = Server::new(native); + let (published, status) = cache.sync_once(&mut server).unwrap().unwrap(); + assert_eq!(published, id); + assert_eq!(cache.snapshot().unwrap(), local); + if matches!(case["case"].as_str().unwrap(), "Split end" | "Split empty") { + accepted += 1; + assert!( + matches!(status, EditStatus::Published { .. }), + "{}: {status:?}", + case["case"] + ); + drop(cache); + cache = Replica::open(&path).unwrap(); + assert_eq!(cache.sync_once(&mut server).unwrap().unwrap().0, dependent); + assert_eq!(server.publications, 2); + let store = Store::parse(&server.visible).unwrap(); + let index = RevisionIndex::parse(&store).unwrap(); + let document = Document::parse(&index).unwrap(); + let space = &document.spaces[&sid]; + let view = &space.revisions[&space.contexts[&ExGuid::default()]]; + let original: ExGuid = + serde_json::from_value(case["original_paragraph"].clone()).unwrap(); + let native_paragraph: ExGuid = + serde_json::from_value(case["new_paragraph"].clone()).unwrap(); + let parent = view + .nodes + .values() + .find(|node| node.children.contains(&original)) + .unwrap(); + let position = parent + .children + .iter() + .position(|id| *id == original) + .unwrap(); + assert_eq!( + &parent.children[position..position + 3], + &[original, split.object(), native_paragraph] + ); + assert!( + matches!(&view.nodes[&native_right].kind, Kind::RichText {text,..} if text.is_empty()) + ); + assert!( + matches!(&view.nodes[&split.text_object()].kind, Kind::RichText {text,..} if text == "Local dependent edit: ") + ); + } else { + conflicts += 1; + assert!( + matches!( + status, + EditStatus::Conflict( + ConflictKind::TextChanged | ConflictKind::StructureChanged + ) + ), + "{}: {status:?}", + case["case"] + ); + assert_eq!(server.publications, 0); + assert_eq!(cache.pending().unwrap(), pending); + drop(cache); + cache = Replica::open(&path).unwrap(); + assert_eq!(cache.snapshot().unwrap(), local); + assert_eq!(cache.pending().unwrap(), pending); + assert_eq!(cache.status(id).unwrap(), Some(status)); + } + } + assert_eq!((accepted, conflicts), (2, 10)); + } + + #[test] + fn native_joins_do_not_acknowledge_or_discard_an_independent_local_branch() { + let source = include_bytes!( + "../../../corpus/paragraph-edit/join-tags/before/notebook/synthetic.one" + ); + let native = include_bytes!( + "../../../corpus/paragraph-edit/join-tags/joined/notebook/synthetic.one" + ); + let store = Store::parse(source).unwrap(); + let index = RevisionIndex::parse(&store).unwrap(); + let document = Document::parse(&index).unwrap(); + let mut cases = 0; + for (sid, page) in document.pages().unwrap() { + let space = &document.spaces[&sid]; + let view = &space.revisions[&space.contexts[&ExGuid::default()]]; + let Kind::Metadata { + title: Some(name), .. + } = &view.nodes[&view.roots[&2]].kind + else { + continue; + }; + if !name.starts_with("Join ") { + continue; + } + cases += 1; + let outline = view.nodes[&page] + .children + .iter() + .find(|id| matches!(view.nodes[id].kind, Kind::Outline { .. })) + .unwrap(); + let children = &view.nodes[outline].children; + let left = view.nodes[&children[0]].content[0]; + let right = view.nodes[&children[1]].content[0]; + let Kind::RichText { + text: left_text, .. + } = &view.nodes[&left].kind + else { + panic!() + }; + let survivor = if left_text.is_empty() { right } else { left }; + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("native-join.sqlite"); + let mut cache = Replica::create(&path, source).unwrap(); + let join = ParagraphJoin::new(left, right, "Offline author").unwrap(); + let id = cache.join(source, sid, &join).unwrap().unwrap(); + let current = cache.snapshot().unwrap(); + cache + .edit_text(¤t, sid, survivor, 0..0, "Local dependent edit: ") + .unwrap() + .unwrap(); + let local = cache.snapshot().unwrap(); + let pending = cache.pending().unwrap(); + let mut server = Server::new(native); + let (conflicted, status) = cache.sync_once(&mut server).unwrap().unwrap(); + assert_eq!(conflicted, id); + assert!( + matches!( + status, + EditStatus::Conflict(ConflictKind::TargetUnavailable) + ), + "{name}: {status:?}" + ); + assert_eq!(server.publications, 0); + assert_eq!(cache.snapshot().unwrap(), local); + assert_eq!(cache.pending().unwrap(), pending); + drop(cache); + cache = Replica::open(&path).unwrap(); + assert_eq!(cache.snapshot().unwrap(), local); + assert_eq!(cache.pending().unwrap(), pending); + assert_eq!(cache.status(id).unwrap(), Some(status)); + } + assert_eq!(cases, 3); + } + #[test] fn split_join_dependencies_reopen_and_rebase_with_remote_text_and_styles() { let (source, sid, left, parent) = fixture(); diff --git a/tools/offline_document_history.py b/tools/offline_document_history.py index e4d7f0af093ea292fb72d6237f27e6837f54aeb8..b6f0d007c4f66e2d6ad7bee21e2d9cfe430fad45 100644 --- a/tools/offline_document_history.py +++ b/tools/offline_document_history.py @@ -17,7 +17,8 @@ def characters(observed): def operation_kinds(events): kinds = tuple(events[0].get('document_kinds', ('insert', 'format'))) - assert kinds in (('insert', 'format'), ('insert', 'format', 'text')), 'Unknown document workload' + assert kinds in (('insert', 'format'), ('insert', 'format', 'text'), + ('insert', 'format', 'text', 'split', 'right_text', 'join')), 'Unknown document workload' return kinds @@ -39,15 +40,18 @@ def document_history(logs, operations): assert [(r['id'], r['revision']) for r in reopened] == [(r['id'], r['revision']) for r in receipts], 'Document receipts changed across reopen' linked = {} for intent, receipt in zip(edits, receipts, strict=True): + changed_targets = {intent['object']} + if intent['kind'] == 'split': changed_targets.add(identity(intent['split'], 2)) + if intent['kind'] == 'join': changed_targets = set(intent['joined']) attempts = [row for row in events if row['event'] == 'remote_attempt' and row['revision'] == receipt['revision'] - and set(row.get('document_changes') or {}) == {intent['object']}] + and set(row.get('document_changes') or {}) == changed_targets] if not attempts and intent['kind'] == 'format': attempts = [row for row in events if row['event'] == 'remote_attempt' and row['state'] == 'Unknown' - and set(row.get('document_changes') or {}) == {intent['object']}] + and set(row.get('document_changes') or {}) == changed_targets] assert len(attempts) == 1, 'Document receipt lacks one publication attempt' attempt, = attempts assert attempt['state'] in ('Committed', 'Unknown'), 'Receipt identifies an unpublished document operation' - assert intent['object'] in attempt['documents'] and attempt['document_changes'] == {intent['object']: attempt['documents'][intent['object']]}, 'Document publication changed another target' + 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' successful = [row for row in events if row['event'] == 'remote_attempt' and row['state'] in ('Committed', 'Unknown') and row.get('document_changes') == attempt['document_changes']] @@ -71,6 +75,11 @@ def document_history(logs, operations): else: assert receipt['revision'] == attempt['revision'] linked[intent['id']] = {**attempt, 'acknowledged_us': receipt['at_us'], 'receipt_revision': receipt['revision']} + assert {(row['revision'], row['started_us']) for row in events + if row['event'] == 'remote_attempt' and row['state'] in ('Committed', 'Unknown') + and row.get('document_changes')} == { + (attempt['revision'], attempt['started_us']) for attempt in linked.values() + }, 'A document publication lacks its recorded intent and receipt' for at in range(0, len(edits), len(kinds)): inserted, formatted, *replaced = edits[at:at + len(kinds)] number = inserted['operation'] @@ -102,7 +111,7 @@ def document_history(logs, operations): states = {'insert': {'characters': old, 'attempt': created}, 'format': {'characters': new, 'attempt': changed}} if replaced: - replacement, = replaced + replacement, *boundaries = replaced end = len(text.encode('utf-16-le')) // 2 assert replacement['object'] == target and replacement['space'] == inserted['space'], 'Text edit addresses another object' 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): assert changed['finished_us'] <= attempt['started_us'], 'Text replacement preceded its formatting' assert characters(attempt['documents'][target]) == final, 'Text publication differs from its local intent' states['text'] = {'characters': final, 'attempt': attempt} + paragraph = identity(insertion, 3 if 'Outline' in insertion['placement'] else 1) + for state in states.values(): + state['parts'] = [(paragraph, target, 0, len(state['characters']))] + if replaced and boundaries: + split, right_edit, join = boundaries + assert events[0].get('document_graph') is True, 'Boundary workload requires structural observations' + assert all(row.get('document') == target and row['space'] == inserted['space'] + for row in edits[at:at + len(kinds)]), 'Boundary operation lost its owning document' + intent = split['split'] + right, right_paragraph = identity(intent, 2), identity(intent, 1) + assert split['object'] == intent['text'] == target and intent['author'] == insertion['author'], 'Split addresses another text or author' + assert intent['offset'] == end-2 and split['range'] == [end-2, end-2], 'Split boundary differs from workload' + boundary = len(text)-1 + assert right_edit['object'] == right and right_edit['range'] == [0, 2] and right_edit['replacement'] == 'B🦋', 'Dependent right edit differs from workload' + assert join['object'] == target and join['joined'] == [target, right], 'Join does not retain its original targets' + updated = final[:boundary] + [(char, *final[boundary][1:]) for char in 'B🦋'] + final[boundary+2:] + for row, value, parts in [ + (split, final, [(paragraph, target, 0, boundary), (right_paragraph, right, boundary, len(final))]), + (right_edit, updated, [(paragraph, target, 0, boundary), (right_paragraph, right, boundary, len(updated))]), + (join, updated, [(paragraph, target, 0, len(updated))]), + ]: + attempt = linked[row['id']] + assert list(states.values())[-1]['attempt']['finished_us'] <= attempt['started_us'], 'Boundary publication preceded its dependency' + assert set(attempt['documents']) & {target, right} == {oid for _, oid, _, _ in parts}, 'Boundary publication omitted or resurrected a text object' + for _, oid, start, stop in parts: + assert characters(attempt['documents'][oid]) == value[start:stop], 'Boundary publication differs from its local intent' + states[row['kind']] = {'characters': value, 'parts': parts, 'attempt': attempt} + allocated = {oid for state in states.values() for paragraph, text, _, _ in state['parts'] for oid in (paragraph, text)} + existing = {oid for document in documents.values() for state in document['states'].values() + for paragraph, text, _, _ in state['parts'] for oid in (paragraph, text)} + assert not allocated & existing, 'Two document operations share allocated identities' documents[target] = {'insertion': insertion, 'space': inserted['space'], 'states': states} assert documents, 'No document operations were recorded' - if any(events[0].get('document_graph') for events in logs.values()): - for events in logs.values(): + structural = any(events[0].get('document_graph') for events in logs.values()) + allowed_texts = {oid for document in documents.values() for state in document['states'].values() + for _, oid, _, _ in state['parts']} + for actor, events in logs.items(): + if structural: assert events[0].get('document_graph') is True, 'Client omitted structural observations' - before = None - for row in events: - if row['event'] not in ('read', 'document_read', 'remote_attempt'): continue - observed = row.get('documents') - assert isinstance(observed, dict), 'Snapshot omitted document text' - graph = row.get('document_graph') - assert isinstance(graph, dict), 'Snapshot omitted document graph' - expected = {} - for target, document in documents.items(): - if target not in observed: continue + before, previous = None, {} + reads = [row for row in events if row['event'] in ('read', 'document_read') and row.get('documents') is not None] + assert reads, f'{actor} did not observe document snapshots' + assert any(row['documents'] for row in reads), f'{actor} never observed a created document object' + for row in events: + if row['event'] not in ('read', 'document_read', 'remote_attempt'): continue + is_read = row['event'] != 'remote_attempt' + observed = row.get('documents') + assert isinstance(observed, dict), 'Snapshot omitted document text' + assert set(previous) <= set(observed) <= allowed_texts, 'Reader lost an object or observed an unrecorded insertion' + graph = row.get('document_graph') + if structural: assert isinstance(graph, dict), 'Snapshot omitted document graph' + expected_graph = {} + for target, document in documents.items(): + states = list(document['states'].values()) + known_texts = {oid for state in states for _, oid, _, _ in state['parts']} + known_paragraphs = {oid for state in states for oid, _, _, _ in state['parts']} + if is_read and row['started_us'] > states[0]['attempt']['acknowledged_us']: + assert target in observed, 'Reader missed an acknowledged insertion' + if target not in observed: + assert not known_texts & observed.keys(), 'Snapshot omitted the original paragraph text' + continue + actual = {oid: characters(observed[oid]) for oid in known_texts & observed.keys()} + matches = [i for i, state in enumerate(states) + if actual == {oid: state['characters'][start:stop] for _, oid, start, stop in state['parts']} + and (not structural or known_paragraphs & graph.keys() == {oid for oid, _, _, _ in state['parts']})] + assert matches, 'Reader observed partial or invented document content' + if is_read: + matches = [i for i in matches if i >= previous.get(target, 0)] + assert matches, 'Reader reverted document content' + matches = [i for i in matches if row['finished_us'] >= states[i]['attempt']['started_us']] + assert matches, 'Reader observed future document content' + acknowledged = max((i for i, state in enumerate(states) if row['started_us'] > state['attempt']['acknowledged_us']), default=0) + matches = [i for i in matches if i >= acknowledged] + assert matches, 'Reader missed acknowledged document content' + previous[target] = min(matches) + if structural: insertion = document['insertion'] outline = insertion['parent'] - paragraph = identity(insertion, 1) if 'Outline' in insertion['placement']: - outline = paragraph - paragraph = identity(insertion, 3) - expected[outline] = {'parent': insertion['parent'], 'children': [], 'content': [], - 'child_level': 1, 'position': insertion['placement']['Outline']} - assert outline in expected, 'Snapshot omitted the inserted outline' - expected[outline]['children'].append(paragraph) - expected[paragraph] = {'parent': outline, 'children': [], 'content': [target], - 'child_level': 1, 'position': None} - assert graph == expected, 'Snapshot contains partial, reordered or invented paragraph structure' - if row['event'] == 'remote_attempt': + outline = identity(insertion, 1) + expected_graph[outline] = {'parent': insertion['parent'], 'children': [], 'content': [], + 'child_level': 1, 'position': insertion['placement']['Outline']} + assert outline in expected_graph, 'Snapshot omitted the inserted outline' + for paragraph, oid, _, _ in states[min(matches)]['parts']: + expected_graph[outline]['children'].append(paragraph) + expected_graph[paragraph] = {'parent': outline, 'children': [], 'content': [oid], + 'child_level': 1, 'position': None} + if structural: + assert graph == expected_graph, 'Snapshot contains partial, reordered or invented paragraph structure' + if not is_read: assert before is not None, 'Publication omitted its observed source' for image, delta in [('documents', 'document_changes'), ('document_graph', 'document_graph_changes')]: expected_delta = {key: row[image].get(key) for key in before[image].keys() | row[image].keys() if before[image].get(key) != row[image].get(key)} assert row.get(delta) == expected_delta, 'Publication diff omitted or invented a changed object' - else: - before = row - for actor, events in logs.items(): - previous = {} - reads = [row for row in events if row['event'] in ('read', 'document_read') and row.get('documents') is not None] - assert reads, f'{actor} did not observe document snapshots' - assert any(row['documents'] for row in reads), f'{actor} never observed a created document object' - for read in reads: - observed = read['documents'] - assert set(previous) <= set(observed) <= set(documents), 'Reader lost an object or observed an unrecorded insertion' - for target, document in documents.items(): - states = list(document['states'].values()) - if read['started_us'] > states[0]['attempt']['acknowledged_us']: - assert target in observed, 'Reader missed an acknowledged insertion' - if target not in observed: continue - actual = characters(observed[target]) - matches = [i for i, state in enumerate(states) if actual == state['characters']] - assert len(matches) == 1, 'Reader observed partial or invented document content' - current, = matches - assert current >= previous.get(target, 0), 'Reader reverted document content' - assert read['finished_us'] >= states[current]['attempt']['started_us'], 'Reader observed future document content' - for i, state in enumerate(states): - if read['started_us'] > state['attempt']['acknowledged_us']: - assert current >= i, 'Reader missed acknowledged document content' - previous[target] = current + if is_read: before = row return documents @@ -184,10 +230,15 @@ def verify_model(model, documents): insertion = document['insertion'] if 'Paragraph' in insertion['placement']: children[insertion['parent']].append(identity(insertion, 1)) + retired = {oid for document in documents.values() for state in document['states'].values() + for paragraph, text, _, _ in state['parts'] for oid in (paragraph, text)} - { + oid for document in documents.values() for paragraph, text, _, _ in list(document['states'].values())[-1]['parts'] + for oid in (paragraph, text)} found = set() for sid, _, revision, page in ordered_pages(model): nodes = revision['nodes'] for target, node in walk(revision, page): + assert target not in retired, 'Final model resurrected a retired paragraph or text' if target not in documents: continue assert target not in found, 'Inserted text is reachable twice' found.add(target) @@ -198,6 +249,7 @@ def verify_model(model, documents): object_id = identity(insertion, 1) assert object_id in nodes[insertion['parent']]['children'], 'Insertion lost its parent' paragraph = identity(insertion, 3) if 'Outline' in insertion['placement'] else object_id + assert not nodes[paragraph]['children'], 'Final paragraph gained an unrecorded child' assert nodes[paragraph]['content'] == [target], 'Inserted paragraph content changed' if 'Outline' in insertion['placement']: position = insertion['placement']['Outline'] diff --git a/tools/test_offline_document_history.py b/tools/test_offline_document_history.py index c7a1e68997c547fd8fb7938b02d71968dc6ad918..76aace55ff7c9990d1d64d4baf94bbbf7cd32fd0 100644 --- a/tools/test_offline_document_history.py +++ b/tools/test_offline_document_history.py @@ -1,4 +1,7 @@ import copy +import json +from pathlib import Path +import subprocess import unittest import uuid from unittest.mock import patch @@ -6,8 +9,11 @@ from unittest.mock import patch from offline_document_history import document_history, identity, verify_model, verify_native -def history(text_edits=False, graph=False): +def history(text_edits=False, graph=False, boundaries=False): + if boundaries: text_edits, graph = True, True logs = {'w0': [{'event': 'ready', 'document_operations': True, 'document_kinds': ['insert', 'format', 'text'] if text_edits else ['insert', 'format']}]} + if boundaries: logs['w0'][0]['document_kinds'] += ['split', 'right_text', 'join'] + kinds = logs['w0'][0]['document_kinds'] events = logs['w0'] if graph: events[0]['document_graph'] = True structure = {} @@ -19,9 +25,13 @@ def history(text_edits=False, graph=False): 'placement': {'Outline': {'x': 144, 'y': 144}} if operation == 0 else {'Paragraph': {'before': None}}} previous = insertion target = identity(insertion, 2) - for step, kind in enumerate(('insert', 'format', 'text') if text_edits else ('insert', 'format')): - timestamp = 10 + operation*30 + step*10 - local_id = operation*3 + step + 2 + end = len(insertion['text'].encode('utf-16-le'))//2 + split = {'guid': list(uuid.UUID(int=100+operation).bytes_le), 'text': target, + 'offset': end-2, 'author': insertion['author'], 'created': 1} + right, right_paragraph = identity(split, 2), identity(split, 1) + for step, kind in enumerate(kinds): + timestamp = 10 + operation*max(3, len(kinds))*10 + step*10 + local_id = operation*max(3, len(kinds)) + step + 2 event = {'event': 'local_document_commit', 'id': local_id, 'operation': operation, 'kind': kind, 'space': 'space', 'object': target, 'text': insertion['text'], 'insertion': insertion if step == 0 else None, 'range': [1, len(insertion['text'].encode('utf-16-le'))//2-2], @@ -30,7 +40,12 @@ def history(text_edits=False, graph=False): if kind == 'text': end = len(insertion['text'].encode('utf-16-le'))//2 event.update(range=[end-3, end], replacement=' e\u0301🐈') + if boundaries: event['document'] = target + if kind == 'split': event.update(split=split, range=[end-2, end-2]) + if kind == 'right_text': event.update(object=right, range=[0, 2], replacement='B🦋') + if kind == 'join': event['joined'] = [target, right] events.append(event) + prior_observed = copy.deepcopy(observed) prior_structure = copy.deepcopy(structure) if graph: events.append({'event': 'read', 'started_us': timestamp-1, 'finished_us': timestamp, @@ -44,26 +59,45 @@ def history(text_edits=False, graph=False): structure[parent]['children'].append(paragraph) structure[paragraph] = {'parent': parent, 'children': [], 'content': [target], 'child_level': 1, 'position': None} - value = insertion['text'].replace(' 🦀', ' e\u0301🐈') if kind == 'text' else insertion['text'] - runs = [] - for index, char in enumerate(value): - selected = index > 0 and (kind == 'text' or (kind == 'format' and index < len(value)-1)) - runs.append({'text': char, 'bold': selected, 'size': 18+operation if selected else 11, - 'color': 0x563412 if selected else 0xff000000}) - observed[target] = {'text': value, 'runs': runs} + if kind == 'split': + boundary = len(insertion['text'])-1 + old = observed[target] + left_runs, right_runs = old['runs'][:boundary], old['runs'][boundary:] + observed[target] = {'text': old['text'][:boundary], 'runs': left_runs} + observed[right] = {'text': old['text'][boundary:], 'runs': right_runs} + structure[parent]['children'].insert(structure[parent]['children'].index(paragraph)+1, right_paragraph) + structure[right_paragraph] = {**copy.deepcopy(structure[paragraph]), 'content': [right]} + elif kind == 'right_text': + runs = [{**observed[right]['runs'][0], 'text': c} for c in 'B🦋'] + observed[right]['runs'][2:] + observed[right] = {'text': 'B🦋🐈', 'runs': runs} + elif kind == 'join': + observed[target]['text'] += observed[right]['text'] + observed[target]['runs'] += observed[right]['runs'] + del observed[right] + structure[parent]['children'].remove(right_paragraph) + del structure[right_paragraph] + else: + value = insertion['text'].replace(' 🦀', ' e\u0301🐈') if kind == 'text' else insertion['text'] + runs = [] + for index, char in enumerate(value): + selected = index > 0 and (kind == 'text' or (kind == 'format' and index < len(value)-1)) + runs.append({'text': char, 'bold': selected, 'size': 18+operation if selected else 11, + 'color': 0x563412 if selected else 0xff000000}) + observed[target] = {'text': value, 'runs': runs} revision = f'revision-{local_id}' events.append({'event': 'remote_attempt', 'revision': revision, 'state': 'Committed', - 'document_changes': {target: copy.deepcopy(observed[target])}, + 'document_changes': {key: copy.deepcopy(observed.get(key)) for key in observed.keys() | prior_observed.keys() + if observed.get(key) != prior_observed.get(key)}, 'documents': copy.deepcopy(observed), 'started_us': timestamp, 'finished_us': timestamp+1}) if graph: events[-1].update(document_graph=copy.deepcopy(structure), document_graph_changes={ - key: copy.deepcopy(value) for key, value in structure.items() if prior_structure.get(key) != value}) + key: copy.deepcopy(structure.get(key)) for key in structure.keys() | prior_structure.keys() if prior_structure.get(key) != structure.get(key)}) events.append({'event': 'document_receipt', 'id': local_id, 'revision': revision, 'at_us': timestamp+2}) if text_edits: events.append({'event': 'read', 'started_us': timestamp+3, 'finished_us': timestamp+4, 'documents': copy.deepcopy(observed), **({'document_graph': copy.deepcopy(structure)} if graph else {})}) events.extend({'event': 'reopened_document_receipt', 'id': row['id'], 'revision': row['revision']} for row in list(events) if row['event'] == 'document_receipt') - read = {'event': 'read', 'started_us': 100, 'finished_us': 101, 'documents': observed, **({'document_graph': structure} if graph else {})} + 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 {})} events.extend([read, {'event': 'done'}]) logs['r0'] = [{'event': 'ready', **({'document_graph': True} if graph else {})}, copy.deepcopy(read), {'event': 'done'}] return logs @@ -106,6 +140,94 @@ class DocumentHistoryTests(unittest.TestCase): with self.subTest(mutation=mutation), self.assertRaises(AssertionError): document_history(logs, 2) + def test_production_cache_workload_matches_the_independent_history_model(self): + result = subprocess.run(['cargo', 'test', '--locked', '-p', 'onestore-offline', '--features', 'smb', + '--example', 'smb_offline_client', + 'tests::document_workload_retains_dependencies_and_receipts_across_reopen', + '--', '--exact', '--nocapture'], + cwd=Path(__file__).resolve().parent.parent, capture_output=True, text=True, timeout=120) + self.assertEqual(result.returncode, 0, result.stdout + result.stderr) + events = [json.loads(line) for line in result.stdout.splitlines() if line.startswith('{')] + documents = document_history({'w0': events}, 2) + self.assertEqual(sum(len(document['states']) for document in documents.values()), 12) + self.assertEqual(sum(row['event'] == 'remote_attempt' and row['state'] == 'Unknown' for row in events), 4) + self.assertEqual(sum(row['event'] == 'remote_confirm' and row['state'] == 'Committed' for row in events), 4) + + def test_boundary_histories_retain_every_intermediate_graph_and_retired_identity(self): + baseline = history(boundaries=True) + documents = document_history(baseline, 2) + for document in documents.values(): + self.assertEqual(list(document['states']), ['insert', 'format', 'text', 'split', 'right_text', 'join']) + self.assertEqual(''.join(c for c, *_ in document['states']['join']['characters']), + document['insertion']['text'].replace('🦀', 'B🦋🐈')) + for mutation in ('missing-suffix', 'half-join', 'resurrected', 'reorder', 'missing-removal', + 'wrong-split', 'wrong-dependent-target', 'wrong-join', 'stale', 'future', 'duplicate-attempt'): + logs = copy.deepcopy(baseline) + split = next(row for row in logs['w0'] if row.get('kind') == 'split') + right, right_paragraph = identity(split['split'], 2), identity(split['split'], 1) + split_read = next(row for row in logs['w0'] if row['event'] == 'read' and row['started_us'] == 43) + joined_read = next(row for row in logs['w0'] if row['event'] == 'read' and row['started_us'] == 63) + if mutation == 'missing-suffix': del split_read['documents'][right] + elif mutation == 'half-join': + joined_read['documents'][right] = copy.deepcopy(split_read['documents'][right]) + elif mutation == 'resurrected': + joined_read['document_graph'][right_paragraph] = copy.deepcopy(split_read['document_graph'][right_paragraph]) + elif mutation == 'reorder': + next(node for node in split_read['document_graph'].values() if node['position'])['children'].reverse() + elif mutation == 'missing-removal': + join_attempt = next(row for row in logs['w0'] if row['event'] == 'remote_attempt' and row['started_us'] == 60) + del join_attempt['document_changes'][right] + elif mutation == 'wrong-split': split['split']['offset'] -= 1 + elif mutation == 'wrong-dependent-target': + next(row for row in logs['w0'] if row.get('kind') == 'right_text')['object'] = split['object'] + elif mutation == 'wrong-join': + next(row for row in logs['w0'] if row.get('kind') == 'join')['joined'].reverse() + elif mutation == 'stale': + joined_read.update(documents=copy.deepcopy(split_read['documents']), document_graph=copy.deepcopy(split_read['document_graph'])) + elif mutation == 'future': split_read.update(started_us=0, finished_us=1) + else: + attempt = copy.deepcopy(next(row for row in logs['w0'] if row['event'] == 'remote_attempt' and row['started_us'] == 40)) + attempt.update(state='Unknown', revision='duplicate') + logs['w0'].insert(1, attempt) + with self.subTest(mutation=mutation), self.assertRaises(AssertionError): + document_history(logs, 2) + + def test_no_reader_phase_can_omit_an_acknowledged_text_or_paragraph(self): + baseline = history(boundaries=True) + for index, row in enumerate(baseline['w0']): + if row['event'] != 'read': continue + for field in ('documents', 'document_graph'): + for target in row[field]: + changed = copy.deepcopy(baseline) + del changed['w0'][index][field][target] + with self.subTest(read=index, field=field, target=target), self.assertRaises(AssertionError): + document_history(changed, 2) + extra = copy.deepcopy(next(row for row in baseline['w0'] if row['event'] == 'remote_attempt')) + extra.update(revision='unreceipted', started_us=1000, finished_us=1001, + document_changes={'unrecorded': {'text': 'extra'}}) + baseline['w0'].insert(-1, extra) + with self.assertRaisesRegex(AssertionError, 'lacks its recorded intent and receipt'): + document_history(baseline, 2) + + def test_unknown_boundary_receipts_require_the_original_revision(self): + for kind in ('split', 'join'): + logs = history(boundaries=True) + events = logs['w0'] + intent = next(row for row in events if row.get('kind') == kind) + receipt = next(row for row in events if row['event'] == 'document_receipt' and row['id'] == intent['id']) + attempt = next(row for row in events if row['event'] == 'remote_attempt' and row['revision'] == receipt['revision']) + attempt['state'] = 'Unknown' + confirmation = {'event': 'remote_confirm', 'state': 'Committed', + 'started_us': attempt['finished_us'], 'finished_us': receipt['at_us'], + 'revisions': {'space': [receipt['revision']]}} + events.insert(events.index(receipt), confirmation) + document_history(logs, 2) + confirmation['revisions']['space'] = ['later-revision'] + receipt['revision'] = 'later-revision' + next(row for row in events if row['event'] == 'reopened_document_receipt' and row['id'] == intent['id'])['revision'] = 'later-revision' + with self.subTest(kind=kind), self.assertRaises(AssertionError): + document_history(logs, 2) + def test_cross_run_text_requires_complete_ordered_states_and_exact_receipts(self): logs = history(True) documents = document_history(logs, 2) diff --git a/tools/verify_offline.py b/tools/verify_offline.py index 11016d4ea26eeabc6ee8b9e57282fd57c27ff9d9..837db2499902aeaa28912cfce3000e9c42b2ccae 100644 --- a/tools/verify_offline.py +++ b/tools/verify_offline.py @@ -32,7 +32,7 @@ def verify(output, max_gap=120): edits = {event['id']: event for event in events if event['event'] in ('local_commit', 'local_document_commit')} receipts = {event['id']: event for event in events if event['event'] in ('remote_receipt', 'document_receipt')} attempts = {event['revision']: event for event in events if event['event'] == 'remote_attempt'} - publication_starts = {id: documents[edit['object']]['states'][edit['kind']]['attempt']['started_us'] if edit['event'] == 'local_document_commit' + publication_starts = {id: documents[edit.get('document', edit['object'])]['states'][edit['kind']]['attempt']['started_us'] if edit['event'] == 'local_document_commit' else attempts[receipts[id]['revision']]['started_us'] for id, edit in edits.items()} 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()]) assert queues[actor] >= 2, 'Writer did not establish a durable local queue before publication' -- 2.54.0