Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/rollover.rs

19.3 KiB, 80 runs

created by r1870400018:18290, 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//! Ozone data file rollover and supersession accounting test.
2//!
3//! A zone seals its live data file when the next record would push it past
4//! `data_file_max_bytes`, and the writer bot asks its zone bot for the next
5//! file number. Every record written into a file is recorded in that file's
6//! `FileState.dmap`, and when the record is superseded by a later write of the
7//! same key the entry is flipped from `Cur` to `Old`, which is what tells the
8//! garbage collector how much of the file is reclaimable.
9//!
10//! This test drives many rollovers with a tiny file size limit and repeatedly
11//! overwrites a small set of keys, so that supersessions land in files that
12//! have already been sealed. It then asserts the accounting invariant:
13//!
14//! * one `dmap` entry per record written,
15//! * exactly one `Cur` entry per distinct key,
16//! * every other entry `Old`,
17//! * an old-record counter that agrees with the map,
18//! * not one error logged by any bot.
19//!
20//! When the zone hands a writer a file number that is already in use, the
21//! receiving file bot replaces that file's `FileState` wholesale and the
22//! entries recorded so far are lost. Every later supersession of one of those
23//! records then fails in `FileState::register_old` with "a data entry starting
24//! at position N in the FileState was not found", garbage is never registered,
25//! and the store grows without bound. The invariant above catches that.
26//!
27//! The remaining phases cover a restart, which takes over the incomplete live
28//! file left behind and rebuilds the caches from the index files, and garbage
29//! collection, measured against the same churn run with collection disabled.
30//!
31//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
32//! Anthropic Claude
33
34use oxedyne_fe2o3_core::{
35 prelude::*,
36 alt::Override,
37};
38use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
39use oxedyne_fe2o3_hash::{
40 csum::ChecksumScheme,
41 hash::HashScheme,
42};
43use oxedyne_fe2o3_iop_db::api::{
44 Database,
45 RestSchemesOverride,
46};
47use oxedyne_fe2o3_jdat::prelude::*;
48use oxedyne_fe2o3_o3db_sync::{
49 base::{
50 cfg::OzoneConfig,
51 constant,
52 index::WorkerInd,
53 },
54 data::core::RestSchemesInput,
55 file::state::{
56 DataState,
57 FileStateMap,
58 },
59 test::setup,
60 O3db,
61};
62
63use std::{
64 collections::BTreeMap,
65 fs,
66 path::{
67 Path,
68 PathBuf,
69 },
70 thread,
71 time::Duration,
72};
73
74const NKEYS: usize = 40; // distinct keys churned through the rollover
75const NOVER: usize = 12; // overwrites of each key after its first write
76const GC_ROUNDS: usize = 40; // overwrite rounds in the collection phase
77
78pub fn test_rollover(_filter: &'static str) -> Outcome<()> {
79
80 let db_root = res!(canonical_dir("./test_db_rollover"));
81
82 let enckey = [0x7bu8; 32];
83 let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..]));
84 let crc32 = ChecksumScheme::new_crc32();
85 let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> =
86 RestSchemesOverride::default()
87 .set_encrypter(Override::Default(aes_gcm.clone()));
88 let schms2 = Some(&schms2);
89 let user = setup::Uid::default();
90
91 let schms_input = RestSchemesInput::new(
92 Some(aes_gcm.clone()),
93 None::<HashScheme>,
94 None::<HashScheme>,
95 Some(crc32.clone()),
96 );
97
98 // One writer bot per zone is the simple case: a rollover can only ever
99 // collide with the writer's own live file.
100 res!(supersession_accounting(
101 "one writer per zone",
102 &db_root,
103 &schms_input,
104 schms2,
105 user,
106 1,
107 ));
108 // Two writer bots per zone is the shipped default, and is the case where a
109 // reused file number can be handed to a writer other than the one that
110 // already holds it.
111 res!(supersession_accounting(
112 "two writers per zone",
113 &db_root,
114 &schms_input,
115 schms2,
116 user,
117 2,
118 ));
119
120 // A restart takes over the incomplete live file left behind, which is the
121 // other way the zone counter could be left at a number already in use.
122 res!(restart_keeps_file_numbers_unique(
123 &db_root,
124 &schms_input,
125 schms2,
126 user,
127 ));
128
129 res!(gc_reclaims_sealed_files(
130 &db_root,
131 &schms_input,
132 schms2,
133 user,
134 ));
135
136 test!(sync_log::stream(), "Rollover test passed.");
137 Ok(())
138}
139
140/// Tiny data files, so rollovers happen every few records, and garbage
141/// collection under the caller's control.
142fn rollover_cfg(nwbots: u16) -> Outcome<OzoneConfig> {
143 let mut cfg = res!(setup::default_cfg());
144 cfg.num_zones = 1;
145 cfg.num_cbots_per_zone = 2;
146 cfg.num_fbots_per_zone = 2;
147 cfg.num_igbots_per_zone = 2;
148 cfg.num_wbots_per_zone = nwbots;
149 cfg.data_file_max_bytes = 4_000;
150 cfg.rest_chunk_threshold = 3_000;
151 cfg.rest_chunk_bytes = 1_000;
152 cfg.zone_overrides = mapdat!{
153 1u16 => mapdat!{ "dir" => "", "max_size" => 100_000_000u64 },
154 }.get_map().unwrap();
155 Ok(cfg)
156}
157
158/// Writes `NKEYS` keys `NOVER + 1` times each through many rollovers, then
159/// checks that every record written is still accounted for in some file state,
160/// and that exactly one record per key is current.
161fn supersession_accounting(
162 label: &str,
163 db_root: &PathBuf,
164 schms_input: &RestSchemesInput<
165 EncryptionScheme,
166 HashScheme,
167 HashScheme,
168 ChecksumScheme,
169 >,
170 schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>,
171 user: setup::Uid,
172 nwbots: u16,
173)
174 -> Outcome<()>
175{
176 test!(sync_log::stream(), "+--- rollover: {} ---", label);
177
178 let cfg = res!(rollover_cfg(nwbots));
179 let db = res!(setup::start_db(
180 db_root.clone(),
181 Some(cfg.clone()),
182 schms_input.clone(),
183 None,
184 false, // gc off: nothing is retired, so the accounting is exact.
185 true, // wipe: start from an empty zone directory.
186 ));
187
188 thread::sleep(Duration::from_millis(200));
189
190 let mut nwrites = 0usize;
191 for round in 0..(NOVER + 1) {
192 for k in 0..NKEYS {
193 res!(db.insert(
194 dat!(fmt!("rk{:03}", k)),
195 dat!(fmt!("rv{:03}r{:02}", k, round)),
196 user,
197 schms2,
198 ));
199 nwrites += 1;
200 }
201 }
202
203 // Let the trailing UpdateData and ScheduleOld messages drain.
204 thread::sleep(Duration::from_millis(500));
205
206 // Read every key back, to confirm the rollover left the data readable.
207 for k in 0..NKEYS {
208 let key = dat!(fmt!("rk{:03}", k));
209 match res!(db.get(&key, schms2)) {
210 Some((val, _meta)) => {
211 let expected = dat!(fmt!("rv{:03}r{:02}", k, NOVER));
212 if val != expected {
213 return Err(err!(
214 "{}: key {} round-tripped as {:?}, expected {:?}.",
215 label, k, val, expected;
216 Test, Mismatch));
217 }
218 },
219 None => return Err(err!(
220 "{}: key {} not found after {} writes.", label, k, nwrites;
221 Test, Missing)),
222 }
223 }
224
225 let states = res!(db.api().collect_file_states(constant::USER_REQUEST_WAIT));
226 let (nfiles, ncur, nold, noldcnt) = count_data_states(&states);
227
228 test!(sync_log::stream(),
229 "{}: {} writes over {} tracked files: {} current, {} old, {} tracked \
230 total, {} counted old.",
231 label, nwrites, nfiles, ncur, nold, ncur + nold, noldcnt);
232
233 if ncur + nold != nwrites {
234 return Err(err!(
235 "{}: {} records were written but the file states track {} \
236 ({} current + {} old). Entries lost from a FileState can never be \
237 flagged old, so their bytes can never be reclaimed.",
238 label, nwrites, ncur + nold, ncur, nold;
239 Test, Mismatch, Data));
240 }
241 if ncur != NKEYS {
242 return Err(err!(
243 "{}: {} distinct keys were written but {} records are flagged \
244 current.", label, NKEYS, ncur;
245 Test, Mismatch, Data));
246 }
247 if noldcnt != nold {
248 return Err(err!(
249 "{}: {} records are flagged old in the record maps but the \
250 old-record counters total {}. The two disagree, so the \
251 reclaimable byte total the collector trusts is wrong.",
252 label, nold, noldcnt;
253 Test, Mismatch, Data));
254 }
255 res!(assert_no_bot_errors(&db, label));
256
257 res!(db.shutdown());
258 thread::sleep(Duration::from_millis(200));
259
260 test!(sync_log::stream(), "+--- rollover: {} : passed ---", label);
261 Ok(())
262}
263
264/// Churns, shuts the database down, restarts it over the files left behind and
265/// churns again. On restart a writer bot takes over the incomplete live file,
266/// so the zone live file counter must be left above it: if the counter is
267/// wound back below a file already in use, the first rollover after the restart
268/// hands that number out again and the file's record entries are lost.
269fn restart_keeps_file_numbers_unique(
270 db_root: &PathBuf,
271 schms_input: &RestSchemesInput<
272 EncryptionScheme,
273 HashScheme,
274 HashScheme,
275 ChecksumScheme,
276 >,
277 schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>,
278 user: setup::Uid,
279)
280 -> Outcome<()>
281{
282 test!(sync_log::stream(), "+--- rollover: restart ---");
283
284 let cfg = res!(rollover_cfg(1));
285
286 // First run, from empty.
287 let db = res!(setup::start_db(
288 db_root.clone(),
289 Some(cfg.clone()),
290 schms_input.clone(),
291 None,
292 false, // gc off.
293 true, // wipe.
294 ));
295 thread::sleep(Duration::from_millis(200));
296 let n1 = res!(churn(&db, "sk", user, schms2, NOVER + 1));
297 thread::sleep(Duration::from_millis(500));
298 res!(assert_no_bot_errors(&db, "restart: first run"));
299 res!(db.shutdown());
300 thread::sleep(Duration::from_millis(500));
301
302 // Second run, over the files the first left behind.
303 let db = res!(setup::start_db(
304 db_root.clone(),
305 Some(cfg.clone()),
306 schms_input.clone(),
307 None,
308 false, // gc off.
309 false, // no wipe: this is the point of the phase.
310 ));
311 thread::sleep(Duration::from_millis(500));
312 res!(assert_no_bot_errors(&db, "restart: after reload"));
313 let n2 = res!(churn(&db, "sk", user, schms2, NOVER + 1));
314 thread::sleep(Duration::from_millis(500));
315
316 let states = res!(db.api().collect_file_states(constant::USER_REQUEST_WAIT));
317 let (nfiles, ncur, nold, noldcnt) = count_data_states(&states);
318
319 test!(sync_log::stream(),
320 "restart: {} writes before and {} after, {} tracked files: {} current, \
321 {} old, {} tracked total, {} counted old.",
322 n1, n2, nfiles, ncur, nold, ncur + nold, noldcnt);
323
324 if ncur + nold != n1 + n2 {
325 return Err(err!(
326 "restart: {} records were written across the two runs but the file \
327 states track {} ({} current + {} old).",
328 n1 + n2, ncur + nold, ncur, nold;
329 Test, Mismatch, Data));
330 }
331 if ncur != NKEYS {
332 return Err(err!(
333 "restart: {} distinct keys were written but {} records are flagged \
334 current.", NKEYS, ncur;
335 Test, Mismatch, Data));
336 }
337 if noldcnt != nold {
338 return Err(err!(
339 "restart: {} records are flagged old in the record maps but the \
340 old-record counters total {}.", nold, noldcnt;
341 Test, Mismatch, Data));
342 }
343 res!(assert_no_bot_errors(&db, "restart: second run"));
344
345 res!(db.shutdown());
346 thread::sleep(Duration::from_millis(200));
347
348 test!(sync_log::stream(), "+--- rollover: restart : passed ---");
349 Ok(())
350}
351
352fn churn<
353 ENC: oxedyne_fe2o3_iop_crypto::enc::Encrypter + 'static,
354 KH: oxedyne_fe2o3_iop_hash::api::Hasher + 'static,
355 PR: oxedyne_fe2o3_iop_hash::api::Hasher + 'static,
356 CS: oxedyne_fe2o3_iop_hash::csum::Checksummer + 'static,
357>(
358 db: &O3db<{ setup::UID_LEN }, setup::Uid, ENC, KH, PR, CS>,
359 prefix: &str,
360 user: setup::Uid,
361 schms2: Option<&RestSchemesOverride<ENC, KH>>,
362 rounds: usize,
363)
364 -> Outcome<usize>
365{
366 let mut nwrites = 0;
367 for round in 0..rounds {
368 for k in 0..NKEYS {
369 res!(db.insert(
370 dat!(fmt!("{}{:03}", prefix, k)),
371 dat!(fmt!("{}v{:03}r{:02}", prefix, k, round)),
372 user,
373 schms2,
374 ));
375 nwrites += 1;
376 }
377 }
378 Ok(nwrites)
379}
380
381/// Heavy overwrite churn over many sealed data files, run twice: once with
382/// garbage collection off and once with it on. The store with collection on
383/// must be a fraction of the size of the store without it. The run with
384/// collection off is the oracle: it is the same workload measured with the only
385/// mechanism that can reclaim bytes disabled, so the comparison cannot be
386/// satisfied by anything other than reclamation actually happening.
387fn gc_reclaims_sealed_files(
388 db_root: &PathBuf,
389 schms_input: &RestSchemesInput<
390 EncryptionScheme,
391 HashScheme,
392 HashScheme,
393 ChecksumScheme,
394 >,
395 schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>,
396 user: setup::Uid,
397)
398 -> Outcome<()>
399{
400 test!(sync_log::stream(), "+--- rollover: garbage collection ---");
401
402 let uncollected = res!(churn_and_measure(
403 db_root, schms_input, schms2, user, false));
404 let collected = res!(churn_and_measure(
405 db_root, schms_input, schms2, user, true));
406
407 test!(sync_log::stream(),
408 "gc: {} bytes of data files with collection off, {} with it on.",
409 uncollected, collected);
410
411 // The live set is a fortieth of everything written, so a working collector
412 // should leave well under a quarter of the uncollected store. Anything at
413 // or above it means sealed files are not being reclaimed.
414 if collected * 4 >= uncollected {
415 return Err(err!(
416 "Garbage collection reclaimed little or nothing: the same churn \
417 left {} bytes of data files with collection on against {} with it \
418 off.", collected, uncollected;
419 Test, Mismatch, Data));
420 }
421
422 test!(sync_log::stream(), "+--- rollover: garbage collection : passed ---");
423 Ok(())
424}
425
426/// The database is freshly wiped first, so the returned byte total covers only
427/// this run's churn.
428fn churn_and_measure(
429 db_root: &PathBuf,
430 schms_input: &RestSchemesInput<
431 EncryptionScheme,
432 HashScheme,
433 HashScheme,
434 ChecksumScheme,
435 >,
436 schms2: Option<&RestSchemesOverride<EncryptionScheme, HashScheme>>,
437 user: setup::Uid,
438 gc_on: bool,
439)
440 -> Outcome<u64>
441{
442 let label = if gc_on { "collection on" } else { "collection off" };
443 let cfg = res!(rollover_cfg(1));
444 let db = res!(setup::start_db(
445 db_root.clone(),
446 Some(cfg.clone()),
447 schms_input.clone(),
448 None,
449 gc_on,
450 true, // wipe.
451 ));
452
453 thread::sleep(Duration::from_millis(200));
454
455 // Enough churn that the superseded copies vastly outweigh the live set.
456 let mut nwrites = 0usize;
457 for round in 0..GC_ROUNDS {
458 for k in 0..NKEYS {
459 res!(db.insert(
460 dat!(fmt!("gk{:03}", k)),
461 dat!(fmt!("gv{:03}r{:02}", k, round)),
462 user,
463 schms2,
464 ));
465 nwrites += 1;
466 }
467 }
468
469 // Give the gbots time to finish transcribing.
470 thread::sleep(Duration::from_secs(2));
471
472 let bytes = res!(zone_data_bytes(&db_root));
473 let states = res!(db.api().collect_file_states(constant::USER_REQUEST_WAIT));
474 let (nfiles, ncur, nold, _) = count_data_states(&states);
475
476 test!(sync_log::stream(),
477 "gc {}: {} writes, {} bytes of data files remain over {} tracked \
478 files, {} current and {} old records tracked.",
479 label, nwrites, bytes, nfiles, ncur, nold);
480
481 if nold > nwrites - NKEYS {
482 return Err(err!(
483 "gc {}: more records are flagged old ({}) than were superseded ({}).",
484 label, nold, nwrites - NKEYS;
485 Test, Mismatch, Data));
486 }
487 res!(assert_no_bot_errors(&db, label));
488
489 res!(db.shutdown());
490 thread::sleep(Duration::from_millis(200));
491
492 Ok(bytes)
493}
494
495/// Fails if any bot in the database has logged an error. The flag-as-old
496/// failure this test exists for is logged by a file bot rather than returned
497/// to the caller, so this is how the log storm itself is asserted away.
498fn assert_no_bot_errors<
499 ENC: oxedyne_fe2o3_iop_crypto::enc::Encrypter + 'static,
500 KH: oxedyne_fe2o3_iop_hash::api::Hasher + 'static,
501 PR: oxedyne_fe2o3_iop_hash::api::Hasher + 'static,
502 CS: oxedyne_fe2o3_iop_hash::csum::Checksummer + 'static,
503>(
504 db: &O3db<{ setup::UID_LEN }, setup::Uid, ENC, KH, PR, CS>,
505 label: &str,
506)
507 -> Outcome<()>
508{
509 let (errs, nbots) = res!(db.api().bot_error_count(constant::USER_REQUEST_WAIT));
510 test!(sync_log::stream(), "{}: {} bots reported {} errors.", label, nbots, errs);
511 if errs > 0 {
512 return Err(err!(
513 "{}: the {} bots of the database logged {} errors during the run. \
514 A supersession that cannot be registered is logged, not returned, \
515 so any error here means the rollover bookkeeping is still wrong.",
516 label, nbots, errs;
517 Test, Mismatch, Data));
518 }
519 Ok(())
520}
521
522/// Across every file bot in every zone: the number of tracked files, the
523/// current and old record counts, and the old count the file states track.
524fn count_data_states(
525 states: &BTreeMap<WorkerInd, FileStateMap>,
526)
527 -> (usize, usize, usize, usize)
528{
529 let mut nfiles = 0;
530 let mut ncur = 0;
531 let mut nold = 0;
532 let mut noldcnt = 0;
533 for (_wind, fstates) in states {
534 for (_fnum, fstat) in fstates.map() {
535 nfiles += 1;
536 noldcnt += fstat.get_old_count();
537 for (_start, dstat) in fstat.data_map() {
538 match dstat {
539 DataState::Cur => ncur += 1,
540 DataState::Old => nold += 1,
541 }
542 }
543 }
544 }
545 (nfiles, ncur, nold, noldcnt)
546}
547
548fn zone_data_bytes(db_root: &Path) -> Outcome<u64> {
549 let mut total = 0u64;
550 res!(walk_data_files(db_root, &mut total));
551 Ok(total)
552}
553
554fn walk_data_files(dir: &Path, total: &mut u64) -> Outcome<()> {
555 for entry in res!(fs::read_dir(dir)) {
556 let entry = res!(entry);
557 let path = entry.path();
558 if path.is_dir() {
559 res!(walk_data_files(&path, total));
560 } else if path.extension().map(|e| e == constant::DATA_FILE_EXT).unwrap_or(false) {
561 let meta = res!(entry.metadata());
562 *total += meta.len();
563 }
564 }
565 Ok(())
566}
567
568/// Creates the directory if it does not exist.
569fn canonical_dir(p: &str) -> Outcome<PathBuf> {
570 match Path::new(p).canonicalize() {
571 Ok(path) => Ok(path),
572 Err(_) => {
573 res!(fs::create_dir_all(p));
574 match Path::new(p).canonicalize() {
575 Ok(path) => Ok(path),
576 Err(e) => Err(err!(e, "Cannot canonicalise {:?}.", p; IO, Path)),
577 }
578 },
579 }
580}