oxedyne/fe2o3/fe2o3_o3db_sync/tests/close_order.rs
13.0 KiB, 1 run
created by r1870400018:61839, 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 close finishes the bots in the order they depend on one another: the readers, scanners and |
| 2 | //! collectors that ask the cache and file bots for things, then the writers, whose syncers release |
| 3 | //! records to the cache bots, and only then the cache and file bots. Each check holds one bot or |
| 4 | //! barrier up past the shutdown's three seconds and says what the callers waiting on it must be |
| 5 | //! told. The hooks are process-wide, hence a binary of its own with a single test. Found by QA |
| 6 | //! of lane o3i, 2026-09-24. |
| 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::{ |
| 18 | cfg::OzoneConfig, |
| 19 | constant, |
| 20 | }, |
| 21 | comm::msg::OzoneMsg, |
| 22 | data::core::RestSchemesInput, |
| 23 | test::{ |
| 24 | hooks, |
| 25 | setup::{ |
| 26 | self, |
| 27 | Uid, |
| 28 | UID_LEN, |
| 29 | }, |
| 30 | }, |
| 31 | }; |
| 32 | |
| 33 | use std::{ |
| 34 | collections::BTreeMap, |
| 35 | path::{ |
| 36 | Path, |
| 37 | PathBuf, |
| 38 | }, |
| 39 | sync::{ |
| 40 | Arc, |
| 41 | Mutex, |
| 42 | atomic::{ |
| 43 | AtomicBool, |
| 44 | Ordering, |
| 45 | }, |
| 46 | }, |
| 47 | thread, |
| 48 | time::{ |
| 49 | Duration, |
| 50 | Instant, |
| 51 | }, |
| 52 | }; |
| 53 | |
| 54 | type TestDb = O3db< |
| 55 | { UID_LEN }, |
| 56 | Uid, |
| 57 | (), |
| 58 | HashScheme, |
| 59 | HashScheme, |
| 60 | ChecksumScheme, |
| 61 | >; |
| 62 | |
| 63 | const HOLD: Duration = Duration::from_millis(4_000); // past the shutdown's 3 s |
| 64 | const READERS: u64 = 8; // threads, one read each queued |
| 65 | const KEYS: u64 = 20; |
| 66 | const CLOSE_WITHIN: Duration = Duration::from_secs(8); // one read timing out adds 5 s |
| 67 | const DEADLINE: Duration = Duration::from_secs(15); // stands in for the 120 s one |
| 68 | const TOLD_WITHIN: Duration = Duration::from_secs(10); |
| 69 | |
| 70 | #[test] |
| 71 | fn main() -> Outcome<()> { |
| 72 | log_set_level!("warn"); |
| 73 | // Whatever fails, the next check must not inherit the fault. |
| 74 | let reads = close_under_read_load_answers_the_reads_queued(); |
| 75 | hooks::set_insert_delay(Duration::ZERO); |
| 76 | let failed = close_tells_a_write_its_barrier_failed(); |
| 77 | hooks::set_barrier_failure(false); |
| 78 | hooks::set_barrier_delay(Duration::ZERO); |
| 79 | let slow = close_confirms_a_write_behind_a_slow_barrier(); |
| 80 | hooks::set_barrier_delay(Duration::ZERO); |
| 81 | log_finish_wait!(); |
| 82 | let failed: Vec<Error<ErrTag>> = [reads, failed, slow].into_iter() |
| 83 | .filter_map(|r| r.err()) |
| 84 | .collect(); |
| 85 | match failed.len() { |
| 86 | 0 => Ok(()), |
| 87 | 1 => match failed.into_iter().next() { |
| 88 | Some(e) => Err(e), |
| 89 | None => Ok(()), // unreachable |
| 90 | }, |
| 91 | n => Err(err!( |
| 92 | "{} checks failed: {:?}", n, failed; |
| 93 | Test)), |
| 94 | } |
| 95 | } |
| 96 | |
| 97 | /// The store is closed while eight reads wait at its one reader bot behind a cache bot held by an |
| 98 | /// insert. Every read asked before the close is answered with its value, and the close takes |
| 99 | /// about as long as the hold. The cache bots used to be finished while reads were still queued |
| 100 | /// at the reader bot, and once nothing was drained any more each of those reads waited out |
| 101 | /// `BOT_REQUEST_TIMEOUT` on a cache bot that had ended: 38.65 s for this close, where it took |
| 102 | /// 4.70 s when the queue was drained. |
| 103 | fn close_under_read_load_answers_the_reads_queued() -> Outcome<()> { |
| 104 | let root = res!(fresh("./test_db_close_order_reads")); |
| 105 | let db = res!(open(&root, res!(config()))); |
| 106 | for i in 0..KEYS { |
| 107 | res!(db.insert(key(i), dat!(i), Uid::default(), None)); |
| 108 | } |
| 109 | |
| 110 | hooks::set_insert_delay(HOLD); |
| 111 | let writing = db.clone(); |
| 112 | let writer = res!(thread::Builder::new().name(fmt!("close order writer")).spawn(move || { |
| 113 | // Its answer is not the point: it holds the cache bot. |
| 114 | let _ = writing.insert(key(KEYS), dat!(KEYS), Uid::default(), None); |
| 115 | })); |
| 116 | thread::sleep(Duration::from_millis(50)); |
| 117 | |
| 118 | let stop = Arc::new(AtomicBool::new(false)); |
| 119 | let asked = Arc::new(Mutex::new(Vec::new())); // (asked at, key, answer) |
| 120 | let mut readers = Vec::new(); |
| 121 | for r in 0..READERS { |
| 122 | let (db, stop, asked) = (db.clone(), stop.clone(), asked.clone()); |
| 123 | readers.push(res!(thread::Builder::new().name(fmt!("close order reader {}", r)).spawn( |
| 124 | move || { |
| 125 | let mut i = r; |
| 126 | while !stop.load(Ordering::Relaxed) { |
| 127 | let at = Instant::now(); |
| 128 | let answer = match db.get(&key(i % KEYS), None) { |
| 129 | Ok(Some((v, _))) if v == dat!(i % KEYS) => Ok(()), |
| 130 | Ok(Some((v, _))) => Err(fmt!("the value {:?}", v)), |
| 131 | Ok(None) => Err(fmt!("nothing")), |
| 132 | Err(e) => Err(fmt!("{}", e)), |
| 133 | }; |
| 134 | if let Ok(mut asked) = asked.lock() { |
| 135 | asked.push((at, i % KEYS, answer)); |
| 136 | } |
| 137 | i += 1; |
| 138 | } |
| 139 | }))); |
| 140 | } |
| 141 | thread::sleep(Duration::from_millis(300)); |
| 142 | |
| 143 | let begun = Instant::now(); |
| 144 | let closed = db.close(); |
| 145 | let took = begun.elapsed(); |
| 146 | stop.store(true, Ordering::Relaxed); |
| 147 | hooks::set_insert_delay(Duration::ZERO); |
| 148 | for reader in readers { |
| 149 | let _ = reader.join(); |
| 150 | } |
| 151 | let _ = writer.join(); |
| 152 | res!(closed); |
| 153 | let asked = match asked.lock() { |
| 154 | Ok(asked) => asked.clone(), |
| 155 | Err(_) => return Err(err!("A reader thread panicked holding the answers."; Test, Thread)), |
| 156 | }; |
| 157 | let queued: Vec<_> = asked.iter().filter(|(at, _, _)| *at < begun).collect(); |
| 158 | msg!("A close with {} reads queued behind a cache bot held for {:?} took {:?}.", |
| 159 | queued.len(), HOLD, took); |
| 160 | let wrong: Vec<_> = queued.iter() |
| 161 | .filter_map(|(_, k, a)| a.as_ref().err().map(|e| fmt!("key {}: {}", k, e))) |
| 162 | .collect(); |
| 163 | if !wrong.is_empty() { |
| 164 | return Err(err!( |
| 165 | "Of {} reads asked before the store was closed, which took {:?}, {} were not \ |
| 166 | answered with their value: {:?}", queued.len(), took, wrong.len(), wrong; |
| 167 | Test, Missing)); |
| 168 | } |
| 169 | if took > CLOSE_WITHIN { |
| 170 | return Err(err!( |
| 171 | "A close with {} reads queued behind a cache bot held for {:?} took {:?}, where it \ |
| 172 | should take about as long as the hold and at most {:?}.", |
| 173 | queued.len(), HOLD, took, CLOSE_WITHIN; |
| 174 | Test, Timeout)); |
| 175 | } |
| 176 | Ok(()) |
| 177 | } |
| 178 | |
| 179 | /// The store is closed while a written record waits on a barrier that fails, after the |
| 180 | /// shutdown's three seconds. Its caller is told the barrier failed, tagged `Unconfirmed`, as |
| 181 | /// soon as it does. The failure is told by the cache bot the syncer releases the record to, and |
| 182 | /// the cache bots used to be finished once the three seconds ran out, so it was never told: its |
| 183 | /// caller waited out the durability deadline and heard a timeout. |
| 184 | fn close_tells_a_write_its_barrier_failed() -> Outcome<()> { |
| 185 | let root = res!(fresh("./test_db_close_order_failed")); |
| 186 | let mut cfg = res!(config()); |
| 187 | cfg.sync_on_write = true; |
| 188 | let db = res!(open(&root, cfg)); |
| 189 | res!(db.insert(key(1), dat!(1u64), Uid::default(), None)); |
| 190 | |
| 191 | hooks::set_barrier_delay(HOLD); |
| 192 | hooks::set_barrier_failure(true); |
| 193 | let writer = res!(write_during_close(&db, key(2), dat!(2u64))); |
| 194 | thread::sleep(Duration::from_millis(100)); |
| 195 | let begun = Instant::now(); |
| 196 | let closed = db.close(); |
| 197 | let took = begun.elapsed(); |
| 198 | let heard = match writer.join() { |
| 199 | Ok(heard) => heard, |
| 200 | Err(_) => return Err(err!("The writing thread panicked."; Test, Thread)), |
| 201 | }; |
| 202 | hooks::set_barrier_failure(false); |
| 203 | hooks::set_barrier_delay(Duration::ZERO); |
| 204 | msg!("A close with a write behind a barrier failing after {:?} took {:?}; the write heard \ |
| 205 | after {:?}.", HOLD, took, heard.0); |
| 206 | res!(closed); |
| 207 | match heard { |
| 208 | (_, Ok(())) => Err(err!( |
| 209 | "A write whose barrier failed as the store closed was confirmed."; |
| 210 | Test, Unexpected)), |
| 211 | (took, Err(e)) => { |
| 212 | if !e.tags().contains(&ErrTag::Unconfirmed) || e.tags().contains(&ErrTag::Timeout) { |
| 213 | return Err(err!( |
| 214 | "A write whose barrier failed as the store closed was told after {:?}, tagged \ |
| 215 | {:?}, where it should hear the barrier failed, tagged Unconfirmed and not \ |
| 216 | Timeout: {}", took, e.tags(), e; |
| 217 | Test, Mismatch)); |
| 218 | } |
| 219 | if took > TOLD_WITHIN { |
| 220 | return Err(err!( |
| 221 | "A write whose barrier failed as the store closed was told after {:?}, where \ |
| 222 | the barrier failed after {:?}.", took, HOLD; |
| 223 | Test, Timeout)); |
| 224 | } |
| 225 | Ok(()) |
| 226 | }, |
| 227 | } |
| 228 | } |
| 229 | |
| 230 | /// The store is closed while a written record waits on a barrier that succeeds, after the |
| 231 | /// shutdown's three seconds. Its caller is told it is durable, and it is there after a restart. |
| 232 | /// The cache bots used to be finished once the three seconds ran out, so the record the syncer |
| 233 | /// released afterwards was never answered, and its caller waited out the durability deadline. |
| 234 | fn close_confirms_a_write_behind_a_slow_barrier() -> Outcome<()> { |
| 235 | let root = res!(fresh("./test_db_close_order_slow")); |
| 236 | let mut cfg = res!(config()); |
| 237 | cfg.sync_on_write = true; |
| 238 | let db = res!(open(&root, cfg.clone())); |
| 239 | res!(db.insert(key(1), dat!(1u64), Uid::default(), None)); |
| 240 | |
| 241 | hooks::set_barrier_delay(HOLD); |
| 242 | let writer = res!(write_during_close(&db, key(3), dat!(3u64))); |
| 243 | thread::sleep(Duration::from_millis(100)); |
| 244 | let begun = Instant::now(); |
| 245 | let closed = db.close(); |
| 246 | let took = begun.elapsed(); |
| 247 | let heard = match writer.join() { |
| 248 | Ok(heard) => heard, |
| 249 | Err(_) => return Err(err!("The writing thread panicked."; Test, Thread)), |
| 250 | }; |
| 251 | hooks::set_barrier_delay(Duration::ZERO); |
| 252 | msg!("A close with a write behind a barrier succeeding after {:?} took {:?}; the write heard \ |
| 253 | after {:?}.", HOLD, took, heard.0); |
| 254 | res!(closed); |
| 255 | match heard { |
| 256 | (took, Err(e)) => return Err(err!( |
| 257 | "A write whose slow barrier succeeded as the store closed was told after {:?}: {}", |
| 258 | took, e; |
| 259 | Test, Missing)), |
| 260 | (took, Ok(())) if took > TOLD_WITHIN => return Err(err!( |
| 261 | "A write whose slow barrier succeeded as the store closed was confirmed after {:?}, \ |
| 262 | where the barrier took {:?}.", took, HOLD; |
| 263 | Test, Timeout)), |
| 264 | _ => (), |
| 265 | } |
| 266 | let db = res!(open(&root, cfg)); |
| 267 | match res!(db.get(&key(3), None)) { |
| 268 | Some((v, _)) => { |
| 269 | res!(db.close()); |
| 270 | req!(v, dat!(3u64), "A write confirmed as the store closed."); |
| 271 | Ok(()) |
| 272 | }, |
| 273 | None => { |
| 274 | let _ = db.close(); // the check has failed already, and says why |
| 275 | Err(err!( |
| 276 | "A write confirmed as the store closed is missing after a restart."; |
| 277 | Test, Missing)) |
| 278 | }, |
| 279 | } |
| 280 | } |
| 281 | |
| 282 | /// Writes on a thread of its own, returning how long it waited for its final answer and what |
| 283 | /// that answer was. |
| 284 | fn write_during_close( |
| 285 | db: &TestDb, |
| 286 | k: Dat, |
| 287 | v: Dat, |
| 288 | ) |
| 289 | -> Outcome<thread::JoinHandle<(Duration, Outcome<()>)>> |
| 290 | { |
| 291 | let writing = db.clone(); |
| 292 | Ok(res!(thread::Builder::new().name(fmt!("close order writer")).spawn(move || { |
| 293 | let begun = Instant::now(); |
| 294 | let result = write(&writing, k, v); |
| 295 | (begun.elapsed(), result) |
| 296 | }))) |
| 297 | } |
| 298 | |
| 299 | /// Writes and waits `DEADLINE` for durability. The final answer's error is passed on as it is, |
| 300 | /// since `Error::tags` reports only the outermost error's tags. |
| 301 | fn write(db: &TestDb, k: Dat, v: Dat) -> Outcome<()> { |
| 302 | let resp = res!(db.api().store(k, v, Uid::default())); |
| 303 | let n = match res!(resp.recv_timeout(constant::USER_REQUEST_TIMEOUT)) { |
| 304 | OzoneMsg::Chunks(n) => n, |
| 305 | msg => return Err(err!( |
| 306 | "Expected the record count, received {:?}.", msg; Test, Unexpected)), |
| 307 | }; |
| 308 | ok!(resp.recv_write_acks(n, constant::USER_REQUEST_TIMEOUT, DEADLINE)); |
| 309 | Ok(()) |
| 310 | } |
| 311 | |
| 312 | fn key(i: u64) -> Dat { |
| 313 | dat!(fmt!("close order key {:03}", i)) |
| 314 | } |
| 315 | |
| 316 | /// One zone with one bot of each kind, so every read queues at the one reader bot and every |
| 317 | /// record reaches the one cache bot. |
| 318 | fn config() -> Outcome<OzoneConfig> { |
| 319 | let mut cfg = res!(setup::default_cfg()); |
| 320 | cfg.num_zones = 1; |
| 321 | cfg.num_cbots_per_zone = 1; |
| 322 | cfg.num_fbots_per_zone = 1; |
| 323 | cfg.num_igbots_per_zone = 1; |
| 324 | cfg.num_rbots_per_zone = 1; |
| 325 | cfg.num_wbots_per_zone = 1; |
| 326 | cfg.zone_overrides = BTreeMap::new(); |
| 327 | Ok(cfg) |
| 328 | } |
| 329 | |
| 330 | fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> { |
| 331 | RestSchemesInput::new( |
| 332 | None::<()>, |
| 333 | None::<HashScheme>, |
| 334 | None::<HashScheme>, |
| 335 | Some(ChecksumScheme::new_crc32()), |
| 336 | ) |
| 337 | } |
| 338 | |
| 339 | /// An empty directory of this test's own. |
| 340 | fn fresh(dir: &str) -> Outcome<PathBuf> { |
| 341 | let _ = std::fs::remove_dir_all(dir); // absent the first time |
| 342 | res!(std::fs::create_dir_all(dir)); |
| 343 | Ok(res!(Path::new(dir).canonicalize())) |
| 344 | } |
| 345 | |
| 346 | fn open(root: &Path, cfg: OzoneConfig) -> Outcome<TestDb> { |
| 347 | let mut db = res!(TestDb::new(root.to_path_buf(), Some(cfg), schemes(), Uid::default())); |
| 348 | res!(db.start("test")); |
| 349 | Ok(db) |
| 350 | } |