| 1 | #[path = "../../onestore/examples/support/concurrent.rs"] |
| 2 | mod concurrent; |
| 3 | |
| 4 | use notebook::smb::{Client, Credentials}; |
| 5 | use notebook::{EditStatus, Error, Remote, Replica, SmbRemote}; |
| 6 | use onestore::{CommitError, CommitState, ExGuid, RevisionIndex, Stamp, Store, Transaction}; |
| 7 | use serde_json::json; |
| 8 | use std::{ |
| 9 | env, |
| 10 | io::{self, Write}, |
| 11 | path::{Path, PathBuf}, |
| 12 | sync::Arc, |
| 13 | thread, |
| 14 | time::{Duration, Instant, SystemTime, UNIX_EPOCH}, |
| 15 | }; |
| 16 | |
| 17 | mod support { |
| 18 | pub mod view; |
| 19 | } |
| 20 | use concurrent::{DocumentView, document_view}; |
| 21 | use support::view::{cached, view}; |
| 22 | |
| 23 | impl DocumentView { |
| 24 | fn changes(&self, before: &Self) -> Self { |
| 25 | let difference = |
| 26 | |old: &serde_json::Map<String, serde_json::Value>, |
| 27 | new: &serde_json::Map<String, serde_json::Value>| { |
| 28 | old.keys() |
| 29 | .chain(new.keys()) |
| 30 | .filter(|id| old.get(*id) != new.get(*id)) |
| 31 | .map(|id| (id.clone(), new.get(id).cloned().unwrap_or_default())) |
| 32 | .collect() |
| 33 | }; |
| 34 | Self { |
| 35 | texts: difference(&before.texts, &self.texts), |
| 36 | graph: difference(&before.graph, &self.graph), |
| 37 | } |
| 38 | } |
| 39 | } |
| 40 | |
| 41 | fn now() -> u128 { |
| 42 | SystemTime::now() |
| 43 | .duration_since(UNIX_EPOCH) |
| 44 | .unwrap() |
| 45 | .as_micros() |
| 46 | } |
| 47 | |
| 48 | #[derive(Clone)] |
| 49 | enum Pause { |
| 50 | Outage(PathBuf), |
| 51 | FormatReply(PathBuf), |
| 52 | } |
| 53 | |
| 54 | struct Traced<R> { |
| 55 | remote: R, |
| 56 | before: Option<(String, Option<DocumentView>)>, |
| 57 | /// The last image read, which publications apply to. |
| 58 | read: Vec<u8>, |
| 59 | pause: Option<Pause>, |
| 60 | documents: bool, |
| 61 | } |
| 62 | |
| 63 | impl<R: Remote> Remote for Traced<R> { |
| 64 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 65 | let started = now(); |
| 66 | let bytes = self.remote.read()?; |
| 67 | let observed = view(&bytes).map_err(|error| io::Error::other(error.to_string()))?; |
| 68 | let documents = if self.documents { |
| 69 | Some(document_view(&bytes).map_err(io::Error::other)?) |
| 70 | } else { |
| 71 | None |
| 72 | }; |
| 73 | println!( |
| 74 | "{}", |
| 75 | 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)}) |
| 76 | ); |
| 77 | self.before = Some((observed.text, documents)); |
| 78 | self.read.clone_from(&bytes); |
| 79 | Ok(bytes) |
| 80 | } |
| 81 | fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { |
| 82 | // Publication reads nothing while the stamp holds; observe an unseen base once. |
| 83 | if Stamp::of(&self.read).ok().as_ref() != Some(transaction.base()) { |
| 84 | self.read().map_err(|error| CommitError { |
| 85 | state: CommitState::NotCommitted, |
| 86 | error, |
| 87 | })?; |
| 88 | } |
| 89 | let mut image = self.read.clone(); |
| 90 | transaction.apply(&mut image).map_err(|error| CommitError { |
| 91 | state: CommitState::NotCommitted, |
| 92 | error: io::Error::other(error.to_string()), |
| 93 | })?; |
| 94 | let after = view(&image).map_err(|error| CommitError { |
| 95 | state: CommitState::NotCommitted, |
| 96 | error: io::Error::other(error.to_string()), |
| 97 | })?; |
| 98 | let documents = if self.documents { |
| 99 | Some(document_view(&image).map_err(|error| CommitError { |
| 100 | state: CommitState::NotCommitted, |
| 101 | error: io::Error::other(error), |
| 102 | })?) |
| 103 | } else { |
| 104 | None |
| 105 | }; |
| 106 | let changes = if let Some(after) = &documents { |
| 107 | let before = self |
| 108 | .before |
| 109 | .as_ref() |
| 110 | .and_then(|(_, documents)| documents.as_ref()) |
| 111 | .ok_or_else(|| CommitError { |
| 112 | state: CommitState::NotCommitted, |
| 113 | error: io::Error::other("Document publication has no observed source"), |
| 114 | })?; |
| 115 | Some(after.changes(before)) |
| 116 | } else { |
| 117 | None |
| 118 | }; |
| 119 | let pause = match &self.pause { |
| 120 | Some(Pause::Outage(marker)) => Some(( |
| 121 | marker, |
| 122 | marker.parent().unwrap().join("offline-outage-resumed"), |
| 123 | "outage", |
| 124 | )), |
| 125 | Some(Pause::FormatReply(marker)) |
| 126 | if changes.as_ref().is_some_and(|changes| { |
| 127 | changes.texts.len() == 1 |
| 128 | && self |
| 129 | .before |
| 130 | .as_ref() |
| 131 | .and_then(|(_, before)| before.as_ref()) |
| 132 | .is_some_and(|before| { |
| 133 | changes.texts.keys().all(|id| before.texts.contains_key(id)) |
| 134 | }) |
| 135 | }) => |
| 136 | { |
| 137 | Some((marker, marker.with_extension("resume"), "format")) |
| 138 | } |
| 139 | _ => None, |
| 140 | }; |
| 141 | if let Some((marker, resumed, kind)) = pause |
| 142 | && !resumed.exists() |
| 143 | { |
| 144 | std::fs::write(marker, after.revision.to_string()).map_err(|error| CommitError { |
| 145 | state: CommitState::NotCommitted, |
| 146 | error, |
| 147 | })?; |
| 148 | println!( |
| 149 | "{}", |
| 150 | json!({"event":"publication_paused", "revision":after.revision.to_string(), "kind":kind, "at_us":now()}) |
| 151 | ); |
| 152 | let deadline = Instant::now() + Duration::from_secs(120); |
| 153 | while !resumed.exists() { |
| 154 | if Instant::now() >= deadline { |
| 155 | return Err(CommitError { |
| 156 | state: CommitState::NotCommitted, |
| 157 | error: io::Error::new( |
| 158 | io::ErrorKind::TimedOut, |
| 159 | "Offline publication barrier timed out", |
| 160 | ), |
| 161 | }); |
| 162 | } |
| 163 | thread::sleep(Duration::from_millis(10)); |
| 164 | } |
| 165 | } |
| 166 | let started = now(); |
| 167 | let result = self.remote.publish(transaction); |
| 168 | let finished = now(); |
| 169 | println!( |
| 170 | "{}", |
| 171 | json!({"event":"remote_attempt", "started_us":started, "finished_us":finished, |
| 172 | "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, |
| 173 | "state":format!("{:?}", result.as_ref().map_or_else(|error| error.state, |_| CommitState::Committed)), |
| 174 | "documents":documents.as_ref().map(|view| &view.texts), "document_graph":documents.as_ref().map(|view| &view.graph), |
| 175 | "document_changes":changes.as_ref().map(|view| &view.texts), "document_graph_changes":changes.as_ref().map(|view| &view.graph)}) |
| 176 | ); |
| 177 | if result |
| 178 | .as_ref() |
| 179 | .is_err_and(|error| error.state == CommitState::Unknown) |
| 180 | && let Some(Pause::FormatReply(marker)) = &self.pause |
| 181 | && marker.with_extension("isolate").exists() |
| 182 | { |
| 183 | std::fs::write(marker.with_extension("isolate"), serde_json::to_vec(&json!({"space":after.space.to_string(), "revision":after.revision.to_string(), "after_us":finished})).unwrap()) |
| 184 | .map_err(|error| CommitError { state:CommitState::Unknown, error })?; |
| 185 | println!( |
| 186 | "{}", |
| 187 | json!({"event":"confirmation_paused", "revision":after.revision.to_string(), "at_us":now()}) |
| 188 | ); |
| 189 | let deadline = Instant::now() + Duration::from_secs(120); |
| 190 | while !marker.with_extension("confirmation-resume").exists() { |
| 191 | if Instant::now() >= deadline { |
| 192 | return Err(CommitError { |
| 193 | state: CommitState::Unknown, |
| 194 | error: io::Error::new( |
| 195 | io::ErrorKind::TimedOut, |
| 196 | "Offline confirmation barrier timed out", |
| 197 | ), |
| 198 | }); |
| 199 | } |
| 200 | thread::sleep(Duration::from_millis(10)); |
| 201 | } |
| 202 | } |
| 203 | if result.is_ok() { |
| 204 | self.read = image; |
| 205 | self.before = Some((after.text, documents)); |
| 206 | } |
| 207 | result |
| 208 | } |
| 209 | fn stamp(&mut self) -> io::Result<Stamp> { |
| 210 | self.remote.stamp() |
| 211 | } |
| 212 | fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> { |
| 213 | let captured = (|| -> Result<_, Box<dyn std::error::Error>> { |
| 214 | // A confirmation follows the read of its image, except for edits that changed nothing. |
| 215 | if Stamp::of(&self.read).ok().as_ref() != Some(base) { |
| 216 | self.read = self.remote.read()?; |
| 217 | } |
| 218 | let snapshot = &self.read[..]; |
| 219 | let store = Store::parse(snapshot)?; |
| 220 | let index = RevisionIndex::parse(&store)?; |
| 221 | let revisions = index |
| 222 | .spaces |
| 223 | .iter() |
| 224 | .map(|(id, space)| { |
| 225 | ( |
| 226 | id.to_string(), |
| 227 | space |
| 228 | .revisions |
| 229 | .keys() |
| 230 | .map(ToString::to_string) |
| 231 | .collect::<Vec<_>>(), |
| 232 | ) |
| 233 | }) |
| 234 | .collect::<std::collections::BTreeMap<_, _>>(); |
| 235 | let current = index |
| 236 | .spaces |
| 237 | .iter() |
| 238 | .filter_map(|(id, space)| { |
| 239 | space |
| 240 | .labels |
| 241 | .get(&(ExGuid::default(), 1)) |
| 242 | .map(|revision| (id.to_string(), revision.to_string())) |
| 243 | }) |
| 244 | .collect::<std::collections::BTreeMap<_, _>>(); |
| 245 | let path = env::var_os("ONESTORE_OFFLINE_CONFIRM_DIR").map(|directory| { |
| 246 | PathBuf::from(directory).join(format!("{}-{}.one", std::process::id(), now())) |
| 247 | }); |
| 248 | if let Some(path) = &path { |
| 249 | std::fs::OpenOptions::new() |
| 250 | .write(true) |
| 251 | .create_new(true) |
| 252 | .open(path)? |
| 253 | .write_all(snapshot)?; |
| 254 | } |
| 255 | Ok((revisions, current, path)) |
| 256 | })(); |
| 257 | let (revisions, current, capture) = captured.map_err(|error| CommitError { |
| 258 | state: CommitState::NotCommitted, |
| 259 | error: io::Error::other(error.to_string()), |
| 260 | })?; |
| 261 | let started = now(); |
| 262 | let result = self.remote.confirm(base); |
| 263 | println!( |
| 264 | "{}", |
| 265 | json!({"event":"remote_confirm", "started_us":started, "finished_us":now(), "revisions":revisions, "current_revisions":current, "capture":capture.as_ref().and_then(|path| path.file_name()).map(|name| name.to_string_lossy()), "text":self.before.as_ref().map(|(text,_)|text), |
| 266 | "state":format!("{:?}", result.as_ref().map_or_else(|error| error.state, |_| CommitState::Committed)), "error":result.as_ref().err().map(|error|error.error.to_string())}) |
| 267 | ); |
| 268 | result |
| 269 | } |
| 270 | } |
| 271 | |
| 272 | /// The edit appending `token` at `at` of a text. |
| 273 | fn append(space: ExGuid, text: ExGuid, at: u32, token: &str) -> onestore::op::Edit { |
| 274 | onestore::op::Edit { |
| 275 | at: u64::try_from(now()).unwrap_or_default() * 10 + 116_444_736_000_000_000, |
| 276 | ops: vec![onestore::op::Op::Page { |
| 277 | space, |
| 278 | op: onestore::op::PageOp::Text { |
| 279 | text, |
| 280 | range: at..at, |
| 281 | with: token.to_owned(), |
| 282 | }, |
| 283 | }], |
| 284 | } |
| 285 | } |
| 286 | |
| 287 | fn main() -> Result<(), Box<dyn std::error::Error>> { |
| 288 | let args: Vec<_> = env::args().skip(1).collect(); |
| 289 | if args.len() != 7 |
| 290 | || !["read", "write"].contains(&args[0].as_str()) |
| 291 | || !args[2].bytes().all(|b| b.is_ascii_alphanumeric()) |
| 292 | { |
| 293 | return Err("Expected read|write FILE ACTOR OPERATIONS START_FILE STOP_FILE SEED".into()); |
| 294 | } |
| 295 | let documents = env::var_os("ONESTORE_OFFLINE_DOCUMENTS").is_some(); |
| 296 | let address = env::var("ONESTORE_SMB_LAB")?; |
| 297 | let share = env::var("ONESTORE_SMB_SHARE")?; |
| 298 | let initial = Client::connect( |
| 299 | &address, |
| 300 | &share, |
| 301 | Credentials::default(), |
| 302 | Duration::from_secs(5), |
| 303 | )?; |
| 304 | if args[0] == "read" { |
| 305 | return concurrent::run( |
| 306 | &args, |
| 307 | |path| initial.read(path, 256 * 1024 * 1024), |
| 308 | |_, _, _, _, _, _| unreachable!(), |
| 309 | ); |
| 310 | } |
| 311 | let timeout: u64 = env::var("ONESTORE_CLIENT_TIMEOUT_MS") |
| 312 | .unwrap_or_else(|_| "600000".into()) |
| 313 | .parse()?; |
| 314 | let deadline = Instant::now() |
| 315 | .checked_add(Duration::from_millis(timeout)) |
| 316 | .ok_or("Invalid timeout")?; |
| 317 | let operations: usize = args[3].parse()?; |
| 318 | let mut seed: u64 = args[6].parse()?; |
| 319 | if operations == 0 { |
| 320 | return Err("Expected positive operations".into()); |
| 321 | } |
| 322 | let source = initial.read(&args[1], 256 * 1024 * 1024)?; |
| 323 | drop(initial); |
| 324 | let cache_path = Path::new(&args[4]) |
| 325 | .parent() |
| 326 | .ok_or("Missing workload directory")? |
| 327 | .join(format!("{}.sqlite", args[2])); |
| 328 | let cache = Arc::new(Replica::create(&cache_path, &source)?); |
| 329 | println!( |
| 330 | "{}", |
| 331 | json!({"event":"ready", "pid":std::process::id(), "actor":args[2], "offline":true, "document_operations":false,"document_graph":documents,"document_kinds":[]}) |
| 332 | ); |
| 333 | while !Path::new(&args[4]).exists() { |
| 334 | if Instant::now() >= deadline { |
| 335 | return Err("Start barrier timed out".into()); |
| 336 | } |
| 337 | thread::sleep(Duration::from_millis(5)); |
| 338 | } |
| 339 | let path = args[1].clone(); |
| 340 | let outage = env::var_os("ONESTORE_OFFLINE_OUTAGE_DIR").map(PathBuf::from); |
| 341 | let pause = outage |
| 342 | .as_ref() |
| 343 | .map(|directory| Pause::Outage(directory.join(format!("offline-paused-{}", args[2])))) |
| 344 | .or_else(|| { |
| 345 | env::var_os("ONESTORE_OFFLINE_FORMAT_REPLY_DIR").map(|directory| { |
| 346 | Pause::FormatReply( |
| 347 | PathBuf::from(directory).join(format!("offline-paused-{}", args[2])), |
| 348 | ) |
| 349 | }) |
| 350 | }); |
| 351 | let (fatal_tx, fatal_rx) = std::sync::mpsc::channel(); |
| 352 | let worker = cache.start_sync(Duration::from_millis(50), move || { |
| 353 | let client = Client::connect(&address, &share, Credentials::default(), Duration::from_secs(5))?; |
| 354 | println!("{}", json!({"event":"transport_connected", "at_us":now()})); |
| 355 | Ok(Traced { remote: SmbRemote::new(client, &path, 256 * 1024 * 1024), before:None, read:Vec::new(), pause:pause.clone(), documents }) |
| 356 | }, move |result| { |
| 357 | if let Err(error) = result { |
| 358 | println!("{}", json!({"event":"sync_error", "error":error.to_string(), "at_us":now()})); |
| 359 | if !matches!(error, Error::RemoteIo(_) | Error::Remote(_)) && !matches!(error, Error::Io(error) if error.kind() == io::ErrorKind::WouldBlock) { |
| 360 | let _ = fatal_tx.send(error.to_string()); |
| 361 | } |
| 362 | } |
| 363 | })?; |
| 364 | let mut ids = Vec::new(); |
| 365 | let mut tokens = Vec::new(); |
| 366 | let result = (|| -> Result<(), Box<dyn std::error::Error>> { |
| 367 | let mut generated = 0; |
| 368 | let mut received = 0; |
| 369 | loop { |
| 370 | match fatal_rx.try_recv() { |
| 371 | Ok(error) => return Err(error.into()), |
| 372 | Err(std::sync::mpsc::TryRecvError::Disconnected) => { |
| 373 | return Err("Worker exited before completion".into()); |
| 374 | } |
| 375 | Err(std::sync::mpsc::TryRecvError::Empty) => {} |
| 376 | } |
| 377 | while received < ids.len() { |
| 378 | match cache.status(ids[received])? { |
| 379 | Some(EditStatus::Published { revision }) => { |
| 380 | println!( |
| 381 | "{}", |
| 382 | json!({"event":"remote_receipt", "id":ids[received], "revision":revision.to_string(), "at_us":now()}) |
| 383 | ); |
| 384 | received += 1; |
| 385 | } |
| 386 | Some(_) => break, |
| 387 | None => return Err("Local intent disappeared".into()), |
| 388 | } |
| 389 | } |
| 390 | if Instant::now() >= deadline { |
| 391 | return Err("Offline workload timed out; cache retained".into()); |
| 392 | } |
| 393 | let pending = cache.pending()?; |
| 394 | let capacity = if outage |
| 395 | .as_ref() |
| 396 | .is_some_and(|directory| !directory.join("offline-outage-down").exists()) |
| 397 | { |
| 398 | 1 |
| 399 | } else { |
| 400 | 8 |
| 401 | }; |
| 402 | if generated < operations && pending.len() < capacity { |
| 403 | let (space, object, text) = cached(&cache)?; |
| 404 | let at = u32::try_from(text.encode_utf16().count())?; |
| 405 | let token = format!(" [{}:{}]", args[2], generated); |
| 406 | let started = now(); |
| 407 | match cache.apply("Offline document writer", append(space, object, at, &token)) { |
| 408 | Ok(id) => { |
| 409 | println!( |
| 410 | "{}", |
| 411 | json!({"event":"local_commit", "id":id, "operation":generated, "space":space.to_string(), "object":object.to_string(), "before":text, "token":token, "started_us":started, "finished_us":now()}) |
| 412 | ); |
| 413 | ids.push(id); |
| 414 | tokens.push(token); |
| 415 | generated += 1; |
| 416 | } |
| 417 | other => return Err(format!("Unexpected local result: {other:?}").into()), |
| 418 | } |
| 419 | } |
| 420 | if received == ids.len() && cache.pending()?.is_empty() { |
| 421 | // Concurrent appends meet at the end of the text: a merge keeps the remote's |
| 422 | // text and this writer's on a conflict page. Append the published tokens |
| 423 | // the text lacks after it again. |
| 424 | let (space, object, text) = cached(&cache)?; |
| 425 | let missing: Vec<String> = tokens |
| 426 | .iter() |
| 427 | .filter(|token| !text.contains(token.as_str())) |
| 428 | .cloned() |
| 429 | .collect(); |
| 430 | let mut at = u32::try_from(text.encode_utf16().count())?; |
| 431 | for token in &missing { |
| 432 | ids.push( |
| 433 | cache.apply("Offline document writer", append(space, object, at, token))?, |
| 434 | ); |
| 435 | at += u32::try_from(token.encode_utf16().count())?; |
| 436 | } |
| 437 | if !missing.is_empty() { |
| 438 | println!( |
| 439 | "{}", |
| 440 | json!({"event":"reviewed_append", "remote":text, "tokens":missing, "at_us":now()}) |
| 441 | ); |
| 442 | } |
| 443 | } |
| 444 | if generated == operations && received == ids.len() { |
| 445 | if !cache.pending()?.is_empty() { |
| 446 | return Err("Acknowledged queue did not drain".into()); |
| 447 | } |
| 448 | break; |
| 449 | } |
| 450 | seed = seed |
| 451 | .wrapping_mul(6364136223846793005) |
| 452 | .wrapping_add(1442695040888963407); |
| 453 | thread::sleep(Duration::from_millis(1 + (seed >> 32) % 7)); |
| 454 | } |
| 455 | Ok(()) |
| 456 | })(); |
| 457 | let stopped = worker.stop(); |
| 458 | result?; |
| 459 | stopped?; |
| 460 | drop(cache); |
| 461 | let reopened = Replica::open(&cache_path)?; |
| 462 | if !reopened.pending()?.is_empty() { |
| 463 | return Err("Pending edits reappeared after reopen".into()); |
| 464 | } |
| 465 | for id in ids { |
| 466 | let Some(EditStatus::Published { revision }) = reopened.status(id)? else { |
| 467 | return Err("Receipt did not survive reopen".into()); |
| 468 | }; |
| 469 | println!( |
| 470 | "{}", |
| 471 | json!({"event":"reopened_receipt", "id":id, "revision":revision.to_string()}) |
| 472 | ); |
| 473 | } |
| 474 | println!( |
| 475 | "{}", |
| 476 | json!({"event":"done", "operations":operations, "at_us":now()}) |
| 477 | ); |
| 478 | Ok(()) |
| 479 | } |