oxedyne/fe2o3/fe2o3_o3db_sync/tests/wrong_record.rs
15.2 KiB, 1 run
created by r1870400018:61843, 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 | //! Reads, writes and garbage collection all at once, with every record the same size, and every |
| 2 | //! value judged against something outside the store: the key it was asked for, and the versions |
| 3 | //! its writer issued and had acknowledged. A read that raced a collection returned another key's |
| 4 | //! value, or an older version of its own, about once in three collections (QA of lane sto, |
| 5 | //! 2026-09-23, whose stress harness this is); with collection off there were none. Each writer |
| 6 | //! thread owns its keys, so versions and stamps agree. The long variant is ignored by default. |
| 7 | |
| 8 | use oxedyne_fe2o3_core::prelude::*; |
| 9 | use oxedyne_fe2o3_hash::{ |
| 10 | csum::ChecksumScheme, |
| 11 | hash::HashScheme, |
| 12 | }; |
| 13 | use oxedyne_fe2o3_iop_db::api::Database; |
| 14 | use oxedyne_fe2o3_jdat::prelude::*; |
| 15 | use oxedyne_fe2o3_o3db_sync::{ |
| 16 | O3db, |
| 17 | base::cfg::OzoneConfig, |
| 18 | comm::response::Wait, |
| 19 | data::core::RestSchemesInput, |
| 20 | test::setup::{ |
| 21 | self, |
| 22 | Uid, |
| 23 | UID_LEN, |
| 24 | }, |
| 25 | }; |
| 26 | |
| 27 | use std::{ |
| 28 | collections::BTreeMap, |
| 29 | path::{ |
| 30 | Path, |
| 31 | PathBuf, |
| 32 | }, |
| 33 | sync::{ |
| 34 | Arc, |
| 35 | atomic::{ |
| 36 | AtomicBool, |
| 37 | AtomicU64, |
| 38 | Ordering, |
| 39 | }, |
| 40 | }, |
| 41 | thread, |
| 42 | time::{ |
| 43 | Duration, |
| 44 | Instant, |
| 45 | }, |
| 46 | }; |
| 47 | |
| 48 | type TestDb = O3db< |
| 49 | { UID_LEN }, |
| 50 | Uid, |
| 51 | (), |
| 52 | HashScheme, |
| 53 | HashScheme, |
| 54 | ChecksumScheme, |
| 55 | >; |
| 56 | |
| 57 | fn env_u64(name: &str, dflt: u64) -> u64 { |
| 58 | match std::env::var(name) { |
| 59 | Ok(s) => match s.parse::<u64>() { |
| 60 | Ok(n) => n, |
| 61 | Err(_) => dflt, |
| 62 | }, |
| 63 | Err(_) => dflt, |
| 64 | } |
| 65 | } |
| 66 | |
| 67 | #[derive(Default)] |
| 68 | struct Counts { |
| 69 | writes: AtomicU64, |
| 70 | write_errs: AtomicU64, |
| 71 | reads: AtomicU64, |
| 72 | wrong_key: AtomicU64, |
| 73 | stale: AtomicU64, |
| 74 | phantom: AtomicU64, |
| 75 | none: AtomicU64, |
| 76 | read_errs: AtomicU64, |
| 77 | bad_shape: AtomicU64, |
| 78 | shrinks: AtomicU64, |
| 79 | removed: AtomicU64, |
| 80 | clears: AtomicU64, |
| 81 | } |
| 82 | |
| 83 | impl Counts { |
| 84 | fn line(&self) -> String { |
| 85 | fmt!("writes {} write_errs {} reads {} wrong_key {} stale {} phantom {} none {} read_errs {} \ |
| 86 | bad_shape {} | file shrinks {} removed {} clears {}", |
| 87 | self.writes.load(Ordering::Relaxed), |
| 88 | self.write_errs.load(Ordering::Relaxed), |
| 89 | self.reads.load(Ordering::Relaxed), |
| 90 | self.wrong_key.load(Ordering::Relaxed), |
| 91 | self.stale.load(Ordering::Relaxed), |
| 92 | self.phantom.load(Ordering::Relaxed), |
| 93 | self.none.load(Ordering::Relaxed), |
| 94 | self.read_errs.load(Ordering::Relaxed), |
| 95 | self.bad_shape.load(Ordering::Relaxed), |
| 96 | self.shrinks.load(Ordering::Relaxed), |
| 97 | self.removed.load(Ordering::Relaxed), |
| 98 | self.clears.load(Ordering::Relaxed), |
| 99 | ) |
| 100 | } |
| 101 | } |
| 102 | |
| 103 | struct Plan { |
| 104 | nkeys: usize, |
| 105 | vbytes: usize, |
| 106 | nw: usize, |
| 107 | nr: usize, |
| 108 | secs: u64, |
| 109 | clear_ms: u64, |
| 110 | phases: u64, |
| 111 | min_gc: u64, // collections a run must see to have tested anything |
| 112 | } |
| 113 | |
| 114 | fn key(i: usize) -> Dat { dat!(fmt!("wrong record key {:04}", i)) } |
| 115 | |
| 116 | fn value(i: usize, ver: u64, len: usize) -> Dat { |
| 117 | let mut v = vec![0u8; len]; |
| 118 | v[0] = (i >> 8) as u8; |
| 119 | v[1] = i as u8; |
| 120 | v[2..10].copy_from_slice(&ver.to_be_bytes()); |
| 121 | for j in 10..len { |
| 122 | v[j] = (i as u8) ^ (ver as u8) ^ (j as u8); |
| 123 | } |
| 124 | Dat::BU32(v) |
| 125 | } |
| 126 | |
| 127 | /// The key index and version a value says it holds, if it has the shape of one of ours. |
| 128 | fn parse(d: &Dat, len: usize) -> Option<(usize, u64)> { |
| 129 | let v = d.bytes_ref()?; |
| 130 | if v.len() != len { return None; } |
| 131 | let i = ((v[0] as usize) << 8) | (v[1] as usize); |
| 132 | let mut b = [0u8; 8]; |
| 133 | b.copy_from_slice(&v[2..10]); |
| 134 | let ver = u64::from_be_bytes(b); |
| 135 | for j in 10..len { |
| 136 | if v[j] != (i as u8) ^ (ver as u8) ^ (j as u8) { return None; } |
| 137 | } |
| 138 | Some((i, ver)) |
| 139 | } |
| 140 | |
| 141 | struct Rng(u64); |
| 142 | impl Rng { |
| 143 | fn next(&mut self) -> u64 { |
| 144 | let mut x = self.0; |
| 145 | x ^= x << 13; |
| 146 | x ^= x >> 7; |
| 147 | x ^= x << 17; |
| 148 | self.0 = x; |
| 149 | x |
| 150 | } |
| 151 | } |
| 152 | |
| 153 | fn config() -> Outcome<OzoneConfig> { |
| 154 | let mut cfg = res!(setup::default_cfg()); |
| 155 | cfg.num_zones = 1; |
| 156 | cfg.num_cbots_per_zone = 2; |
| 157 | cfg.num_fbots_per_zone = 2; |
| 158 | cfg.num_igbots_per_zone = 2; |
| 159 | cfg.num_rbots_per_zone = 2; |
| 160 | cfg.num_wbots_per_zone = 2; |
| 161 | let fbytes = 8_000; |
| 162 | cfg.data_file_max_bytes = fbytes; |
| 163 | cfg.rest_chunk_threshold = fbytes * 7 / 10; // the store refuses over 80%, values are far below |
| 164 | cfg.zone_overrides = BTreeMap::new(); |
| 165 | cfg.sync_on_write = true; |
| 166 | Ok(cfg) |
| 167 | } |
| 168 | |
| 169 | fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> { |
| 170 | RestSchemesInput::new( |
| 171 | None::<()>, |
| 172 | None::<HashScheme>, |
| 173 | None::<HashScheme>, |
| 174 | Some(ChecksumScheme::new_crc32()), |
| 175 | ) |
| 176 | } |
| 177 | |
| 178 | fn open(root: &Path, wipe: bool) -> Outcome<TestDb> { |
| 179 | // Collection off is the control: no fault of any kind in 3.1 M reads (2026-09-23). |
| 180 | let gc_on = env_u64("WRONG_RECORD_GC", 1) == 1; |
| 181 | setup::start_db(root.to_path_buf(), Some(res!(config())), schemes(), None, gc_on, wipe) |
| 182 | } |
| 183 | |
| 184 | fn classify<M>( |
| 185 | c: &Counts, |
| 186 | got: Outcome<Option<(Dat, M)>>, |
| 187 | i: usize, |
| 188 | lo: u64, |
| 189 | hi: u64, |
| 190 | len: usize, |
| 191 | who: &str, |
| 192 | ) { |
| 193 | c.reads.fetch_add(1, Ordering::Relaxed); |
| 194 | match got { |
| 195 | Err(e) => { |
| 196 | let n = c.read_errs.fetch_add(1, Ordering::Relaxed); |
| 197 | // Not a silent fault, but worth seeing: a read that gave up says why. |
| 198 | if n < 5 { msg!("{} read error on key {}: {}", who, i, e); } |
| 199 | }, |
| 200 | Ok(None) => { |
| 201 | if lo > 0 { |
| 202 | let n = c.none.fetch_add(1, Ordering::Relaxed); |
| 203 | if n < 5 { test!(sync_log::stream(), "{} key {} read None, acked {}", who, i, lo); } |
| 204 | } |
| 205 | }, |
| 206 | Ok(Some((d, _))) => match parse(&d, len) { |
| 207 | None => { |
| 208 | let n = c.bad_shape.fetch_add(1, Ordering::Relaxed); |
| 209 | if n < 5 { test!(sync_log::stream(), "{} key {} read a value of no known shape: {:?}", who, i, d); } |
| 210 | }, |
| 211 | Some((j, ver)) => { |
| 212 | if j != i { |
| 213 | let n = c.wrong_key.fetch_add(1, Ordering::Relaxed); |
| 214 | if n < 10 { test!(sync_log::stream(), "{} WRONG KEY: asked {} got key {} version {}", who, i, j, ver); } |
| 215 | } else if ver < lo { |
| 216 | let n = c.stale.fetch_add(1, Ordering::Relaxed); |
| 217 | if n < 10 { test!(sync_log::stream(), "{} STALE: key {} got version {} < acked {}", who, i, ver, lo); } |
| 218 | } else if ver > hi { |
| 219 | let n = c.phantom.fetch_add(1, Ordering::Relaxed); |
| 220 | if n < 10 { test!(sync_log::stream(), "{} PHANTOM: key {} got version {} > issued {}", who, i, ver, hi); } |
| 221 | } |
| 222 | }, |
| 223 | }, |
| 224 | } |
| 225 | } |
| 226 | |
| 227 | fn run_phase( |
| 228 | db: &TestDb, |
| 229 | root: &Path, |
| 230 | plan: &Plan, |
| 231 | issued: &Arc<Vec<AtomicU64>>, |
| 232 | acked: &Arc<Vec<AtomicU64>>, |
| 233 | c: &Arc<Counts>, |
| 234 | ) -> Outcome<()> { |
| 235 | let stop = Arc::new(AtomicBool::new(false)); |
| 236 | let mut hs = Vec::new(); |
| 237 | for w in 0..plan.nw { |
| 238 | let db = db.clone(); |
| 239 | let (issued, acked, c, stop) = (issued.clone(), acked.clone(), c.clone(), stop.clone()); |
| 240 | let (nkeys, nw, vb) = (plan.nkeys, plan.nw, plan.vbytes); |
| 241 | hs.push(res!(thread::Builder::new().name(fmt!("qa writer {}", w)).spawn(move || { |
| 242 | let mine: Vec<usize> = (0..nkeys).filter(|k| k % nw == w).collect(); |
| 243 | let mut rng = Rng(0x9e3779b97f4a7c15 ^ (w as u64 + 1)); |
| 244 | while !stop.load(Ordering::Relaxed) { |
| 245 | let k = mine[(rng.next() as usize) % mine.len()]; |
| 246 | let ver = issued[k].load(Ordering::Relaxed) + 1; |
| 247 | issued[k].store(ver, Ordering::Relaxed); |
| 248 | c.writes.fetch_add(1, Ordering::Relaxed); |
| 249 | match db.insert(key(k), value(k, ver, vb), Uid::default(), None) { |
| 250 | Ok(_) => { |
| 251 | acked[k].store(ver, Ordering::Relaxed); |
| 252 | // Read one's own write back straight away. |
| 253 | let got = db.get(&key(k), None); |
| 254 | let hi = issued[k].load(Ordering::Relaxed); |
| 255 | classify(&c, got, k, ver, hi, vb, "rw"); |
| 256 | }, |
| 257 | Err(e) => { |
| 258 | let n = c.write_errs.fetch_add(1, Ordering::Relaxed); |
| 259 | if n < 5 { test!(sync_log::stream(), "write error on key {} v{}: {}", k, ver, e); } |
| 260 | }, |
| 261 | } |
| 262 | } |
| 263 | }))); |
| 264 | } |
| 265 | for r in 0..plan.nr { |
| 266 | let db = db.clone(); |
| 267 | let (issued, acked, c, stop) = (issued.clone(), acked.clone(), c.clone(), stop.clone()); |
| 268 | let (nkeys, vb) = (plan.nkeys, plan.vbytes); |
| 269 | hs.push(res!(thread::Builder::new().name(fmt!("qa reader {}", r)).spawn(move || { |
| 270 | let mut rng = Rng(0x2545f4914f6cdd1d ^ (r as u64 + 101)); |
| 271 | while !stop.load(Ordering::Relaxed) { |
| 272 | let k = (rng.next() as usize) % nkeys; |
| 273 | let lo = acked[k].load(Ordering::Relaxed); |
| 274 | let got = db.get(&key(k), None); |
| 275 | let hi = issued[k].load(Ordering::Relaxed); |
| 276 | classify(&c, got, k, lo, hi, vb, "r"); |
| 277 | } |
| 278 | }))); |
| 279 | } |
| 280 | // Monitor: clears cached values so reads go to the files, and counts collections by |
| 281 | // watching data files shrink or vanish. |
| 282 | { |
| 283 | let db = db.clone(); |
| 284 | let (c, stop) = (c.clone(), stop.clone()); |
| 285 | let zdir = res!(config()).zone_root(root).join("zone_001"); |
| 286 | let clear_ms = plan.clear_ms; |
| 287 | hs.push(res!(thread::Builder::new().name(fmt!("qa monitor")).spawn(move || { |
| 288 | let mut sizes: BTreeMap<PathBuf, u64> = BTreeMap::new(); |
| 289 | let mut last_clear = Instant::now(); |
| 290 | while !stop.load(Ordering::Relaxed) { |
| 291 | if clear_ms > 0 && last_clear.elapsed() >= Duration::from_millis(clear_ms) { |
| 292 | last_clear = Instant::now(); |
| 293 | if db.api().clear_cache_values(Wait { |
| 294 | max_wait: Duration::from_secs(5), |
| 295 | check_interval: Duration::from_millis(10), |
| 296 | }).is_ok() { |
| 297 | c.clears.fetch_add(1, Ordering::Relaxed); |
| 298 | } |
| 299 | } |
| 300 | let mut now: BTreeMap<PathBuf, u64> = BTreeMap::new(); |
| 301 | if let Ok(rd) = std::fs::read_dir(&zdir) { |
| 302 | for e in rd.flatten() { |
| 303 | let p = e.path(); |
| 304 | let is_dat = match p.extension() { Some(x) => x == "dat", None => false }; |
| 305 | if is_dat { |
| 306 | if let Ok(m) = std::fs::metadata(&p) { now.insert(p, m.len()); } |
| 307 | } |
| 308 | } |
| 309 | } |
| 310 | for (p, s) in &sizes { |
| 311 | match now.get(p) { |
| 312 | Some(n) if n < s => { c.shrinks.fetch_add(1, Ordering::Relaxed); }, |
| 313 | None => { c.removed.fetch_add(1, Ordering::Relaxed); }, |
| 314 | _ => (), |
| 315 | } |
| 316 | } |
| 317 | sizes = now; |
| 318 | thread::sleep(Duration::from_millis(5)); |
| 319 | } |
| 320 | }))); |
| 321 | } |
| 322 | thread::sleep(Duration::from_secs(plan.secs)); |
| 323 | stop.store(true, Ordering::Relaxed); |
| 324 | for h in hs { |
| 325 | if h.join().is_err() { test!(sync_log::stream(), "a thread panicked"); } |
| 326 | } |
| 327 | Ok(()) |
| 328 | } |
| 329 | |
| 330 | /// Every key must read back a version its writer issued no earlier than the last it had |
| 331 | /// acknowledged. After a reopen this is the durability check. |
| 332 | fn verify( |
| 333 | db: &TestDb, |
| 334 | plan: &Plan, |
| 335 | issued: &Arc<Vec<AtomicU64>>, |
| 336 | acked: &Arc<Vec<AtomicU64>>, |
| 337 | label: &str, |
| 338 | ) -> Outcome<(u64, u64, String)> { |
| 339 | let c = Counts::default(); |
| 340 | for k in 0..plan.nkeys { |
| 341 | let lo = acked[k].load(Ordering::Relaxed); |
| 342 | let hi = issued[k].load(Ordering::Relaxed); |
| 343 | let got = db.get(&key(k), None); |
| 344 | classify(&c, got, k, lo, hi, plan.vbytes, label); |
| 345 | } |
| 346 | let bad = c.wrong_key.load(Ordering::Relaxed) + c.stale.load(Ordering::Relaxed) |
| 347 | + c.phantom.load(Ordering::Relaxed) + c.none.load(Ordering::Relaxed) |
| 348 | + c.read_errs.load(Ordering::Relaxed) + c.bad_shape.load(Ordering::Relaxed); |
| 349 | Ok((c.reads.load(Ordering::Relaxed), bad, fmt!("{}: {}", label, c.line()))) |
| 350 | } |
| 351 | |
| 352 | #[test] |
| 353 | fn main() -> Outcome<()> { |
| 354 | run("./test_db_wrong_record", Plan { |
| 355 | nkeys: 40, |
| 356 | vbytes: 100, |
| 357 | nw: 4, |
| 358 | nr: 4, |
| 359 | secs: env_u64("WRONG_RECORD_SECS", 12), |
| 360 | clear_ms: 20, |
| 361 | phases: 2, |
| 362 | min_gc: 20, |
| 363 | }) |
| 364 | } |
| 365 | |
| 366 | /// The stress the fault was found with, three minutes of it. |
| 367 | #[test] |
| 368 | #[ignore] |
| 369 | fn long() -> Outcome<()> { |
| 370 | run("./test_db_wrong_record_long", Plan { |
| 371 | nkeys: 40, |
| 372 | vbytes: 100, |
| 373 | nw: 4, |
| 374 | nr: 4, |
| 375 | secs: env_u64("WRONG_RECORD_SECS", 60), |
| 376 | clear_ms: 20, |
| 377 | phases: 3, |
| 378 | min_gc: 100, |
| 379 | }) |
| 380 | } |
| 381 | |
| 382 | fn run(dir: &str, plan: Plan) -> Outcome<()> { |
| 383 | log_set_level!("error"); |
| 384 | let _ = std::fs::remove_dir_all(dir); // absent the first time |
| 385 | res!(std::fs::create_dir_all(dir)); |
| 386 | let root = res!(Path::new(dir).canonicalize()); |
| 387 | let issued: Arc<Vec<AtomicU64>> = Arc::new((0..plan.nkeys).map(|_| AtomicU64::new(0)).collect()); |
| 388 | let acked: Arc<Vec<AtomicU64>> = Arc::new((0..plan.nkeys).map(|_| AtomicU64::new(0)).collect()); |
| 389 | let mut silent = 0u64; |
| 390 | let mut collections = 0u64; |
| 391 | let mut lines = Vec::new(); |
| 392 | let mut db = res!(open(&root, true)); |
| 393 | for p in 0..plan.phases { |
| 394 | let c = Arc::new(Counts::default()); |
| 395 | let begun = Instant::now(); |
| 396 | res!(run_phase(&db, &root, &plan, &issued, &acked, &c)); |
| 397 | lines.push(fmt!("phase {} ({:?}): {}", p, begun.elapsed(), c.line())); |
| 398 | silent += c.wrong_key.load(Ordering::Relaxed) + c.stale.load(Ordering::Relaxed) |
| 399 | + c.phantom.load(Ordering::Relaxed) + c.none.load(Ordering::Relaxed) |
| 400 | + c.bad_shape.load(Ordering::Relaxed); |
| 401 | collections += c.shrinks.load(Ordering::Relaxed) + c.removed.load(Ordering::Relaxed); |
| 402 | thread::sleep(Duration::from_secs(2)); |
| 403 | let (_, bad, line) = res!(verify(&db, &plan, &issued, &acked, &fmt!("phase {} settled", p))); |
| 404 | silent += bad; |
| 405 | lines.push(line); |
| 406 | res!(db.close()); |
| 407 | db = res!(open(&root, false)); |
| 408 | let (_, bad, line) = res!(verify(&db, &plan, &issued, &acked, &fmt!("phase {} reopened", p))); |
| 409 | silent += bad; |
| 410 | lines.push(line); |
| 411 | } |
| 412 | res!(db.close()); |
| 413 | // The counts are the finding, so they are shown whichever way the run goes. |
| 414 | for line in &lines { |
| 415 | msg!("{}", line); |
| 416 | } |
| 417 | log_finish_wait!(); |
| 418 | if silent > 0 { |
| 419 | return Err(err!( |
| 420 | "{} reads returned a value other than the one asked for (another key's, an older \ |
| 421 | version, one never written, or none), over {} collections: {:?}", |
| 422 | silent, collections, lines; |
| 423 | Test, Mismatch)); |
| 424 | } |
| 425 | if collections < plan.min_gc { |
| 426 | return Err(err!( |
| 427 | "Only {} collections ran, fewer than the {} this run needs to have tested reads \ |
| 428 | racing them: {:?}", collections, plan.min_gc, lines; |
| 429 | Test, Missing)); |
| 430 | } |
| 431 | Ok(()) |
| 432 | } |