1use onestore::{
2 CommitState, ExGuid, RevisionIndex, Store,
3 document::{Document, Kind},
4};
5use serde_json::json;
6use std::{
7 fs,
8 io::{self, Write},
9 path::Path,
10 thread,
11 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
12};
13
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> {
21 let store = Store::parse(bytes)?;
22 let index = RevisionIndex::parse(&store)?;
23 index.validate_current()?;
24 let document = Document::parse(&index)?;
25 let mut observed = DocumentView::default();
26 for (sid, page) in document.pages()? {
27 let space = &document.spaces[&sid];
28 let revision = &space.revisions[&space.contexts[&ExGuid::default()]];
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 }
80 }
81 }
82 }
83 Ok(observed)
84}
85
86pub fn run(
87 args: &[String],
88 mut read: impl FnMut(&str) -> io::Result<Vec<u8>>,
89 mut commit: impl FnMut(
90 &str,
91 &[u8],
92 ExGuid,
93 ExGuid,
94 std::ops::Range<u32>,
95 &str,
96 ) -> Result<(), onestore::CommitError>,
97) -> Result<(), Box<dyn std::error::Error>> {
98 if args.len() != 7 || !["read", "write", "edit"].contains(&args[0].as_str()) {
99 return Err(
100 "Expected read|write|edit FILE ACTOR OPERATIONS START_FILE STOP_FILE SEED.".into(),
101 );
102 }
103 let mut random: u64 = args[6].parse()?;
104 let operations: usize = args[3].parse()?;
105 if operations == 0 {
106 return Err("Choose at least one operation.".into());
107 }
108 let timeout = match std::env::var("ONESTORE_CLIENT_TIMEOUT_MS") {
109 Ok(value) => value.parse::<u64>()?,
110 Err(std::env::VarError::NotPresent) => 600_000,
111 Err(error) => return Err(error.into()),
112 };
113 if timeout == 0 {
114 return Err("Choose a positive client timeout.".into());
115 }
116 let deadline = Instant::now()
117 .checked_add(Duration::from_millis(timeout))
118 .ok_or("Client timeout exceeds the clock range.")?;
119 let mut output = io::stdout().lock();
120 let mut log = |event: serde_json::Value| -> io::Result<()> {
121 writeln!(output, "{event}")?;
122 output.flush()
123 };
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 )?;
128 while !Path::new(&args[4]).exists() {
129 if Instant::now() > deadline {
130 return Err("Start barrier timed out.".into());
131 }
132 thread::sleep(Duration::from_millis(5));
133 }
134 let maintenance = std::env::var_os("ONESTORE_MAINTENANCE_DIR").map(std::path::PathBuf::from);
135 let mut completed = 0;
136 let mut attempts = 0;
137 while completed < operations || (args[0] == "read" && !Path::new(&args[5]).exists()) {
138 if Instant::now() > deadline {
139 return Err("Concurrent client timed out.".into());
140 }
141 if let Some(control) = &maintenance
142 && !control.join("resume").exists()
143 && ((args[0] != "read" && completed == operations / 2)
144 || (args[0] == "read" && control.join("pause").exists()))
145 {
146 fs::write(control.join(format!("paused-{}", args[2])), b"paused")?;
147 log(json!({"event": "paused", "completed": completed}))?;
148 while !control.join("resume").exists() {
149 if Instant::now() > deadline {
150 return Err("Maintenance pause timed out.".into());
151 }
152 thread::sleep(Duration::from_millis(10));
153 }
154 log(json!({"event": "resumed", "completed": completed}))?;
155 }
156 attempts += 1;
157 random = random
158 .wrapping_mul(6364136223846793005)
159 .wrapping_add(1442695040888963407);
160 thread::sleep(Duration::from_millis((random >> 32) % 7));
161 let started = SystemTime::now().duration_since(UNIX_EPOCH)?.as_micros();
162 let source = match read(&args[1]) {
163 Ok(source) => source,
164 Err(error)
165 if [
166 io::ErrorKind::WouldBlock,
167 io::ErrorKind::ResourceBusy,
168 io::ErrorKind::PermissionDenied,
169 io::ErrorKind::NotFound,
170 ]
171 .contains(&error.kind()) =>
172 {
173 log(
174 json!({"event": "read_busy", "attempt": attempts, "kind": format!("{:?}", error.kind())}),
175 )?;
176 thread::sleep(Duration::from_millis(100));
177 continue;
178 }
179 Err(error) => {
180 log(
181 json!({"event": "read_error", "attempt": attempts, "kind": format!("{:?}", error.kind())}),
182 )?;
183 return Err(error.into());
184 }
185 };
186 let preserve = |error: onestore::Error| {
187 let path = Path::new(&args[4])
188 .parent()
189 .unwrap()
190 .join(format!("invalid-{}-{attempts}.one", std::process::id()));
191 if let Err(failure) = fs::write(&path, &source) {
192 eprintln!(
193 "Could not save invalid snapshot {}: {failure}",
194 path.display()
195 );
196 }
197 error
198 };
199 let store = Store::parse(&source).map_err(preserve)?;
200 if !store.checksum_mismatches.is_empty() {
201 return Err(preserve(onestore::Error {
202 offset: store.checksum_mismatches[0],
203 message: "A reader observed transaction checksum damage.",
204 })
205 .into());
206 }
207 let index = RevisionIndex::parse(&store).map_err(preserve)?;
208 index.validate_current().map_err(preserve)?;
209 let document = Document::parse(&index).map_err(preserve)?;
210 let mut targets = Vec::new();
211 for (sid, page) in document.pages().map_err(preserve)? {
212 let space = &document.spaces[&sid];
213 let revision = &space.revisions[&space.contexts[&ExGuid::default()]];
214 let mut pending = vec![page];
215 let mut seen = std::collections::BTreeSet::new();
216 while let Some(oid) = pending.pop() {
217 if !seen.insert(oid) {
218 continue;
219 }
220 let node = &revision.nodes[&oid];
221 pending.extend(
222 node.children
223 .iter()
224 .chain(&node.content)
225 .chain(&node.structure)
226 .copied(),
227 );
228 if let Kind::RichText { text, .. } = &node.kind
229 && text.starts_with("Concurrent edits:")
230 {
231 revision.text_runs(oid).map_err(preserve)?;
232 targets.push((sid, oid, text));
233 }
234 }
235 }
236 let [(sid, oid, text)] = targets.as_slice() else {
237 return Err(preserve(onestore::Error {
238 offset: 0,
239 message: "Expected one concurrent-edit paragraph.",
240 })
241 .into());
242 };
243 let observed = if documents {
244 Some(document_view(&source).map_err(preserve)?)
245 } else {
246 None
247 };
248 let read_finished = SystemTime::now().duration_since(UNIX_EPOCH)?.as_micros();
249 log(
250 json!({"event": "read", "attempt": attempts, "started_us": started, "finished_us": read_finished,
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)}),
252 )?;
253 if args[0] == "read" {
254 completed += 1;
255 continue;
256 }
257 let token = format!(" [{}:{}]", args[2], completed);
258 let offset = u32::try_from(text.encode_utf16().count())?;
259 let mut range = offset..offset;
260 let mut replacement = token.clone();
261 if args[0] == "edit" {
262 let prefix = "Concurrent edits:";
263 let mut boundaries = vec![u32::try_from(prefix.encode_utf16().count())?];
264 for character in text[prefix.len()..].chars() {
265 boundaries.push(boundaries.last().unwrap() + character.len_utf16() as u32);
266 }
267 let first = ((random >> 16) % boundaries.len() as u64) as usize;
268 let second = ((random >> 40) % boundaries.len() as u64) as usize;
269 range = boundaries[first.min(second)]..boundaries[first.max(second)];
270 replacement = format!(" café 🦀{token}");
271 }
272 log(
273 json!({"event": "intent", "attempt": attempts, "operation": completed,
274 "source_transaction": store.header.transaction_count, "before": text,
275 "range": [range.start, range.end], "replacement": replacement, "token": token}),
276 )?;
277 thread::sleep(Duration::from_millis((random >> 48) % 13));
278 let commit_started = SystemTime::now().duration_since(UNIX_EPOCH)?.as_micros();
279 let result = commit(&args[1], &source, *sid, *oid, range, &replacement);
280 let finished = SystemTime::now().duration_since(UNIX_EPOCH)?.as_micros();
281 match result {
282 Ok(()) => {
283 log(
284 json!({"event": "commit", "attempt": attempts, "operation": completed, "token": token,
285 "started_us": commit_started, "finished_us": finished, "source_transaction": store.header.transaction_count}),
286 )?;
287 completed += 1;
288 }
289 Err(error)
290 if error.state == CommitState::NotCommitted
291 && [
292 io::ErrorKind::WouldBlock,
293 io::ErrorKind::ResourceBusy,
294 io::ErrorKind::PermissionDenied,
295 io::ErrorKind::NotFound,
296 ]
297 .contains(&error.error.kind()) =>
298 {
299 log(
300 json!({"event": "retry", "attempt": attempts, "started_us": commit_started,
301 "finished_us": finished, "kind": format!("{:?}", error.error.kind())}),
302 )?;
303 }
304 Err(error) => {
305 log(
306 json!({"event": "commit_error", "attempt": attempts, "operation": completed,
307 "token": token, "state": format!("{:?}", error.state), "kind": format!("{:?}", error.error.kind()),
308 "started_us": commit_started, "finished_us": finished}),
309 )?;
310 return Err(error.into());
311 }
312 }
313 }
314 log(json!({"event": "done", "completed": completed, "attempts": attempts}))?;
315 Ok(())
316}
317
318#[cfg(test)]
319mod tests {
320 use super::*;
321 use onestore::{
322 TextAttribute,
323 document::{Format, Layout},
324 op::{Edit, Op, PageOp},
325 page::{
326 Outline, PageObject, PageParagraph, Paragraph, ParagraphContent, TextObject,
327 text::new_id,
328 },
329 };
330
331 fn paragraph(text: &str, level: u32) -> PageParagraph {
332 PageParagraph {
333 id: new_id().unwrap(),
334 parent: None,
335 level,
336 style: None,
337 format: Format::default(),
338 content: ParagraphContent::Text(TextObject {
339 id: new_id().unwrap(),
340 date_field: None,
341 text: Paragraph::new(
342 text.into(),
343 Format {
344 font: Some("Calibri".into()),
345 font_size: Some(11.0),
346 language: Some(0x409),
347 ..Format::default()
348 },
349 ),
350 tags: Vec::new(),
351 }),
352 lists: Vec::new(),
353 tags: Vec::new(),
354 media: Default::default(),
355 collapsed: false,
356 }
357 }
358
359 fn edited(image: &[u8], space: ExGuid, ops: Vec<PageOp>) -> Vec<u8> {
360 let arena = onestore::Arena::default();
361 let mut section = onestore::Section::open(&arena, image.to_vec()).unwrap();
362 let ops = ops.into_iter().map(|op| Op::Page { space, op }).collect();
363 section
364 .apply(
365 "Author",
366 &Edit {
367 at: 134_000_000_000_000_000,
368 ops,
369 },
370 )
371 .unwrap();
372 section.seal().unwrap();
373 section.image()
374 }
375
376 #[test]
377 fn observation_follows_split_suffixes_and_moved_children_through_active_ancestry() {
378 let source = onestore::create_section("observation.one", "Original", "Author").unwrap();
379 assert_eq!(document_view(&source).unwrap(), DocumentView::default());
380 let store = Store::parse(&source).unwrap();
381 let index = RevisionIndex::parse(&store).unwrap();
382 let document = Document::parse(&index).unwrap();
383 let (space, page) = document.pages().unwrap()[0];
384 let first = paragraph("Document w0:0 🦀", 1);
385 let (paragraph_id, text) = (first.id, first.text().unwrap().id);
386 let outline_id = new_id().unwrap();
387 let inserted = edited(
388 &source,
389 space,
390 vec![
391 PageOp::Add {
392 object: PageObject::Outline(Outline {
393 id: outline_id,
394 title: false,
395 min_width: None,
396 layout: Layout {
397 x: Some(144.0),
398 y: Some(216.0),
399 ..Default::default()
400 },
401 indents: Vec::new(),
402 paragraphs: vec![first],
403 unsupported: Vec::new(),
404 }),
405 before: None,
406 },
407 PageOp::Format {
408 text,
409 range: 0..16,
410 set: vec![TextAttribute::Bold(true)],
411 clear: Vec::new(),
412 },
413 ],
414 );
415 let before = document_view(&inserted).unwrap();
416 assert_eq!(before.texts.len(), 1);
417 assert_eq!(before.graph.len(), 2);
418 let outline = outline_id.to_string();
419 let paragraph_name = paragraph_id.to_string();
420 assert_eq!(before.graph[&outline]["children"], json!([paragraph_name]));
421 let child = paragraph("Unmarked child", 2);
422 let child_id = child.id;
423 let with_child = edited(
424 &inserted,
425 space,
426 vec![PageOp::Insert {
427 container: paragraph_id,
428 before: None,
429 paragraphs: vec![child],
430 }],
431 );
432 let before = document_view(&with_child).unwrap();
433 assert_eq!(before.texts.len(), 2);
434 assert_eq!(before.graph.len(), 3);
435 assert_eq!(
436 before.graph[&outline],
437 json!({"parent":page.to_string(), "children":[paragraph_name], "content":[], "child_level":1, "position":{"x":144.0,"y":216.0}})
438 );
439 for offset in [14, 16] {
440 let (split, right) = (new_id().unwrap(), new_id().unwrap());
441 let split_edit = edited(
442 &with_child,
443 space,
444 vec![PageOp::Split {
445 text,
446 at: offset,
447 paragraph: split,
448 right,
449 lists: Vec::new(),
450 }],
451 );
452 let observed = document_view(&split_edit).unwrap();
453 assert_eq!(observed.texts.len(), 3);
454 assert_eq!(observed.graph.len(), 4);
455 assert_eq!(
456 observed.texts[&right.to_string()]["text"],
457 if offset == 14 { "🦀" } else { "" }
458 );
459 assert_eq!(
460 observed.graph[&outline]["children"],
461 json!([paragraph_name, split.to_string()])
462 );
463 assert_eq!(observed.graph[&paragraph_name]["children"], json!([]));
464 assert_eq!(
465 observed.graph[&split.to_string()]["children"],
466 json!([child_id.to_string()])
467 );
468 assert_eq!(
469 observed.graph[&child_id.to_string()]["parent"],
470 split.to_string()
471 );
472 let joined = edited(&split_edit, space, vec![PageOp::Join { left: text, right }]);
473 let joined = document_view(&joined).unwrap();
474 assert_eq!(joined.graph, before.graph);
475 let characters = |view: &DocumentView| {
476 view.texts
477 .iter()
478 .map(|(id, text)| {
479 let runs = text["runs"]
480 .as_array()
481 .unwrap()
482 .iter()
483 .flat_map(|run| {
484 run["text"]
485 .as_str()
486 .unwrap()
487 .chars()
488 .map(|c| json!([c, run["bold"], run["size"], run["color"]]))
489 })
490 .collect::<Vec<_>>();
491 (id.clone(), (text["text"].clone(), runs))
492 })
493 .collect::<std::collections::BTreeMap<_, _>>()
494 };
495 assert_eq!(characters(&joined), characters(&before));
496 }
497 }
498}