oxedyne/fe2o3/fe2o3_o3db_sync/src/test/hooks.rs
5.3 KiB, 7 runs
created by r1870400018:61243, 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 | //! Delays and failures a test puts in the store's way, so that a slow or failing disk, or a slow |
| 2 | //! start, can be reproduced on a machine that has none of them. Each is off until a test sets it, |
| 3 | //! and each is process-wide, so a test that sets one runs in a test binary of its own. |
| 4 | |
| 5 | use std::{ |
| 6 | sync::atomic::{ |
| 7 | AtomicBool, |
| 8 | AtomicU64, |
| 9 | Ordering, |
| 10 | }, |
| 11 | thread, |
| 12 | time::Duration, |
| 13 | }; |
| 14 | |
| 15 | static BARRIER_DELAY_MS: AtomicU64 = AtomicU64::new(0); // before each barrier |
| 16 | static PUBLISH_DELAY_MS: AtomicU64 = AtomicU64::new(0); // before channels are handed over |
| 17 | static COLLECT_DELAY_MS: AtomicU64 = AtomicU64::new(0); // before each garbage collection |
| 18 | static FORWARD_DELAY_MS: AtomicU64 = AtomicU64::new(0); // before a supersession is forwarded |
| 19 | static INSERT_DELAY_MS: AtomicU64 = AtomicU64::new(0); // before each cache bot insert |
| 20 | static SUP_PANICS: AtomicBool = AtomicBool::new(false); // the supervisor panics starting up |
| 21 | static BARRIER_FAILS: AtomicBool = AtomicBool::new(false); // every durability barrier fails |
| 22 | static BARRIERS_FAILED: AtomicU64 = AtomicU64::new(0); // failed by the switch above |
| 23 | static SYNCER_STOPS: AtomicBool = AtomicBool::new(false); // syncers stop after their next batch |
| 24 | static SYNCERS_STOPPED: AtomicU64 = AtomicU64::new(0); // stopped by the switch above |
| 25 | static PAIR_HAND_FAILS: AtomicBool = AtomicBool::new(false); // new live pairs cannot be handed over |
| 26 | |
| 27 | /// Holds every durability barrier this long before it syncs, as an fsync queued behind the rest |
| 28 | /// of a busy disk's writes would be held. |
| 29 | pub fn set_barrier_delay(d: Duration) { |
| 30 | BARRIER_DELAY_MS.store(millis(d), Ordering::Relaxed); |
| 31 | } |
| 32 | |
| 33 | /// Holds the supervisor this long before it hands the database's channels to the handle that |
| 34 | /// started it, as a starved machine spawning a few dozen threads would. |
| 35 | pub fn set_publish_delay(d: Duration) { |
| 36 | PUBLISH_DELAY_MS.store(millis(d), Ordering::Relaxed); |
| 37 | } |
| 38 | |
| 39 | /// Holds every garbage collection this long before it reads the file, as a large file on a busy |
| 40 | /// disk would take, so that writes can be made to land while a file is being collected. |
| 41 | pub fn set_collect_delay(d: Duration) { |
| 42 | COLLECT_DELAY_MS.store(millis(d), Ordering::Relaxed); |
| 43 | } |
| 44 | |
| 45 | /// Holds a file bot this long before it forwards a supersession to the file bot of the superseded |
| 46 | /// record's file, as a file bot behind a long queue would, so that the supersession can reach a |
| 47 | /// file after a collection of it has finished. |
| 48 | pub fn set_forward_delay(d: Duration) { |
| 49 | FORWARD_DELAY_MS.store(millis(d), Ordering::Relaxed); |
| 50 | } |
| 51 | |
| 52 | /// Holds every cache bot insert this long, as a cache bot behind a long queue would, so that a |
| 53 | /// shutdown's time can run out with written records still queued at it. |
| 54 | pub fn set_insert_delay(d: Duration) { |
| 55 | INSERT_DELAY_MS.store(millis(d), Ordering::Relaxed); |
| 56 | } |
| 57 | |
| 58 | /// Makes the supervisor panic once it has brought the bots up, before the database is ready, as a |
| 59 | /// fault in its start-up would. |
| 60 | pub fn set_supervisor_panic(on: bool) { |
| 61 | SUP_PANICS.store(on, Ordering::Relaxed); |
| 62 | } |
| 63 | |
| 64 | /// Makes every durability barrier fail without syncing, as a disk that has started returning |
| 65 | /// write-back errors would, and counts each barrier it fails. |
| 66 | pub fn set_barrier_failure(on: bool) { |
| 67 | BARRIER_FAILS.store(on, Ordering::Relaxed); |
| 68 | } |
| 69 | |
| 70 | /// How many durability barriers `set_barrier_failure` has failed so far. |
| 71 | pub fn barriers_failed() -> u64 { |
| 72 | BARRIERS_FAILED.load(Ordering::Relaxed) |
| 73 | } |
| 74 | |
| 75 | /// Makes each writer's syncer stop once it has released the records it holds, as one that had |
| 76 | /// panicked would, and counts the syncers it stops. |
| 77 | pub fn set_syncer_stop(on: bool) { |
| 78 | SYNCER_STOPS.store(on, Ordering::Relaxed); |
| 79 | } |
| 80 | |
| 81 | /// How many syncers `set_syncer_stop` has stopped so far. |
| 82 | pub fn syncers_stopped() -> u64 { |
| 83 | SYNCERS_STOPPED.load(Ordering::Relaxed) |
| 84 | } |
| 85 | |
| 86 | /// Makes handing a writer's new live pair to its syncer fail, as running out of file descriptors |
| 87 | /// to duplicate the pair's with would. |
| 88 | pub fn set_pair_hand_failure(on: bool) { |
| 89 | PAIR_HAND_FAILS.store(on, Ordering::Relaxed); |
| 90 | } |
| 91 | |
| 92 | pub(crate) fn barrier_delay() { |
| 93 | pause(&BARRIER_DELAY_MS); |
| 94 | } |
| 95 | |
| 96 | pub(crate) fn publish_delay() { |
| 97 | pause(&PUBLISH_DELAY_MS); |
| 98 | } |
| 99 | |
| 100 | pub(crate) fn collect_delay() { |
| 101 | pause(&COLLECT_DELAY_MS); |
| 102 | } |
| 103 | |
| 104 | pub(crate) fn forward_delay() { |
| 105 | pause(&FORWARD_DELAY_MS); |
| 106 | } |
| 107 | |
| 108 | pub(crate) fn insert_delay() { |
| 109 | pause(&INSERT_DELAY_MS); |
| 110 | } |
| 111 | |
| 112 | pub(crate) fn supervisor_panic() { |
| 113 | if SUP_PANICS.load(Ordering::Relaxed) { |
| 114 | panic!("The supervisor panicked starting up (test::hooks::set_supervisor_panic)."); |
| 115 | } |
| 116 | } |
| 117 | |
| 118 | /// Is this syncer to stop now? Counted when it is. |
| 119 | pub(crate) fn syncer_stops() -> bool { |
| 120 | let stops = SYNCER_STOPS.load(Ordering::Relaxed); |
| 121 | if stops { |
| 122 | SYNCERS_STOPPED.fetch_add(1, Ordering::Relaxed); |
| 123 | } |
| 124 | stops |
| 125 | } |
| 126 | |
| 127 | pub(crate) fn pair_hand_fails() -> bool { |
| 128 | PAIR_HAND_FAILS.load(Ordering::Relaxed) |
| 129 | } |
| 130 | |
| 131 | /// Is the disk to fail this sync? Counted when it is. |
| 132 | pub(crate) fn sync_fails() -> bool { |
| 133 | let fails = BARRIER_FAILS.load(Ordering::Relaxed); |
| 134 | if fails { |
| 135 | BARRIERS_FAILED.fetch_add(1, Ordering::Relaxed); |
| 136 | } |
| 137 | fails |
| 138 | } |
| 139 | |
| 140 | fn millis(d: Duration) -> u64 { |
| 141 | d.as_millis().min(u64::MAX as u128) as u64 |
| 142 | } |
| 143 | |
| 144 | fn pause(ms: &AtomicU64) { |
| 145 | let ms = ms.load(Ordering::Relaxed); |
| 146 | if ms > 0 { |
| 147 | thread::sleep(Duration::from_millis(ms)); |
| 148 | } |
| 149 | } |