| 1 | use super::*; |
| 2 | use serde_json::{Value, json}; |
| 3 | use std::{collections::BTreeMap, path::Path, process::Child}; |
| 4 | |
| 5 | #[derive(Clone)] |
| 6 | struct Snapshot { |
| 7 | content: BTreeMap<ExGuid, String>, |
| 8 | modified: BTreeMap<(ExGuid, ExGuid), u32>, |
| 9 | } |
| 10 | |
| 11 | fn snapshot(bytes: &[u8]) -> Snapshot { |
| 12 | let store = Store::parse(bytes).unwrap(); |
| 13 | assert!(store.checksum_mismatches.is_empty()); |
| 14 | let index = RevisionIndex::parse(&store).unwrap(); |
| 15 | index.validate_current().unwrap(); |
| 16 | Document::parse(&index).unwrap(); |
| 17 | let mut content = BTreeMap::new(); |
| 18 | let mut modified = BTreeMap::new(); |
| 19 | for (sid, space) in &index.spaces { |
| 20 | let revision = index |
| 21 | .resolve(*sid, space.labels[&(ExGuid::default(), 1)]) |
| 22 | .unwrap(); |
| 23 | let mut objects = BTreeMap::new(); |
| 24 | for (oid, object) in &revision.objects { |
| 25 | if let Some(onestore::FileDataReference::Internal(guid)) = |
| 26 | object.file_reference().unwrap() |
| 27 | { |
| 28 | store.file_data(guid).unwrap(); |
| 29 | } |
| 30 | let data = match object.data { |
| 31 | onestore::ObjectData::Properties(bytes) => { |
| 32 | let mut properties = onestore::PropertySets::parse(bytes).unwrap(); |
| 33 | for property in &mut properties.sets[0] { |
| 34 | if property.id == 0x14001d7a { |
| 35 | let onestore::Value::Bytes(value) = property.value else { |
| 36 | panic!() |
| 37 | }; |
| 38 | assert!( |
| 39 | modified |
| 40 | .insert( |
| 41 | (*sid, *oid), |
| 42 | u32::from_le_bytes(value.try_into().unwrap()) |
| 43 | ) |
| 44 | .is_none() |
| 45 | ); |
| 46 | property.value = onestore::Value::Bytes(&[0; 4]); |
| 47 | } |
| 48 | } |
| 49 | format!("{properties:?}") |
| 50 | } |
| 51 | other => format!("{other:?}"), |
| 52 | }; |
| 53 | objects.insert( |
| 54 | *oid, |
| 55 | ( |
| 56 | object.jcid, |
| 57 | object.reference_count, |
| 58 | data, |
| 59 | format!("{:?}", object.references().unwrap()), |
| 60 | ), |
| 61 | ); |
| 62 | } |
| 63 | content.insert(*sid, format!("{:?} {objects:?}", revision.roots)); |
| 64 | } |
| 65 | Snapshot { content, modified } |
| 66 | } |
| 67 | |
| 68 | fn stamp() -> u32 { |
| 69 | (SystemTime::now() |
| 70 | .duration_since(UNIX_EPOCH) |
| 71 | .unwrap() |
| 72 | .as_secs() |
| 73 | - 315532800) |
| 74 | .try_into() |
| 75 | .unwrap() |
| 76 | } |
| 77 | |
| 78 | fn published( |
| 79 | actual: &Snapshot, |
| 80 | old: &Snapshot, |
| 81 | new: &Snapshot, |
| 82 | time: std::ops::RangeInclusive<u32>, |
| 83 | ) -> bool { |
| 84 | if actual.content == old.content { |
| 85 | assert_eq!( |
| 86 | actual.modified, old.modified, |
| 87 | "Unpublished modification time changed" |
| 88 | ); |
| 89 | return false; |
| 90 | } |
| 91 | assert_eq!( |
| 92 | actual.content, new.content, |
| 93 | "Partial or unexpected publication" |
| 94 | ); |
| 95 | assert!(actual.modified.keys().eq(new.modified.keys())); |
| 96 | for (key, value) in &actual.modified { |
| 97 | if old.modified.get(key) != new.modified.get(key) { |
| 98 | assert!( |
| 99 | time.contains(value), |
| 100 | "Modification time outside the commit interval" |
| 101 | ); |
| 102 | } else { |
| 103 | assert_eq!( |
| 104 | *value, new.modified[key], |
| 105 | "Unrelated modification time changed" |
| 106 | ); |
| 107 | } |
| 108 | } |
| 109 | true |
| 110 | } |
| 111 | |
| 112 | #[test] |
| 113 | fn publication_oracle_bounds_changed_timestamps_and_preserves_others() { |
| 114 | let id = ExGuid::default(); |
| 115 | let other = ExGuid { n: 1, ..id }; |
| 116 | let old = Snapshot { |
| 117 | content: BTreeMap::from([(id, "old".into())]), |
| 118 | modified: BTreeMap::from([((id, id), 10), ((id, other), 7)]), |
| 119 | }; |
| 120 | let new = Snapshot { |
| 121 | content: BTreeMap::from([(id, "new".into())]), |
| 122 | modified: BTreeMap::from([((id, id), 20), ((id, other), 7)]), |
| 123 | }; |
| 124 | let mut actual = new.clone(); |
| 125 | actual.modified.insert((id, id), 30); |
| 126 | assert!(published(&actual, &old, &new, 25..=35)); |
| 127 | assert!(!published(&old, &old, &new, 25..=35)); |
| 128 | for (key, value) in [((id, id), 24), ((id, id), 36), ((id, other), 30)] { |
| 129 | let mut changed = actual.clone(); |
| 130 | changed.modified.insert(key, value); |
| 131 | assert!(std::panic::catch_unwind(|| published(&changed, &old, &new, 25..=35)).is_err()); |
| 132 | } |
| 133 | actual.content = old.content.clone(); |
| 134 | assert!(std::panic::catch_unwind(|| published(&actual, &old, &new, 25..=35)).is_err()); |
| 135 | } |
| 136 | |
| 137 | struct Proxy(Child); |
| 138 | impl Drop for Proxy { |
| 139 | fn drop(&mut self) { |
| 140 | let _ = self.0.kill(); |
| 141 | let _ = self.0.wait(); |
| 142 | } |
| 143 | } |
| 144 | |
| 145 | fn records(output: &Path) -> Vec<Value> { |
| 146 | let data = fs::read_to_string(output.join("proxy.jsonl")).unwrap(); |
| 147 | data.rsplit_once('\n') |
| 148 | .map_or("", |(complete, _)| complete) |
| 149 | .lines() |
| 150 | .map(|line| serde_json::from_str(line).unwrap()) |
| 151 | .collect() |
| 152 | } |
| 153 | |
| 154 | fn configure(output: &Path, state: Value) -> usize { |
| 155 | fs::write(output.join("control.tmp"), state.to_string()).unwrap(); |
| 156 | fs::rename(output.join("control.tmp"), output.join("control.json")).unwrap(); |
| 157 | let deadline = Instant::now() + Duration::from_secs(5); |
| 158 | loop { |
| 159 | let events = records(output); |
| 160 | if let Some(index) = events |
| 161 | .iter() |
| 162 | .rposition(|event| event.get("control") == Some(&state)) |
| 163 | { |
| 164 | return index + 1; |
| 165 | } |
| 166 | assert!( |
| 167 | Instant::now() < deadline, |
| 168 | "proxy did not acknowledge its control state" |
| 169 | ); |
| 170 | std::thread::sleep(Duration::from_millis(5)); |
| 171 | } |
| 172 | } |
| 173 | |
| 174 | fn proxy(output: &Path) -> (Proxy, String) { |
| 175 | let address = std::env::var("ONESTORE_SMB_LAB").unwrap(); |
| 176 | let (host, port) = address.rsplit_once(':').unwrap(); |
| 177 | let mut proxy = Proxy( |
| 178 | std::process::Command::new("python3") |
| 179 | .arg(Path::new(env!("CARGO_MANIFEST_DIR")).join("../../tools/smb-proxy.py")) |
| 180 | .arg(output.join("control.json")) |
| 181 | .args(["--port", "0", "--server", host, "--server-port", port]) |
| 182 | .stdout(fs::File::create(output.join("proxy.jsonl")).unwrap()) |
| 183 | .stderr(fs::File::create(output.join("proxy.stderr")).unwrap()) |
| 184 | .spawn() |
| 185 | .unwrap(), |
| 186 | ); |
| 187 | let deadline = Instant::now() + Duration::from_secs(5); |
| 188 | let port = loop { |
| 189 | if let Some(port) = records(output) |
| 190 | .iter() |
| 191 | .find_map(|event| event["listening"].as_u64()) |
| 192 | { |
| 193 | break port; |
| 194 | } |
| 195 | assert!(proxy.0.try_wait().unwrap().is_none()); |
| 196 | assert!(Instant::now() < deadline, "proxy did not start"); |
| 197 | std::thread::sleep(Duration::from_millis(5)); |
| 198 | }; |
| 199 | (proxy, format!("127.0.0.1:{port}")) |
| 200 | } |
| 201 | |
| 202 | #[test] |
| 203 | #[ignore = "requires an owned Samba share and a new ONESTORE_SMB_EVIDENCE directory"] |
| 204 | fn live_asset_loss() { |
| 205 | let output = std::path::PathBuf::from(std::env::var("ONESTORE_SMB_EVIDENCE").unwrap()); |
| 206 | fs::create_dir(&output).unwrap(); |
| 207 | let (_proxy, address) = proxy(&output); |
| 208 | let observer = client(); |
| 209 | let bytes: Vec<_> = (0..1048577).map(|index| (index % 251) as u8).collect(); |
| 210 | let path = format!( |
| 211 | "asset-loss-{}.onebin", |
| 212 | SystemTime::now() |
| 213 | .duration_since(UNIX_EPOCH) |
| 214 | .unwrap() |
| 215 | .as_nanos() |
| 216 | ); |
| 217 | create(&observer, &path, &bytes); |
| 218 | for (command, occurrence) in [(8, 1), (8, 9), (8, 17), (8, 18), (6, 1)] { |
| 219 | for direction in ["request", "response"] { |
| 220 | let reader = Client::connect( |
| 221 | &address, |
| 222 | "agent", |
| 223 | Credentials::default(), |
| 224 | Duration::from_secs(5), |
| 225 | ) |
| 226 | .unwrap(); |
| 227 | // Response occurrences count only replies with the selected status. |
| 228 | let begin = configure( |
| 229 | &output, |
| 230 | json!({"cut":command,"occurrence":if command == 8 && occurrence == 18 && direction == "response" { 1 } else { occurrence }, |
| 231 | "direction":direction,"status":if command == 8 && occurrence == 18 { "0xc0000011" } else { "0x0" }}), |
| 232 | ); |
| 233 | assert!( |
| 234 | reader.read_asset(&path, bytes.len()).is_err(), |
| 235 | "{command}/{occurrence}/{direction}" |
| 236 | ); |
| 237 | assert_eq!( |
| 238 | reader.read_asset(&path, bytes.len()).unwrap_err().kind(), |
| 239 | io::ErrorKind::NotConnected |
| 240 | ); |
| 241 | assert_eq!( |
| 242 | records(&output)[begin..] |
| 243 | .iter() |
| 244 | .filter(|row| row.get("cut").is_some()) |
| 245 | .count(), |
| 246 | 1 |
| 247 | ); |
| 248 | configure( |
| 249 | &output, |
| 250 | json!({"phase":format!("{command}-{occurrence}-{direction}-reconnected")}), |
| 251 | ); |
| 252 | let reconnected = Client::connect( |
| 253 | &address, |
| 254 | "agent", |
| 255 | Credentials::default(), |
| 256 | Duration::from_secs(5), |
| 257 | ) |
| 258 | .unwrap(); |
| 259 | assert_eq!(reconnected.read_asset(&path, bytes.len()).unwrap(), bytes); |
| 260 | } |
| 261 | } |
| 262 | assert_eq!(observer.read_asset(&path, bytes.len()).unwrap(), bytes); |
| 263 | } |
| 264 | |
| 265 | #[test] |
| 266 | #[ignore = "requires an owned Samba share and a new ONESTORE_SMB_EVIDENCE directory"] |
| 267 | fn live_message_loss() { |
| 268 | let output = std::path::PathBuf::from(std::env::var("ONESTORE_SMB_EVIDENCE").unwrap()); |
| 269 | fs::create_dir(&output).unwrap(); |
| 270 | fs::create_dir(output.join("interrupted")).unwrap(); |
| 271 | fs::create_dir(output.join("recovered")).unwrap(); |
| 272 | fs::create_dir(output.join("source")).unwrap(); |
| 273 | let (_proxy, proxied) = proxy(&output); |
| 274 | let observer = client(); |
| 275 | let prefix = SystemTime::now() |
| 276 | .duration_since(UNIX_EPOCH) |
| 277 | .unwrap() |
| 278 | .as_nanos(); |
| 279 | let fixtures = [ |
| 280 | ( |
| 281 | "chunked", |
| 282 | onestore::create_section("fault.one", "Before café 🦀", "Fault test").unwrap(), |
| 283 | ), |
| 284 | ( |
| 285 | "carry", |
| 286 | fs::read( |
| 287 | Path::new(env!("CARGO_MANIFEST_DIR")) |
| 288 | .join("../../corpus/append/round-01/tx-255/notebook/synthetic.one"), |
| 289 | ) |
| 290 | .unwrap(), |
| 291 | ), |
| 292 | ]; |
| 293 | let mut results = Vec::new(); |
| 294 | for (fixture, source) in fixtures { |
| 295 | fs::write( |
| 296 | output.join("source").join(format!("{fixture}.one")), |
| 297 | &source, |
| 298 | ) |
| 299 | .unwrap(); |
| 300 | let (_, _, before) = text(&source); |
| 301 | let after = if fixture == "chunked" { |
| 302 | "After café 🦀 ".repeat(6000) |
| 303 | } else { |
| 304 | "After café 🦀".to_owned() |
| 305 | }; |
| 306 | let suffix = " [reconnected]"; |
| 307 | fs::write( |
| 308 | output.join(format!("{fixture}-intent.json")), |
| 309 | serde_json::to_vec(&json!({"before":before,"after":after,"suffix":suffix})).unwrap(), |
| 310 | ) |
| 311 | .unwrap(); |
| 312 | let old = snapshot(&source); |
| 313 | let deadline = Instant::now() + Duration::from_secs(2); |
| 314 | while old.modified.values().any(|value| *value >= stamp()) { |
| 315 | assert!( |
| 316 | Instant::now() < deadline, |
| 317 | "source timestamps did not precede the test" |
| 318 | ); |
| 319 | std::thread::sleep(Duration::from_millis(5)); |
| 320 | } |
| 321 | let everything = 0..before.encode_utf16().count() as u32; |
| 322 | let new = snapshot(&applied( |
| 323 | &source, |
| 324 | &replaced(&source, everything.clone(), &after), |
| 325 | )); |
| 326 | let baseline_path = format!("fault-{prefix}-{fixture}-baseline.one"); |
| 327 | create(&observer, &baseline_path, &source); |
| 328 | configure(&output, json!({"phase":format!("{fixture}-setup")})); |
| 329 | let baseline = Client::connect( |
| 330 | &proxied, |
| 331 | "agent", |
| 332 | Credentials::default(), |
| 333 | Duration::from_secs(5), |
| 334 | ) |
| 335 | .unwrap(); |
| 336 | let start = configure(&output, json!({"phase":format!("{fixture}-baseline")})); |
| 337 | let began = stamp(); |
| 338 | baseline |
| 339 | .commit_transaction( |
| 340 | &baseline_path, |
| 341 | &replaced(&source, everything.clone(), &after), |
| 342 | ) |
| 343 | .unwrap(); |
| 344 | let events = records(&output); |
| 345 | let mut occurrences = BTreeMap::new(); |
| 346 | let mut cuts = Vec::new(); |
| 347 | for event in &events[start..] { |
| 348 | let Some(direction) = event["direction"].as_str() else { |
| 349 | continue; |
| 350 | }; |
| 351 | if event["status"] == "0x103" { |
| 352 | continue; |
| 353 | } |
| 354 | let command = event["command"].as_u64().unwrap(); |
| 355 | let status = event["status"].as_str().map(str::to_owned); |
| 356 | let count = occurrences |
| 357 | .entry((direction.to_owned(), command, status.clone())) |
| 358 | .or_insert(0); |
| 359 | *count += 1; |
| 360 | cuts.push((direction.to_owned(), command, status, *count)); |
| 361 | } |
| 362 | assert!( |
| 363 | cuts.iter() |
| 364 | .any(|(direction, command, _, count)| direction == "response" |
| 365 | && *command == 7 |
| 366 | && *count == 3) |
| 367 | ); |
| 368 | assert!(published( |
| 369 | &snapshot(&observer.read(&baseline_path, 1 << 20).unwrap()), |
| 370 | &old, |
| 371 | &new, |
| 372 | began..=stamp() |
| 373 | )); |
| 374 | drop(baseline); |
| 375 | for (case, (direction, command, status, occurrence)) in cuts.iter().enumerate() { |
| 376 | let name = format!("{fixture}-{case:03}"); |
| 377 | let path = format!("fault-{prefix}-{name}.one"); |
| 378 | create(&observer, &path, &source); |
| 379 | configure(&output, json!({"phase":format!("{name}-setup")})); |
| 380 | let interrupted = Client::connect( |
| 381 | &proxied, |
| 382 | "agent", |
| 383 | Credentials::default(), |
| 384 | Duration::from_secs(5), |
| 385 | ) |
| 386 | .unwrap(); |
| 387 | let start = configure( |
| 388 | &output, |
| 389 | json!({"phase":name, "direction":direction, "cut":command, "status":status, "occurrence":occurrence}), |
| 390 | ); |
| 391 | let began = stamp(); |
| 392 | let transaction = replaced(&source, everything.clone(), &after); |
| 393 | let error = interrupted |
| 394 | .commit_transaction(&path, &transaction) |
| 395 | .unwrap_err(); |
| 396 | assert_eq!( |
| 397 | interrupted.read(&path, 1 << 20).unwrap_err().kind(), |
| 398 | io::ErrorKind::NotConnected |
| 399 | ); |
| 400 | let events = records(&output); |
| 401 | assert_eq!( |
| 402 | events[start..] |
| 403 | .iter() |
| 404 | .filter(|event| event.get("cut").is_some()) |
| 405 | .count(), |
| 406 | 1 |
| 407 | ); |
| 408 | assert!( |
| 409 | !events |
| 410 | .iter() |
| 411 | .any(|event| event.get("trace_error").is_some()) |
| 412 | ); |
| 413 | let fresh = client(); |
| 414 | let deadline = Instant::now() + Duration::from_secs(5); |
| 415 | loop { |
| 416 | match fresh |
| 417 | .open(&path, true) |
| 418 | .and_then(|file| file.coordinate(&path, true, &[])) |
| 419 | { |
| 420 | Ok(file) => { |
| 421 | file.close().unwrap(); |
| 422 | break; |
| 423 | } |
| 424 | Err(error) |
| 425 | if error.kind() == io::ErrorKind::WouldBlock |
| 426 | && Instant::now() < deadline => |
| 427 | { |
| 428 | std::thread::sleep(Duration::from_millis(5)) |
| 429 | } |
| 430 | Err(error) => panic!("{name}: retired session retained exclusion: {error}"), |
| 431 | } |
| 432 | } |
| 433 | let saved = fresh.read(&path, 1 << 20).unwrap(); |
| 434 | fs::write( |
| 435 | output.join("interrupted").join(format!("{name}.one")), |
| 436 | &saved, |
| 437 | ) |
| 438 | .unwrap(); |
| 439 | let visible = published(&snapshot(&saved), &old, &new, began..=stamp()); |
| 440 | match error.state { |
| 441 | CommitState::NotCommitted => assert!(!visible, "{name}"), |
| 442 | CommitState::Committed => assert!(visible, "{name}"), |
| 443 | CommitState::Unknown => {} |
| 444 | } |
| 445 | if !visible { |
| 446 | fresh |
| 447 | .commit_transaction(&path, &replaced(&saved, everything.clone(), &after)) |
| 448 | .unwrap(); |
| 449 | } |
| 450 | let saved = fresh.read(&path, 1 << 20).unwrap(); |
| 451 | assert!( |
| 452 | published(&snapshot(&saved), &old, &new, began..=stamp()), |
| 453 | "{name}: replay did not publish" |
| 454 | ); |
| 455 | let end = after.encode_utf16().count() as u32; |
| 456 | fresh |
| 457 | .commit_transaction(&path, &replaced(&saved, end..end, suffix)) |
| 458 | .unwrap(); |
| 459 | let recovered = fresh.read(&path, 1 << 20).unwrap(); |
| 460 | assert_eq!(text(&recovered).2, format!("{after}{suffix}")); |
| 461 | fs::write( |
| 462 | output.join("recovered").join(format!("{name}.one")), |
| 463 | recovered, |
| 464 | ) |
| 465 | .unwrap(); |
| 466 | results.push(json!({"case":name,"path":path,"direction":direction,"command":command,"status":status,"occurrence":occurrence, |
| 467 | "state":format!("{:?}",error.state),"visible":if visible {"after"} else {"before"}, |
| 468 | "retired_session_rejected":true,"exclusion_released":true,"fresh_commit_succeeded":true})); |
| 469 | fs::write( |
| 470 | output.join("results.json"), |
| 471 | serde_json::to_vec_pretty(&results).unwrap(), |
| 472 | ) |
| 473 | .unwrap(); |
| 474 | } |
| 475 | } |
| 476 | for state in ["NotCommitted", "Unknown", "Committed"] { |
| 477 | assert!(results.iter().any(|result| result["state"] == state)); |
| 478 | } |
| 479 | assert!( |
| 480 | results |
| 481 | .iter() |
| 482 | .any(|result| result["state"] == "Unknown" && result["visible"] == "before") |
| 483 | ); |
| 484 | assert!( |
| 485 | results |
| 486 | .iter() |
| 487 | .any(|result| result["state"] == "Unknown" && result["visible"] == "after") |
| 488 | ); |
| 489 | println!("{} message-loss cases passed", results.len()); |
| 490 | } |