Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/gc_stale_floc.rs

37.6 KiB, 7 runs

created by r1870400018:38151, 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 key carried through a garbage collection must keep reading back, every time.
2//!
3//! # The fault
4//!
5//! Garbage collection transcribes a sealed data file to a temporary, then renames the temporary
6//! over it (`InitGarbageBot::collect_file`, step 7). The rename replaces the path, so every
7//! surviving record now sits at a new offset in a new inode, and the old inode is unlinked. Two
8//! caches have to follow that.
9//!
10//! The location cache does follow it: `cache_data_file` sends each carried record's new location
11//! to its cbot, which re-anchors the entry (`Cache::reanchor`), and any location still
12//! in flight is mapped through the file state's ephemeral old -> new move map.
13//!
14//! The reader's file cache does not. `ReaderBot::get_file` keeps an open `File` per file number
15//! for `constant::FILE_CACHE_EXPIRY_SECS` -- fifteen minutes -- and nothing invalidates it when a
16//! collection renames that file away. A reader that had read the file before the collection goes
17//! on reading the unlinked old inode, at the NEW offsets. Those offsets are not record boundaries
18//! in the old file, so the bytes are not a record, and the checksum verify in `ReaderBot::read`
19//! fails: `api.rs get_wait` -> `bot_reader.rs` -> `fe2o3_hash csum.rs` -> `[Checksum][Input
20//! Mismatch]`. That is the production symptom exactly, and it explains its shape -- key-specific
21//! (only keys in a collected file), intermittent (only readers that touched the file first), on a
22//! store that is byte-perfect on disk, cleared by a cold start, and re-made by the next collection.
23//!
24//! The `postgc` flag is the signal that would have caught it, and it is discarded: the fbot sets it
25//! when a location came through the move map, the rbot passes it up through `Value` and
26//! `Responder::recv_daticle`, and `api.rs get_wait` drops it on the floor.
27//!
28//! # What the tests do
29//!
30//! Each case writes one survivor key behind a few fillers -- behind, because a survivor at the head
31//! of its file is transcribed to the head of the new file and never moves -- then fills the zone's
32//! data files, supersedes enough fillers to push the survivor's file past the 30% old-data
33//! collection trigger, and reads the survivor back. Cached values are cleared first
34//! (`clear_cache_values`), so every read must go to the file: this is the production shape, where
35//! values are jettisoned under memory pressure or never loaded and only locations are held.
36//!
37//! - `no_compaction` -- the control: the same writes and churn with collection off. Nothing
38//! shrinks and every read is clean.
39//! - `some_survive` -- half the fillers superseded, so the compacted file keeps several records.
40//! - `only_survivor` -- every filler superseded, so the survivor is the one record carried over.
41//! Neither of these two reads the survivor until the collection has
42//! finished, so the reader opens the new inode and both pass: that is what
43//! localises the fault to a handle opened beforehand rather than to the
44//! location or to the file on disk.
45//! - `read_during_compaction` -- the repro. A reader hammers the survivor while the supersession
46//! burst drives collection, which is the shape a live gateway is in. The
47//! reads are clean until the rename, fail from the rename onwards, still
48//! fail once everything has settled, and are clean again after a restart.
49//!
50//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
51//! Anthropic Claude
52
53use oxedyne_fe2o3_core::{
54 prelude::*,
55 alt::Override,
56};
57use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
58use oxedyne_fe2o3_hash::{
59 csum::ChecksumScheme,
60 hash::HashScheme,
61};
62use oxedyne_fe2o3_iop_db::api::{
63 Database,
64 RestSchemesOverride,
65};
66use oxedyne_fe2o3_jdat::prelude::*;
67use oxedyne_fe2o3_o3db_sync::{
68 O3db,
69 base::constant,
70 comm::response::Wait,
71 data::core::RestSchemesInput,
72 prelude::{
73 Checksummer,
74 Encrypter,
75 Hasher,
76 },
77 test::setup,
78};
79
80use std::{
81 collections::BTreeMap,
82 fs,
83 path::{
84 Path,
85 PathBuf,
86 },
87 sync::{
88 Arc,
89 Mutex,
90 atomic::{
91 AtomicBool,
92 AtomicUsize,
93 Ordering,
94 },
95 },
96 thread,
97 time::{
98 Duration,
99 Instant,
100 },
101};
102
103const VALUE_BYTES: usize = 200; // well under the chunking threshold: one record per value
104const NPRE: usize = 5; // fillers written ahead of the survivor in its own data file
105const NFILL: usize = 60; // enough fillers to seal several data files behind the survivor
106const NREADS: usize = 6; // a first read plus five more, so a read-once fault is visible
107const GC_TIMEOUT: Duration = Duration::from_secs(30);
108
109/// Which supersession pattern the case drives, and hence whether a collection runs at all.
110#[derive(Clone, Copy, Debug)]
111enum Case {
112 /// Half the fillers superseded: the compacted file keeps several records beside the survivor.
113 SomeSurvive,
114 /// Every filler superseded: the survivor is the only record carried into the new file.
115 OnlySurvivor,
116 /// The control. Collection is off, so nothing moves and every read must be clean.
117 NoCompaction,
118}
119
120impl Case {
121 fn name(&self) -> &'static str {
122 match self {
123 Self::SomeSurvive => "some_survive",
124 Self::OnlySurvivor => "only_survivor",
125 Self::NoCompaction => "no_compaction",
126 }
127 }
128
129 fn gc_on(&self) -> bool {
130 match self {
131 Self::SomeSurvive |
132 Self::OnlySurvivor => true,
133 Self::NoCompaction => false,
134 }
135 }
136
137 /// Is the filler at this index superseded during the churn phase? Every filler written ahead
138 /// of the survivor is superseded whatever the case, because those are the records whose
139 /// removal shifts the survivor to a new offset: a survivor at the head of its file is carried
140 /// to the head of the transcribed file and never moves, which would prove nothing.
141 fn supersedes(&self, i: usize) -> bool {
142 if i < NPRE {
143 return true;
144 }
145 match self {
146 Self::SomeSurvive => i % 2 == 0,
147 Self::OnlySurvivor => true,
148 Self::NoCompaction => i % 2 == 0,
149 }
150 }
151}
152
153fn wait() -> Wait {
154 Wait {
155 max_wait: Duration::from_secs(30),
156 check_interval: constant::CHECK_INTERVAL,
157 }
158}
159
160/// A byte string of the given size, filled deterministically so each version is distinct on disk.
161fn value_of(seed: u8, len: usize) -> Dat {
162 let mut v = vec![0u8; len];
163 for (i, b) in v.iter_mut().enumerate() {
164 *b = seed.wrapping_add((i % 251) as u8);
165 }
166 Dat::BU32(v)
167}
168
169fn survivor_key() -> Dat { dat!("gcrace:survivor") }
170fn survivor_val() -> Dat { value_of(0x5a, VALUE_BYTES) }
171fn filler_key(i: usize) -> Dat { dat!(fmt!("gcrace:fill:{:04}", i)) }
172
173/// Every data file under `root`, with its current length.
174fn data_file_sizes(root: &Path) -> Outcome<BTreeMap<PathBuf, u64>> {
175 let mut found = BTreeMap::new();
176 let mut stack = vec![root.to_path_buf()];
177 while let Some(dir) = stack.pop() {
178 for entry in res!(fs::read_dir(&dir)) {
179 let path = res!(entry).path();
180 if path.is_dir() {
181 stack.push(path);
182 } else if path.extension().map(|e| e == constant::DATA_FILE_EXT).unwrap_or(false) {
183 let len = res!(fs::metadata(&path)).len();
184 found.insert(path, len);
185 }
186 }
187 }
188 Ok(found)
189}
190
191/// A compaction rewrites a sealed data file in place, so it shows up as a file present before the
192/// churn whose length has fallen. Polls until one is seen or the timeout expires.
193fn wait_for_compaction(
194 root: &Path,
195 before: &BTreeMap<PathBuf, u64>,
196 timeout: Duration,
197)
198 -> Outcome<Option<(PathBuf, u64, u64)>>
199{
200 let start = Instant::now();
201 loop {
202 let now = res!(data_file_sizes(root));
203 for (path, was) in before {
204 if let Some(is) = now.get(path) {
205 if is < was {
206 return Ok(Some((path.clone(), *was, *is)));
207 }
208 }
209 }
210 if start.elapsed() > timeout {
211 return Ok(None);
212 }
213 thread::sleep(Duration::from_millis(100));
214 }
215}
216
217/// Reads the survivor once, insisting on the exact bytes written. The read index is carried so
218/// that a failure says which read of the sequence broke: the reads before a collection are clean
219/// and the reads after it are not, which is what names the rename as the moment of damage.
220fn read_survivor<
221 ENC: Encrypter + 'static,
222 KH: Hasher + 'static,
223 PR: Hasher + 'static,
224 CS: Checksummer + 'static,
225>(
226 db: &O3db<{ setup::UID_LEN }, setup::Uid, ENC, KH, PR, CS>,
227 schms2: Option<&RestSchemesOverride<ENC, KH>>,
228 label: &str,
229 n: usize,
230)
231 -> Outcome<()>
232{
233 match db.get(&survivor_key(), schms2) {
234 Err(e) => Err(err!(e,
235 "[{}] Read {} of {} of the survivor key failed. A store whose records are all \
236 intact on disk can only fail a read this way by landing somewhere that is not a \
237 record boundary -- a location that was not re-anchored after the collection, or a \
238 file handle still open on the inode the collection renamed away.",
239 label, n, NREADS;
240 Test, Data)),
241 Ok(None) => Err(err!(
242 "[{}] Read {} of {} of the survivor key found nothing, but it was written and \
243 never deleted.", label, n, NREADS;
244 Test, Missing, Data)),
245 Ok(Some((got, _))) => if got == survivor_val() {
246 Ok(())
247 } else {
248 Err(err!(
249 "[{}] Read {} of {} of the survivor key returned different bytes from those \
250 written.", label, n, NREADS;
251 Test, Invalid, Data))
252 },
253 }
254}
255
256fn run_case(case: Case) -> Outcome<()> {
257
258 let label = case.name();
259 let dirname = fmt!("./test_db_gc_stale_floc_{}", label);
260 let db_root = res!(canonical_dir(&dirname));
261
262 let enckey = [0x5cu8; 32];
263 let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..]));
264 let crc32 = ChecksumScheme::new_crc32();
265 let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> =
266 RestSchemesOverride::default()
267 .set_encrypter(Override::Default(aes_gcm.clone()));
268 let schms2 = Some(&schms2);
269 let user = setup::Uid::default();
270
271 let schms_input = RestSchemesInput::new(
272 Some(aes_gcm.clone()),
273 None::<HashScheme>,
274 None::<HashScheme>,
275 Some(crc32.clone()),
276 );
277
278 // One zone and one bot of each kind, so the survivor's routing is fixed and a compaction is
279 // not spread over zones the test would have to reason about. The data file is small enough
280 // that the fillers seal several files behind the survivor, and the values stay well under the
281 // chunking threshold so every write takes the single-record path. Writes are synced so that
282 // the collector's "has this file drained?" check passes promptly rather than on a timer.
283 let mut cfg = res!(setup::default_cfg());
284 cfg.num_zones = 1;
285 cfg.num_cbots_per_zone = 1;
286 cfg.num_fbots_per_zone = 1;
287 cfg.num_igbots_per_zone = 1;
288 cfg.num_rbots_per_zone = 1;
289 cfg.num_wbots_per_zone = 1;
290 cfg.data_file_max_bytes = 4_000;
291 cfg.rest_chunk_threshold = 1_500;
292 cfg.rest_chunk_bytes = 64;
293 cfg.init_load_caches = true;
294 cfg.sync_on_write = true;
295 cfg.zone_overrides = DaticleMap::new();
296
297 test!(sync_log::stream(), "+--- gc stale floc: {} ---", label);
298
299 let db = res!(setup::start_db(
300 db_root.clone(),
301 Some(cfg.clone()),
302 schms_input.clone(),
303 None,
304 case.gc_on(),
305 true, // wipe
306 ));
307 thread::sleep(Duration::from_millis(500));
308
309 // 1. A few fillers go in ahead of the survivor, so the survivor sits part way into the zone's
310 // first data file rather than at its head. Superseding those puts the survivor at a
311 // different offset in the transcribed file, which is the whole subject of the test.
312 for i in 0..NPRE {
313 res!(db.insert(filler_key(i), value_of(i as u8, VALUE_BYTES), user, schms2));
314 }
315 res!(db.insert(survivor_key(), survivor_val(), user, schms2));
316
317 // 2. The rest of the fillers seal that file and several after it.
318 for i in NPRE..NFILL {
319 res!(db.insert(filler_key(i), value_of(i as u8, VALUE_BYTES), user, schms2));
320 }
321 thread::sleep(Duration::from_millis(500));
322
323 // 3. Clear cached values, so from here a read of any key must go to its file. This is the
324 // production shape: only locations are held, and a stale one cannot hide behind a cached
325 // copy of the value.
326 res!(db.api().clear_cache_values(wait()));
327
328 let before = res!(data_file_sizes(&db_root));
329 if before.len() < 3 {
330 let _ = db.shutdown();
331 return Err(err!(
332 "[{}] Only {} data files were written; the configuration should have sealed \
333 several behind the survivor, so this case would prove nothing.", label, before.len();
334 Test, Size));
335 }
336
337 // 4. Supersede the fillers. Each overwrite flags the previous record old in its sealed file,
338 // and once a file passes the 30% old-data trigger the collector transcribes it, carrying
339 // the survivor to a new offset.
340 for i in 0..NFILL {
341 if case.supersedes(i) {
342 res!(db.insert(filler_key(i), value_of((i + 128) as u8, VALUE_BYTES), user, schms2));
343 }
344 }
345
346 // 5. Wait for a compaction to actually land before reading anything.
347 let compacted = res!(wait_for_compaction(
348 &db_root,
349 &before,
350 if case.gc_on() { GC_TIMEOUT } else { Duration::from_secs(3) },
351 ));
352 match (case.gc_on(), &compacted) {
353 (true, None) => {
354 let _ = db.shutdown();
355 return Err(err!(
356 "[{}] No data file shrank within {:?}, so no compaction ran and this case \
357 would prove nothing about a location that survived one.", label, GC_TIMEOUT;
358 Test, Missing));
359 },
360 (true, Some((path, was, is))) => test!(sync_log::stream(),
361 "[{}] Compacted {:?}: {} -> {} bytes.", label, path, was, is),
362 (false, Some((path, was, is))) => {
363 let _ = db.shutdown();
364 return Err(err!(
365 "[{}] The control ran with collection off, but {:?} still shrank from {} to \
366 {} bytes.", label, path, was, is;
367 Test, Invalid));
368 },
369 (false, None) => test!(sync_log::stream(),
370 "[{}] Control: no data file shrank, as expected with collection off.", label),
371 }
372 thread::sleep(Duration::from_millis(500));
373
374 // Clear again, so the reads below cannot be answered from a value cached by the churn.
375 res!(db.api().clear_cache_values(wait()));
376
377 // 6. Read the survivor repeatedly. No read has touched this file before now, so the reader
378 // opens it fresh and these reads are clean -- which is the half of the evidence that says
379 // the location re-anchor works and the transcribed file is right.
380 let mut first_failure: Option<(usize, Error<ErrTag>)> = None;
381 for n in 1..=NREADS {
382 match read_survivor(&db, schms2, label, n) {
383 Ok(()) => test!(sync_log::stream(), "[{}] Read {} of {}: clean.", label, n, NREADS),
384 Err(e) => {
385 test!(sync_log::stream(), "[{}] Read {} of {}: FAILED.", label, n, NREADS);
386 if first_failure.is_none() {
387 first_failure = Some((n, e));
388 }
389 },
390 }
391 }
392
393 let _ = db.shutdown();
394 thread::sleep(Duration::from_millis(300));
395
396 match first_failure {
397 None => {
398 test!(sync_log::stream(),
399 "[{}] All {} reads of the survivor returned the bytes written.", label, NREADS);
400 Ok(())
401 },
402 Some((n, e)) => Err(err!(e,
403 "[{}] The survivor key stopped reading back at read {} of {}.", label, n, NREADS;
404 Test, Data)),
405 }
406}
407
408/// A reader hammering one key while a supersession burst drives collection. This is the shape a
409/// live gateway is in, and it is the one that matters: the reader opens the survivor's data file
410/// before the collection renames a transcribed copy over it, and from the rename onwards it is
411/// reading an unlinked inode at offsets that belong to the file that replaced it. The reads after
412/// everything settles say whether the damage lasts, and the reads after a restart say whether the
413/// store on disk was ever at fault.
414fn read_during_compaction() -> Outcome<()> {
415
416 let label = "read_during_compaction";
417 let db_root = res!(canonical_dir("./test_db_gc_stale_floc_concurrent"));
418
419 let enckey = [0x5cu8; 32];
420 let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..]));
421 let crc32 = ChecksumScheme::new_crc32();
422 let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> =
423 RestSchemesOverride::default()
424 .set_encrypter(Override::Default(aes_gcm.clone()));
425 let schms2 = Some(&schms2);
426 let user = setup::Uid::default();
427
428 let schms_input = RestSchemesInput::new(
429 Some(aes_gcm.clone()),
430 None::<HashScheme>,
431 None::<HashScheme>,
432 Some(crc32.clone()),
433 );
434
435 let mut cfg = res!(setup::default_cfg());
436 cfg.num_zones = 1;
437 cfg.num_cbots_per_zone = 1;
438 cfg.num_fbots_per_zone = 1;
439 cfg.num_igbots_per_zone = 1;
440 cfg.num_rbots_per_zone = 2;
441 cfg.num_wbots_per_zone = 1;
442 cfg.data_file_max_bytes = 4_000;
443 cfg.rest_chunk_threshold = 1_500;
444 cfg.rest_chunk_bytes = 64;
445 cfg.init_load_caches = true;
446 cfg.sync_on_write = true;
447 cfg.zone_overrides = DaticleMap::new();
448
449 test!(sync_log::stream(), "+--- gc stale floc: {} ---", label);
450
451 let db = res!(setup::start_db(
452 db_root.clone(),
453 Some(cfg.clone()),
454 schms_input.clone(),
455 None,
456 true, // gc on
457 true, // wipe
458 ));
459 thread::sleep(Duration::from_millis(500));
460
461 for i in 0..NPRE {
462 res!(db.insert(filler_key(i), value_of(i as u8, VALUE_BYTES), user, schms2));
463 }
464 res!(db.insert(survivor_key(), survivor_val(), user, schms2));
465 for i in NPRE..NFILL {
466 res!(db.insert(filler_key(i), value_of(i as u8, VALUE_BYTES), user, schms2));
467 }
468 thread::sleep(Duration::from_millis(500));
469 res!(db.api().clear_cache_values(wait()));
470
471 let before = res!(data_file_sizes(&db_root));
472
473 let stop = Arc::new(AtomicBool::new(false));
474 let reads = Arc::new(AtomicUsize::new(0));
475 // Every fault is recorded and the reader carries on, because whether the damage is one blip
476 // per collection or a key that stays unreadable is the difference between an in-flight
477 // location and a cached one, and only a reader that keeps going can tell them apart.
478 let faults: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
479
480 let outcome = thread::scope(|scope| {
481 let stop_r = Arc::clone(&stop);
482 let reads_r = Arc::clone(&reads);
483 let faults_r = Arc::clone(&faults);
484 let db_r = &db;
485 scope.spawn(move || {
486 while !stop_r.load(Ordering::Relaxed) {
487 let n = reads_r.fetch_add(1, Ordering::Relaxed) + 1;
488 let fault = match db_r.get(&survivor_key(), schms2) {
489 Ok(Some((got, _))) => if got == survivor_val() {
490 None
491 } else {
492 Some(fmt!("Read {} returned different bytes from those written.", n))
493 },
494 Ok(None) => Some(fmt!(
495 "Read {} found nothing under a key that was never deleted.", n)),
496 Err(e) => Some(fmt!("Read {} failed: {}", n, e)),
497 };
498 if let Some(msg) = fault {
499 let mut slot = lock_mutex_thread!(faults_r, "concurrent reader");
500 slot.push(msg);
501 }
502 thread::sleep(Duration::from_millis(2));
503 }
504 });
505
506 // Drive the collection underneath the reader.
507 let mut result = Ok(());
508 for i in 0..NFILL {
509 if let Err(e) = db.insert(
510 filler_key(i),
511 value_of((i + 128) as u8, VALUE_BYTES),
512 user,
513 schms2,
514 ) {
515 result = Err(e);
516 break;
517 }
518 }
519 thread::sleep(Duration::from_secs(3));
520 stop.store(true, Ordering::Relaxed);
521 result
522 });
523 res!(outcome);
524
525 let compacted = res!(wait_for_compaction(&db_root, &before, Duration::from_secs(5)));
526 let nreads = reads.load(Ordering::Relaxed);
527
528 let fault_list = {
529 let slot = lock_mutex!(faults);
530 slot.clone()
531 };
532
533 // Once everything has settled, read the survivor again. A one-off race on a location that was
534 // in flight during the transcription leaves nothing behind, so these reads would be clean; a
535 // cached handle on the renamed-away inode keeps failing until it expires or the process ends.
536 let mut after_settled = Vec::new();
537 for n in 1..=NREADS {
538 if let Err(e) = read_survivor(&db, schms2, label, n) {
539 after_settled.push(fmt!("{}", e));
540 }
541 }
542 // The file states carry the move map, so a dump here says whether an unconsumed old -> new
543 // entry is still sitting on the file the survivor lives in.
544 res!(db.api().dump_file_states(wait()));
545
546 // Every read must hand its file back. The fbot will not collect a file whose reader count is
547 // above zero, so a read that returned early without sending its completion leaves a count that
548 // never comes down and a file that can never be collected again -- a burst of failures here
549 // was seen to strand a count in the thousands. Nothing is reading by now, so every count must
550 // be zero.
551 let mut stranded = Vec::new();
552 for (wind, fstates) in res!(db.api().collect_file_states(wait())) {
553 for (fnum, fstat) in fstates.map() {
554 if fstat.readers() != 0 {
555 stranded.push(fmt!("{} file {} holds {} reader(s)", wind, fnum, fstat.readers()));
556 }
557 }
558 }
559
560 let _ = db.shutdown();
561 thread::sleep(Duration::from_millis(300));
562
563 // The same store reopened: the caches are rebuilt from the index files, so a fault that
564 // clears here was in memory and not on disk. This is the cold start that clears the
565 // production symptom, and it is what says the store itself is intact.
566 let db2 = res!(setup::start_db(
567 db_root.clone(),
568 Some(cfg.clone()),
569 schms_input.clone(),
570 None,
571 false, // gc off: nothing should move during the check
572 false, // keep what is there
573 ));
574 thread::sleep(Duration::from_millis(500));
575 res!(db2.api().clear_cache_values(wait()));
576 let mut after_restart = Vec::new();
577 for n in 1..=NREADS {
578 if let Err(e) = read_survivor(&db2, schms2, label, n) {
579 after_restart.push(fmt!("{}", e));
580 }
581 }
582 let _ = db2.shutdown();
583 thread::sleep(Duration::from_millis(300));
584 test!(sync_log::stream(),
585 "[{}] After a restart, {} of {} reads of the survivor failed.",
586 label, after_restart.len(), NREADS);
587
588 match compacted {
589 None => return Err(err!(
590 "[{}] No data file shrank, so no compaction ran under the reader and this case \
591 would prove nothing.", label;
592 Test, Missing)),
593 Some((path, was, is)) => test!(sync_log::stream(),
594 "[{}] Compacted {:?}: {} -> {} bytes, under {} concurrent reads, {} of which \
595 failed. {} of {} reads after everything settled failed.",
596 label, path, was, is, nreads, fault_list.len(), after_settled.len(), NREADS),
597 }
598
599 for msg in fault_list.iter().take(3) {
600 test!(sync_log::stream(), "[{}] During collection: {}", label, msg);
601 }
602 for msg in after_settled.iter().take(3) {
603 test!(sync_log::stream(), "[{}] After settling: {}", label, msg);
604 }
605
606 if !fault_list.is_empty() || !after_settled.is_empty() {
607 return Err(err!(
608 "[{}] {} of {} reads of the survivor broke while collection was running, and {} of \
609 {} after it settled. First: {}",
610 label, fault_list.len(), nreads, after_settled.len(), NREADS,
611 match fault_list.first() {
612 Some(msg) => msg.clone(),
613 None => match after_settled.first() {
614 Some(msg) => msg.clone(),
615 None => fmt!("none"),
616 },
617 };
618 Test, Data));
619 }
620
621 if !stranded.is_empty() {
622 return Err(err!(
623 "[{}] Every read has finished, but {} file state(s) still hold a reader count: {}. \
624 A read that returns early without telling its fbot strands the count, and the fbot \
625 will never collect that file again.", label, stranded.len(), stranded.join("; ");
626 Test, Data, Mismatch));
627 }
628
629 test!(sync_log::stream(),
630 "[{}] {} reads of the survivor during collection all returned the bytes written, and \
631 every file state is back to zero readers.", label, nreads);
632 Ok(())
633}
634
635const NDEL: usize = 6; // chunked keys deleted during the burst
636const CHUNKED_BYTES: usize = 1_600; // over the 1_500 chunk threshold, so each value is chunked
637
638fn del_key(i: usize) -> Dat { dat!(fmt!("gcrace:del:{:04}", i)) }
639fn del_val(i: usize) -> Dat { value_of((i as u8).wrapping_mul(7).wrapping_add(3), CHUNKED_BYTES) }
640fn chunked_survivor_key() -> Dat { dat!("gcrace:chunksurv") }
641fn chunked_survivor_val() -> Dat { value_of(0xa5, CHUNKED_BYTES) }
642
643/// A delete of a chunked key travels the same read path as a get: `delete_using_responder` calls
644/// `reclaim_chunks_on_delete`, which fetches the bunch key through a reader bot before it can
645/// tombstone the chunk records it names. If that bunch-key record lives in a file a collection is
646/// renaming underneath the reader, the fetch is exposed to exactly the stale-handle race that
647/// `read_during_compaction` drives -- and a failed fetch there fails the delete with `[Checksum]`.
648/// This case deletes a set of chunked keys while a supersession burst compacts the files their
649/// records sit in, and insists every delete succeeds and the key then reads back absent. A reader
650/// hammers a separate chunked survivor throughout, both to keep the collection racing and because a
651/// chunked get is itself a fan-out of reads over bunch and chunk records.
652fn delete_during_compaction() -> Outcome<()> {
653
654 let label = "delete_during_compaction";
655 let db_root = res!(canonical_dir("./test_db_gc_stale_floc_delete"));
656
657 let enckey = [0x5cu8; 32];
658 let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..]));
659 let crc32 = ChecksumScheme::new_crc32();
660 let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> =
661 RestSchemesOverride::default()
662 .set_encrypter(Override::Default(aes_gcm.clone()));
663 let schms2 = Some(&schms2);
664 let user = setup::Uid::default();
665
666 let schms_input = RestSchemesInput::new(
667 Some(aes_gcm.clone()),
668 None::<HashScheme>,
669 None::<HashScheme>,
670 Some(crc32.clone()),
671 );
672
673 let mut cfg = res!(setup::default_cfg());
674 cfg.num_zones = 1;
675 cfg.num_cbots_per_zone = 1;
676 cfg.num_fbots_per_zone = 1;
677 cfg.num_igbots_per_zone = 1;
678 cfg.num_rbots_per_zone = 2;
679 cfg.num_wbots_per_zone = 1;
680 cfg.data_file_max_bytes = 4_000;
681 cfg.rest_chunk_threshold = 1_500;
682 cfg.rest_chunk_bytes = 64;
683 cfg.init_load_caches = true;
684 cfg.sync_on_write = true;
685 cfg.zone_overrides = DaticleMap::new();
686
687 test!(sync_log::stream(), "+--- gc stale floc: {} ---", label);
688
689 let db = res!(setup::start_db(
690 db_root.clone(),
691 Some(cfg.clone()),
692 schms_input.clone(),
693 None,
694 true, // gc on
695 true, // wipe
696 ));
697 thread::sleep(Duration::from_millis(500));
698
699 // A few fillers ahead of the chunked records, so those records are carried to a new offset by
700 // the collection rather than sitting at the head and never moving.
701 for i in 0..NPRE {
702 res!(db.insert(filler_key(i), value_of(i as u8, VALUE_BYTES), user, schms2));
703 }
704 // The chunked survivor the reader hammers throughout, and the chunked keys that get deleted mid
705 // burst. Each is asserted to have actually chunked, or the delete would never reach the
706 // reclaim fetch this case exists to exercise.
707 let (_, ns) = res!(db.insert(chunked_survivor_key(), chunked_survivor_val(), user, schms2));
708 if ns < 2 {
709 let _ = db.shutdown();
710 return Err(err!(
711 "[{}] The chunked survivor stored in {} chunk(s); the value must exceed the chunk \
712 threshold so the delete path fetches a bunch key.", label, ns;
713 Test, Size));
714 }
715 for i in 0..NDEL {
716 let (_, nc) = res!(db.insert(del_key(i), del_val(i), user, schms2));
717 if nc < 2 {
718 let _ = db.shutdown();
719 return Err(err!(
720 "[{}] Delete-target {} stored in {} chunk(s); it must chunk so its delete drives \
721 reclaim_chunks_on_delete.", label, i, nc;
722 Test, Size));
723 }
724 }
725 for i in NPRE..NFILL {
726 res!(db.insert(filler_key(i), value_of(i as u8, VALUE_BYTES), user, schms2));
727 }
728 thread::sleep(Duration::from_millis(500));
729 res!(db.api().clear_cache_values(wait()));
730
731 let before = res!(data_file_sizes(&db_root));
732
733 let stop = Arc::new(AtomicBool::new(false));
734 let reads = Arc::new(AtomicUsize::new(0));
735 let read_faults: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
736
737 let outcome = thread::scope(|scope| {
738 let stop_r = Arc::clone(&stop);
739 let reads_r = Arc::clone(&reads);
740 let faults_r = Arc::clone(&read_faults);
741 let db_r = &db;
742 scope.spawn(move || {
743 while !stop_r.load(Ordering::Relaxed) {
744 let n = reads_r.fetch_add(1, Ordering::Relaxed) + 1;
745 let fault = match db_r.get(&chunked_survivor_key(), schms2) {
746 Ok(Some((got, _))) => if got == chunked_survivor_val() {
747 None
748 } else {
749 Some(fmt!("Read {} of the chunked survivor returned different bytes.", n))
750 },
751 Ok(None) => Some(fmt!(
752 "Read {} of the chunked survivor found nothing under a live key.", n)),
753 Err(e) => Some(fmt!("Read {} of the chunked survivor failed: {}", n, e)),
754 };
755 if let Some(msg) = fault {
756 let mut slot = lock_mutex_thread!(faults_r, "concurrent reader");
757 slot.push(msg);
758 }
759 thread::sleep(Duration::from_millis(2));
760 }
761 });
762
763 // Drive the collection with a supersession burst, and delete the chunked keys through it.
764 // The deletes are spaced across the burst so some land while a file holding a bunch-key
765 // record is mid-rename -- the moment the reclaim fetch is exposed to the stale handle.
766 let mut result = Ok(());
767 let mut delete_faults: Vec<String> = Vec::new();
768 let mut next_del = 0usize;
769 for i in 0..NFILL {
770 if let Err(e) = db.insert(
771 filler_key(i),
772 value_of((i + 128) as u8, VALUE_BYTES),
773 user,
774 schms2,
775 ) {
776 result = Err(e);
777 break;
778 }
779 // Roughly one delete every NFILL/NDEL fillers, so they are strewn through the burst.
780 if next_del < NDEL && i >= (next_del * NFILL) / NDEL {
781 match db.delete(&del_key(next_del), user, schms2) {
782 Ok(_) => (),
783 Err(e) => delete_faults.push(fmt!(
784 "Delete of {:?} during collection failed: {}", del_key(next_del), e)),
785 }
786 next_del += 1;
787 }
788 }
789 // Any deletes not yet issued (short burst) go out now, still before settling.
790 while next_del < NDEL {
791 if let Err(e) = db.delete(&del_key(next_del), user, schms2) {
792 delete_faults.push(fmt!(
793 "Delete of {:?} during collection failed: {}", del_key(next_del), e));
794 }
795 next_del += 1;
796 }
797 thread::sleep(Duration::from_secs(3));
798 stop.store(true, Ordering::Relaxed);
799 result.map(|()| delete_faults)
800 });
801 let delete_faults = res!(outcome);
802
803 let compacted = res!(wait_for_compaction(&db_root, &before, Duration::from_secs(5)));
804 let nreads = reads.load(Ordering::Relaxed);
805
806 let read_fault_list = {
807 let slot = lock_mutex!(read_faults);
808 slot.clone()
809 };
810
811 // Every deleted key must now read back absent: the tombstone superseded the bunch key and the
812 // reclaim tombstoned each chunk, so a get reconstructs nothing. A key that still reads a value
813 // would mean the delete's write landed but its reclaim fetch was skipped or lost.
814 let mut still_present = Vec::new();
815 for i in 0..NDEL {
816 match db.get(&del_key(i), schms2) {
817 Ok(None) => (),
818 Ok(Some(_)) => still_present.push(fmt!("{:?} still returns a value", del_key(i))),
819 Err(e) => still_present.push(fmt!("read-back of {:?} failed: {}", del_key(i), e)),
820 }
821 }
822
823 // The chunked survivor, never deleted, must still read back its bytes.
824 let mut survivor_faults = Vec::new();
825 for n in 1..=NREADS {
826 match db.get(&chunked_survivor_key(), schms2) {
827 Ok(Some((got, _))) => if got != chunked_survivor_val() {
828 survivor_faults.push(fmt!("post-settle read {} returned different bytes", n));
829 },
830 Ok(None) => survivor_faults.push(fmt!("post-settle read {} found nothing", n)),
831 Err(e) => survivor_faults.push(fmt!("post-settle read {} failed: {}", n, e)),
832 }
833 }
834
835 // No read may have stranded a reader count: a fetch that returned early on a checksum mismatch
836 // without reporting completion leaves a file uncollectable. This holds for a delete's internal
837 // fetch exactly as for a plain read.
838 let mut stranded = Vec::new();
839 for (wind, fstates) in res!(db.api().collect_file_states(wait())) {
840 for (fnum, fstat) in fstates.map() {
841 if fstat.readers() != 0 {
842 stranded.push(fmt!("{} file {} holds {} reader(s)", wind, fnum, fstat.readers()));
843 }
844 }
845 }
846
847 let _ = db.shutdown();
848 thread::sleep(Duration::from_millis(300));
849
850 match compacted {
851 None => return Err(err!(
852 "[{}] No data file shrank, so no compaction ran under the deletes and this case \
853 would prove nothing.", label;
854 Test, Missing)),
855 Some((path, was, is)) => test!(sync_log::stream(),
856 "[{}] Compacted {:?}: {} -> {} bytes, under {} concurrent reads and {} deletes; {} \
857 deletes failed, {} reads failed.",
858 label, path, was, is, nreads, NDEL,
859 delete_faults.len(), read_fault_list.len()),
860 }
861
862 for msg in delete_faults.iter().take(3) {
863 test!(sync_log::stream(), "[{}] Delete fault: {}", label, msg);
864 }
865 for msg in read_fault_list.iter().take(3) {
866 test!(sync_log::stream(), "[{}] Read fault: {}", label, msg);
867 }
868
869 if !delete_faults.is_empty() {
870 return Err(err!(
871 "[{}] {} of {} chunked-key deletes failed while their files were being compacted. \
872 First: {}", label, delete_faults.len(), NDEL, delete_faults[0];
873 Test, Data));
874 }
875 if !still_present.is_empty() {
876 return Err(err!(
877 "[{}] {} deleted key(s) did not read back absent: {}.",
878 label, still_present.len(), still_present.join("; ");
879 Test, Data, Mismatch));
880 }
881 if !read_fault_list.is_empty() || !survivor_faults.is_empty() {
882 return Err(err!(
883 "[{}] The chunked survivor broke: {} read(s) failed during collection and {} after.",
884 label, read_fault_list.len(), survivor_faults.len();
885 Test, Data));
886 }
887 if !stranded.is_empty() {
888 return Err(err!(
889 "[{}] Every read has finished, but {} file state(s) still hold a reader count: {}.",
890 label, stranded.len(), stranded.join("; ");
891 Test, Data, Mismatch));
892 }
893
894 test!(sync_log::stream(),
895 "[{}] {} chunked-key deletes all succeeded under a live collection, each read back \
896 absent, the chunked survivor stayed intact across {} reads, and every file state is back \
897 to zero readers.", label, NDEL, nreads);
898 Ok(())
899}
900
901/// Creates the directory if it does not exist.
902fn canonical_dir(p: &str) -> Outcome<PathBuf> {
903 match Path::new(p).canonicalize() {
904 Ok(path) => Ok(path),
905 Err(_) => {
906 res!(fs::create_dir_all(p));
907 match Path::new(p).canonicalize() {
908 Ok(path) => Ok(path),
909 Err(e) => Err(err!(e, "Cannot canonicalise {:?}.", p; IO, Path)),
910 }
911 },
912 }
913}
914
915pub fn test_gc_stale_floc(_filter: &'static str) -> Outcome<()> {
916 res!(run_case(Case::NoCompaction));
917 res!(run_case(Case::SomeSurvive));
918 res!(run_case(Case::OnlySurvivor));
919 res!(read_during_compaction());
920 res!(delete_during_compaction());
921 Ok(())
922}
923
924#[test]
925fn main() -> Outcome<()> {
926 log_set_level!("debug");
927 let outcome = test_gc_stale_floc("all");
928 log_finish_wait!();
929 outcome
930}