| 1 | use onestore::{ |
| 2 | CommitState, ExGuid, RevisionIndex, Store, |
| 3 | document::{Document, Kind}, |
| 4 | }; |
| 5 | use serde_json::json; |
| 6 | use 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)] |
| 15 | pub struct DocumentView { |
| 16 | pub texts: serde_json::Map<String, serde_json::Value>, |
| 17 | pub graph: serde_json::Map<String, serde_json::Value>, |
| 18 | } |
| 19 | |
| 20 | pub 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 | |
| 86 | pub 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)] |
| 319 | mod 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 | } |