| 1 | #[path = "support/current.rs"] |
| 2 | mod current; |
| 3 | use current::current; |
| 4 | |
| 5 | #[path = "support/checkpoint.rs"] |
| 6 | mod checkpoint; |
| 7 | #[path = "support/ops.rs"] |
| 8 | mod ops; |
| 9 | #[path = "support/trace.rs"] |
| 10 | mod trace; |
| 11 | |
| 12 | use onestore::read_snapshot; |
| 13 | use onestore::{ |
| 14 | ExGuid, RevisionIndex, Store, |
| 15 | document::{Document, Kind}, |
| 16 | }; |
| 17 | use std::{fs, io}; |
| 18 | use trace::{Event, Trace}; |
| 19 | |
| 20 | /// Types `with` over `range` of the text object in `source` and commits it to `io`. |
| 21 | fn commit_text( |
| 22 | io: &mut impl onestore::CommitIo, |
| 23 | source: &[u8], |
| 24 | space: ExGuid, |
| 25 | text: ExGuid, |
| 26 | range: std::ops::Range<u32>, |
| 27 | with: &str, |
| 28 | ) -> Result<(), onestore::CommitError> { |
| 29 | let op = onestore::op::PageOp::Text { |
| 30 | text, |
| 31 | range, |
| 32 | with: with.into(), |
| 33 | }; |
| 34 | ops::transaction(source, "Author", vec![onestore::op::Op::Page { space, op }]) |
| 35 | .unwrap() |
| 36 | .unwrap() |
| 37 | .commit(io) |
| 38 | } |
| 39 | |
| 40 | #[test] |
| 41 | fn storage_inspection_preserves_opaque_images_without_claiming_edit_readiness() { |
| 42 | for path in [ |
| 43 | "native-encrypted/encrypted-01/notebook/synthetic.one", |
| 44 | "native-protected-boundaries/notebook/synthetic.one", |
| 45 | "native-encrypted/cold-encrypted-02/notebook/Open Notebook.one", |
| 46 | "malformed/native-inflight.one", |
| 47 | ] { |
| 48 | let source = fs::read(format!("../../corpus/{path}")).unwrap(); |
| 49 | for block in [1, 17, 65536] { |
| 50 | let mut read = |offset: u64, output: &mut [u8]| { |
| 51 | let offset = usize::try_from(offset).unwrap(); |
| 52 | let count = output |
| 53 | .len() |
| 54 | .min(block) |
| 55 | .min(source.len().saturating_sub(offset)); |
| 56 | output[..count].copy_from_slice(&source[offset..offset + count]); |
| 57 | Ok(count) |
| 58 | }; |
| 59 | assert_eq!( |
| 60 | onestore::read_storage_snapshot(&mut read, source.len()).unwrap(), |
| 61 | Some(source.clone()) |
| 62 | ); |
| 63 | assert!(read_snapshot(&mut read, source.len()).unwrap().is_none()); |
| 64 | } |
| 65 | let mut headers = 0; |
| 66 | let changed = onestore::read_storage_snapshot( |
| 67 | |offset, output| { |
| 68 | let offset = usize::try_from(offset).unwrap(); |
| 69 | let count = output.len().min(source.len().saturating_sub(offset)); |
| 70 | output[..count].copy_from_slice(&source[offset..offset + count]); |
| 71 | if offset == 0 { |
| 72 | headers += 1; |
| 73 | if headers == 2 { |
| 74 | output[0] ^= 1; |
| 75 | } |
| 76 | } |
| 77 | Ok(count) |
| 78 | }, |
| 79 | source.len(), |
| 80 | ) |
| 81 | .unwrap(); |
| 82 | assert!(changed.is_none()); |
| 83 | let mut broken = source.clone(); |
| 84 | let offset = usize::try_from(Store::parse(&source).unwrap().header.root.offset).unwrap(); |
| 85 | broken[offset] ^= 1; |
| 86 | assert_eq!( |
| 87 | onestore::read_storage_snapshot( |
| 88 | |offset, output| { |
| 89 | let offset = usize::try_from(offset).unwrap(); |
| 90 | let count = output.len().min(broken.len().saturating_sub(offset)); |
| 91 | output[..count].copy_from_slice(&broken[offset..offset + count]); |
| 92 | Ok(count) |
| 93 | }, |
| 94 | broken.len() |
| 95 | ) |
| 96 | .unwrap_err() |
| 97 | .kind(), |
| 98 | std::io::ErrorKind::InvalidData |
| 99 | ); |
| 100 | } |
| 101 | } |
| 102 | |
| 103 | fn target(bytes: &[u8]) -> (ExGuid, ExGuid, u32) { |
| 104 | let store = Store::parse(bytes).unwrap(); |
| 105 | let index = RevisionIndex::parse(&store).unwrap(); |
| 106 | let doc = Document::parse(&index).unwrap(); |
| 107 | doc.spaces |
| 108 | .iter() |
| 109 | .find_map(|(sid, space)| { |
| 110 | space.revisions[&space.contexts[&ExGuid::default()]] |
| 111 | .nodes |
| 112 | .iter() |
| 113 | .find_map(|(oid, node)| { |
| 114 | if let Kind::RichText { text, .. } = &node.kind |
| 115 | && (text.starts_with("Fictitious") || text.starts_with("Transaction")) |
| 116 | { |
| 117 | return Some((*sid, *oid, text.encode_utf16().count() as u32)); |
| 118 | } |
| 119 | None |
| 120 | }) |
| 121 | }) |
| 122 | .unwrap() |
| 123 | } |
| 124 | |
| 125 | #[test] |
| 126 | fn version_cached_readers_cannot_miss_a_completed_publication() { |
| 127 | for path in [ |
| 128 | "native/20260905-05/snapshots/03-format-unicode/notebook/synthetic.one", |
| 129 | "append/round-01/tx-255/notebook/synthetic.one", |
| 130 | ] { |
| 131 | let source = fs::read(format!("../../corpus/{path}")).unwrap(); |
| 132 | let (sid, oid, end) = target(&source); |
| 133 | let mut trace = Trace { |
| 134 | bytes: source.clone(), |
| 135 | events: Vec::new(), |
| 136 | }; |
| 137 | commit_text(&mut trace, &source, sid, oid, end..end, " [cached reader]").unwrap(); |
| 138 | let final_content = current(&trace.bytes); |
| 139 | let mut visible = source; |
| 140 | for event in &trace.events { |
| 141 | let Event::Write(offset, bytes) = event else { |
| 142 | continue; |
| 143 | }; |
| 144 | visible.resize(visible.len().max(offset + bytes.len()), 0); |
| 145 | for (index, byte) in bytes.iter().enumerate() { |
| 146 | visible[offset + index] = *byte; |
| 147 | if (212..252).contains(&(offset + index)) { |
| 148 | let cached_content = current(&visible); |
| 149 | let refreshed = if visible[212..252] == trace.bytes[212..252] { |
| 150 | cached_content |
| 151 | } else { |
| 152 | current(&trace.bytes) |
| 153 | }; |
| 154 | assert_eq!( |
| 155 | refreshed, |
| 156 | final_content, |
| 157 | "{path}: cached header at {}", |
| 158 | offset + index |
| 159 | ); |
| 160 | } |
| 161 | } |
| 162 | } |
| 163 | } |
| 164 | } |
| 165 | |
| 166 | #[test] |
| 167 | fn published_snapshots_survive_interleaved_commit_io() { |
| 168 | let cases = [ |
| 169 | ( |
| 170 | "unicode", |
| 171 | "native/20260905-05/snapshots/03-format-unicode/notebook/synthetic.one", |
| 172 | ), |
| 173 | ( |
| 174 | "attachment", |
| 175 | "native/20260905-05/snapshots/06-attachment/notebook/synthetic.one", |
| 176 | ), |
| 177 | ( |
| 178 | "rollover-256", |
| 179 | "append/round-01/tx-255/notebook/synthetic.one", |
| 180 | ), |
| 181 | ( |
| 182 | "rollover-65536", |
| 183 | "append/round-01/tx-65535/notebook/synthetic.one", |
| 184 | ), |
| 185 | ( |
| 186 | "checkpoint", |
| 187 | "native/20260905-05/snapshots/03-format-unicode/notebook/synthetic.one", |
| 188 | ), |
| 189 | ]; |
| 190 | for (name, path) in cases { |
| 191 | let mut source = fs::read(format!("../../corpus/{path}")).unwrap(); |
| 192 | let (sid, oid, _) = target(&source); |
| 193 | if name == "checkpoint" { |
| 194 | source = checkpoint::pending(&source, sid, oid); |
| 195 | } |
| 196 | let (_, _, end) = target(&source); |
| 197 | let before = current(&source); |
| 198 | let mut trace = Trace { |
| 199 | bytes: source.clone(), |
| 200 | events: Vec::new(), |
| 201 | }; |
| 202 | commit_text(&mut trace, &source, sid, oid, end..end, " [reader café 🦀]").unwrap(); |
| 203 | let after = current(&trace.bytes); |
| 204 | assert_ne!(before, after); |
| 205 | let mut writes = Vec::new(); |
| 206 | for event in &trace.events { |
| 207 | if let Event::Write(offset, bytes) = event { |
| 208 | let piece = if *offset < 1024 { 17 } else { 4096 }; |
| 209 | writes.extend( |
| 210 | bytes |
| 211 | .chunks(piece) |
| 212 | .enumerate() |
| 213 | .map(|(i, bytes)| (offset + i * piece, bytes)), |
| 214 | ); |
| 215 | } |
| 216 | } |
| 217 | let mut accepted = [0; 2]; |
| 218 | let mut retried = 0; |
| 219 | let mut interleaved = 0; |
| 220 | for run in 0..writes.len() + 1 + 512 { |
| 221 | let paused = run <= writes.len(); |
| 222 | let mut step = if paused { run } else { 0 }; |
| 223 | let mut visible = source.clone(); |
| 224 | for &(offset, bytes) in &writes[..step] { |
| 225 | visible.resize(visible.len().max(offset + bytes.len()), 0); |
| 226 | visible[offset..offset + bytes.len()].copy_from_slice(bytes); |
| 227 | } |
| 228 | let mut seed = run as u64 + 1; |
| 229 | let mut read_calls = 0; |
| 230 | let mut overlapped = false; |
| 231 | let read_limit = [17, 193, 4096, 65536][run % 4]; |
| 232 | let result = read_snapshot( |
| 233 | |offset, output| { |
| 234 | seed ^= seed << 13; |
| 235 | seed ^= seed >> 7; |
| 236 | seed ^= seed << 17; |
| 237 | if !paused { |
| 238 | let count = if seed.is_multiple_of(11) { |
| 239 | writes.len() |
| 240 | } else { |
| 241 | (seed % 4) as usize |
| 242 | }; |
| 243 | let end = writes.len().min(step + count); |
| 244 | for &(at, bytes) in &writes[step..end] { |
| 245 | visible.resize(visible.len().max(at + bytes.len()), 0); |
| 246 | visible[at..at + bytes.len()].copy_from_slice(bytes); |
| 247 | } |
| 248 | overlapped |= read_calls > 0 && end > step; |
| 249 | step = end; |
| 250 | } |
| 251 | read_calls += 1; |
| 252 | let offset = offset as usize; |
| 253 | let count = output |
| 254 | .len() |
| 255 | .min(read_limit) |
| 256 | .min(visible.len().saturating_sub(offset)); |
| 257 | if count != 0 { |
| 258 | output[..count].copy_from_slice(&visible[offset..offset + count]); |
| 259 | } |
| 260 | Ok(count) |
| 261 | }, |
| 262 | trace.bytes.len(), |
| 263 | ); |
| 264 | interleaved += usize::from(overlapped); |
| 265 | if paused && (run == 0 || run == writes.len()) { |
| 266 | assert!( |
| 267 | matches!(&result, Ok(Some(_))), |
| 268 | "{name}: quiescent run {run}" |
| 269 | ); |
| 270 | } |
| 271 | match result { |
| 272 | Ok(Some(bytes)) => { |
| 273 | let checked = std::panic::catch_unwind(|| { |
| 274 | let observed = current(&bytes); |
| 275 | assert!( |
| 276 | observed == before || observed == after, |
| 277 | "{name}: run {run}, write step {step}" |
| 278 | ); |
| 279 | if run == writes.len() { |
| 280 | assert_eq!(observed, after); |
| 281 | } |
| 282 | observed == after |
| 283 | }); |
| 284 | match checked { |
| 285 | Ok(new) => accepted[usize::from(new)] += 1, |
| 286 | Err(failure) => { |
| 287 | let time = std::time::SystemTime::now() |
| 288 | .duration_since(std::time::UNIX_EPOCH) |
| 289 | .unwrap() |
| 290 | .as_nanos(); |
| 291 | let path = std::path::PathBuf::from(format!( |
| 292 | "../../evidence/m9/read-interleaving-failure-{name}-{run}-{time}" |
| 293 | )); |
| 294 | fs::create_dir_all(&path).unwrap(); |
| 295 | fs::write(path.join("source.one"), &source).unwrap(); |
| 296 | fs::write(path.join("observed.one"), &bytes).unwrap(); |
| 297 | fs::write( |
| 298 | path.join("replay.json"), |
| 299 | serde_json::to_vec(&serde_json::json!({ |
| 300 | "run": run, "read_limit": read_limit, "writes": writes, |
| 301 | })) |
| 302 | .unwrap(), |
| 303 | ) |
| 304 | .unwrap(); |
| 305 | eprintln!("Replay: {}", path.display()); |
| 306 | std::panic::resume_unwind(failure); |
| 307 | } |
| 308 | } |
| 309 | } |
| 310 | Ok(None) => retried += 1, |
| 311 | Err(error) |
| 312 | if matches!( |
| 313 | error.kind(), |
| 314 | io::ErrorKind::UnexpectedEof | io::ErrorKind::InvalidData |
| 315 | ) => |
| 316 | { |
| 317 | retried += 1 |
| 318 | } |
| 319 | Err(error) => panic!("{name}: run {run}: {error}"), |
| 320 | } |
| 321 | } |
| 322 | assert!( |
| 323 | accepted[0] > 0 && accepted[1] > 0 && interleaved > 0 && retried > 0, |
| 324 | "{name}" |
| 325 | ); |
| 326 | println!( |
| 327 | "{name}: old={}, new={}, retry={retried}, overlap={interleaved}", |
| 328 | accepted[0], accepted[1] |
| 329 | ); |
| 330 | } |
| 331 | } |
| 332 | |
| 333 | #[test] |
| 334 | fn snapshot_rejects_short_io_and_unbounded_allocation() { |
| 335 | let source = onestore::create_section("test.one", "read", "test").unwrap(); |
| 336 | let mut calls = 0; |
| 337 | let result = read_snapshot( |
| 338 | |offset, output| { |
| 339 | calls += 1; |
| 340 | if calls == 1 { |
| 341 | return Err(io::ErrorKind::Interrupted.into()); |
| 342 | } |
| 343 | let offset = offset as usize; |
| 344 | let count = output.len().min(7).min(source.len().saturating_sub(offset)); |
| 345 | if count != 0 { |
| 346 | output[..count].copy_from_slice(&source[offset..offset + count]); |
| 347 | } |
| 348 | Ok(count) |
| 349 | }, |
| 350 | source.len(), |
| 351 | ) |
| 352 | .unwrap() |
| 353 | .unwrap(); |
| 354 | assert_eq!(current(&source), current(&result)); |
| 355 | assert!(calls > source.len() / 7); |
| 356 | let error = read_snapshot(|_, output| Ok(output.len() + 1), source.len()).unwrap_err(); |
| 357 | assert_eq!(error.kind(), io::ErrorKind::InvalidData); |
| 358 | let error = read_snapshot(|_, _| Ok(0), source.len()).unwrap_err(); |
| 359 | assert_eq!(error.kind(), io::ErrorKind::UnexpectedEof); |
| 360 | let error = read_snapshot( |
| 361 | |_, output| { |
| 362 | output.copy_from_slice(&source[..1024]); |
| 363 | Ok(1024) |
| 364 | }, |
| 365 | 1023, |
| 366 | ) |
| 367 | .unwrap_err(); |
| 368 | assert_eq!(error.kind(), io::ErrorKind::InvalidData); |
| 369 | } |
| 370 | |
| 371 | #[test] |
| 372 | fn native_tail_truncation_before_header_update_is_retried() { |
| 373 | let compact = onestore::create_section("test.one", "Transaction retained", "test").unwrap(); |
| 374 | let mut expanded = compact.clone(); |
| 375 | expanded.resize(compact.len() + 216, 0); |
| 376 | let length = expanded.len() as u64; |
| 377 | expanded[196..204].copy_from_slice(&length.to_le_bytes()); |
| 378 | assert_eq!(current(&expanded), current(&compact)); |
| 379 | for block in [17, 193, 65536] { |
| 380 | for trigger in [0, 1024usize.div_ceil(block)] { |
| 381 | let mut visible = expanded.clone(); |
| 382 | let mut calls = 0; |
| 383 | let result = read_snapshot( |
| 384 | |offset, output| { |
| 385 | if calls == trigger { |
| 386 | visible.truncate(compact.len()); |
| 387 | } |
| 388 | calls += 1; |
| 389 | let offset = offset as usize; |
| 390 | let count = output |
| 391 | .len() |
| 392 | .min(block) |
| 393 | .min(visible.len().saturating_sub(offset)); |
| 394 | if count != 0 { |
| 395 | output[..count].copy_from_slice(&visible[offset..offset + count]); |
| 396 | } |
| 397 | Ok(count) |
| 398 | }, |
| 399 | expanded.len(), |
| 400 | ) |
| 401 | .unwrap(); |
| 402 | assert!( |
| 403 | result.is_none(), |
| 404 | "accepted a snapshot with a stale expected length" |
| 405 | ); |
| 406 | } |
| 407 | } |
| 408 | let result = read_snapshot( |
| 409 | |offset, output| { |
| 410 | let offset = offset as usize; |
| 411 | let count = output.len().min(compact.len().saturating_sub(offset)); |
| 412 | if count != 0 { |
| 413 | output[..count].copy_from_slice(&compact[offset..offset + count]); |
| 414 | } |
| 415 | Ok(count) |
| 416 | }, |
| 417 | expanded.len(), |
| 418 | ) |
| 419 | .unwrap() |
| 420 | .unwrap(); |
| 421 | assert_eq!(result, compact); |
| 422 | } |
| 423 | |
| 424 | #[test] |
| 425 | fn unpublished_append_remains_available_for_retry() { |
| 426 | let source = onestore::create_section("test.one", "Transaction before", "test").unwrap(); |
| 427 | let (sid, oid, end) = target(&source); |
| 428 | let mut trace = Trace { |
| 429 | bytes: source.clone(), |
| 430 | events: Vec::new(), |
| 431 | }; |
| 432 | commit_text(&mut trace, &source, sid, oid, end..end, " abandoned").unwrap(); |
| 433 | let Event::Write(offset, append) = &trace.events[0] else { |
| 434 | panic!() |
| 435 | }; |
| 436 | assert_eq!(*offset, source.len()); |
| 437 | for length in [1, 17, append.len()] { |
| 438 | let mut persisted = source.clone(); |
| 439 | persisted.extend_from_slice(&append[..length]); |
| 440 | let snapshot = read_snapshot( |
| 441 | |offset, output| { |
| 442 | let offset = offset as usize; |
| 443 | let count = output.len().min(persisted.len().saturating_sub(offset)); |
| 444 | output[..count].copy_from_slice(&persisted[offset..offset + count]); |
| 445 | Ok(count) |
| 446 | }, |
| 447 | persisted.len(), |
| 448 | ) |
| 449 | .unwrap() |
| 450 | .unwrap(); |
| 451 | assert_eq!(snapshot, persisted); |
| 452 | assert_eq!(current(&snapshot), current(&source)); |
| 453 | let mut retry = Trace { |
| 454 | bytes: persisted, |
| 455 | events: Vec::new(), |
| 456 | }; |
| 457 | commit_text(&mut retry, &snapshot, sid, oid, end..end, " retry").unwrap(); |
| 458 | assert_ne!(current(&retry.bytes), current(&source)); |
| 459 | } |
| 460 | } |
| 461 | |
| 462 | #[test] |
| 463 | fn header_comparison_cannot_replace_maintenance_exclusion() { |
| 464 | let mut source = onestore::create_section("test.one", "AAAA BBBB", "test").unwrap(); |
| 465 | source[212..228].fill(9); |
| 466 | source[236..252].fill(7); |
| 467 | let encoded: Vec<_> = "AAAA BBBB" |
| 468 | .encode_utf16() |
| 469 | .flat_map(u16::to_le_bytes) |
| 470 | .collect(); |
| 471 | let offset = source |
| 472 | .windows(encoded.len()) |
| 473 | .position(|bytes| bytes == encoded) |
| 474 | .unwrap(); |
| 475 | let mut changed = source.clone(); |
| 476 | let replacement: Vec<_> = "ZZZZ YYYY" |
| 477 | .encode_utf16() |
| 478 | .flat_map(u16::to_le_bytes) |
| 479 | .collect(); |
| 480 | changed[offset..offset + replacement.len()].copy_from_slice(&replacement); |
| 481 | let before = current(&source); |
| 482 | let after = current(&changed); |
| 483 | assert_ne!(before, after); |
| 484 | let mut visible = &source; |
| 485 | let result = read_snapshot( |
| 486 | |at, output| { |
| 487 | let at = at as usize; |
| 488 | if at >= offset + 8 { |
| 489 | visible = &changed; |
| 490 | } |
| 491 | let count = output.len().min(2).min(visible.len().saturating_sub(at)); |
| 492 | if count != 0 { |
| 493 | output[..count].copy_from_slice(&visible[at..at + count]); |
| 494 | } |
| 495 | Ok(count) |
| 496 | }, |
| 497 | source.len(), |
| 498 | ) |
| 499 | .unwrap() |
| 500 | .unwrap(); |
| 501 | let hybrid = current(&result); |
| 502 | assert_ne!(hybrid, before); |
| 503 | assert_ne!(hybrid, after); |
| 504 | } |