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 | |
| 53 | use oxedyne_fe2o3_core::{ |
| 54 | prelude::*, |
| 55 | alt::Override, |
| 56 | }; |
| 57 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 58 | use oxedyne_fe2o3_hash::{ |
| 59 | csum::ChecksumScheme, |
| 60 | hash::HashScheme, |
| 61 | }; |
| 62 | use oxedyne_fe2o3_iop_db::api::{ |
| 63 | Database, |
| 64 | RestSchemesOverride, |
| 65 | }; |
| 66 | use oxedyne_fe2o3_jdat::prelude::*; |
| 67 | use 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 | |
| 80 | use 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 | |
| 103 | const VALUE_BYTES: usize = 200; // well under the chunking threshold: one record per value |
| 104 | const NPRE: usize = 5; // fillers written ahead of the survivor in its own data file |
| 105 | const NFILL: usize = 60; // enough fillers to seal several data files behind the survivor |
| 106 | const NREADS: usize = 6; // a first read plus five more, so a read-once fault is visible |
| 107 | const 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)] |
| 111 | enum 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 | |
| 120 | impl 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 | |
| 153 | fn 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. |
| 161 | fn 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 | |
| 169 | fn survivor_key() -> Dat { dat!("gcrace:survivor") } |
| 170 | fn survivor_val() -> Dat { value_of(0x5a, VALUE_BYTES) } |
| 171 | fn filler_key(i: usize) -> Dat { dat!(fmt!("gcrace:fill:{:04}", i)) } |
| 172 | |
| 173 | /// Every data file under `root`, with its current length. |
| 174 | fn 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. |
| 193 | fn 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. |
| 220 | fn 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 | |
| 256 | fn 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. |
| 414 | fn 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 | |
| 635 | const NDEL: usize = 6; // chunked keys deleted during the burst |
| 636 | const CHUNKED_BYTES: usize = 1_600; // over the 1_500 chunk threshold, so each value is chunked |
| 637 | |
| 638 | fn del_key(i: usize) -> Dat { dat!(fmt!("gcrace:del:{:04}", i)) } |
| 639 | fn del_val(i: usize) -> Dat { value_of((i as u8).wrapping_mul(7).wrapping_add(3), CHUNKED_BYTES) } |
| 640 | fn chunked_survivor_key() -> Dat { dat!("gcrace:chunksurv") } |
| 641 | fn 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. |
| 652 | fn 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. |
| 902 | fn 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 | |
| 915 | pub 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] |
| 925 | fn main() -> Outcome<()> { |
| 926 | log_set_level!("debug"); |
| 927 | let outcome = test_gc_stale_floc("all"); |
| 928 | log_finish_wait!(); |
| 929 | outcome |
| 930 | } |