oxedyne/fe2o3/fe2o3_o3db_sync/tests/scan_under_gc.rs
13.7 KiB, 36 runs
created by r1870400018:21762, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! Ozone scan-under-garbage-collection integration test. |
| 2 | //! |
| 3 | //! A scan is a foreground request with a person waiting on it; garbage |
| 4 | //! collection is background work with no deadline at all. This test |
| 5 | //! holds the two apart. |
| 6 | //! |
| 7 | //! It builds a store whose files are all past the collection trigger |
| 8 | //! while collection is switched off, so no work has been done yet and |
| 9 | //! none is queued. A quiet scan then measures what the walk costs with |
| 10 | //! nothing in its way. Collection is switched on and a short burst of |
| 11 | //! writes dispatches the whole backlog at once, and the same scan is |
| 12 | //! issued again. Two properties are checked: |
| 13 | //! |
| 14 | //! 1. **A concurrent scan sees every key, and does not fail.** |
| 15 | //! Collection rewrites a file's index. A scanner walking that index |
| 16 | //! at the wrong moment can see a file with no index at all and drop |
| 17 | //! every key whose current value lives in it, or read a half-written |
| 18 | //! record and fail outright. Entry count and outcome are checked on |
| 19 | //! every scan. This is the sharp half of the test: before the index |
| 20 | //! rebuild was corrected it failed here every run. |
| 21 | //! |
| 22 | //! 2. **Latency does not track the collector's backlog.** The scan |
| 23 | //! issued while the backlog is being worked through costs about what |
| 24 | //! the quiet one did. This half is a tripwire rather than a proof. |
| 25 | //! A one-shot backlog can only ever cost about as much as one full |
| 26 | //! scan -- the collector's work over every file and the scan's walk |
| 27 | //! over every file are both a fixed cost per record, so the whole |
| 28 | //! backlog is roughly one scan's worth -- and the allowance below is |
| 29 | //! set well clear of ordinary noise. What it catches is a change |
| 30 | //! that puts the scan back on a queue behind unbounded background |
| 31 | //! work; what it cannot show, at any size a test can afford, is the |
| 32 | //! multi-second wait a store under sustained churn produces. |
| 33 | //! |
| 34 | //! A burst is used rather than a sustained write storm because it makes |
| 35 | //! the backlog a known quantity. Under a storm the collector is idle |
| 36 | //! most of the time -- work arrives for it in proportion to the write |
| 37 | //! rate, and it serves that work about as fast as several writers can |
| 38 | //! produce it -- so a queue forms only by luck, and the test would then |
| 39 | //! measure the machine rather than the database. |
| 40 | //! |
| 41 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 42 | //! Anthropic Claude |
| 43 | |
| 44 | use oxedyne_fe2o3_core::{ |
| 45 | prelude::*, |
| 46 | alt::Override, |
| 47 | }; |
| 48 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 49 | use oxedyne_fe2o3_hash::{ |
| 50 | csum::ChecksumScheme, |
| 51 | hash::HashScheme, |
| 52 | }; |
| 53 | use oxedyne_fe2o3_iop_db::api::{ |
| 54 | Database, |
| 55 | RestSchemesOverride, |
| 56 | ScanOpts, |
| 57 | }; |
| 58 | use oxedyne_fe2o3_jdat::prelude::*; |
| 59 | use oxedyne_fe2o3_o3db_sync::{ |
| 60 | data::core::RestSchemesInput, |
| 61 | test::setup, |
| 62 | }; |
| 63 | |
| 64 | use std::{ |
| 65 | path::Path, |
| 66 | thread, |
| 67 | time::{ |
| 68 | Duration, |
| 69 | Instant, |
| 70 | }, |
| 71 | }; |
| 72 | |
| 73 | /// The dimensions of one run. |
| 74 | struct Shape { |
| 75 | dir: &'static str, // directory name for this run's store |
| 76 | keys: usize, // distinct keys held in the store |
| 77 | // Value payload size in bytes, before encoding. Large, so that a |
| 78 | // collection moves far more than a scan reads. |
| 79 | val_bytes: usize, |
| 80 | file_bytes: u64, // and so the size of one collection's work |
| 81 | } |
| 82 | |
| 83 | // The store the crate's own test suite can afford. |
| 84 | const QUICK: Shape = Shape { |
| 85 | dir: "./test_db_scan_under_gc", |
| 86 | keys: 6_000, |
| 87 | val_bytes: 20_000, |
| 88 | file_bytes: 8_000_000, |
| 89 | }; |
| 90 | |
| 91 | // A store several times larger, where the backlog is deep enough to cost more |
| 92 | // than the user request deadline outright rather than only running late. Writes |
| 93 | // some hundreds of megabytes. |
| 94 | const DEEP: Shape = Shape { |
| 95 | dir: "./test_db_scan_stall", |
| 96 | keys: 30_000, |
| 97 | val_bytes: 20_000, |
| 98 | file_bytes: 8_000_000, |
| 99 | }; |
| 100 | |
| 101 | // Fraction of the keys superseded before collection is switched on. It has to |
| 102 | // clear OLD_DATA_PERCENT_GC_TRIGGER, which is measured against the maximum data |
| 103 | // file size rather than the file's own size. |
| 104 | const OLD_FRAC: f64 = 0.5; |
| 105 | |
| 106 | const NUM_BASELINE: usize = 3; // quiet scans establishing the latency baseline |
| 107 | const NUM_LOADED: usize = 5; // scans issued while the backlog is worked off |
| 108 | |
| 109 | /// The spacing of the keys superseded to set the backlog going: about two per |
| 110 | /// file, which reaches every file holding original |
| 111 | /// records while keeping the burst far shorter than the collection work |
| 112 | /// it dispatches. A finer stride lets the collector work the queue down |
| 113 | /// as fast as it is filled, and the test then measures nothing. |
| 114 | fn trigger_stride(shape: &Shape) -> usize { |
| 115 | let per_file = (shape.file_bytes as usize) / shape.val_bytes; |
| 116 | std::cmp::max(1, per_file / 2) |
| 117 | } |
| 118 | |
| 119 | /// Runs the scan-under-collection test as its own integration binary. |
| 120 | /// |
| 121 | /// It is kept out of `tests/main.rs` deliberately: it is the slowest |
| 122 | /// test in the crate and the only one that wants the disk to itself, so |
| 123 | /// `cargo test -p oxedyne_fe2o3_o3db_sync --test scan_under_gc` can run |
| 124 | /// it alone. |
| 125 | #[test] |
| 126 | fn main() -> Outcome<()> { |
| 127 | |
| 128 | // Logging every stored key costs several times what the writes do, |
| 129 | // and this test writes a great many of them. `test` is the lowest |
| 130 | // level that still carries the test's own reporting. |
| 131 | log_set_level!("test"); |
| 132 | |
| 133 | let outcome = test_scan_under_gc(&QUICK); |
| 134 | |
| 135 | log_finish_wait!(); |
| 136 | |
| 137 | outcome |
| 138 | } |
| 139 | |
| 140 | /// The same test against a store deep enough to reproduce the |
| 141 | /// production symptom rather than only the mechanism behind it. It |
| 142 | /// writes for some minutes, so it is ignored by default: |
| 143 | /// |
| 144 | /// ```ignore |
| 145 | /// cargo test -p oxedyne_fe2o3_o3db_sync --test scan_under_gc -- --ignored --nocapture |
| 146 | /// ``` |
| 147 | #[test] |
| 148 | #[ignore] |
| 149 | fn deep() -> Outcome<()> { |
| 150 | |
| 151 | log_set_level!("test"); |
| 152 | |
| 153 | let outcome = test_scan_under_gc(&DEEP); |
| 154 | |
| 155 | log_finish_wait!(); |
| 156 | |
| 157 | outcome |
| 158 | } |
| 159 | |
| 160 | pub fn test_scan_under_gc(shape: &Shape) -> Outcome<()> { |
| 161 | |
| 162 | let db_root = res!(Path::new(shape.dir).canonicalize().or_else(|_| { |
| 163 | ok!(std::fs::create_dir_all(shape.dir)); |
| 164 | Path::new(shape.dir).canonicalize() |
| 165 | })); |
| 166 | |
| 167 | // Fixed key so the test is deterministic. |
| 168 | let enckey = [0x37u8; 32]; |
| 169 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..])); |
| 170 | let crc32 = ChecksumScheme::new_crc32(); |
| 171 | let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> = |
| 172 | RestSchemesOverride::default() |
| 173 | .set_encrypter(Override::Default(aes_gcm.clone())); |
| 174 | let schms2 = Some(&schms2); |
| 175 | let user = setup::Uid::default(); |
| 176 | |
| 177 | let schms_input = RestSchemesInput::new( |
| 178 | Some(aes_gcm.clone()), |
| 179 | None::<HashScheme>, |
| 180 | None::<HashScheme>, |
| 181 | Some(crc32.clone()), |
| 182 | ); |
| 183 | |
| 184 | let mut cfg = res!(setup::default_cfg()); |
| 185 | // One zone, and one bot of each kind in it, so there is exactly one |
| 186 | // queue per role and no ambiguity about what a scan is waiting for. |
| 187 | cfg.num_zones = 1; |
| 188 | cfg.num_cbots_per_zone = 1; |
| 189 | cfg.num_fbots_per_zone = 1; |
| 190 | cfg.num_wbots_per_zone = 1; |
| 191 | // One collector is the sharpest form of the defect: a scan sent to |
| 192 | // the collector's pool is then guaranteed to sit behind whatever is |
| 193 | // already queued there. It is a legal production setting, and with |
| 194 | // more collectors the defect merely becomes intermittent. |
| 195 | cfg.num_igbots_per_zone = 1; |
| 196 | cfg.data_file_max_bytes = shape.file_bytes; |
| 197 | // Values are well under the file size, so chunking stays out of the |
| 198 | // way and each file holds whole records only. |
| 199 | cfg.rest_chunk_threshold = shape.file_bytes / 8; |
| 200 | cfg.rest_chunk_bytes = shape.file_bytes / 40; |
| 201 | // Small enough that the cache jettisons values, so a scan cannot be |
| 202 | // served from a fully resident cache by accident. |
| 203 | cfg.cache_size_limit_bytes = 40_000_000; |
| 204 | cfg.zone_overrides = mapdat!{ |
| 205 | 1u16 => mapdat!{ "dir" => "", "max_size" => 8_000_000_000u64 }, |
| 206 | }.get_map().unwrap(); |
| 207 | |
| 208 | test!(sync_log::stream(), "+---------------------------------------------+"); |
| 209 | test!(sync_log::stream(), "| SCAN UNDER GARBAGE COLLECTION TEST |"); |
| 210 | test!(sync_log::stream(), "+---------------------------------------------+"); |
| 211 | test!(sync_log::stream(), "{} keys of {} bytes in {} byte files.", |
| 212 | shape.keys, shape.val_bytes, shape.file_bytes); |
| 213 | |
| 214 | // Collection stays off while the store is built, so that nothing is |
| 215 | // collected until the test chooses the moment. |
| 216 | let mut db = res!(setup::start_db( |
| 217 | db_root.clone(), |
| 218 | Some(cfg.clone()), |
| 219 | schms_input.clone(), |
| 220 | None, |
| 221 | false, // gc off |
| 222 | true, // wipe |
| 223 | )); |
| 224 | thread::sleep(Duration::from_secs(1)); |
| 225 | |
| 226 | let payload = |i: usize, round: usize| -> Dat { |
| 227 | let seed = fmt!("k{:06}r{:02}:", i, round); |
| 228 | let mut s = String::with_capacity(shape.val_bytes); |
| 229 | while s.len() < shape.val_bytes { |
| 230 | s.push_str(&seed); |
| 231 | } |
| 232 | s.truncate(shape.val_bytes); |
| 233 | Dat::Str(s) |
| 234 | }; |
| 235 | |
| 236 | // 1. Populate. |
| 237 | test!(sync_log::stream(), "Populating {} keys.", shape.keys); |
| 238 | let t = Instant::now(); |
| 239 | for i in 0..shape.keys { |
| 240 | res!(db.insert(dat!(fmt!("rec:{:06}", i)), payload(i, 0), user, schms2)); |
| 241 | } |
| 242 | test!(sync_log::stream(), "Populated in {:?}.", t.elapsed()); |
| 243 | |
| 244 | // 2. Supersede every second key, with collection still off. Half of |
| 245 | // every file is now garbage -- well past the trigger -- and the |
| 246 | // other half is untouched, so a later write to any of those keys |
| 247 | // still lands on the file it was written to and can set its |
| 248 | // collection going. |
| 249 | let num_old = ((shape.keys as f64) * OLD_FRAC) as usize; |
| 250 | test!(sync_log::stream(), "Superseding {} keys with collection off.", num_old); |
| 251 | let t = Instant::now(); |
| 252 | let mut i = 0; |
| 253 | while i < shape.keys { |
| 254 | res!(db.insert(dat!(fmt!("rec:{:06}", i)), payload(i, 1), user, schms2)); |
| 255 | i += 2; |
| 256 | } |
| 257 | test!(sync_log::stream(), "Superseded in {:?}.", t.elapsed()); |
| 258 | thread::sleep(Duration::from_secs(2)); |
| 259 | |
| 260 | // 3. Quiet baseline. The collector has done nothing and has nothing |
| 261 | // queued, so this is the cost of the walk alone. |
| 262 | let mut worst_quiet = Duration::ZERO; |
| 263 | for _ in 0..NUM_BASELINE { |
| 264 | let t = Instant::now(); |
| 265 | let entries = res!(db.scan(&ScanOpts::all(), schms2)); |
| 266 | let dt = t.elapsed(); |
| 267 | if dt > worst_quiet { |
| 268 | worst_quiet = dt; |
| 269 | } |
| 270 | if entries.len() != shape.keys { |
| 271 | return Err(err!( |
| 272 | "Quiet scan returned {} entries, expected {}.", |
| 273 | entries.len(), shape.keys; |
| 274 | Test, Mismatch)); |
| 275 | } |
| 276 | } |
| 277 | test!(sync_log::stream(), "Worst quiet scan latency {:?}.", worst_quiet); |
| 278 | |
| 279 | // 4. Switch collection on and dispatch the backlog. Each of these |
| 280 | // writes supersedes a record in a different file, and the file |
| 281 | // bot answers each by handing one collection to the zone's |
| 282 | // collector. They return as soon as the message is sent, so by |
| 283 | // the end of the burst the collector's queue holds the lot. |
| 284 | res!(ok!(db.updated_api()).activate_gc(true)); |
| 285 | test!(sync_log::stream(), "Collection enabled; dispatching the backlog."); |
| 286 | let stride = trigger_stride(shape); |
| 287 | let t = Instant::now(); |
| 288 | let mut triggers = 0; |
| 289 | let mut i = 1; |
| 290 | while i < shape.keys { |
| 291 | res!(db.insert(dat!(fmt!("rec:{:06}", i)), payload(i, 2), user, schms2)); |
| 292 | triggers += 1; |
| 293 | i += stride; |
| 294 | } |
| 295 | test!(sync_log::stream(), "{} triggers dispatched in {:?}.", triggers, t.elapsed()); |
| 296 | |
| 297 | // 5. Scan while the collector works through it. This is an operator |
| 298 | // refreshing a view. |
| 299 | let mut worst_loaded = Duration::ZERO; |
| 300 | let mut failures: Vec<String> = Vec::new(); |
| 301 | for _ in 0..NUM_LOADED { |
| 302 | let t = Instant::now(); |
| 303 | match db.scan(&ScanOpts::all(), schms2) { |
| 304 | Err(e) => { |
| 305 | let dt = t.elapsed(); |
| 306 | if dt > worst_loaded { |
| 307 | worst_loaded = dt; |
| 308 | } |
| 309 | failures.push(fmt!("scan failed after {:?}: {}", dt, e)); |
| 310 | }, |
| 311 | Ok(entries) => { |
| 312 | let dt = t.elapsed(); |
| 313 | if dt > worst_loaded { |
| 314 | worst_loaded = dt; |
| 315 | } |
| 316 | if entries.len() != shape.keys { |
| 317 | failures.push(fmt!( |
| 318 | "scan during collection returned {} entries, expected {}", |
| 319 | entries.len(), shape.keys)); |
| 320 | } |
| 321 | }, |
| 322 | } |
| 323 | } |
| 324 | test!(sync_log::stream(), "Worst scan latency during collection {:?}, {} misbehaved.", |
| 325 | worst_loaded, failures.len()); |
| 326 | for f in failures.iter().take(5) { |
| 327 | test!(sync_log::stream(), " {}", f); |
| 328 | } |
| 329 | |
| 330 | thread::sleep(Duration::from_secs(5)); |
| 331 | res!(db.shutdown()); |
| 332 | |
| 333 | // 6. Judgement. |
| 334 | if !failures.is_empty() { |
| 335 | return Err(err!( |
| 336 | "{} scan(s) issued during garbage collection did not behave: first was '{}'.", |
| 337 | failures.len(), failures[0]; |
| 338 | Test, Mismatch)); |
| 339 | } |
| 340 | // A scan served independently of the collector costs what the walk |
| 341 | // costs. The allowance is loose enough for the disk contention a |
| 342 | // running collection genuinely causes, and far tighter than a |
| 343 | // backlog of collections. |
| 344 | let allowance = worst_quiet * 3 + Duration::from_millis(200); |
| 345 | if worst_loaded > allowance { |
| 346 | return Err(err!( |
| 347 | "Worst scan latency during collection was {:?}, more than the allowance \ |
| 348 | of {:?} derived from the worst quiet latency of {:?}. The scan is \ |
| 349 | waiting on the collector, and on a larger store that wait passes the \ |
| 350 | user request deadline.", |
| 351 | worst_loaded, allowance, worst_quiet; |
| 352 | Test, Excessive)); |
| 353 | } |
| 354 | |
| 355 | test!(sync_log::stream(), |
| 356 | "Scan under garbage collection test passed: worst during collection {:?}, \ |
| 357 | worst quiet {:?}.", worst_loaded, worst_quiet); |
| 358 | Ok(()) |
| 359 | } |