oxedyne/fe2o3/fe2o3_o3db_sync/tests/store_timeouts.rs
17.4 KiB, 11 runs
created by r1870400018:61247, 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 | //! On a busy machine Oregami's store reported writes as failed that went on to land, and came up |
| 2 | //! answering nothing (2026-09-23). Each fault is reproduced here without the machine having to be |
| 3 | //! busy: `test::hooks` holds every durability barrier, or the supervisor, for longer than the |
| 4 | //! deadline that used to expire. The hooks are process-wide, which is why this is a test binary |
| 5 | //! of its own with a single test. |
| 6 | |
| 7 | use oxedyne_fe2o3_core::prelude::*; |
| 8 | use oxedyne_fe2o3_hash::{ |
| 9 | csum::ChecksumScheme, |
| 10 | hash::HashScheme, |
| 11 | }; |
| 12 | use oxedyne_fe2o3_iop_db::api::Database; |
| 13 | use oxedyne_fe2o3_jdat::prelude::*; |
| 14 | use oxedyne_fe2o3_o3db_sync::{ |
| 15 | O3db, |
| 16 | base::{ |
| 17 | cfg::OzoneConfig, |
| 18 | constant, |
| 19 | }, |
| 20 | comm::msg::OzoneMsg, |
| 21 | data::core::RestSchemesInput, |
| 22 | test::{ |
| 23 | hooks, |
| 24 | setup::{ |
| 25 | self, |
| 26 | Uid, |
| 27 | UID_LEN, |
| 28 | }, |
| 29 | }, |
| 30 | }; |
| 31 | |
| 32 | use std::{ |
| 33 | collections::BTreeMap, |
| 34 | os::unix::fs::PermissionsExt, |
| 35 | path::{ |
| 36 | Path, |
| 37 | PathBuf, |
| 38 | }, |
| 39 | thread, |
| 40 | time::{ |
| 41 | Duration, |
| 42 | Instant, |
| 43 | }, |
| 44 | }; |
| 45 | |
| 46 | type TestDb = O3db< |
| 47 | { UID_LEN }, |
| 48 | Uid, |
| 49 | (), |
| 50 | HashScheme, |
| 51 | HashScheme, |
| 52 | ChecksumScheme, |
| 53 | >; |
| 54 | |
| 55 | #[test] |
| 56 | fn main() -> Outcome<()> { |
| 57 | log_set_level!("warn"); |
| 58 | let outcome = run(); |
| 59 | // Whatever failed, the next binary must not inherit a slow store. |
| 60 | hooks::set_barrier_delay(Duration::ZERO); |
| 61 | hooks::set_publish_delay(Duration::ZERO); |
| 62 | hooks::set_supervisor_panic(false); |
| 63 | log_finish_wait!(); |
| 64 | outcome |
| 65 | } |
| 66 | |
| 67 | /// Every check runs, whichever fail, and all that fail are reported. |
| 68 | fn run() -> Outcome<()> { |
| 69 | let slow_start = slow_start_hands_over_live_channels(); |
| 70 | hooks::set_publish_delay(Duration::ZERO); |
| 71 | let slow_barrier = slow_barrier_still_acknowledges(); |
| 72 | hooks::set_barrier_delay(Duration::ZERO); |
| 73 | let deadline = durability_deadline_reports_written(); |
| 74 | hooks::set_barrier_delay(Duration::ZERO); |
| 75 | let writer = writer_failure_reaches_the_caller(); |
| 76 | let zone = failed_zone_fails_the_start(); |
| 77 | let panicked = panicked_supervisor_stops_its_bots(); |
| 78 | hooks::set_supervisor_panic(false); |
| 79 | let failed: Vec<Error<ErrTag>> = [slow_start, slow_barrier, deadline, writer, zone, panicked] |
| 80 | .into_iter() |
| 81 | .filter_map(|r| r.err()) |
| 82 | .collect(); |
| 83 | match failed.len() { |
| 84 | 0 => Ok(()), |
| 85 | 1 => match failed.into_iter().next() { |
| 86 | Some(e) => Err(e), |
| 87 | None => Ok(()), // unreachable |
| 88 | }, |
| 89 | n => Err(err!( |
| 90 | "{} checks failed: {:?}", n, failed; |
| 91 | Test)), |
| 92 | } |
| 93 | } |
| 94 | |
| 95 | /// The supervisor takes longer to hand over its channels than the one second `start` used to |
| 96 | /// sleep. The handle must hold the live channels when `start` returns, with no `updated_api` |
| 97 | /// and no sleep: before, every bot but the supervisor was unreachable, and a ping heard one of |
| 98 | /// them. |
| 99 | fn slow_start_hands_over_live_channels() -> Outcome<()> { |
| 100 | let delay = Duration::from_millis(2_500); |
| 101 | let root = res!(fresh("./test_db_store_timeouts_start")); |
| 102 | hooks::set_publish_delay(delay); |
| 103 | let begun = Instant::now(); |
| 104 | let opened = open(&root, res!(config())); |
| 105 | hooks::set_publish_delay(Duration::ZERO); |
| 106 | let db = res!(opened); |
| 107 | let took = begun.elapsed(); |
| 108 | if took < delay { |
| 109 | return Err(err!( |
| 110 | "start returned after {:?}, before the supervisor had handed over its channels \ |
| 111 | {:?} in.", took, delay; |
| 112 | Test, Timeout)); |
| 113 | } |
| 114 | |
| 115 | // Straight through the handle, as Oregami's wake does. A ping goes to every bot the handle |
| 116 | // holds a channel for and fails unless each of them answers. |
| 117 | let (_, pongs) = res!(db.api().ping_bots(constant::USER_REQUEST_WAIT)); |
| 118 | if pongs.len() < 2 { |
| 119 | return Err(err!( |
| 120 | "Only {} bot answered a ping as soon as start returned.", pongs.len(); |
| 121 | Test, Missing)); |
| 122 | } |
| 123 | |
| 124 | let key = dat!("reachable after a slow start"); |
| 125 | res!(db.insert(key.clone(), dat!(7u8), Uid::default(), None)); |
| 126 | match res!(db.get(&key, None)) { |
| 127 | Some((v, _)) => req!(v, dat!(7u8), "The value written after a slow start."), |
| 128 | None => return Err(err!( |
| 129 | "A value written after a slow start could not be read back."; Test, Missing)), |
| 130 | } |
| 131 | res!(db.close()); |
| 132 | Ok(()) |
| 133 | } |
| 134 | |
| 135 | /// Every durability barrier takes longer than the user request deadline. A write used to fail |
| 136 | /// at that deadline and land a moment later; now it is answered written at once and durable |
| 137 | /// when the barrier completes. Writes queued behind the barrier share it, and none of them is |
| 138 | /// held by it before being written. |
| 139 | fn slow_barrier_still_acknowledges() -> Outcome<()> { |
| 140 | let delay = constant::USER_REQUEST_TIMEOUT + Duration::from_secs(1); |
| 141 | let root = res!(fresh("./test_db_store_timeouts_barrier")); |
| 142 | let mut cfg = res!(config()); |
| 143 | cfg.sync_on_write = true; |
| 144 | let db = res!(open(&root, cfg)); |
| 145 | |
| 146 | hooks::set_barrier_delay(delay); |
| 147 | let begun = Instant::now(); |
| 148 | let mut writers = Vec::new(); |
| 149 | for i in 0..3u8 { |
| 150 | let db = db.clone(); |
| 151 | writers.push(res!(thread::Builder::new() |
| 152 | .name(fmt!("writer {}", i)) |
| 153 | .spawn(move || -> Outcome<Duration> { |
| 154 | let begun = Instant::now(); |
| 155 | res!(db.insert(key(i), dat!(i), Uid::default(), None)); |
| 156 | Ok(begun.elapsed()) |
| 157 | }))); |
| 158 | } |
| 159 | for (i, writer) in writers.into_iter().enumerate() { |
| 160 | match writer.join() { |
| 161 | Ok(Ok(took)) => if took < delay { |
| 162 | return Err(err!( |
| 163 | "Write {} was acknowledged in {:?}, before the {:?} barrier it waits on \ |
| 164 | could have completed.", i, took, delay; |
| 165 | Test, Unexpected)); |
| 166 | }, |
| 167 | Ok(Err(e)) => return Err(err!(e, |
| 168 | "Write {} failed behind a {:?} barrier, which the user request deadline of \ |
| 169 | {:?} does not measure.", i, delay, constant::USER_REQUEST_TIMEOUT; |
| 170 | Test, Write)), |
| 171 | Err(_) => return Err(err!("Writer thread {} panicked.", i; Test, Thread)), |
| 172 | } |
| 173 | } |
| 174 | hooks::set_barrier_delay(Duration::ZERO); |
| 175 | test!(sync_log::stream(), "Three writes behind {:?} barriers took {:?}.", delay, begun.elapsed()); |
| 176 | for i in 0..3u8 { |
| 177 | match res!(db.get(&key(i), None)) { |
| 178 | Some((v, _)) => req!(v, dat!(i), "The value of a slowly synced write."), |
| 179 | None => return Err(err!("Write {} is not readable.", i; Test, Missing)), |
| 180 | } |
| 181 | } |
| 182 | res!(db.close()); |
| 183 | |
| 184 | // Reopened, to show the writes are in the files and not just the cache. |
| 185 | let db = res!(open(&root, res!(config()))); |
| 186 | for i in 0..3u8 { |
| 187 | match res!(db.get(&key(i), None)) { |
| 188 | Some((v, _)) => req!(v, dat!(i), "The value of a slowly synced write, reopened."), |
| 189 | None => return Err(err!("Write {} did not survive a reopen.", i; Test, Missing)), |
| 190 | } |
| 191 | } |
| 192 | res!(db.close()); |
| 193 | Ok(()) |
| 194 | } |
| 195 | |
| 196 | /// When the durability deadline itself expires, the caller is told the write is in the files |
| 197 | /// but not yet confirmed durable, and that is true: the value is readable once the barrier |
| 198 | /// completes. The deadline is passed in short here, since the real one is two minutes. |
| 199 | fn durability_deadline_reports_written() -> Outcome<()> { |
| 200 | let delay = Duration::from_secs(4); |
| 201 | let durability = Duration::from_secs(2); |
| 202 | let root = res!(fresh("./test_db_store_timeouts_deadline")); |
| 203 | let mut cfg = res!(config()); |
| 204 | cfg.sync_on_write = true; |
| 205 | let db = res!(open(&root, cfg)); |
| 206 | |
| 207 | hooks::set_barrier_delay(delay); |
| 208 | let begun = Instant::now(); |
| 209 | let resp = res!(db.api().store(key(9), dat!(9u8), Uid::default())); |
| 210 | let n = match res!(resp.recv_timeout(constant::USER_REQUEST_TIMEOUT)) { |
| 211 | OzoneMsg::Chunks(n) => n, |
| 212 | msg => return Err(err!("Expected the record count, received {:?}.", msg; |
| 213 | Test, Unexpected)), |
| 214 | }; |
| 215 | let outcome = resp.recv_write_acks(n, constant::USER_REQUEST_TIMEOUT, durability); |
| 216 | let took = begun.elapsed(); |
| 217 | hooks::set_barrier_delay(Duration::ZERO); |
| 218 | match outcome { |
| 219 | Ok(acks) => return Err(err!( |
| 220 | "A write behind a {:?} barrier was confirmed durable within {:?}: {:?}.", |
| 221 | delay, durability, acks; |
| 222 | Test, Unexpected)), |
| 223 | Err(e) => { |
| 224 | let text = fmt!("{:?}", e); |
| 225 | if !text.contains("not confirmed durable") || !text.contains("has not failed") { |
| 226 | return Err(err!( |
| 227 | "An expired durability deadline must say the write is in the files but not \ |
| 228 | yet durable, and it said: {}", text; |
| 229 | Test, Mismatch)); |
| 230 | } |
| 231 | if took >= delay { |
| 232 | return Err(err!( |
| 233 | "The durability deadline of {:?} took {:?} to expire.", durability, took; |
| 234 | Test, Timeout)); |
| 235 | } |
| 236 | if !e.tags().contains(&ErrTag::Unconfirmed) { |
| 237 | return Err(err!( |
| 238 | "An expired durability deadline is tagged {:?}, which does not say \ |
| 239 | Unconfirmed: a caller cannot tell it from a write that never landed without \ |
| 240 | reading the words.", e.tags(); |
| 241 | Test, Mismatch)); |
| 242 | } |
| 243 | }, |
| 244 | } |
| 245 | // What the error said: the write lands when the disk completes it. |
| 246 | thread::sleep(delay.saturating_sub(begun.elapsed()) + Duration::from_millis(500)); |
| 247 | match res!(db.get(&key(9), None)) { |
| 248 | Some((v, _)) => req!(v, dat!(9u8), "A write reported as not yet durable."), |
| 249 | None => return Err(err!( |
| 250 | "A write reported as written but not yet durable never became readable."; |
| 251 | Test, Missing)), |
| 252 | } |
| 253 | res!(db.close()); |
| 254 | Ok(()) |
| 255 | } |
| 256 | |
| 257 | /// A writer that fails says why, where it used to leave the caller to time out on nothing. The |
| 258 | /// zone directory is made read-only, so the write that needs a new live file cannot have one. |
| 259 | fn writer_failure_reaches_the_caller() -> Outcome<()> { |
| 260 | let root = res!(fresh("./test_db_store_timeouts_writer")); |
| 261 | let mut cfg = res!(config()); |
| 262 | cfg.num_zones = 1; |
| 263 | let db = res!(open(&root, cfg.clone())); |
| 264 | let big = || dat!(vec![0x5au8; 1_000]); // two of these overflow a 2,000 byte live file |
| 265 | |
| 266 | res!(db.insert(key(1), big(), Uid::default(), None)); |
| 267 | let zdir = res!(zone_dir(&root, &cfg)); |
| 268 | res!(std::fs::set_permissions(&zdir, std::fs::Permissions::from_mode(0o555))); |
| 269 | let begun = Instant::now(); |
| 270 | let outcome = db.insert(key(2), big(), Uid::default(), None); |
| 271 | let took = begun.elapsed(); |
| 272 | res!(std::fs::set_permissions(&zdir, std::fs::Permissions::from_mode(0o755))); |
| 273 | match outcome { |
| 274 | Ok(_) => return Err(err!( |
| 275 | "A write needing a live file in a read-only zone directory succeeded."; Test, Unexpected)), |
| 276 | Err(e) => { |
| 277 | let text = fmt!("{:?}", e); |
| 278 | if text.contains("Failed to receive a message via responder") { |
| 279 | return Err(err!( |
| 280 | "A writer's failure reached the caller as a timeout: {}", text; Test, Mismatch)); |
| 281 | } |
| 282 | if !text.contains("new live file") || !text.contains("ermission denied") { |
| 283 | return Err(err!( |
| 284 | "A writer's failure must name its cause, and it said: {}", text; Test, Mismatch)); |
| 285 | } |
| 286 | if took >= constant::USER_REQUEST_TIMEOUT { |
| 287 | return Err(err!( |
| 288 | "The writer's failure took {:?} to arrive.", took; Test, Timeout)); |
| 289 | } |
| 290 | }, |
| 291 | } |
| 292 | // The directory is writable again, and so is the store. |
| 293 | res!(db.insert(key(3), big(), Uid::default(), None)); |
| 294 | match res!(db.get(&key(3), None)) { |
| 295 | Some((v, _)) => req!(v, big(), "A write after the zone recovered."), |
| 296 | None => return Err(err!("A write after the zone recovered is missing."; Test, Missing)), |
| 297 | } |
| 298 | res!(db.close()); |
| 299 | Ok(()) |
| 300 | } |
| 301 | |
| 302 | /// A zone that cannot be initialised fails the start. It used to be logged, and the store came |
| 303 | /// up with that zone's writer holding no directory, so its writes landed in the working directory |
| 304 | /// of whatever process had opened it. The failed start returns once the bots it brought up have |
| 305 | /// stopped: it returned with 32 threads running where there had been 4, and a caller that opens |
| 306 | /// the directory again at once, as Oregami's forge does on its next request, could have two sets |
| 307 | /// of bots over one directory. |
| 308 | fn failed_zone_fails_the_start() -> Outcome<()> { |
| 309 | let root = res!(fresh("./test_db_store_timeouts_zone")); |
| 310 | let mut cfg = res!(config()); |
| 311 | cfg.zone_overrides = res!(mapdat!{ |
| 312 | 1u16 => mapdat!{ |
| 313 | "dir" => "../test_db_store_timeouts_no_such_container", |
| 314 | "max_size" => 1_000_000u64, |
| 315 | }, |
| 316 | }.get_map().ok_or_else(|| err!("The zone override is not a map."; Test, Bug))); |
| 317 | let before = res!(threads()); |
| 318 | let begun = Instant::now(); |
| 319 | let mut db = res!(TestDb::new(root.clone(), Some(cfg), schemes(), Uid::default())); |
| 320 | match db.start("test") { |
| 321 | Ok(_) => return Err(err!( |
| 322 | "A store whose zone directory does not exist started."; Test, Unexpected)), |
| 323 | Err(e) => { |
| 324 | let text = fmt!("{:?}", e); |
| 325 | if !text.contains("does not exist and must be created") { |
| 326 | return Err(err!( |
| 327 | "The failed start must name the missing directory, and it said: {}", text; |
| 328 | Test, Mismatch)); |
| 329 | } |
| 330 | }, |
| 331 | } |
| 332 | let took = begun.elapsed(); |
| 333 | if took >= constant::USER_REQUEST_TIMEOUT { |
| 334 | return Err(err!("The failed start took {:?} to report.", took; Test, Timeout)); |
| 335 | } |
| 336 | // A thread that has let go of its bot ends a moment later, so the count gets that moment. |
| 337 | let at_return = res!(threads()); |
| 338 | let settled = Instant::now() + Duration::from_millis(100); |
| 339 | let mut after = at_return; |
| 340 | while after > before && Instant::now() < settled { |
| 341 | thread::sleep(Duration::from_millis(5)); |
| 342 | after = res!(threads()); |
| 343 | } |
| 344 | if after > before { |
| 345 | return Err(err!( |
| 346 | "A failed start returned with {} threads running, {} a moment later, where there were \ |
| 347 | {} before it: the bots it brought up were still running.", at_return, after, before; |
| 348 | Test, Unexpected)); |
| 349 | } |
| 350 | // Nothing is left running for a close to wait on. |
| 351 | let begun = Instant::now(); |
| 352 | res!(db.close()); |
| 353 | if begun.elapsed() >= Duration::from_secs(1) { |
| 354 | return Err(err!( |
| 355 | "Closing a store whose start failed took {:?}.", begun.elapsed(); Test, Timeout)); |
| 356 | } |
| 357 | Ok(()) |
| 358 | } |
| 359 | |
| 360 | /// The supervisor panics once it has brought the bots up, before the database is ready. The |
| 361 | /// start fails and returns once those bots have stopped. It returned at once, with every bot |
| 362 | /// still running and nothing left that would stop them, so the directory could be opened again |
| 363 | /// beside them. |
| 364 | fn panicked_supervisor_stops_its_bots() -> Outcome<()> { |
| 365 | let root = res!(fresh("./test_db_store_timeouts_panic")); |
| 366 | let cfg = res!(config()); |
| 367 | let before = res!(threads()); |
| 368 | hooks::set_supervisor_panic(true); |
| 369 | let mut db = res!(TestDb::new(root.clone(), Some(cfg), schemes(), Uid::default())); |
| 370 | let started = db.start("test"); |
| 371 | hooks::set_supervisor_panic(false); |
| 372 | if started.is_ok() { |
| 373 | return Err(err!( |
| 374 | "A store whose supervisor panicked while starting it started."; Test, Unexpected)); |
| 375 | } |
| 376 | // A thread that has let go of its bot ends a moment later, so the count gets that moment. |
| 377 | let at_return = res!(threads()); |
| 378 | let settled = Instant::now() + Duration::from_millis(100); |
| 379 | let mut after = at_return; |
| 380 | while after > before && Instant::now() < settled { |
| 381 | thread::sleep(Duration::from_millis(5)); |
| 382 | after = res!(threads()); |
| 383 | } |
| 384 | if after > before { |
| 385 | return Err(err!( |
| 386 | "A start whose supervisor panicked returned with {} threads running, {} a moment \ |
| 387 | later, where there were {} before it: the bots it brought up were still running.", |
| 388 | at_return, after, before; |
| 389 | Test, Unexpected)); |
| 390 | } |
| 391 | let begun = Instant::now(); |
| 392 | res!(db.close()); |
| 393 | if begun.elapsed() >= Duration::from_secs(1) { |
| 394 | return Err(err!( |
| 395 | "Closing a store whose start failed took {:?}.", begun.elapsed(); Test, Timeout)); |
| 396 | } |
| 397 | Ok(()) |
| 398 | } |
| 399 | |
| 400 | fn key(i: u8) -> Dat { |
| 401 | dat!(fmt!("store timeouts key {}", i)) |
| 402 | } |
| 403 | |
| 404 | /// The threads this process is running. |
| 405 | fn threads() -> Outcome<usize> { |
| 406 | Ok(res!(std::fs::read_dir("/proc/self/task")).count()) |
| 407 | } |
| 408 | |
| 409 | /// The crate's test configuration, with every zone inside the test's own directory. |
| 410 | fn config() -> Outcome<OzoneConfig> { |
| 411 | let mut cfg = res!(setup::default_cfg()); |
| 412 | cfg.zone_overrides = BTreeMap::new(); |
| 413 | Ok(cfg) |
| 414 | } |
| 415 | |
| 416 | fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> { |
| 417 | RestSchemesInput::new( |
| 418 | None::<()>, |
| 419 | None::<HashScheme>, |
| 420 | None::<HashScheme>, |
| 421 | Some(ChecksumScheme::new_crc32()), |
| 422 | ) |
| 423 | } |
| 424 | |
| 425 | /// An empty directory of this test's own. |
| 426 | fn fresh(dir: &str) -> Outcome<PathBuf> { |
| 427 | let _ = std::fs::remove_dir_all(dir); |
| 428 | res!(std::fs::create_dir_all(dir)); |
| 429 | Ok(res!(Path::new(dir).canonicalize())) |
| 430 | } |
| 431 | |
| 432 | fn open(root: &Path, cfg: OzoneConfig) -> Outcome<TestDb> { |
| 433 | let mut db = res!(TestDb::new(root.to_path_buf(), Some(cfg), schemes(), Uid::default())); |
| 434 | res!(db.start("test")); |
| 435 | Ok(db) |
| 436 | } |
| 437 | |
| 438 | /// The directory of the first zone, which is the only one when there is one. |
| 439 | fn zone_dir(root: &Path, cfg: &OzoneConfig) -> Outcome<PathBuf> { |
| 440 | let zdir = cfg.zone_root(root).join("zone_001"); |
| 441 | if !zdir.is_dir() { |
| 442 | return Err(err!("Expected a zone directory at {:?}.", zdir; Test, Missing)); |
| 443 | } |
| 444 | Ok(zdir) |
| 445 | } |