oxedyne/fe2o3/fe2o3_o3db_sync/tests/gc_reevaluate.rs
15.1 KiB, 1 run
created by r1870400018:61245, 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 file whose garbage the collector passed over once kept it for good (2026-09-23). Two |
| 2 | //! faults, found because the online orphan sweep reclaimed a different fraction of what it retired |
| 3 | //! on every run, and on some runs too little: |
| 4 | //! |
| 5 | //! - A file bot applied a supersession it handled itself, rather than one sent to it by another |
| 6 | //! file bot, straight to its file even while the file was being collected. The collection's |
| 7 | //! result then replaced the file's state, so the record was carried into the rewritten file as |
| 8 | //! current, with a move entry nothing would clear, and a file holding a move entry is never |
| 9 | //! collected again. |
| 10 | //! - Collection was considered only when a supersession arrived. A file that passed the trigger |
| 11 | //! while it was live, when it could not be collected, was not looked at again when it was |
| 12 | //! sealed. |
| 13 | //! |
| 14 | //! Each is reproduced on one zone with one file bot, so that every supersession is one the bot |
| 15 | //! handles itself, and judged by the size of the data file on disk. The second is checked twice, |
| 16 | //! once with the file drained when it is sealed and once with its last records still on their way, |
| 17 | //! which only their landing can follow up. `test::hooks` holds collections open, so that |
| 18 | //! supersessions land during one or a file can be measured before one, and holds durability |
| 19 | //! barriers; the hooks are process-wide, which is why this is a test binary of its own. |
| 20 | |
| 21 | use oxedyne_fe2o3_core::prelude::*; |
| 22 | use oxedyne_fe2o3_hash::{ |
| 23 | csum::ChecksumScheme, |
| 24 | hash::HashScheme, |
| 25 | }; |
| 26 | use oxedyne_fe2o3_iop_db::api::Database; |
| 27 | use oxedyne_fe2o3_jdat::prelude::*; |
| 28 | use oxedyne_fe2o3_o3db_sync::{ |
| 29 | O3db, |
| 30 | base::{ |
| 31 | cfg::OzoneConfig, |
| 32 | constant, |
| 33 | }, |
| 34 | comm::response::Wait, |
| 35 | data::core::RestSchemesInput, |
| 36 | file::{ |
| 37 | core::FileType, |
| 38 | zdir::ZoneDir, |
| 39 | }, |
| 40 | test::{ |
| 41 | hooks, |
| 42 | setup::{ |
| 43 | self, |
| 44 | Uid, |
| 45 | UID_LEN, |
| 46 | }, |
| 47 | }, |
| 48 | }; |
| 49 | |
| 50 | use std::{ |
| 51 | collections::BTreeMap, |
| 52 | path::{ |
| 53 | Path, |
| 54 | PathBuf, |
| 55 | }, |
| 56 | thread, |
| 57 | time::{ |
| 58 | Duration, |
| 59 | Instant, |
| 60 | }, |
| 61 | }; |
| 62 | |
| 63 | type TestDb = O3db< |
| 64 | { UID_LEN }, |
| 65 | Uid, |
| 66 | (), |
| 67 | HashScheme, |
| 68 | HashScheme, |
| 69 | ChecksumScheme, |
| 70 | >; |
| 71 | |
| 72 | const FILE_BYTES: u64 = 16_000; // the collection trigger is 30% of this, 4,800 bytes |
| 73 | const VALUE_BYTES: usize = 1_000; // five superseded records pass the trigger, four do not |
| 74 | const NKEYS: usize = 10; // most of the first file |
| 75 | const NFILL: usize = 8; // seals the first file part way through |
| 76 | const HOLD: Duration = Duration::from_millis(1_500); |
| 77 | const SETTLE: Duration = Duration::from_secs(120); // a collection's syncs on a busy disk |
| 78 | |
| 79 | #[test] |
| 80 | fn main() -> Outcome<()> { |
| 81 | log_set_level!("warn"); |
| 82 | // Whatever failed, the next check and the next binary must not inherit a slow collector. |
| 83 | let during = supersession_during_collection_is_kept(); |
| 84 | hooks::set_collect_delay(Duration::ZERO); |
| 85 | let sealed = garbage_made_while_live_is_collected_once_sealed(); |
| 86 | hooks::set_collect_delay(Duration::ZERO); |
| 87 | let landed = garbage_is_collected_when_the_last_write_lands(); |
| 88 | hooks::set_collect_delay(Duration::ZERO); |
| 89 | hooks::set_barrier_delay(Duration::ZERO); |
| 90 | log_finish_wait!(); |
| 91 | let failed: Vec<Error<ErrTag>> = [during, sealed, landed].into_iter() |
| 92 | .filter_map(|r| r.err()) |
| 93 | .collect(); |
| 94 | match failed.len() { |
| 95 | 0 => Ok(()), |
| 96 | 1 => match failed.into_iter().next() { |
| 97 | Some(e) => Err(e), |
| 98 | None => Ok(()), // unreachable |
| 99 | }, |
| 100 | n => Err(err!( |
| 101 | "{} of the 3 checks failed: {:?}", n, failed; |
| 102 | Test)), |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | /// Half of ten records in a sealed file are superseded, which starts its collection, and the other |
| 107 | /// half while that collection is held open. All ten must be reclaimed: the second half by a second |
| 108 | /// collection, once the first has finished and the supersessions it held back are applied. Before, |
| 109 | /// the second half was lost and the file kept those five records for good. |
| 110 | fn supersession_during_collection_is_kept() -> Outcome<()> { |
| 111 | let root = res!(fresh("./test_db_gc_reevaluate_during")); |
| 112 | let db = res!(open(&root)); |
| 113 | let f1 = res!(data_file(&root, 1)); |
| 114 | |
| 115 | for i in 0..NKEYS { |
| 116 | res!(db.insert(key(i), value(i, VALUE_BYTES), Uid::default(), None)); |
| 117 | } |
| 118 | let kbytes = size(&f1); |
| 119 | for i in 0..NFILL { |
| 120 | res!(db.insert(filler(i), value(i, VALUE_BYTES), Uid::default(), None)); |
| 121 | } |
| 122 | let sealed = res!(sealed_size(&root, &f1)); |
| 123 | |
| 124 | hooks::set_collect_delay(HOLD); |
| 125 | for i in 0..(NKEYS / 2) { |
| 126 | res!(db.insert(key(i), value(i + 100, VALUE_BYTES), Uid::default(), None)); |
| 127 | } |
| 128 | res!(wait_for_collection(&db, 1)); |
| 129 | for i in (NKEYS / 2)..NKEYS { |
| 130 | res!(db.insert(key(i), value(i + 100, VALUE_BYTES), Uid::default(), None)); |
| 131 | } |
| 132 | hooks::set_collect_delay(Duration::ZERO); |
| 133 | |
| 134 | let want = res!(garbage_removed(sealed, kbytes)); |
| 135 | let got = settle_to(&f1, want); |
| 136 | if got != want { |
| 137 | let _ = db.close(); |
| 138 | return Err(err!( |
| 139 | "Data file 1 held {} bytes of ten superseded records and {} of fillers when sealed. \ |
| 140 | With every record superseded, half of them while it was being collected, it should \ |
| 141 | have settled at {} bytes, but after {:?} it is {}: {} bytes of superseded records \ |
| 142 | were never reclaimed.", |
| 143 | kbytes, want, want, SETTLE, got, got.saturating_sub(want); |
| 144 | Test, Mismatch, Size)); |
| 145 | } |
| 146 | for i in 0..NKEYS { |
| 147 | match res!(db.get(&key(i), None)) { |
| 148 | Some((v, _)) => req!(v, value(i + 100, VALUE_BYTES), "A value carried through two collections."), |
| 149 | None => return Err(err!("Key {} is missing after two collections.", i; Test, Missing)), |
| 150 | } |
| 151 | } |
| 152 | for i in 0..NFILL { |
| 153 | match res!(db.get(&filler(i), None)) { |
| 154 | Some((v, _)) => req!(v, value(i, VALUE_BYTES), "A filler carried through two collections."), |
| 155 | None => return Err(err!("Filler {} is missing after two collections.", i; Test, Missing)), |
| 156 | } |
| 157 | } |
| 158 | res!(db.close()); |
| 159 | test!(sync_log::stream(), "Supersessions made during a collection were kept: {} -> {} bytes.", sealed, got); |
| 160 | Ok(()) |
| 161 | } |
| 162 | |
| 163 | /// Ten records are superseded while their file is still live, which pushes it past the collection |
| 164 | /// trigger when it cannot be collected, and the file is then sealed with nothing in it superseded |
| 165 | /// again. The ten must still be reclaimed. Before, nothing looked at the file after the seal. |
| 166 | fn garbage_made_while_live_is_collected_once_sealed() -> Outcome<()> { |
| 167 | let root = res!(fresh("./test_db_gc_reevaluate_sealed")); |
| 168 | let db = res!(open(&root)); |
| 169 | let f1 = res!(data_file(&root, 1)); |
| 170 | |
| 171 | for i in 0..NKEYS { |
| 172 | res!(db.insert(key(i), value(i, VALUE_BYTES), Uid::default(), None)); |
| 173 | } |
| 174 | let kbytes = size(&f1); |
| 175 | for i in 0..NKEYS { |
| 176 | res!(db.insert(key(i), value(i, 8), Uid::default(), None)); |
| 177 | } |
| 178 | // The seal may start the collection at once, so it is held until the file has been measured. |
| 179 | hooks::set_collect_delay(HOLD); |
| 180 | for i in 0..NFILL { |
| 181 | res!(db.insert(filler(i), value(i, VALUE_BYTES), Uid::default(), None)); |
| 182 | } |
| 183 | let sealed = res!(sealed_size(&root, &f1)); |
| 184 | hooks::set_collect_delay(Duration::ZERO); |
| 185 | |
| 186 | let want = res!(garbage_removed(sealed, kbytes)); |
| 187 | let got = settle_to(&f1, want); |
| 188 | if got != want { |
| 189 | let _ = db.close(); |
| 190 | return Err(err!( |
| 191 | "Data file 1 was sealed at {} bytes, {} of them ten records superseded while it was \ |
| 192 | live. It should have settled at {} bytes, but after {:?} it is {}.", |
| 193 | sealed, kbytes, want, SETTLE, got; |
| 194 | Test, Mismatch, Size)); |
| 195 | } |
| 196 | for i in 0..NKEYS { |
| 197 | match res!(db.get(&key(i), None)) { |
| 198 | Some((v, _)) => req!(v, value(i, 8), "A value written while its file was live."), |
| 199 | None => return Err(err!("Key {} is missing after the collection.", i; Test, Missing)), |
| 200 | } |
| 201 | } |
| 202 | res!(db.close()); |
| 203 | test!(sync_log::stream(), "Garbage made while live was collected once sealed: {} -> {} bytes.", sealed, got); |
| 204 | Ok(()) |
| 205 | } |
| 206 | |
| 207 | /// As the previous check, but the file's last records are still on their way to their file bot |
| 208 | /// when it is sealed: every durability barrier is held, so the seal finds the file not yet drained |
| 209 | /// and cannot collect it. The last of those records landing must. Before, nothing looked at the |
| 210 | /// file again. |
| 211 | fn garbage_is_collected_when_the_last_write_lands() -> Outcome<()> { |
| 212 | let root = res!(fresh("./test_db_gc_reevaluate_landed")); |
| 213 | let mut cfg = res!(config()); |
| 214 | cfg.sync_on_write = true; |
| 215 | let db = res!(setup::start_db(root.clone(), Some(cfg), schemes(), None, true, true)); |
| 216 | let f1 = res!(data_file(&root, 1)); |
| 217 | |
| 218 | for i in 0..NKEYS { |
| 219 | res!(db.insert(key(i), value(i, VALUE_BYTES), Uid::default(), None)); |
| 220 | } |
| 221 | let kbytes = size(&f1); |
| 222 | for i in 0..NKEYS { |
| 223 | res!(db.insert(key(i), value(i, 8), Uid::default(), None)); |
| 224 | } |
| 225 | // The fillers are sent together, so the one that seals the file is written while the records |
| 226 | // before it wait on their barrier. The collection is held until the file has been measured. |
| 227 | hooks::set_barrier_delay(HOLD); |
| 228 | hooks::set_collect_delay(HOLD); |
| 229 | let mut resps = Vec::new(); |
| 230 | for i in 0..NFILL { |
| 231 | resps.push(res!(db.api().store(filler(i), value(i, VALUE_BYTES), Uid::default()))); |
| 232 | } |
| 233 | // Measured once sealed, while its last records still wait on their barrier. |
| 234 | let sealed = res!(wait_until_sealed(&root, &f1)); |
| 235 | for resp in resps { |
| 236 | res!(resp.recv_store_ack()); |
| 237 | } |
| 238 | hooks::set_barrier_delay(Duration::ZERO); |
| 239 | hooks::set_collect_delay(Duration::ZERO); |
| 240 | |
| 241 | let want = res!(garbage_removed(sealed, kbytes)); |
| 242 | let got = settle_to(&f1, want); |
| 243 | if got != want { |
| 244 | let _ = db.close(); |
| 245 | return Err(err!( |
| 246 | "Data file 1 was sealed at {} bytes before its last records had landed, {} of them \ |
| 247 | ten records superseded while it was live. It should have settled at {} bytes once \ |
| 248 | they landed, but after {:?} it is {}.", |
| 249 | sealed, kbytes, want, SETTLE, got; |
| 250 | Test, Mismatch, Size)); |
| 251 | } |
| 252 | for i in 0..NFILL { |
| 253 | match res!(db.get(&filler(i), None)) { |
| 254 | Some((v, _)) => req!(v, value(i, VALUE_BYTES), "A filler that sealed a collected file."), |
| 255 | None => return Err(err!("Filler {} is missing after the collection.", i; Test, Missing)), |
| 256 | } |
| 257 | } |
| 258 | res!(db.close()); |
| 259 | test!(sync_log::stream(), "Garbage was collected when the last write landed: {} -> {} bytes.", sealed, got); |
| 260 | Ok(()) |
| 261 | } |
| 262 | |
| 263 | fn key(i: usize) -> Dat { dat!(fmt!("gc reevaluate key {:02}", i)) } |
| 264 | fn filler(i: usize) -> Dat { dat!(fmt!("gc reevaluate filler {:02}", i)) } |
| 265 | |
| 266 | /// A byte string of the given length, distinct for each seed. |
| 267 | fn value(seed: usize, len: usize) -> Dat { |
| 268 | let mut v = vec![0u8; len]; |
| 269 | for (i, b) in v.iter_mut().enumerate() { |
| 270 | *b = (seed as u8).wrapping_add((i % 251) as u8); |
| 271 | } |
| 272 | Dat::BU32(v) |
| 273 | } |
| 274 | |
| 275 | /// One zone and one bot of each kind, with values well under the chunking threshold, so that every |
| 276 | /// write is one record in the first zone and every supersession is handled by the one file bot. |
| 277 | fn config() -> Outcome<OzoneConfig> { |
| 278 | let mut cfg = res!(setup::default_cfg()); |
| 279 | cfg.num_zones = 1; |
| 280 | cfg.num_cbots_per_zone = 1; |
| 281 | cfg.num_fbots_per_zone = 1; |
| 282 | cfg.num_igbots_per_zone = 1; |
| 283 | cfg.num_rbots_per_zone = 1; |
| 284 | cfg.num_wbots_per_zone = 1; |
| 285 | cfg.data_file_max_bytes = FILE_BYTES; |
| 286 | cfg.rest_chunk_threshold = 3_000; // well over VALUE_BYTES |
| 287 | cfg.zone_overrides = BTreeMap::new(); |
| 288 | Ok(cfg) |
| 289 | } |
| 290 | |
| 291 | fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> { |
| 292 | RestSchemesInput::new( |
| 293 | None::<()>, |
| 294 | None::<HashScheme>, |
| 295 | None::<HashScheme>, |
| 296 | Some(ChecksumScheme::new_crc32()), |
| 297 | ) |
| 298 | } |
| 299 | |
| 300 | /// An empty directory of this test's own. |
| 301 | fn fresh(dir: &str) -> Outcome<PathBuf> { |
| 302 | let _ = std::fs::remove_dir_all(dir); |
| 303 | res!(std::fs::create_dir_all(dir)); |
| 304 | Ok(res!(Path::new(dir).canonicalize())) |
| 305 | } |
| 306 | |
| 307 | fn open(root: &Path) -> Outcome<TestDb> { |
| 308 | setup::start_db( |
| 309 | root.to_path_buf(), |
| 310 | Some(res!(config())), |
| 311 | schemes(), |
| 312 | None, |
| 313 | true, // collection on |
| 314 | true, // wipe |
| 315 | ) |
| 316 | } |
| 317 | |
| 318 | fn data_file(root: &Path, fnum: u32) -> Outcome<PathBuf> { |
| 319 | let cfg = res!(config()); |
| 320 | Ok(cfg.zone_root(root).join("zone_001").join(ZoneDir::relative_file_path(&FileType::Data, fnum))) |
| 321 | } |
| 322 | |
| 323 | /// A file that has been collected away entirely is no longer there. |
| 324 | fn size(path: &Path) -> u64 { |
| 325 | match std::fs::metadata(path) { |
| 326 | Ok(m) => m.len(), |
| 327 | Err(_) => 0, |
| 328 | } |
| 329 | } |
| 330 | |
| 331 | /// The size of data file 1 once the writer has moved on to file 2, after which it cannot grow. |
| 332 | fn sealed_size(root: &Path, f1: &Path) -> Outcome<u64> { |
| 333 | let f2 = res!(data_file(root, 2)); |
| 334 | if !f2.is_file() { |
| 335 | return Err(err!( |
| 336 | "The fillers did not seal data file 1: there is no data file 2 at {:?}.", f2; |
| 337 | Test, Missing)); |
| 338 | } |
| 339 | Ok(size(f1)) |
| 340 | } |
| 341 | |
| 342 | /// Waits for the writer to move on to data file 2, then gives the size of data file 1. |
| 343 | fn wait_until_sealed(root: &Path, f1: &Path) -> Outcome<u64> { |
| 344 | let f2 = res!(data_file(root, 2)); |
| 345 | let begun = Instant::now(); |
| 346 | while !f2.is_file() { |
| 347 | if begun.elapsed() >= SETTLE { |
| 348 | return Err(err!( |
| 349 | "The fillers did not seal data file 1 within {:?}: there is no data file 2 at {:?}.", |
| 350 | SETTLE, f2; |
| 351 | Test, Timeout)); |
| 352 | } |
| 353 | thread::sleep(Duration::from_millis(5)); |
| 354 | } |
| 355 | Ok(size(f1)) |
| 356 | } |
| 357 | |
| 358 | /// The size a sealed file should settle at once its superseded records are gone. |
| 359 | fn garbage_removed(sealed: u64, garbage: u64) -> Outcome<u64> { |
| 360 | match sealed.checked_sub(garbage) { |
| 361 | Some(want) => Ok(want), |
| 362 | None => Err(err!( |
| 363 | "Data file 1 was sealed at {} bytes, less than the {} its first records took.", |
| 364 | sealed, garbage; |
| 365 | Test, Invalid, Size)), |
| 366 | } |
| 367 | } |
| 368 | |
| 369 | /// Waits until the file bot has handed file `fnum` to a collector. |
| 370 | fn wait_for_collection(db: &TestDb, fnum: u32) -> Outcome<()> { |
| 371 | let begun = Instant::now(); |
| 372 | while begun.elapsed() < SETTLE { |
| 373 | let states = res!(db.api().collect_file_states(Wait { |
| 374 | max_wait: constant::USER_REQUEST_TIMEOUT, |
| 375 | check_interval: constant::CHECK_INTERVAL, |
| 376 | })); |
| 377 | for (_, shard) in &states { |
| 378 | if let Some(fstat) = shard.map().get(&fnum) { |
| 379 | if fstat.gc_active() { |
| 380 | return Ok(()); |
| 381 | } |
| 382 | } |
| 383 | } |
| 384 | thread::sleep(Duration::from_millis(20)); |
| 385 | } |
| 386 | Err(err!( |
| 387 | "Superseding half the records in data file {} did not start its collection within {:?}.", |
| 388 | fnum, SETTLE; |
| 389 | Test, Timeout)) |
| 390 | } |
| 391 | |
| 392 | /// Waits for the file to reach the given size, returning the size it has when it does or when |
| 393 | /// the wait runs out. |
| 394 | fn settle_to(path: &Path, want: u64) -> u64 { |
| 395 | let begun = Instant::now(); |
| 396 | loop { |
| 397 | let got = size(path); |
| 398 | if got == want || begun.elapsed() >= SETTLE { |
| 399 | return got; |
| 400 | } |
| 401 | thread::sleep(Duration::from_millis(100)); |
| 402 | } |
| 403 | } |