1use super::*;
2use serde_json::{Value, json};
3use std::{collections::BTreeMap, path::Path, process::Child};
4
5#[derive(Clone)]
6struct Snapshot {
7 content: BTreeMap<ExGuid, String>,
8 modified: BTreeMap<(ExGuid, ExGuid), u32>,
9}
10
11fn 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
68fn stamp() -> u32 {
69 (SystemTime::now()
70 .duration_since(UNIX_EPOCH)
71 .unwrap()
72 .as_secs()
73 - 315532800)
74 .try_into()
75 .unwrap()
76}
77
78fn 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]
113fn 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
137struct Proxy(Child);
138impl Drop for Proxy {
139 fn drop(&mut self) {
140 let _ = self.0.kill();
141 let _ = self.0.wait();
142 }
143}
144
145fn 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
154fn 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
174fn 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"]
204fn 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"]
267fn 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}