authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-08 09:16:18-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-08 09:16:18-07:00
logc0d4cbf0d0de5d6a8f0db4455c2165bb952802af
treec050308b9d24958cb4d388fab9392f9aa2333158
parent4cd338a8b4d4e4acf3fc1f4e48cf94ac7596e5ea
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

test: observe paragraph structure in concurrent notebook histories

Follow generated outline ancestry so readers retain split suffixes and nested children. Record ordered paragraph graphs with text observations and include removed identities in publication differences. Validate every advertised graph and candidate delta against independent insertion intents; preserve older capture verification. Exercise split/join observation, child transfer and missing or invented structure. Affected example tests, 132 Python tests and workspace Clippy pass. Assisted-by: gpt-6-astra

4 files changed, 309 insertions(+), 33 deletions(-)

crates/onestore-offline/examples/smb_offline_client.rs+69-16
......@@ -20,9 +20,27 @@ use std::{
2020mod support {
2121 pub mod view;
2222}
23use concurrent::document_view;
23use concurrent::{DocumentView, document_view};
2424use support::view::view;
2525
26impl DocumentView {
27 fn changes(&self, before: &Self) -> Self {
28 let difference =
29 |old: &serde_json::Map<String, serde_json::Value>,
30 new: &serde_json::Map<String, serde_json::Value>| {
31 old.keys()
32 .chain(new.keys())
33 .filter(|id| old.get(*id) != new.get(*id))
34 .map(|id| (id.clone(), new.get(id).cloned().unwrap_or_default()))
35 .collect()
36 };
37 Self {
38 texts: difference(&before.texts, &self.texts),
39 graph: difference(&before.graph, &self.graph),
40 }
41 }
42}
43
2644const DOCUMENT_OPERATIONS: [&str; 3] = ["insert", "format", "text"];
2745
2846fn now() -> u128 {
......@@ -40,7 +58,7 @@ enum Pause {
4058
4159struct Traced {
4260 remote: SmbRemote,
43 before: Option<(String, Option<serde_json::Value>)>,
61 before: Option<(String, Option<DocumentView>)>,
4462 pause: Option<Pause>,
4563 documents: bool,
4664}
......@@ -57,7 +75,7 @@ impl Remote for Traced {
5775 };
5876 println!(
5977 "{}",
60 json!({"event":"read", "started_us":started, "finished_us":now(), "text":observed.text, "documents":documents})
78 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)})
6179 );
6280 self.before = Some((observed.text, documents));
6381 Ok(bytes)
......@@ -84,15 +102,7 @@ impl Remote for Traced {
84102 state: CommitState::NotCommitted,
85103 error: io::Error::other("Document publication has no observed source"),
86104 })?;
87 Some(
88 after
89 .as_object()
90 .unwrap()
91 .iter()
92 .filter(|(id, value)| before.get(*id) != Some(*value))
93 .map(|(id, value)| (id.clone(), value.clone()))
94 .collect::<serde_json::Map<_, _>>(),
95 )
105 Some(after.changes(before))
96106 } else {
97107 None
98108 };
......@@ -104,12 +114,14 @@ impl Remote for Traced {
104114 )),
105115 Some(Pause::FormatReply(marker))
106116 if changes.as_ref().is_some_and(|changes| {
107 changes.len() == 1
117 changes.texts.len() == 1
108118 && self
109119 .before
110120 .as_ref()
111121 .and_then(|(_, before)| before.as_ref())
112 .is_some_and(|before| changes.keys().all(|id| before.get(id).is_some()))
122 .is_some_and(|before| {
123 changes.texts.keys().all(|id| before.texts.contains_key(id))
124 })
113125 }) =>
114126 {
115127 Some((marker, marker.with_extension("resume"), "format"))
......@@ -149,7 +161,8 @@ impl Remote for Traced {
149161 json!({"event":"remote_attempt", "started_us":started, "finished_us":finished,
150162 "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,
151163 "state":format!("{:?}", result.as_ref().map_or_else(|error| error.state, |_| CommitState::Committed)),
152 "documents":documents, "document_changes":changes})
164 "documents":documents.as_ref().map(|view| &view.texts), "document_graph":documents.as_ref().map(|view| &view.graph),
165 "document_changes":changes.as_ref().map(|view| &view.texts), "document_graph_changes":changes.as_ref().map(|view| &view.graph)})
153166 );
154167 if result
155168 .as_ref()
......@@ -417,7 +430,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
417430 let cache = Arc::new(Replica::create(&cache_path, &source)?);
418431 println!(
419432 "{}",
420 json!({"event":"ready", "pid":std::process::id(), "actor":args[2], "offline":true, "document_operations":documents,"document_kinds":if documents {DOCUMENT_OPERATIONS.as_slice()} else {&[]}})
433 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 {&[]}})
421434 );
422435 while !Path::new(&args[4]).exists() {
423436 if Instant::now() >= deadline {
......@@ -615,3 +628,43 @@ fn append_review_requires_a_unique_ordered_history_and_an_absent_new_token() {
615628 None
616629 );
617630}
631
632#[cfg(test)]
633mod tests {
634 use super::*;
635
636 #[test]
637 fn publication_differences_retain_removed_text_and_graph_identities() {
638 let before = DocumentView {
639 texts: serde_json::from_value(json!({"left":{"text":"ab"}, "right":{"text":"cd"}, "untouched":{"text":"ef"}})).unwrap(),
640 graph: serde_json::from_value(json!({"outline":{"children":["a","b"]}, "a":{"content":["left"]}, "b":{"content":["right"]}})).unwrap(),
641 };
642 let after = DocumentView {
643 texts: serde_json::from_value(
644 json!({"left":{"text":"abcd"}, "untouched":{"text":"ef"}}),
645 )
646 .unwrap(),
647 graph: serde_json::from_value(
648 json!({"outline":{"children":["a"]}, "a":{"content":["left"]}}),
649 )
650 .unwrap(),
651 };
652 let changes = after.changes(&before);
653 assert_eq!(
654 changes.texts,
655 serde_json::from_value::<serde_json::Map<_, _>>(
656 json!({"left":{"text":"abcd"}, "right":null})
657 )
658 .unwrap()
659 );
660 assert_eq!(
661 changes.graph,
662 serde_json::from_value::<serde_json::Map<_, _>>(
663 json!({"outline":{"children":["a"]}, "b":null})
664 )
665 .unwrap()
666 );
667 assert_eq!(before.changes(&after).texts["right"], before.texts["right"]);
668 assert_eq!(before.changes(&before), DocumentView::default());
669 }
670}
crates/onestore/examples/support/concurrent.rs+161-13
......@@ -11,25 +11,76 @@ use std::{
1111 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
1212};
1313
14pub fn document_view(bytes: &[u8]) -> Result<serde_json::Value, onestore::Error> {
14#[derive(Debug, Default, PartialEq)]
15pub struct DocumentView {
16 pub texts: serde_json::Map<String, serde_json::Value>,
17 pub graph: serde_json::Map<String, serde_json::Value>,
18}
19
20pub fn document_view(bytes: &[u8]) -> Result<DocumentView, onestore::Error> {
1521 let store = Store::parse(bytes)?;
1622 let index = RevisionIndex::parse(&store)?;
1723 index.validate_current()?;
1824 let document = Document::parse(&index)?;
19 let mut texts = serde_json::Map::new();
20 for (sid, _) in document.pages()? {
25 let mut observed = DocumentView::default();
26 for (sid, page) in document.pages()? {
2127 let space = &document.spaces[&sid];
2228 let revision = &space.revisions[&space.contexts[&ExGuid::default()]];
23 for (id, node) in &revision.nodes {
24 if let Kind::RichText { text, .. } = &node.kind
25 && text.starts_with("Document w")
26 {
27 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::<Vec<_>>();
28 texts.insert(id.to_string(), json!({"text":text,"runs":runs}));
29 for outline in &revision.nodes[&page].children {
30 if !matches!(revision.nodes[outline].kind, Kind::Outline { .. }) {
31 continue;
32 }
33 let mut pending = vec![(*outline, page)];
34 let mut parents = std::collections::BTreeMap::new();
35 while let Some((id, parent)) = pending.pop() {
36 if parents.insert(id, parent).is_some() {
37 return Err(onestore::Error {
38 offset: 0,
39 message: "Document observation contains a repeated object",
40 });
41 }
42 let node = &revision.nodes[&id];
43 pending.extend(
44 node.children
45 .iter()
46 .chain(&node.content)
47 .map(|child| (*child, id)),
48 );
49 }
50 // The workload keeps its marker in the left paragraph through boundary edits.
51 if !parents.keys().any(|id| {
52 matches!(&revision.nodes[id].kind,
53 Kind::RichText { text, .. } if text.starts_with("Document w"))
54 }) {
55 continue;
56 }
57 for (id, parent) in parents {
58 let key = id.to_string();
59 if observed.texts.contains_key(&key) || observed.graph.contains_key(&key) {
60 return Err(onestore::Error {
61 offset: 0,
62 message: "Document observation contains a repeated object",
63 });
64 }
65 let node = &revision.nodes[&id];
66 if let Kind::RichText { text, .. } = &node.kind {
67 let runs = revision.text_runs(id)?.into_iter().map(|run| json!({
68 "text": run.text, "bold": run.format.bold.unwrap_or(false),
69 "size": run.format.font_size, "color": run.format.color.unwrap_or(0xff000000)
70 })).collect::<Vec<_>>();
71 observed.texts.insert(key, json!({"text":text,"runs":runs}));
72 } else {
73 observed.graph.insert(key, json!({
74 "parent": parent.to_string(), "children": node.children.iter().map(ToString::to_string).collect::<Vec<_>>(),
75 "content": node.content.iter().map(ToString::to_string).collect::<Vec<_>>(),
76 "child_level": node.child_level,
77 "position": if id == *outline { Some(json!({"x":node.layout.x,"y":node.layout.y})) } else { None }
78 }));
79 }
2980 }
3081 }
3182 }
32 Ok(texts.into())
83 Ok(observed)
3384}
3485
3586pub fn run(
......@@ -70,14 +121,16 @@ pub fn run(
70121 writeln!(output, "{event}")?;
71122 output.flush()
72123 };
73 log(json!({"event": "ready", "pid": std::process::id(), "actor": args[2]}))?;
124 let documents = std::env::var_os("ONESTORE_OFFLINE_DOCUMENTS").is_some();
125 log(
126 json!({"event": "ready", "pid": std::process::id(), "actor": args[2], "document_graph": documents}),
127 )?;
74128 while !Path::new(&args[4]).exists() {
75129 if Instant::now() > deadline {
76130 return Err("Start barrier timed out.".into());
77131 }
78132 thread::sleep(Duration::from_millis(5));
79133 }
80 let documents = std::env::var_os("ONESTORE_OFFLINE_DOCUMENTS").is_some();
81134 let maintenance = std::env::var_os("ONESTORE_MAINTENANCE_DIR").map(std::path::PathBuf::from);
82135 let mut completed = 0;
83136 let mut attempts = 0;
......@@ -187,10 +240,15 @@ pub fn run(
187240 })
188241 .into());
189242 };
243 let observed = if documents {
244 Some(document_view(&source).map_err(preserve)?)
245 } else {
246 None
247 };
190248 let read_finished = SystemTime::now().duration_since(UNIX_EPOCH)?.as_micros();
191249 log(
192250 json!({"event": "read", "attempt": attempts, "started_us": started, "finished_us": read_finished,
193 "transaction": store.header.transaction_count, "text": text, "documents":if documents {Some(document_view(&source).map_err(preserve)?)}else{None}}),
251 "transaction": store.header.transaction_count, "text": text, "documents":observed.as_ref().map(|view| &view.texts), "document_graph":observed.as_ref().map(|view| &view.graph)}),
194252 )?;
195253 if args[0] == "read" {
196254 completed += 1;
......@@ -256,3 +314,93 @@ pub fn run(
256314 log(json!({"event": "done", "completed": completed, "attempts": attempts}))?;
257315 Ok(())
258316}
317
318#[cfg(test)]
319mod tests {
320 use super::*;
321 use onestore::{Insertion, ParagraphJoin, ParagraphSplit, PreparedEdit, TextAttribute};
322
323 #[test]
324 fn observation_follows_split_suffixes_and_moved_children_through_active_ancestry() {
325 let source = onestore::create_section("observation.one", "Original", "Author").unwrap();
326 assert_eq!(document_view(&source).unwrap(), DocumentView::default());
327 let store = Store::parse(&source).unwrap();
328 let index = RevisionIndex::parse(&store).unwrap();
329 let document = Document::parse(&index).unwrap();
330 let (space, page) = document.pages().unwrap()[0];
331 let insertion = Insertion::outline(page, 144.0, 216.0, "Document w0:0 🦀", "Author")
332 .unwrap()
333 .with_formatting(0..16, &[TextAttribute::Bold(true)])
334 .unwrap();
335 let inserted = PreparedEdit::insert(&source, space, &insertion).unwrap();
336 let before = document_view(inserted.as_bytes()).unwrap();
337 assert_eq!(before.texts.len(), 1);
338 assert_eq!(before.graph.len(), 2);
339 let outline = insertion.object().to_string();
340 let paragraph = before.graph[&outline]["children"][0]
341 .as_str()
342 .unwrap()
343 .to_owned();
344 let child =
345 Insertion::paragraph(paragraph.parse().unwrap(), None, "Unmarked child", "Author")
346 .unwrap();
347 let with_child = PreparedEdit::insert(inserted.as_bytes(), space, &child).unwrap();
348 let before = document_view(with_child.as_bytes()).unwrap();
349 assert_eq!(before.texts.len(), 2);
350 assert_eq!(before.graph.len(), 3);
351 assert_eq!(
352 before.graph[&outline],
353 json!({"parent":page.to_string(), "children":[paragraph], "content":[], "child_level":1, "position":{"x":144.0,"y":216.0}})
354 );
355 for offset in [14, 16] {
356 let split = ParagraphSplit::new(insertion.text_object(), offset, "Author").unwrap();
357 let split_edit = PreparedEdit::split(with_child.as_bytes(), space, &split).unwrap();
358 let observed = document_view(split_edit.as_bytes()).unwrap();
359 assert_eq!(observed.texts.len(), 3);
360 assert_eq!(observed.graph.len(), 4);
361 assert_eq!(
362 observed.texts[&split.text_object().to_string()]["text"],
363 if offset == 14 { "🦀" } else { "" }
364 );
365 assert_eq!(
366 observed.graph[&outline]["children"],
367 json!([paragraph, split.object().to_string()])
368 );
369 assert_eq!(observed.graph[&paragraph]["children"], json!([]));
370 assert_eq!(
371 observed.graph[&split.object().to_string()]["children"],
372 json!([child.object().to_string()])
373 );
374 assert_eq!(
375 observed.graph[&child.object().to_string()]["parent"],
376 split.object().to_string()
377 );
378 let join =
379 ParagraphJoin::new(insertion.text_object(), split.text_object(), "Author").unwrap();
380 let joined = PreparedEdit::join(split_edit.as_bytes(), space, &join).unwrap();
381 let joined = document_view(joined.as_bytes()).unwrap();
382 assert_eq!(joined.graph, before.graph);
383 let characters = |view: &DocumentView| {
384 view.texts
385 .iter()
386 .map(|(id, text)| {
387 let runs = text["runs"]
388 .as_array()
389 .unwrap()
390 .iter()
391 .flat_map(|run| {
392 run["text"]
393 .as_str()
394 .unwrap()
395 .chars()
396 .map(|c| json!([c, run["bold"], run["size"], run["color"]]))
397 })
398 .collect::<Vec<_>>();
399 (id.clone(), (text["text"].clone(), runs))
400 })
401 .collect::<std::collections::BTreeMap<_, _>>()
402 };
403 assert_eq!(characters(&joined), characters(&before));
404 }
405 }
406}
tools/offline_document_history.py+34
......@@ -114,6 +114,40 @@ def document_history(logs, operations):
114114 states['text'] = {'characters': final, 'attempt': attempt}
115115 documents[target] = {'insertion': insertion, 'space': inserted['space'], 'states': states}
116116 assert documents, 'No document operations were recorded'
117 if any(events[0].get('document_graph') for events in logs.values()):
118 for events in logs.values():
119 assert events[0].get('document_graph') is True, 'Client omitted structural observations'
120 before = None
121 for row in events:
122 if row['event'] not in ('read', 'document_read', 'remote_attempt'): continue
123 observed = row.get('documents')
124 assert isinstance(observed, dict), 'Snapshot omitted document text'
125 graph = row.get('document_graph')
126 assert isinstance(graph, dict), 'Snapshot omitted document graph'
127 expected = {}
128 for target, document in documents.items():
129 if target not in observed: continue
130 insertion = document['insertion']
131 outline = insertion['parent']
132 paragraph = identity(insertion, 1)
133 if 'Outline' in insertion['placement']:
134 outline = paragraph
135 paragraph = identity(insertion, 3)
136 expected[outline] = {'parent': insertion['parent'], 'children': [], 'content': [],
137 'child_level': 1, 'position': insertion['placement']['Outline']}
138 assert outline in expected, 'Snapshot omitted the inserted outline'
139 expected[outline]['children'].append(paragraph)
140 expected[paragraph] = {'parent': outline, 'children': [], 'content': [target],
141 'child_level': 1, 'position': None}
142 assert graph == expected, 'Snapshot contains partial, reordered or invented paragraph structure'
143 if row['event'] == 'remote_attempt':
144 assert before is not None, 'Publication omitted its observed source'
145 for image, delta in [('documents', 'document_changes'), ('document_graph', 'document_graph_changes')]:
146 expected_delta = {key: row[image].get(key) for key in before[image].keys() | row[image].keys()
147 if before[image].get(key) != row[image].get(key)}
148 assert row.get(delta) == expected_delta, 'Publication diff omitted or invented a changed object'
149 else:
150 before = row
117151 for actor, events in logs.items():
118152 previous = {}
119153 reads = [row for row in events if row['event'] in ('read', 'document_read') and row.get('documents') is not None]
tools/test_offline_document_history.py+45-4
......@@ -6,9 +6,11 @@ from unittest.mock import patch
66from offline_document_history import document_history, identity, verify_model, verify_native
77
88
9def history(text_edits=False):
9def history(text_edits=False, graph=False):
1010 logs = {'w0': [{'event': 'ready', 'document_operations': True, 'document_kinds': ['insert', 'format', 'text'] if text_edits else ['insert', 'format']}]}
1111 events = logs['w0']
12 if graph: events[0]['document_graph'] = True
13 structure = {}
1214 observed = {}
1315 previous = None
1416 for operation in range(2):
......@@ -29,6 +31,19 @@ def history(text_edits=False):
2931 end = len(insertion['text'].encode('utf-16-le'))//2
3032 event.update(range=[end-3, end], replacement=' e\u0301🐈')
3133 events.append(event)
34 prior_structure = copy.deepcopy(structure)
35 if graph:
36 events.append({'event': 'read', 'started_us': timestamp-1, 'finished_us': timestamp,
37 'documents': copy.deepcopy(observed), 'document_graph': prior_structure})
38 if step == 0:
39 parent, paragraph = insertion['parent'], identity(insertion, 1)
40 if operation == 0:
41 parent, paragraph = paragraph, identity(insertion, 3)
42 structure[parent] = {'parent': insertion['parent'], 'children': [], 'content': [],
43 'child_level': 1, 'position': {'x': 144, 'y': 144}}
44 structure[parent]['children'].append(paragraph)
45 structure[paragraph] = {'parent': parent, 'children': [], 'content': [target],
46 'child_level': 1, 'position': None}
3247 value = insertion['text'].replace(' 🦀', ' e\u0301🐈') if kind == 'text' else insertion['text']
3348 runs = []
3449 for index, char in enumerate(value):
......@@ -40,14 +55,17 @@ def history(text_edits=False):
4055 events.append({'event': 'remote_attempt', 'revision': revision, 'state': 'Committed',
4156 'document_changes': {target: copy.deepcopy(observed[target])},
4257 'documents': copy.deepcopy(observed), 'started_us': timestamp, 'finished_us': timestamp+1})
58 if graph:
59 events[-1].update(document_graph=copy.deepcopy(structure), document_graph_changes={
60 key: copy.deepcopy(value) for key, value in structure.items() if prior_structure.get(key) != value})
4361 events.append({'event': 'document_receipt', 'id': local_id, 'revision': revision, 'at_us': timestamp+2})
4462 if text_edits:
45 events.append({'event': 'read', 'started_us': timestamp+3, 'finished_us': timestamp+4, 'documents': copy.deepcopy(observed)})
63 events.append({'event': 'read', 'started_us': timestamp+3, 'finished_us': timestamp+4, 'documents': copy.deepcopy(observed), **({'document_graph': copy.deepcopy(structure)} if graph else {})})
4664 events.extend({'event': 'reopened_document_receipt', 'id': row['id'], 'revision': row['revision']}
4765 for row in list(events) if row['event'] == 'document_receipt')
48 read = {'event': 'read', 'started_us': 100, 'finished_us': 101, 'documents': observed}
66 read = {'event': 'read', 'started_us': 100, 'finished_us': 101, 'documents': observed, **({'document_graph': structure} if graph else {})}
4967 events.extend([read, {'event': 'done'}])
50 logs['r0'] = [{'event': 'ready'}, copy.deepcopy(read), {'event': 'done'}]
68 logs['r0'] = [{'event': 'ready', **({'document_graph': True} if graph else {})}, copy.deepcopy(read), {'event': 'done'}]
5169 return logs
5270
5371
......@@ -65,6 +83,29 @@ class DocumentHistoryTests(unittest.TestCase):
6583 paragraphs[0][1][1]['font_size'] = 19
6684 with self.assertRaisesRegex(AssertionError, 'font size'): verify_native(paragraphs, expected)
6785
86 def test_structural_snapshots_and_publication_differences_are_complete(self):
87 baseline = history(text_edits=True, graph=True)
88 self.assertEqual(len(document_history(baseline, 2)), 2)
89 for mutation in ('missing-graph', 'missing-client', 'reorder', 'duplicate', 'reparent',
90 'missing-content', 'invented', 'level', 'position', 'delta'):
91 logs = copy.deepcopy(baseline)
92 row = logs['r0'][1]
93 outline, first, second = list(row['document_graph'])
94 graph = row['document_graph']
95 if mutation == 'missing-graph': del row['document_graph']
96 elif mutation == 'missing-client': del logs['r0'][0]['document_graph']
97 elif mutation == 'reorder': graph[outline]['children'].reverse()
98 elif mutation == 'duplicate': graph[outline]['children'].append(first)
99 elif mutation == 'reparent': graph[second]['parent'] = first
100 elif mutation == 'missing-content': graph[first]['content'] = []
101 elif mutation == 'invented': graph['unrecorded'] = copy.deepcopy(graph[first])
102 elif mutation == 'level': graph[first]['child_level'] = 2
103 elif mutation == 'position': graph[outline]['position']['x'] = 145
104 else:
105 next(row for row in logs['w0'] if row['event'] == 'remote_attempt')['document_graph_changes'] = {}
106 with self.subTest(mutation=mutation), self.assertRaises(AssertionError):
107 document_history(logs, 2)
108
68109 def test_cross_run_text_requires_complete_ordered_states_and_exact_receipts(self):
69110 logs = history(True)
70111 documents = document_history(logs, 2)