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'