Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/record_identity.rs

18.2 KiB, 1 run

created by r1870400018:61841, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! A collection moves the records it keeps, and with records of one size a new offset can equal an
2//! old one that something still refers to. These checks read and supersede records around a
3//! collection and judge every answer by the key and version the test wrote. Found by QA of lane
4//! sto, 2026-09-23 (F1, F2). `test::hooks` holds collections open and delays a file bot, and the
5//! hooks are process-wide, which is why this is a test binary of its own with a single test.
6
7use oxedyne_fe2o3_core::prelude::*;
8use oxedyne_fe2o3_hash::{
9 csum::ChecksumScheme,
10 hash::HashScheme,
11};
12use oxedyne_fe2o3_data::time::Timestamp;
13use oxedyne_fe2o3_iop_db::api::{
14 Database,
15 Meta,
16};
17use oxedyne_fe2o3_jdat::prelude::*;
18use oxedyne_fe2o3_o3db_sync::{
19 O3db,
20 base::{
21 cfg::OzoneConfig,
22 constant,
23 },
24 comm::response::Wait,
25 data::{
26 cache::Cache,
27 core::RestSchemesInput,
28 },
29 file::{
30 core::FileType,
31 floc::FileLocation,
32 zdir::ZoneDir,
33 },
34 test::{
35 hooks,
36 setup::{
37 self,
38 Uid,
39 UID_LEN,
40 },
41 },
42};
43
44use std::{
45 collections::BTreeMap,
46 path::{
47 Path,
48 PathBuf,
49 },
50 thread,
51 time::{
52 Duration,
53 Instant,
54 },
55};
56
57type TestDb = O3db<
58 { UID_LEN },
59 Uid,
60 (),
61 HashScheme,
62 HashScheme,
63 ChecksumScheme,
64>;
65
66const VALUE_BYTES: usize = 100;
67const HOLD: Duration = Duration::from_millis(1_500); // a collection held open
68const FORWARD: Duration = Duration::from_secs(4); // a supersession held back
69const SETTLE: Duration = Duration::from_secs(10);
70
71#[test]
72fn main() -> Outcome<()> {
73 log_set_level!("warn");
74 // Whatever failed, the next check and the next binary must not inherit a slow bot.
75 let queued = read_queued_behind_a_collection_returns_its_own_record();
76 hooks::set_collect_delay(Duration::ZERO);
77 let pending = new_offset_equal_to_a_pending_old_one_is_not_remapped();
78 hooks::set_collect_delay(Duration::ZERO);
79 hooks::set_forward_delay(Duration::ZERO);
80 let carried = read_of_a_carried_record_leaves_its_move_for_its_supersession();
81 hooks::set_collect_delay(Duration::ZERO);
82 hooks::set_forward_delay(Duration::ZERO);
83 log_finish_wait!();
84 let failed: Vec<Error<ErrTag>> = [queued, pending, carried].into_iter()
85 .filter_map(|r| r.err())
86 .collect();
87 match failed.len() {
88 0 => Ok(()),
89 1 => match failed.into_iter().next() {
90 Some(e) => Err(e),
91 None => Ok(()), // unreachable
92 },
93 n => Err(err!(
94 "{} of the 3 checks failed: {:?}", n, failed;
95 Test)),
96 }
97}
98
99/// Data file 1 holds a to e, a and b are superseded, and a read of c arrives while the collection
100/// that follows is held open. The read waits for the collection and then used the offset c had
101/// before it, which in the rewritten file is where e now starts, and e's record passes its own
102/// checksum: the read returned e's value for c.
103fn read_queued_behind_a_collection_returns_its_own_record() -> Outcome<()> {
104 let len = res!(record_len());
105 let cfg = res!(config(1, 5 * len + len / 2));
106 let root = res!(fresh("./test_db_record_identity_queued"));
107 let db = res!(open(&root, cfg));
108 res!(fill_file_1(&db, &root, len));
109
110 res!(db.insert(key(0), value(0, 2), Uid::default(), None));
111 hooks::set_collect_delay(HOLD);
112 // Two of the five records superseded, past the collection trigger.
113 res!(db.insert(key(1), value(1, 2), Uid::default(), None));
114 res!(wait_for_collection(&db, 1));
115 // The read goes to the file, and reaches its file bot while the collection is held.
116 res!(clear_values(&db));
117 let got = db.get(&key(2), None);
118 hooks::set_collect_delay(Duration::ZERO);
119
120 let result = judge(got, 2, 1, "A read that waited for a collection of its file");
121 res!(db.close());
122 result
123}
124
125/// With a file bot for odd files and one for even, file 1's collection is held open while c is
126/// superseded by a write to file 2, whose file bot is slow to pass the supersession on. So the
127/// collection carries c, and leaves a move entry at c's old offset for the supersession to find
128/// when it arrives. In the rewritten file e starts at that same offset. A read of e was matched
129/// to c's move entry by offset alone and returned c's old value; the supersession, arriving to
130/// find its entry gone, then flagged e's record as old, and the next collection dropped it.
131fn new_offset_equal_to_a_pending_old_one_is_not_remapped() -> Outcome<()> {
132 let len = res!(record_len());
133 let cfg = res!(config(2, 5 * len + len / 2));
134 let root = res!(fresh("./test_db_record_identity_pending"));
135 let db = res!(open(&root, cfg.clone()));
136 res!(fill_file_1(&db, &root, len));
137 let f1 = data_file(&root, 1);
138
139 res!(db.insert(key(0), value(0, 2), Uid::default(), None));
140 hooks::set_collect_delay(HOLD);
141 res!(db.insert(key(1), value(1, 2), Uid::default(), None));
142 res!(wait_for_collection(&db, 1));
143 // Superseded before the collection updates the caches, and passed on after it has finished.
144 hooks::set_forward_delay(FORWARD);
145 let forwarded = Instant::now();
146 res!(db.insert(key(2), value(2, 2), Uid::default(), None));
147 hooks::set_collect_delay(Duration::ZERO);
148 // c, d and e are carried: c to where a was, e to where c was.
149 if !settle_to(&f1, 3 * len) {
150 let _ = db.close();
151 return Err(err!(
152 "Data file 1 did not settle at the three records its collection carries, {} bytes, \
153 within {:?}: it is {} bytes.", 3 * len, SETTLE, size(&f1);
154 Test, Timeout));
155 }
156 res!(clear_values(&db));
157 let early = judge(db.get(&key(4), None), 4, 1,
158 "A read of e at the offset c had before the collection");
159 if forwarded.elapsed() >= FORWARD {
160 let _ = db.close();
161 return Err(err!(
162 "The read of e came {:?} after c's supersession was held back for {:?}, so the check \
163 did not test what it says.", forwarded.elapsed(), FORWARD;
164 Test, Timeout));
165 }
166 // The supersession of c arrives, and is applied to its own record.
167 thread::sleep(FORWARD.saturating_sub(forwarded.elapsed()) + Duration::from_millis(500));
168 hooks::set_forward_delay(Duration::ZERO);
169
170 // d superseded too: a second collection of file 1, which carries e alone.
171 res!(db.insert(key(3), value(3, 2), Uid::default(), None));
172 let _ = settle_to(&f1, len); // judged by the reads below
173 res!(clear_values(&db));
174 let late = judge(db.get(&key(4), None), 4, 1, "A read of e after a second collection");
175 res!(db.close());
176
177 let db = res!(open_again(&root, cfg));
178 let mut after = Vec::new();
179 for (i, ver) in [(0, 2), (1, 2), (2, 2), (3, 2), (4, 1), (5, 1)] {
180 after.push(judge(db.get(&key(i), None), i, ver, "A read after a restart"));
181 }
182 res!(db.close());
183 res!(early);
184 res!(late);
185 for result in after {
186 res!(result);
187 }
188 Ok(())
189}
190
191/// As in the check before, file 1's collection is held open while c is superseded by a write to
192/// file 2 whose supersession is held back, so the collection carries c and leaves its move entry
193/// for that supersession. This time a read of c, given c's old offset before the write, waits for
194/// the collection and is remapped through c's entry. A read that spent the entry left the
195/// supersession to find it gone, and the supersession then flagged the record at c's old offset,
196/// which is e's now, as old: the next collection dropped e, and kept the c it had been sent for.
197fn read_of_a_carried_record_leaves_its_move_for_its_supersession() -> Outcome<()> {
198 let len = res!(record_len());
199 let cfg = res!(config(2, 5 * len + len / 2));
200 let root = res!(fresh("./test_db_record_identity_carried"));
201 let db = res!(open(&root, cfg.clone()));
202 res!(fill_file_1(&db, &root, len));
203 let f1 = data_file(&root, 1);
204
205 res!(db.insert(key(0), value(0, 2), Uid::default(), None));
206 hooks::set_collect_delay(HOLD);
207 res!(db.insert(key(1), value(1, 2), Uid::default(), None));
208 res!(wait_for_collection(&db, 1));
209 // The read takes c's old offset from the cache bot and waits at the file bot.
210 res!(clear_values(&db));
211 let reading = db.clone();
212 let reader = res!(thread::Builder::new().name(fmt!("record identity reader")).spawn(
213 move || reading.get(&key(2), None)));
214 thread::sleep(Duration::from_millis(100));
215 // Superseded before the collection updates the caches, and passed on after it has finished.
216 hooks::set_forward_delay(FORWARD);
217 let forwarded = Instant::now();
218 res!(db.insert(key(2), value(2, 2), Uid::default(), None));
219 hooks::set_collect_delay(Duration::ZERO);
220 let early = match reader.join() {
221 // Asked before the write, so the version it was given, whichever that was.
222 Ok(got) => match judge(got.clone(), 2, 1, "A read of c queued behind the collection") {
223 Ok(()) => Ok(()),
224 Err(_) => judge(got, 2, 2, "A read of c queued behind the collection"),
225 },
226 Err(_) => Err(err!("The reading thread panicked."; Test, Thread)),
227 };
228 if !settle_to(&f1, 3 * len) {
229 let _ = db.close();
230 return Err(err!(
231 "Data file 1 did not settle at the three records its collection carries, {} bytes, \
232 within {:?}: it is {} bytes.", 3 * len, SETTLE, size(&f1);
233 Test, Timeout));
234 }
235 if forwarded.elapsed() >= FORWARD {
236 let _ = db.close();
237 return Err(err!(
238 "The collection ended {:?} after c's supersession was held back for {:?}, so the \
239 check did not test what it says.", forwarded.elapsed(), FORWARD;
240 Test, Timeout));
241 }
242 // The supersession of c arrives, and is applied to its own record.
243 thread::sleep(FORWARD.saturating_sub(forwarded.elapsed()) + Duration::from_millis(500));
244 hooks::set_forward_delay(Duration::ZERO);
245
246 // d superseded too: a second collection of file 1, which carries e alone.
247 res!(db.insert(key(3), value(3, 2), Uid::default(), None));
248 let _ = settle_to(&f1, len); // judged by the reads below
249 res!(clear_values(&db));
250 let late = judge(db.get(&key(4), None), 4, 1, "A read of e after a second collection");
251 res!(db.close());
252
253 let db = res!(open_again(&root, cfg));
254 let mut after = Vec::new();
255 for (i, ver) in [(0, 2), (1, 2), (2, 2), (3, 2), (4, 1), (5, 1)] {
256 after.push(judge(db.get(&key(i), None), i, ver, "A read after a restart"));
257 }
258 res!(db.close());
259 res!(early);
260 res!(late);
261 for result in after {
262 res!(result);
263 }
264 Ok(())
265}
266
267/// A key with two records in one collected file, the older carried because its supersession has
268/// yet to arrive, and the cache naming the newer. The collection's update for the older record
269/// leaves the cached location alone, and the newer record's update moves it and gives back where
270/// it was, for its move entry to be spent. Matched by file alone, the older record's update
271/// moved the location to the older record's new offset, and the entries were spent against the
272/// wrong offsets.
273#[test]
274fn reanchor_moves_only_the_record_the_cache_names() -> Outcome<()> {
275 let mut cache = Cache::<{ UID_LEN }, Uid>::new(None);
276 let k = res!(key(9).as_bytes());
277 let older = Meta { time: res!(Timestamp::now()), user: Uid::default() };
278 thread::sleep(Duration::from_millis(2));
279 let newer = Meta { time: res!(Timestamp::now()), user: Uid::default() };
280 let at = |start: u64| FileLocation { fnum: 1, start, klen: 30, vlen: 100 };
281 res!(cache.insert(k.clone(), None, at(260), older.clone()));
282 res!(cache.insert(k.clone(), None, at(390), newer.clone()));
283 // The collection carries both, the older to 0 and the newer to 130.
284 if let Some(was) = cache.reanchor(&k, &at(0), &older) {
285 return Err(err!(
286 "The update for the older of a key's two records in a collected file moved the \
287 location the cache holds for the newer, from {:?} to offset 0.", was;
288 Test, Mismatch));
289 }
290 match cache.reanchor(&k, &at(130), &newer) {
291 Some(was) if was.start == 390 => Ok(()),
292 other => Err(err!(
293 "The update for the newer of a key's two records in a collected file gave back {:?} \
294 as where the cache had it, where that was offset 390.", other;
295 Test, Mismatch)),
296 }
297}
298
299/// Writes a to e, which fill data file 1, and f, which seals it.
300fn fill_file_1(db: &TestDb, root: &Path, len: u64) -> Outcome<()> {
301 for i in 0..6 {
302 res!(db.insert(key(i), value(i, 1), Uid::default(), None));
303 }
304 let (f1, f2) = (data_file(root, 1), data_file(root, 2));
305 if !f2.is_file() || size(&f1) != 5 * len {
306 return Err(err!(
307 "Data file 1 should hold five records of {} bytes and be sealed by the sixth, but it \
308 is {} bytes and file 2 {}.", len, size(&f1),
309 if f2.is_file() { "exists" } else { "does not exist" };
310 Test, Invalid, Size));
311 }
312 Ok(())
313}
314
315/// The value read must be the one written for this key at this version.
316fn judge(
317 got: Outcome<Option<(Dat, Meta<{ UID_LEN }, Uid>)>>,
318 i: usize,
319 ver: u8,
320 what: &str,
321)
322 -> Outcome<()>
323{
324 match got {
325 Ok(Some((v, _))) => match parse(&v) {
326 Some((j, w)) if j == i && w == ver => Ok(()),
327 Some((j, w)) => Err(err!(
328 "{}: key {} returned the value written for key {} at version {}, where it should \
329 be version {} of its own.", what, i, j, w, ver;
330 Test, Mismatch)),
331 None => Err(err!(
332 "{}: key {} returned a value the test never wrote: {:?}.", what, i, v;
333 Test, Mismatch)),
334 },
335 Ok(None) => Err(err!(
336 "{}: key {} is missing, where it should be version {}.", what, i, ver;
337 Test, Missing)),
338 Err(e) => Err(err!(e,
339 "{}: key {} could not be read.", what, i;
340 Test, Read)),
341 }
342}
343
344/// The length of one record, every record here being that length.
345fn record_len() -> Outcome<u64> {
346 let root = res!(fresh("./test_db_record_identity_probe"));
347 let db = res!(open(&root, res!(config(1, 2_000))));
348 res!(db.insert(key(0), value(0, 1), Uid::default(), None));
349 res!(db.close());
350 let len = size(&data_file(&root, 1));
351 if len == 0 {
352 return Err(err!("The probe record left data file 1 empty."; Test, Missing));
353 }
354 Ok(len)
355}
356
357/// One zone, one writer, one cache bot and one reader, so that every step lands where the check
358/// says, and each record synced as it is written.
359fn config(nf: u16, max: u64) -> Outcome<OzoneConfig> {
360 let mut cfg = res!(setup::default_cfg());
361 cfg.num_zones = 1;
362 cfg.num_cbots_per_zone = 1;
363 cfg.num_fbots_per_zone = nf;
364 cfg.num_igbots_per_zone = 1;
365 cfg.num_rbots_per_zone = 1;
366 cfg.num_wbots_per_zone = 1;
367 cfg.data_file_max_bytes = max;
368 cfg.rest_chunk_threshold = max * 7 / 10;
369 cfg.zone_overrides = BTreeMap::new();
370 cfg.sync_on_write = true;
371 Ok(cfg)
372}
373
374fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> {
375 RestSchemesInput::new(
376 None::<()>,
377 None::<HashScheme>,
378 None::<HashScheme>,
379 Some(ChecksumScheme::new_crc32()),
380 )
381}
382
383fn open(root: &Path, cfg: OzoneConfig) -> Outcome<TestDb> {
384 setup::start_db(root.to_path_buf(), Some(cfg), schemes(), None, true, true)
385}
386
387fn open_again(root: &Path, cfg: OzoneConfig) -> Outcome<TestDb> {
388 setup::start_db(root.to_path_buf(), Some(cfg), schemes(), None, true, false)
389}
390
391/// Every key has the same length, and so every record.
392fn key(i: usize) -> Dat {
393 dat!(fmt!("record identity key {}", i))
394}
395
396/// A value that says which key and version it was written for.
397fn value(i: usize, ver: u8) -> Dat {
398 let mut v = vec![0u8; VALUE_BYTES];
399 v[0] = i as u8;
400 v[1] = ver;
401 for j in 2..VALUE_BYTES {
402 v[j] = (i as u8) ^ ver ^ (j as u8);
403 }
404 Dat::BU32(v)
405}
406
407fn parse(d: &Dat) -> Option<(usize, u8)> {
408 let v = d.bytes_ref()?;
409 if v.len() != VALUE_BYTES {
410 return None;
411 }
412 let (i, ver) = (v[0], v[1]);
413 for j in 2..VALUE_BYTES {
414 if v[j] != i ^ ver ^ (j as u8) {
415 return None;
416 }
417 }
418 Some((i as usize, ver))
419}
420
421fn clear_values(db: &TestDb) -> Outcome<()> {
422 db.api().clear_cache_values(Wait {
423 max_wait: constant::USER_REQUEST_TIMEOUT,
424 check_interval: constant::CHECK_INTERVAL,
425 })
426}
427
428/// Waits until the file bot has handed file `fnum` to a collector.
429fn wait_for_collection(db: &TestDb, fnum: u32) -> Outcome<()> {
430 let begun = Instant::now();
431 while begun.elapsed() < SETTLE {
432 let states = res!(db.api().collect_file_states(Wait {
433 max_wait: constant::USER_REQUEST_TIMEOUT,
434 check_interval: constant::CHECK_INTERVAL,
435 }));
436 for (_, shard) in &states {
437 if let Some(fstat) = shard.map().get(&fnum) {
438 if fstat.gc_active() {
439 return Ok(());
440 }
441 }
442 }
443 thread::sleep(Duration::from_millis(10));
444 }
445 Err(err!(
446 "Superseding two of the five records in data file {} did not start its collection \
447 within {:?}.", fnum, SETTLE;
448 Test, Timeout))
449}
450
451/// Does the file reach the given size before the wait runs out?
452fn settle_to(path: &Path, want: u64) -> bool {
453 let begun = Instant::now();
454 loop {
455 if size(path) == want {
456 return true;
457 }
458 if begun.elapsed() >= SETTLE {
459 return false;
460 }
461 thread::sleep(Duration::from_millis(5));
462 }
463}
464
465fn data_file(root: &Path, fnum: u32) -> PathBuf {
466 let zones = match config(1, 2_000) {
467 Ok(cfg) => cfg.zone_root(root),
468 Err(_) => root.to_path_buf(), // not reached: the fixture was built from the same config
469 };
470 zones.join("zone_001").join(ZoneDir::relative_file_path(&FileType::Data, fnum))
471}
472
473fn size(path: &Path) -> u64 {
474 match std::fs::metadata(path) {
475 Ok(m) => m.len(),
476 Err(_) => 0,
477 }
478}
479
480/// An empty directory of this test's own.
481fn fresh(dir: &str) -> Outcome<PathBuf> {
482 let _ = std::fs::remove_dir_all(dir); // absent the first time
483 res!(std::fs::create_dir_all(dir));
484 Ok(res!(Path::new(dir).canonicalize()))
485}