oxedyne/fe2o3/fe2o3_o3db_sync/tests/syncer_faults.rs
24.4 KiB, 13 runs
created by r1870400018:61249, 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 writer's durability barrier runs on a thread of its own, its syncer, and these are the ways |
| 2 | //! that can go wrong: the disk failing its syncs, the syncer stopping, a new live pair that cannot |
| 3 | //! be handed to it, and a shutdown arriving while it holds records. Each check says what the |
| 4 | //! store must tell its callers, and what must be on disk afterwards. `test::hooks` makes the |
| 5 | //! faults, and the hooks are process-wide, which is why this is a test binary of its own with a |
| 6 | //! single test. Found by QA of lane sto, 2026-09-23. |
| 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::{ |
| 22 | msg::OzoneMsg, |
| 23 | response::Wait, |
| 24 | }, |
| 25 | data::core::RestSchemesInput, |
| 26 | file::floc::FileNum, |
| 27 | test::{ |
| 28 | hooks, |
| 29 | setup::{ |
| 30 | self, |
| 31 | Uid, |
| 32 | UID_LEN, |
| 33 | }, |
| 34 | }, |
| 35 | }; |
| 36 | |
| 37 | use std::{ |
| 38 | collections::BTreeMap, |
| 39 | path::{ |
| 40 | Path, |
| 41 | PathBuf, |
| 42 | }, |
| 43 | thread, |
| 44 | time::{ |
| 45 | Duration, |
| 46 | Instant, |
| 47 | }, |
| 48 | }; |
| 49 | |
| 50 | type TestDb = O3db< |
| 51 | { UID_LEN }, |
| 52 | Uid, |
| 53 | (), |
| 54 | HashScheme, |
| 55 | HashScheme, |
| 56 | ChecksumScheme, |
| 57 | >; |
| 58 | |
| 59 | const PERIOD_MS: u64 = 200; // the interval policy's |
| 60 | const WATCH: Duration = Duration::from_millis(1_500); // a failing disk watched |
| 61 | const ANSWERED_WITHIN: Duration = Duration::from_secs(10); // a closing store's last write |
| 62 | const SLOW_WRITES: u8 = 20; // arriving at a failing disk |
| 63 | const SLOW_BARRIERS: u64 = 6; // enough for them, and more |
| 64 | const EVERY_N: u32 = 2; // `sync_every_n_writes` |
| 65 | |
| 66 | #[test] |
| 67 | fn main() -> Outcome<()> { |
| 68 | log_set_level!("warn"); |
| 69 | // Whatever fails, the next check and the next binary must not inherit the fault. |
| 70 | let failing = failing_disk_is_retried_once_a_period(); |
| 71 | hooks::set_barrier_failure(false); |
| 72 | let stopped = stopped_syncer_refuses_writes_before_appending(); |
| 73 | hooks::set_syncer_stop(false); |
| 74 | let handoff = failed_hand_off_keeps_the_writer_on_its_pair(); |
| 75 | hooks::set_pair_hand_failure(false); |
| 76 | let closing = close_answers_what_the_syncer_holds(); |
| 77 | hooks::set_barrier_delay(Duration::ZERO); |
| 78 | let slowly = slowly_failing_disk_answers_waiting_writes_together(0); |
| 79 | hooks::set_barrier_failure(false); |
| 80 | hooks::set_barrier_delay(Duration::ZERO); |
| 81 | let every_n = slowly_failing_disk_answers_waiting_writes_together(EVERY_N); |
| 82 | hooks::set_barrier_failure(false); |
| 83 | hooks::set_barrier_delay(Duration::ZERO); |
| 84 | let queued = close_answers_what_a_slow_cache_bot_holds(); |
| 85 | hooks::set_barrier_delay(Duration::ZERO); |
| 86 | hooks::set_insert_delay(Duration::ZERO); |
| 87 | log_finish_wait!(); |
| 88 | let failed: Vec<Error<ErrTag>> = [failing, stopped, handoff, closing, slowly, every_n, queued] |
| 89 | .into_iter() |
| 90 | .filter_map(|r| r.err()) |
| 91 | .collect(); |
| 92 | match failed.len() { |
| 93 | 0 => Ok(()), |
| 94 | 1 => match failed.into_iter().next() { |
| 95 | Some(e) => Err(e), |
| 96 | None => Ok(()), // unreachable |
| 97 | }, |
| 98 | n => Err(err!( |
| 99 | "{} checks failed: {:?}", n, failed; |
| 100 | Test)), |
| 101 | } |
| 102 | } |
| 103 | |
| 104 | /// The disk fails every sync, under the interval policy, the library's default. The barrier the |
| 105 | /// policy owed used to be retried as fast as it could fail: a million times in three seconds, on |
| 106 | /// nearly three cores, each failure logged. It is retried once a period. A write made while the |
| 107 | /// disk is failing waits on a barrier of its own and is told it failed, rather than being answered |
| 108 | /// as if its period will make it durable, and the store recovers with the disk. |
| 109 | fn failing_disk_is_retried_once_a_period() -> Outcome<()> { |
| 110 | let root = res!(fresh("./test_db_syncer_faults_failing")); |
| 111 | let mut cfg = res!(config()); |
| 112 | cfg.sync_interval_ms = PERIOD_MS; |
| 113 | let db = res!(open(&root, cfg)); |
| 114 | // The first write is owed a barrier at once, and the disk is still good. |
| 115 | res!(db.insert(key(1), dat!(1u8), Uid::default(), None)); |
| 116 | |
| 117 | hooks::set_barrier_failure(true); |
| 118 | let counted = hooks::barriers_failed(); |
| 119 | // Made within the period of that barrier, so released without one, leaving one owed. Were |
| 120 | // the machine slow enough for the period to pass first, this write would wait on a barrier of |
| 121 | // its own and be told it failed; a barrier is owed either way, so the outcome is not the point. |
| 122 | let _ = db.insert(key(2), dat!(2u8), Uid::default(), None); |
| 123 | thread::sleep(WATCH); |
| 124 | let tries = hooks::barriers_failed() - counted; |
| 125 | let most = 2 * (WATCH.as_millis() as u64 / PERIOD_MS) + 2; |
| 126 | if tries > most { |
| 127 | let _ = db.close(); // the check has failed already, and says why |
| 128 | return Err(err!( |
| 129 | "While the disk failed every sync for {:?}, the owed barrier was tried {} times, where \ |
| 130 | once a {} ms period is at most {}: the syncer retries as fast as the disk fails.", |
| 131 | WATCH, tries, PERIOD_MS, most; |
| 132 | Test, Mismatch)); |
| 133 | } |
| 134 | if tries < 2 { |
| 135 | let _ = db.close(); // the check has failed already, and says why |
| 136 | return Err(err!( |
| 137 | "While the disk failed every sync for {:?}, the owed barrier was tried {} times: a \ |
| 138 | store that stops retrying never makes its writes durable once the disk recovers.", |
| 139 | WATCH, tries; |
| 140 | Test, Missing)); |
| 141 | } |
| 142 | |
| 143 | match db.insert(key(3), dat!(3u8), Uid::default(), None) { |
| 144 | Ok(_) => { |
| 145 | let _ = db.close(); // the check has failed already, and says why |
| 146 | return Err(err!( |
| 147 | "A write made while the disk failed every sync was answered as if its period \ |
| 148 | would make it durable."; |
| 149 | Test, Unexpected)); |
| 150 | }, |
| 151 | Err(e) => { |
| 152 | let text = fmt!("{:?}", e); |
| 153 | if !text.contains("not confirmed durable") { |
| 154 | let _ = db.close(); // the check has failed already, and says why |
| 155 | return Err(err!( |
| 156 | "A write made while the disk failed must be told its barrier failed, and it \ |
| 157 | was told: {}", text; |
| 158 | Test, Mismatch)); |
| 159 | } |
| 160 | // Written, so a caller must be able to tell it from a write that never landed |
| 161 | // without reading the words. |
| 162 | if !e.tags().contains(&ErrTag::Unconfirmed) { |
| 163 | let _ = db.close(); // the check has failed already, and says why |
| 164 | return Err(err!( |
| 165 | "A write whose barrier failed after it was written reached its caller tagged \ |
| 166 | {:?}, which does not say Unconfirmed: a caller cannot tell it from a write that \ |
| 167 | never landed.", e.tags(); |
| 168 | Test, Mismatch)); |
| 169 | } |
| 170 | }, |
| 171 | } |
| 172 | // The record is in the files all the same, and readable. |
| 173 | match res!(db.get(&key(3), None)) { |
| 174 | Some((v, _)) => req!(v, dat!(3u8), "A write whose barrier failed."), |
| 175 | None => return Err(err!( |
| 176 | "A write whose barrier failed is not readable, although it is in the files."; |
| 177 | Test, Missing)), |
| 178 | } |
| 179 | |
| 180 | // The disk recovers, and the store with it. |
| 181 | hooks::set_barrier_failure(false); |
| 182 | thread::sleep(Duration::from_millis(2 * PERIOD_MS)); |
| 183 | res!(db.insert(key(4), dat!(4u8), Uid::default(), None)); |
| 184 | res!(db.close()); |
| 185 | Ok(()) |
| 186 | } |
| 187 | |
| 188 | /// A writer's syncer stops, as one that panicked would. Every later write through that writer is |
| 189 | /// refused before anything reaches the files. Appended first, it was reported failed and came |
| 190 | /// back at the next start. |
| 191 | fn stopped_syncer_refuses_writes_before_appending() -> Outcome<()> { |
| 192 | let root = res!(fresh("./test_db_syncer_faults_stopped")); |
| 193 | let cfg = res!(config()); |
| 194 | let db = res!(open(&root, cfg.clone())); |
| 195 | res!(db.insert(key(11), dat!(11u8), Uid::default(), None)); |
| 196 | |
| 197 | hooks::set_syncer_stop(true); |
| 198 | let counted = hooks::syncers_stopped(); |
| 199 | // Released, and then its syncer stops. |
| 200 | res!(db.insert(key(12), dat!(12u8), Uid::default(), None)); |
| 201 | let begun = Instant::now(); |
| 202 | while hooks::syncers_stopped() == counted { |
| 203 | if begun.elapsed() > constant::USER_REQUEST_TIMEOUT { |
| 204 | let _ = db.close(); // the check has failed already, and says why |
| 205 | return Err(err!("The syncer did not stop when told to."; Test, Timeout)); |
| 206 | } |
| 207 | thread::sleep(Duration::from_millis(5)); |
| 208 | } |
| 209 | hooks::set_syncer_stop(false); |
| 210 | // A moment for its thread to end. |
| 211 | thread::sleep(Duration::from_millis(100)); |
| 212 | |
| 213 | let refused = match db.insert(key(13), dat!(13u8), Uid::default(), None) { |
| 214 | Ok(_) => { |
| 215 | let _ = db.close(); // the check has failed already, and says why |
| 216 | return Err(err!( |
| 217 | "A write through a writer whose syncer had stopped was confirmed."; |
| 218 | Test, Unexpected)); |
| 219 | }, |
| 220 | Err(e) => { |
| 221 | // Refused before anything was written, so a retry is needed, and nothing may say |
| 222 | // otherwise. |
| 223 | if e.tags().contains(&ErrTag::Unconfirmed) { |
| 224 | let _ = db.close(); // the check has failed already, and says why |
| 225 | return Err(err!( |
| 226 | "A write refused before anything was written is tagged Unconfirmed, as a \ |
| 227 | written one would be."; |
| 228 | Test, Mismatch)); |
| 229 | } |
| 230 | fmt!("{:?}", e) |
| 231 | }, |
| 232 | }; |
| 233 | res!(db.close()); |
| 234 | |
| 235 | // Only what was confirmed is there after a restart. |
| 236 | let db = res!(open(&root, cfg)); |
| 237 | for (k, v) in [(11u8, true), (12, true), (13, false)] { |
| 238 | match (res!(db.get(&key(k), None)), v) { |
| 239 | (Some(_), true) | (None, false) => (), |
| 240 | (None, true) => { |
| 241 | let _ = db.close(); // the check has failed already, and says why |
| 242 | return Err(err!( |
| 243 | "Write {}, confirmed before the syncer stopped, is missing after a restart.", k; |
| 244 | Test, Missing)); |
| 245 | }, |
| 246 | (Some(_), false) => { |
| 247 | let _ = db.close(); // the check has failed already, and says why |
| 248 | return Err(err!( |
| 249 | "Write {}, reported failed because the syncer had stopped, came back after a \ |
| 250 | restart: it had been written before it was refused. It was told: {}", |
| 251 | k, refused; |
| 252 | Test, Unexpected)); |
| 253 | }, |
| 254 | } |
| 255 | } |
| 256 | res!(db.close()); |
| 257 | if !refused.contains("refused before anything was written") { |
| 258 | return Err(err!( |
| 259 | "A write refused because the syncer had stopped must say it was refused before \ |
| 260 | anything was written, and it said: {}", refused; |
| 261 | Test, Mismatch)); |
| 262 | } |
| 263 | Ok(()) |
| 264 | } |
| 265 | |
| 266 | /// Handing a writer's new live pair to its syncer fails, as running out of file descriptors would. |
| 267 | /// The writer stays on the pair its syncer holds. It used to switch to the new pair first, so its |
| 268 | /// next records went to a file no barrier covered and its file bot never flagged live, while the |
| 269 | /// file it had left stayed flagged live for good. |
| 270 | fn failed_hand_off_keeps_the_writer_on_its_pair() -> Outcome<()> { |
| 271 | let root = res!(fresh("./test_db_syncer_faults_handoff")); |
| 272 | let cfg = res!(config()); |
| 273 | let db = res!(open(&root, cfg.clone())); |
| 274 | let big = |i: u8| dat!(vec![i; 1_000]); // two of these overflow a 2,000 byte live file |
| 275 | |
| 276 | res!(db.insert(key(21), big(21), Uid::default(), None)); |
| 277 | hooks::set_pair_hand_failure(true); |
| 278 | // This one needs a new live file, whose pair cannot be handed over. |
| 279 | if db.insert(key(22), big(22), Uid::default(), None).is_ok() { |
| 280 | let _ = db.close(); // the check has failed already, and says why |
| 281 | return Err(err!( |
| 282 | "A write needing a live pair that could not be handed to the syncer was confirmed."; |
| 283 | Test, Unexpected)); |
| 284 | } |
| 285 | // Small enough for the file the writer is on. |
| 286 | res!(db.insert(key(23), dat!(23u8), Uid::default(), None)); |
| 287 | hooks::set_pair_hand_failure(false); |
| 288 | // Two rollovers that work. |
| 289 | res!(db.insert(key(24), big(24), Uid::default(), None)); |
| 290 | res!(db.insert(key(25), big(25), Uid::default(), None)); |
| 291 | |
| 292 | // One file is flagged live, the one the writer is on. |
| 293 | let newest = res!(newest_data_file(&root, &cfg)); |
| 294 | let states = res!(db.api().collect_file_states(Wait { |
| 295 | max_wait: constant::USER_REQUEST_TIMEOUT, |
| 296 | check_interval: constant::CHECK_INTERVAL, |
| 297 | })); |
| 298 | let mut live: Vec<FileNum> = Vec::new(); |
| 299 | for (_, shard) in &states { |
| 300 | for (fnum, fstat) in shard.map() { |
| 301 | if fstat.is_live() { |
| 302 | live.push(*fnum); |
| 303 | } |
| 304 | } |
| 305 | } |
| 306 | live.sort(); |
| 307 | if live != vec![newest] { |
| 308 | let _ = db.close(); // the check has failed already, and says why |
| 309 | return Err(err!( |
| 310 | "After a live pair could not be handed to the syncer, files {:?} are flagged live, \ |
| 311 | where the writer is on file {} alone.", live, newest; |
| 312 | Test, Mismatch)); |
| 313 | } |
| 314 | res!(db.close()); |
| 315 | |
| 316 | // After a restart, every confirmed write is there and the refused one is not. |
| 317 | let db = res!(open(&root, cfg)); |
| 318 | for (k, v) in [(21u8, Some(big(21))), (22, None), (23, Some(dat!(23u8))), (24, Some(big(24))), |
| 319 | (25, Some(big(25)))] |
| 320 | { |
| 321 | match (res!(db.get(&key(k), None)), v) { |
| 322 | (Some((got, _)), Some(want)) => req!(got, want, "A write around a failed hand-off."), |
| 323 | (None, None) => (), |
| 324 | (got, want) => { |
| 325 | let _ = db.close(); // the check has failed already, and says why |
| 326 | return Err(err!( |
| 327 | "Write {} around a failed hand-off: after a restart it is {:?}, where it \ |
| 328 | should be {:?}.", k, got.map(|(v, _)| v), want; |
| 329 | Test, Mismatch)); |
| 330 | }, |
| 331 | } |
| 332 | } |
| 333 | res!(db.close()); |
| 334 | Ok(()) |
| 335 | } |
| 336 | |
| 337 | /// The store is closed while a syncer holds a written record, its barrier still running. The |
| 338 | /// record is answered before the store's bots stop. The cache bots used to be stopped first, and |
| 339 | /// nothing counted what a syncer held, so the record it released afterwards was never answered: |
| 340 | /// its caller waited out the durability deadline and was told that a durable write was not |
| 341 | /// confirmed. |
| 342 | fn close_answers_what_the_syncer_holds() -> Outcome<()> { |
| 343 | let root = res!(fresh("./test_db_syncer_faults_close")); |
| 344 | let mut cfg = res!(config()); |
| 345 | cfg.sync_on_write = true; |
| 346 | let db = res!(open(&root, cfg.clone())); |
| 347 | |
| 348 | hooks::set_barrier_delay(Duration::from_millis(1_500)); |
| 349 | let writing = db.clone(); |
| 350 | let writer = res!(thread::Builder::new() |
| 351 | .name(fmt!("syncer faults writer")) |
| 352 | .spawn(move || -> Outcome<()> { |
| 353 | let resp = res!(writing.api().store(key(31), dat!(31u8), Uid::default())); |
| 354 | let n = match res!(resp.recv_timeout(constant::USER_REQUEST_TIMEOUT)) { |
| 355 | OzoneMsg::Chunks(n) => n, |
| 356 | msg => return Err(err!( |
| 357 | "Expected the record count, received {:?}.", msg; Test, Unexpected)), |
| 358 | }; |
| 359 | // A durability deadline short enough to report an unanswered record in seconds. |
| 360 | res!(resp.recv_write_acks(n, constant::USER_REQUEST_TIMEOUT, ANSWERED_WITHIN)); |
| 361 | Ok(()) |
| 362 | })); |
| 363 | // Written by now, and waiting on its barrier. |
| 364 | thread::sleep(Duration::from_millis(300)); |
| 365 | let closed = db.close(); |
| 366 | hooks::set_barrier_delay(Duration::ZERO); |
| 367 | res!(closed); |
| 368 | match writer.join() { |
| 369 | Ok(Ok(())) => (), |
| 370 | Ok(Err(e)) => return Err(err!(e, |
| 371 | "A write its syncer held when the store was closed was not answered."; Test, Missing)), |
| 372 | Err(_) => return Err(err!("The writing thread panicked."; Test, Thread)), |
| 373 | } |
| 374 | |
| 375 | // And it is on disk. |
| 376 | let db = res!(open(&root, cfg)); |
| 377 | match res!(db.get(&key(31), None)) { |
| 378 | Some((v, _)) => req!(v, dat!(31u8), "A write answered as the store closed."), |
| 379 | None => { |
| 380 | let _ = db.close(); // the check has failed already, and says why |
| 381 | return Err(err!( |
| 382 | "A write answered as the store closed is missing after a restart."; Test, Missing)); |
| 383 | }, |
| 384 | } |
| 385 | res!(db.close()); |
| 386 | Ok(()) |
| 387 | } |
| 388 | |
| 389 | /// The disk fails every sync, slowly, under the interval policy, or under `sync_every_n_writes` |
| 390 | /// when `every_n` is not zero, and twenty writes arrive at once. While the last barrier has |
| 391 | /// failed each write waits on a barrier, and those waiting together share one, as they do under |
| 392 | /// `sync_on_write`. Each had a barrier of its own, one after another, so a disk taking its time |
| 393 | /// to fail answered twenty writes ten times slower, and with writes arriving faster than it |
| 394 | /// failed the queue grew without bound. Fixed for the interval policy on 2026-09-23 and for |
| 395 | /// every-n, where a failed barrier leaves every later write due one, on 2026-09-24. |
| 396 | fn slowly_failing_disk_answers_waiting_writes_together(every_n: u32) -> Outcome<()> { |
| 397 | let root = res!(fresh(if every_n == 0 { |
| 398 | "./test_db_syncer_faults_slowly" |
| 399 | } else { |
| 400 | "./test_db_syncer_faults_slowly_every_n" |
| 401 | })); |
| 402 | let mut cfg = res!(config()); |
| 403 | cfg.sync_interval_ms = PERIOD_MS; |
| 404 | cfg.sync_every_n_writes = every_n; |
| 405 | let policy = match every_n { |
| 406 | 0 => fmt!("the interval policy"), |
| 407 | n => fmt!("every {} writes", n), |
| 408 | }; |
| 409 | let db = res!(open(&root, cfg)); |
| 410 | res!(db.insert(key(41), dat!(41u8), Uid::default(), None)); |
| 411 | |
| 412 | hooks::set_barrier_delay(Duration::from_millis(PERIOD_MS)); |
| 413 | hooks::set_barrier_failure(true); |
| 414 | // Owed a barrier, which fails: under the interval policy it is released within the period of |
| 415 | // the barrier before it, and the barrier follows; under every-n it waits on it. |
| 416 | let begun = Instant::now(); |
| 417 | let counted = hooks::barriers_failed(); |
| 418 | let _ = db.insert(key(42), dat!(42u8), Uid::default(), None); |
| 419 | while hooks::barriers_failed() == counted { |
| 420 | if begun.elapsed() > constant::USER_REQUEST_TIMEOUT { |
| 421 | let _ = db.close(); // the check has failed already, and says why |
| 422 | return Err(err!("The barrier owed on a failing disk was never tried."; Test, Timeout)); |
| 423 | } |
| 424 | thread::sleep(Duration::from_millis(5)); |
| 425 | } |
| 426 | let counted = hooks::barriers_failed(); |
| 427 | let mut resps = Vec::new(); |
| 428 | for i in 0..SLOW_WRITES { |
| 429 | resps.push(res!(db.api().store(key(50 + i), dat!(50 + i), Uid::default()))); |
| 430 | } |
| 431 | let mut told = 0; |
| 432 | for resp in resps { |
| 433 | let n = match res!(resp.recv_timeout(constant::USER_REQUEST_TIMEOUT)) { |
| 434 | OzoneMsg::Chunks(n) => n, |
| 435 | msg => { |
| 436 | let _ = db.close(); // the check has failed already, and says why |
| 437 | return Err(err!( |
| 438 | "Expected the record count, received {:?}.", msg; Test, Unexpected)); |
| 439 | }, |
| 440 | }; |
| 441 | if resp.recv_write_acks(n, constant::USER_REQUEST_TIMEOUT, ANSWERED_WITHIN).is_err() { |
| 442 | told += 1; |
| 443 | } |
| 444 | } |
| 445 | let tries = hooks::barriers_failed() - counted; |
| 446 | hooks::set_barrier_failure(false); |
| 447 | hooks::set_barrier_delay(Duration::ZERO); |
| 448 | res!(db.close()); |
| 449 | if told != SLOW_WRITES { |
| 450 | return Err(err!( |
| 451 | "Under {}, {} of {} writes made while the disk failed were told their barrier \ |
| 452 | failed.", policy, told, SLOW_WRITES; |
| 453 | Test, Mismatch)); |
| 454 | } |
| 455 | if tries > SLOW_BARRIERS { |
| 456 | return Err(err!( |
| 457 | "Under {}, {} writes arriving together while the disk failed each sync slowly took {} \ |
| 458 | barriers, where waiting together they share one, and a few more is the most: they \ |
| 459 | were answered one barrier at a time.", policy, SLOW_WRITES, tries; |
| 460 | Test, Mismatch)); |
| 461 | } |
| 462 | Ok(()) |
| 463 | } |
| 464 | |
| 465 | /// The store is closed while its one cache bot works through a queue of written records, and the |
| 466 | /// shutdown's time runs out with one still queued behind the barrier that ends the writers. That |
| 467 | /// record is answered. Its queue used to be drained into the log when the time ran out, which |
| 468 | /// destroyed it, and its caller waited out the durability deadline for a write that had landed. |
| 469 | fn close_answers_what_a_slow_cache_bot_holds() -> Outcome<()> { |
| 470 | let root = res!(fresh("./test_db_syncer_faults_queued")); |
| 471 | let mut cfg = res!(config()); |
| 472 | cfg.num_cbots_per_zone = 1; |
| 473 | cfg.sync_interval_ms = PERIOD_MS; |
| 474 | let db = res!(open(&root, cfg.clone())); |
| 475 | |
| 476 | // The first write waits on the barrier the policy owes it at once, the second behind it, and |
| 477 | // both reach the cache bot as the barrier ends, after the close has begun. |
| 478 | hooks::set_barrier_delay(Duration::from_millis(2_000)); |
| 479 | hooks::set_insert_delay(Duration::from_millis(1_500)); |
| 480 | let mut writers = Vec::new(); |
| 481 | for i in [61u8, 62] { |
| 482 | let writing = db.clone(); |
| 483 | writers.push(res!(thread::Builder::new() |
| 484 | .name(fmt!("syncer faults writer {}", i)) |
| 485 | .spawn(move || -> Outcome<()> { |
| 486 | let resp = res!(writing.api().store(key(i), dat!(i), Uid::default())); |
| 487 | let n = match res!(resp.recv_timeout(constant::USER_REQUEST_TIMEOUT)) { |
| 488 | OzoneMsg::Chunks(n) => n, |
| 489 | msg => return Err(err!( |
| 490 | "Expected the record count, received {:?}.", msg; Test, Unexpected)), |
| 491 | }; |
| 492 | res!(resp.recv_write_acks(n, constant::USER_REQUEST_TIMEOUT, ANSWERED_WITHIN)); |
| 493 | Ok(()) |
| 494 | }))); |
| 495 | thread::sleep(Duration::from_millis(50)); |
| 496 | } |
| 497 | thread::sleep(Duration::from_millis(100)); |
| 498 | let closed = db.close(); |
| 499 | let mut unanswered = Vec::new(); |
| 500 | for (i, writer) in writers.into_iter().enumerate() { |
| 501 | match writer.join() { |
| 502 | Ok(Ok(())) => (), |
| 503 | Ok(Err(e)) => unanswered.push(fmt!("write {}: {}", i + 1, e)), |
| 504 | Err(_) => unanswered.push(fmt!("write {}: the writing thread panicked", i + 1)), |
| 505 | } |
| 506 | } |
| 507 | hooks::set_barrier_delay(Duration::ZERO); |
| 508 | hooks::set_insert_delay(Duration::ZERO); |
| 509 | res!(closed); |
| 510 | if !unanswered.is_empty() { |
| 511 | return Err(err!( |
| 512 | "Writes queued at the cache bot when the store closed were not answered: {:?}", |
| 513 | unanswered; |
| 514 | Test, Missing)); |
| 515 | } |
| 516 | let db = res!(open(&root, cfg)); |
| 517 | for i in [61u8, 62] { |
| 518 | match res!(db.get(&key(i), None)) { |
| 519 | Some((v, _)) => req!(v, dat!(i), "A write answered as the store closed."), |
| 520 | None => { |
| 521 | let _ = db.close(); // the check has failed already, and says why |
| 522 | return Err(err!( |
| 523 | "Write {} answered as the store closed is missing after a restart.", i; |
| 524 | Test, Missing)); |
| 525 | }, |
| 526 | } |
| 527 | } |
| 528 | res!(db.close()); |
| 529 | Ok(()) |
| 530 | } |
| 531 | |
| 532 | /// The highest-numbered data file in the one zone. |
| 533 | fn newest_data_file(root: &Path, cfg: &OzoneConfig) -> Outcome<FileNum> { |
| 534 | let zdir = cfg.zone_root(root).join("zone_001"); |
| 535 | let mut newest = None; |
| 536 | for entry in res!(std::fs::read_dir(&zdir)) { |
| 537 | let path = res!(entry).path(); |
| 538 | if path.extension().map(|e| e == constant::DATA_FILE_EXT).unwrap_or(false) { |
| 539 | let digits: String = path.file_stem() |
| 540 | .map(|s| s.to_string_lossy().chars().filter(|c| c.is_ascii_digit()).collect()) |
| 541 | .unwrap_or_default(); |
| 542 | let fnum: FileNum = res!(digits.parse::<FileNum>()); |
| 543 | newest = Some(newest.map_or(fnum, |n: FileNum| n.max(fnum))); |
| 544 | } |
| 545 | } |
| 546 | newest.ok_or_else(|| err!("There is no data file in {:?}.", zdir; Test, Missing)) |
| 547 | } |
| 548 | |
| 549 | fn key(i: u8) -> Dat { |
| 550 | dat!(fmt!("syncer faults key {}", i)) |
| 551 | } |
| 552 | |
| 553 | /// One zone, so one writer and one syncer, with every zone inside the test's own directory and |
| 554 | /// no durability policy until a check sets one. |
| 555 | fn config() -> Outcome<OzoneConfig> { |
| 556 | let mut cfg = res!(setup::default_cfg()); |
| 557 | cfg.num_zones = 1; |
| 558 | cfg.zone_overrides = BTreeMap::new(); |
| 559 | Ok(cfg) |
| 560 | } |
| 561 | |
| 562 | fn schemes() -> RestSchemesInput<(), HashScheme, HashScheme, ChecksumScheme> { |
| 563 | RestSchemesInput::new( |
| 564 | None::<()>, |
| 565 | None::<HashScheme>, |
| 566 | None::<HashScheme>, |
| 567 | Some(ChecksumScheme::new_crc32()), |
| 568 | ) |
| 569 | } |
| 570 | |
| 571 | /// An empty directory of this test's own. |
| 572 | fn fresh(dir: &str) -> Outcome<PathBuf> { |
| 573 | let _ = std::fs::remove_dir_all(dir); // absent the first time |
| 574 | res!(std::fs::create_dir_all(dir)); |
| 575 | Ok(res!(Path::new(dir).canonicalize())) |
| 576 | } |
| 577 | |
| 578 | fn open(root: &Path, cfg: OzoneConfig) -> Outcome<TestDb> { |
| 579 | let mut db = res!(TestDb::new(root.to_path_buf(), Some(cfg), schemes(), Uid::default())); |
| 580 | res!(db.start("test")); |
| 581 | Ok(db) |
| 582 | } |