diff --git a/crates/onestore-offline/examples/smb_offline_client.rs b/crates/onestore-offline/examples/smb_offline_client.rs index 3b80ea77392b8b240e16e4b7bfd0fc920062349f..54b85c47b1da226bf799d7e2fa37415d5e3371c1 100644 --- a/crates/onestore-offline/examples/smb_offline_client.rs +++ b/crates/onestore-offline/examples/smb_offline_client.rs @@ -20,9 +20,27 @@ use std::{ mod support { pub mod view; } -use concurrent::document_view; +use concurrent::{DocumentView, document_view}; use support::view::view; +impl DocumentView { + fn changes(&self, before: &Self) -> Self { + let difference = + |old: &serde_json::Map, + new: &serde_json::Map| { + old.keys() + .chain(new.keys()) + .filter(|id| old.get(*id) != new.get(*id)) + .map(|id| (id.clone(), new.get(id).cloned().unwrap_or_default())) + .collect() + }; + Self { + texts: difference(&before.texts, &self.texts), + graph: difference(&before.graph, &self.graph), + } + } +} + const DOCUMENT_OPERATIONS: [&str; 3] = ["insert", "format", "text"]; fn now() -> u128 { @@ -40,7 +58,7 @@ enum Pause { struct Traced { remote: SmbRemote, - before: Option<(String, Option)>, + before: Option<(String, Option)>, pause: Option, documents: bool, } @@ -57,7 +75,7 @@ impl Remote for Traced { }; println!( "{}", - json!({"event":"read", "started_us":started, "finished_us":now(), "text":observed.text, "documents":documents}) + json!({"event":"read", "started_us":started, "finished_us":now(), "text":observed.text, "documents":documents.as_ref().map(|view| &view.texts), "document_graph":documents.as_ref().map(|view| &view.graph)}) ); self.before = Some((observed.text, documents)); Ok(bytes) @@ -84,15 +102,7 @@ impl Remote for Traced { state: CommitState::NotCommitted, error: io::Error::other("Document publication has no observed source"), })?; - Some( - after - .as_object() - .unwrap() - .iter() - .filter(|(id, value)| before.get(*id) != Some(*value)) - .map(|(id, value)| (id.clone(), value.clone())) - .collect::>(), - ) + Some(after.changes(before)) } else { None }; @@ -104,12 +114,14 @@ impl Remote for Traced { )), Some(Pause::FormatReply(marker)) if changes.as_ref().is_some_and(|changes| { - changes.len() == 1 + changes.texts.len() == 1 && self .before .as_ref() .and_then(|(_, before)| before.as_ref()) - .is_some_and(|before| changes.keys().all(|id| before.get(id).is_some())) + .is_some_and(|before| { + changes.texts.keys().all(|id| before.texts.contains_key(id)) + }) }) => { Some((marker, marker.with_extension("resume"), "format")) @@ -149,7 +161,8 @@ impl Remote for Traced { json!({"event":"remote_attempt", "started_us":started, "finished_us":finished, "revision":after.revision.to_string(), "space":after.space.to_string(), "object":after.object.to_string(), "before":self.before.as_ref().map(|(text,_)|text), "after":after.text, "state":format!("{:?}", result.as_ref().map_or_else(|error| error.state, |_| CommitState::Committed)), - "documents":documents, "document_changes":changes}) + "documents":documents.as_ref().map(|view| &view.texts), "document_graph":documents.as_ref().map(|view| &view.graph), + "document_changes":changes.as_ref().map(|view| &view.texts), "document_graph_changes":changes.as_ref().map(|view| &view.graph)}) ); if result .as_ref() @@ -417,7 +430,7 @@ fn main() -> Result<(), Box> { let cache = Arc::new(Replica::create(&cache_path, &source)?); println!( "{}", - json!({"event":"ready", "pid":std::process::id(), "actor":args[2], "offline":true, "document_operations":documents,"document_kinds":if documents {DOCUMENT_OPERATIONS.as_slice()} else {&[]}}) + json!({"event":"ready", "pid":std::process::id(), "actor":args[2], "offline":true, "document_operations":documents,"document_graph":documents,"document_kinds":if documents {DOCUMENT_OPERATIONS.as_slice()} else {&[]}}) ); while !Path::new(&args[4]).exists() { if Instant::now() >= deadline { @@ -615,3 +628,43 @@ fn append_review_requires_a_unique_ordered_history_and_an_absent_new_token() { None ); } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn publication_differences_retain_removed_text_and_graph_identities() { + let before = DocumentView { + texts: serde_json::from_value(json!({"left":{"text":"ab"}, "right":{"text":"cd"}, "untouched":{"text":"ef"}})).unwrap(), + graph: serde_json::from_value(json!({"outline":{"children":["a","b"]}, "a":{"content":["left"]}, "b":{"content":["right"]}})).unwrap(), + }; + let after = DocumentView { + texts: serde_json::from_value( + json!({"left":{"text":"abcd"}, "untouched":{"text":"ef"}}), + ) + .unwrap(), + graph: serde_json::from_value( + json!({"outline":{"children":["a"]}, "a":{"content":["left"]}}), + ) + .unwrap(), + }; + let changes = after.changes(&before); + assert_eq!( + changes.texts, + serde_json::from_value::>( + json!({"left":{"text":"abcd"}, "right":null}) + ) + .unwrap() + ); + assert_eq!( + changes.graph, + serde_json::from_value::>( + json!({"outline":{"children":["a"]}, "b":null}) + ) + .unwrap() + ); + assert_eq!(before.changes(&after).texts["right"], before.texts["right"]); + assert_eq!(before.changes(&before), DocumentView::default()); + } +} diff --git a/crates/onestore/examples/support/concurrent.rs b/crates/onestore/examples/support/concurrent.rs index bb28c70c106927eb9b71ce0a7be2056a641be230..d39a781b522c604809473846aff03471fb3a73e2 100644 --- a/crates/onestore/examples/support/concurrent.rs +++ b/crates/onestore/examples/support/concurrent.rs @@ -11,25 +11,76 @@ use std::{ time::{Duration, Instant, SystemTime, UNIX_EPOCH}, }; -pub fn document_view(bytes: &[u8]) -> Result { +#[derive(Debug, Default, PartialEq)] +pub struct DocumentView { + pub texts: serde_json::Map, + pub graph: serde_json::Map, +} + +pub fn document_view(bytes: &[u8]) -> Result { let store = Store::parse(bytes)?; let index = RevisionIndex::parse(&store)?; index.validate_current()?; let document = Document::parse(&index)?; - let mut texts = serde_json::Map::new(); - for (sid, _) in document.pages()? { + let mut observed = DocumentView::default(); + for (sid, page) in document.pages()? { let space = &document.spaces[&sid]; let revision = &space.revisions[&space.contexts[&ExGuid::default()]]; - for (id, node) in &revision.nodes { - if let Kind::RichText { text, .. } = &node.kind - && text.starts_with("Document w") - { - let runs=revision.text_runs(*id)?.into_iter().map(|run|json!({"text":run.text,"bold":run.format.bold.unwrap_or(false),"size":run.format.font_size,"color":run.format.color.unwrap_or(0xff000000)})).collect::>(); - texts.insert(id.to_string(), json!({"text":text,"runs":runs})); + for outline in &revision.nodes[&page].children { + if !matches!(revision.nodes[outline].kind, Kind::Outline { .. }) { + continue; + } + let mut pending = vec![(*outline, page)]; + let mut parents = std::collections::BTreeMap::new(); + while let Some((id, parent)) = pending.pop() { + if parents.insert(id, parent).is_some() { + return Err(onestore::Error { + offset: 0, + message: "Document observation contains a repeated object", + }); + } + let node = &revision.nodes[&id]; + pending.extend( + node.children + .iter() + .chain(&node.content) + .map(|child| (*child, id)), + ); + } + // The workload keeps its marker in the left paragraph through boundary edits. + if !parents.keys().any(|id| { + matches!(&revision.nodes[id].kind, + Kind::RichText { text, .. } if text.starts_with("Document w")) + }) { + continue; + } + for (id, parent) in parents { + let key = id.to_string(); + if observed.texts.contains_key(&key) || observed.graph.contains_key(&key) { + return Err(onestore::Error { + offset: 0, + message: "Document observation contains a repeated object", + }); + } + let node = &revision.nodes[&id]; + if let Kind::RichText { text, .. } = &node.kind { + let runs = revision.text_runs(id)?.into_iter().map(|run| json!({ + "text": run.text, "bold": run.format.bold.unwrap_or(false), + "size": run.format.font_size, "color": run.format.color.unwrap_or(0xff000000) + })).collect::>(); + observed.texts.insert(key, json!({"text":text,"runs":runs})); + } else { + observed.graph.insert(key, json!({ + "parent": parent.to_string(), "children": node.children.iter().map(ToString::to_string).collect::>(), + "content": node.content.iter().map(ToString::to_string).collect::>(), + "child_level": node.child_level, + "position": if id == *outline { Some(json!({"x":node.layout.x,"y":node.layout.y})) } else { None } + })); + } } } } - Ok(texts.into()) + Ok(observed) } pub fn run( @@ -70,14 +121,16 @@ pub fn run( writeln!(output, "{event}")?; output.flush() }; - log(json!({"event": "ready", "pid": std::process::id(), "actor": args[2]}))?; + let documents = std::env::var_os("ONESTORE_OFFLINE_DOCUMENTS").is_some(); + log( + json!({"event": "ready", "pid": std::process::id(), "actor": args[2], "document_graph": documents}), + )?; while !Path::new(&args[4]).exists() { if Instant::now() > deadline { return Err("Start barrier timed out.".into()); } thread::sleep(Duration::from_millis(5)); } - let documents = std::env::var_os("ONESTORE_OFFLINE_DOCUMENTS").is_some(); let maintenance = std::env::var_os("ONESTORE_MAINTENANCE_DIR").map(std::path::PathBuf::from); let mut completed = 0; let mut attempts = 0; @@ -187,10 +240,15 @@ pub fn run( }) .into()); }; + let observed = if documents { + Some(document_view(&source).map_err(preserve)?) + } else { + None + }; let read_finished = SystemTime::now().duration_since(UNIX_EPOCH)?.as_micros(); log( json!({"event": "read", "attempt": attempts, "started_us": started, "finished_us": read_finished, - "transaction": store.header.transaction_count, "text": text, "documents":if documents {Some(document_view(&source).map_err(preserve)?)}else{None}}), + "transaction": store.header.transaction_count, "text": text, "documents":observed.as_ref().map(|view| &view.texts), "document_graph":observed.as_ref().map(|view| &view.graph)}), )?; if args[0] == "read" { completed += 1; @@ -256,3 +314,93 @@ pub fn run( log(json!({"event": "done", "completed": completed, "attempts": attempts}))?; Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + use onestore::{Insertion, ParagraphJoin, ParagraphSplit, PreparedEdit, TextAttribute}; + + #[test] + fn observation_follows_split_suffixes_and_moved_children_through_active_ancestry() { + let source = onestore::create_section("observation.one", "Original", "Author").unwrap(); + assert_eq!(document_view(&source).unwrap(), DocumentView::default()); + let store = Store::parse(&source).unwrap(); + let index = RevisionIndex::parse(&store).unwrap(); + let document = Document::parse(&index).unwrap(); + let (space, page) = document.pages().unwrap()[0]; + let insertion = Insertion::outline(page, 144.0, 216.0, "Document w0:0 🦀", "Author") + .unwrap() + .with_formatting(0..16, &[TextAttribute::Bold(true)]) + .unwrap(); + let inserted = PreparedEdit::insert(&source, space, &insertion).unwrap(); + let before = document_view(inserted.as_bytes()).unwrap(); + assert_eq!(before.texts.len(), 1); + assert_eq!(before.graph.len(), 2); + let outline = insertion.object().to_string(); + let paragraph = before.graph[&outline]["children"][0] + .as_str() + .unwrap() + .to_owned(); + let child = + Insertion::paragraph(paragraph.parse().unwrap(), None, "Unmarked child", "Author") + .unwrap(); + let with_child = PreparedEdit::insert(inserted.as_bytes(), space, &child).unwrap(); + let before = document_view(with_child.as_bytes()).unwrap(); + assert_eq!(before.texts.len(), 2); + assert_eq!(before.graph.len(), 3); + assert_eq!( + before.graph[&outline], + json!({"parent":page.to_string(), "children":[paragraph], "content":[], "child_level":1, "position":{"x":144.0,"y":216.0}}) + ); + for offset in [14, 16] { + let split = ParagraphSplit::new(insertion.text_object(), offset, "Author").unwrap(); + let split_edit = PreparedEdit::split(with_child.as_bytes(), space, &split).unwrap(); + let observed = document_view(split_edit.as_bytes()).unwrap(); + assert_eq!(observed.texts.len(), 3); + assert_eq!(observed.graph.len(), 4); + assert_eq!( + observed.texts[&split.text_object().to_string()]["text"], + if offset == 14 { "🦀" } else { "" } + ); + assert_eq!( + observed.graph[&outline]["children"], + json!([paragraph, split.object().to_string()]) + ); + assert_eq!(observed.graph[¶graph]["children"], json!([])); + assert_eq!( + observed.graph[&split.object().to_string()]["children"], + json!([child.object().to_string()]) + ); + assert_eq!( + observed.graph[&child.object().to_string()]["parent"], + split.object().to_string() + ); + let join = + ParagraphJoin::new(insertion.text_object(), split.text_object(), "Author").unwrap(); + let joined = PreparedEdit::join(split_edit.as_bytes(), space, &join).unwrap(); + let joined = document_view(joined.as_bytes()).unwrap(); + assert_eq!(joined.graph, before.graph); + let characters = |view: &DocumentView| { + view.texts + .iter() + .map(|(id, text)| { + let runs = text["runs"] + .as_array() + .unwrap() + .iter() + .flat_map(|run| { + run["text"] + .as_str() + .unwrap() + .chars() + .map(|c| json!([c, run["bold"], run["size"], run["color"]])) + }) + .collect::>(); + (id.clone(), (text["text"].clone(), runs)) + }) + .collect::>() + }; + assert_eq!(characters(&joined), characters(&before)); + } + } +} diff --git a/tools/offline_document_history.py b/tools/offline_document_history.py index 3edb363b0a141f66068e86d5b6ad21fd4e22c917..e4d7f0af093ea292fb72d6237f27e6837f54aeb8 100644 --- a/tools/offline_document_history.py +++ b/tools/offline_document_history.py @@ -114,6 +114,40 @@ def document_history(logs, operations): states['text'] = {'characters': final, 'attempt': attempt} 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(): + 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 + 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': + 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] diff --git a/tools/test_offline_document_history.py b/tools/test_offline_document_history.py index 1d51a253ee6ecabb4b75550c963bae971a998871..c7a1e68997c547fd8fb7938b02d71968dc6ad918 100644 --- a/tools/test_offline_document_history.py +++ b/tools/test_offline_document_history.py @@ -6,9 +6,11 @@ from unittest.mock import patch from offline_document_history import document_history, identity, verify_model, verify_native -def history(text_edits=False): +def history(text_edits=False, graph=False): logs = {'w0': [{'event': 'ready', 'document_operations': True, 'document_kinds': ['insert', 'format', 'text'] if text_edits else ['insert', 'format']}]} events = logs['w0'] + if graph: events[0]['document_graph'] = True + structure = {} observed = {} previous = None for operation in range(2): @@ -29,6 +31,19 @@ def history(text_edits=False): end = len(insertion['text'].encode('utf-16-le'))//2 event.update(range=[end-3, end], replacement=' e\u0301🐈') events.append(event) + prior_structure = copy.deepcopy(structure) + if graph: + events.append({'event': 'read', 'started_us': timestamp-1, 'finished_us': timestamp, + 'documents': copy.deepcopy(observed), 'document_graph': prior_structure}) + if step == 0: + parent, paragraph = insertion['parent'], identity(insertion, 1) + if operation == 0: + parent, paragraph = paragraph, identity(insertion, 3) + structure[parent] = {'parent': insertion['parent'], 'children': [], 'content': [], + 'child_level': 1, 'position': {'x': 144, 'y': 144}} + 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): @@ -40,14 +55,17 @@ def history(text_edits=False): events.append({'event': 'remote_attempt', 'revision': revision, 'state': 'Committed', 'document_changes': {target: copy.deepcopy(observed[target])}, '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}) 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)}) + 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} + read = {'event': 'read', 'started_us': 100, 'finished_us': 101, 'documents': observed, **({'document_graph': structure} if graph else {})} events.extend([read, {'event': 'done'}]) - logs['r0'] = [{'event': 'ready'}, copy.deepcopy(read), {'event': 'done'}] + logs['r0'] = [{'event': 'ready', **({'document_graph': True} if graph else {})}, copy.deepcopy(read), {'event': 'done'}] return logs @@ -65,6 +83,29 @@ class DocumentHistoryTests(unittest.TestCase): paragraphs[0][1][1]['font_size'] = 19 with self.assertRaisesRegex(AssertionError, 'font size'): verify_native(paragraphs, expected) + def test_structural_snapshots_and_publication_differences_are_complete(self): + baseline = history(text_edits=True, graph=True) + self.assertEqual(len(document_history(baseline, 2)), 2) + for mutation in ('missing-graph', 'missing-client', 'reorder', 'duplicate', 'reparent', + 'missing-content', 'invented', 'level', 'position', 'delta'): + logs = copy.deepcopy(baseline) + row = logs['r0'][1] + outline, first, second = list(row['document_graph']) + graph = row['document_graph'] + if mutation == 'missing-graph': del row['document_graph'] + elif mutation == 'missing-client': del logs['r0'][0]['document_graph'] + elif mutation == 'reorder': graph[outline]['children'].reverse() + elif mutation == 'duplicate': graph[outline]['children'].append(first) + elif mutation == 'reparent': graph[second]['parent'] = first + elif mutation == 'missing-content': graph[first]['content'] = [] + elif mutation == 'invented': graph['unrecorded'] = copy.deepcopy(graph[first]) + elif mutation == 'level': graph[first]['child_level'] = 2 + elif mutation == 'position': graph[outline]['position']['x'] = 145 + else: + next(row for row in logs['w0'] if row['event'] == 'remote_attempt')['document_graph_changes'] = {} + with self.subTest(mutation=mutation), 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)