lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

failpoint.rs (13450B)


      1 //! Private per-operation durability failpoints used by deterministic tests.
      2 
      3 use core::fmt;
      4 
      5 #[cfg(test)]
      6 use std::{
      7     io::{self, Write},
      8     sync::{Arc, Condvar, Mutex},
      9 };
     10 
     11 #[cfg(test)]
     12 const PROCESS_BARRIER_READY: &[u8] = b"\nRSHR_STEP073_READY\n";
     13 
     14 #[cfg(test)]
     15 pub(crate) fn storage_full_error() -> io::Error {
     16     io::Error::from(io::ErrorKind::StorageFull)
     17 }
     18 
     19 /// Closed inventory of durability edges exercised by the crash-boundary harness.
     20 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     21 pub(crate) enum DurabilityFailpoint {
     22     InitializeBeforeCreate,
     23     InitializeAfterCreate,
     24     InitializeBeforeReservationDirectorySync,
     25     InitializeAfterReservationDirectorySync,
     26     InitializeBeforeFileSync,
     27     InitializeAfterFileSync,
     28     InitializeBeforeCommitDirectorySync,
     29     InitializeAfterCommitDirectorySync,
     30     TransactionBeforeBegin,
     31     TransactionAfterBegin,
     32     TransactionBeforeCommit,
     33     TransactionAfterCommit,
     34     BackupBeforeCreate,
     35     BackupAfterCreate,
     36     BackupBeforeCopy,
     37     BackupAfterCopy,
     38     BackupBeforeFileSync,
     39     BackupAfterFileSync,
     40     BackupBeforeDirectorySync,
     41     BackupAfterDirectorySync,
     42     MarkerBeforeCreate,
     43     MarkerAfterCreate,
     44     MarkerBeforeFileSync,
     45     MarkerAfterFileSync,
     46     MarkerBeforeDirectorySync,
     47     MarkerAfterDirectorySync,
     48     MarkerAdvanceBeforeWriteAndFileSync,
     49     MarkerAdvanceAfterWriteAndFileSync,
     50     MarkerAdvanceBeforeReplace,
     51     MarkerAdvanceAfterReplace,
     52     MarkerAdvanceBeforeDirectorySync,
     53     MarkerAdvanceAfterDirectorySync,
     54     RestoreBeforeRetainLiveRename,
     55     RestoreAfterRetainLiveRename,
     56     RestoreBeforeRetainLiveSync,
     57     RestoreAfterRetainLiveSync,
     58     RestoreBeforeInstallStageRename,
     59     RestoreAfterInstallStageRename,
     60     RestoreBeforeInstallStageSync,
     61     RestoreAfterInstallStageSync,
     62     CloseBeforeDrain,
     63     CloseAfterDrain,
     64     CloseBeforeCheckpoint,
     65     CloseAfterCheckpoint,
     66     CloseBeforeConnectionClose,
     67     CloseAfterConnectionClose,
     68     CloseBeforeAuthorityRelease,
     69     CloseAfterAuthorityRelease,
     70 }
     71 
     72 impl DurabilityFailpoint {
     73     #[cfg(test)]
     74     pub(crate) const ALL: [Self; 48] = [
     75         Self::InitializeBeforeCreate,
     76         Self::InitializeAfterCreate,
     77         Self::InitializeBeforeReservationDirectorySync,
     78         Self::InitializeAfterReservationDirectorySync,
     79         Self::InitializeBeforeFileSync,
     80         Self::InitializeAfterFileSync,
     81         Self::InitializeBeforeCommitDirectorySync,
     82         Self::InitializeAfterCommitDirectorySync,
     83         Self::TransactionBeforeBegin,
     84         Self::TransactionAfterBegin,
     85         Self::TransactionBeforeCommit,
     86         Self::TransactionAfterCommit,
     87         Self::BackupBeforeCreate,
     88         Self::BackupAfterCreate,
     89         Self::BackupBeforeCopy,
     90         Self::BackupAfterCopy,
     91         Self::BackupBeforeFileSync,
     92         Self::BackupAfterFileSync,
     93         Self::BackupBeforeDirectorySync,
     94         Self::BackupAfterDirectorySync,
     95         Self::MarkerBeforeCreate,
     96         Self::MarkerAfterCreate,
     97         Self::MarkerBeforeFileSync,
     98         Self::MarkerAfterFileSync,
     99         Self::MarkerBeforeDirectorySync,
    100         Self::MarkerAfterDirectorySync,
    101         Self::MarkerAdvanceBeforeWriteAndFileSync,
    102         Self::MarkerAdvanceAfterWriteAndFileSync,
    103         Self::MarkerAdvanceBeforeReplace,
    104         Self::MarkerAdvanceAfterReplace,
    105         Self::MarkerAdvanceBeforeDirectorySync,
    106         Self::MarkerAdvanceAfterDirectorySync,
    107         Self::RestoreBeforeRetainLiveRename,
    108         Self::RestoreAfterRetainLiveRename,
    109         Self::RestoreBeforeRetainLiveSync,
    110         Self::RestoreAfterRetainLiveSync,
    111         Self::RestoreBeforeInstallStageRename,
    112         Self::RestoreAfterInstallStageRename,
    113         Self::RestoreBeforeInstallStageSync,
    114         Self::RestoreAfterInstallStageSync,
    115         Self::CloseBeforeDrain,
    116         Self::CloseAfterDrain,
    117         Self::CloseBeforeCheckpoint,
    118         Self::CloseAfterCheckpoint,
    119         Self::CloseBeforeConnectionClose,
    120         Self::CloseAfterConnectionClose,
    121         Self::CloseBeforeAuthorityRelease,
    122         Self::CloseAfterAuthorityRelease,
    123     ];
    124 }
    125 
    126 /// Disabled in ordinary builds; tests may arm one edge on one owned controller.
    127 #[derive(Clone, Default)]
    128 pub(crate) struct DurabilityFailpoints {
    129     #[cfg(test)]
    130     state: Arc<Mutex<TestState>>,
    131 }
    132 
    133 #[cfg(test)]
    134 #[derive(Default)]
    135 struct TestState {
    136     armed: Option<DurabilityFailpoint>,
    137     fired: bool,
    138     reached: Vec<DurabilityFailpoint>,
    139     observations: Vec<(DurabilityFailpoint, u8)>,
    140     process_barrier: Option<TestProcessBarrier>,
    141 }
    142 
    143 #[cfg(test)]
    144 #[derive(Clone)]
    145 struct TestProcessBarrier {
    146     point: DurabilityFailpoint,
    147     occurrence: u8,
    148     seen: u8,
    149     gate: Arc<TestProcessBarrierGate>,
    150 }
    151 
    152 #[cfg(test)]
    153 #[derive(Default)]
    154 struct TestProcessBarrierGate {
    155     state: Mutex<TestProcessBarrierGateState>,
    156     changed: Condvar,
    157 }
    158 
    159 #[cfg(test)]
    160 #[derive(Default)]
    161 struct TestProcessBarrierGateState {
    162     ready: bool,
    163     released: bool,
    164 }
    165 
    166 impl DurabilityFailpoints {
    167     pub(crate) fn hit(&self, point: DurabilityFailpoint) -> Result<(), DurabilityFailpointError> {
    168         #[cfg(test)]
    169         {
    170             let (injected, process_gate) = {
    171                 let mut state = self.state.lock().map_err(|_| DurabilityFailpointError)?;
    172                 if state.reached.len() < DurabilityFailpoint::ALL.len() {
    173                     state.reached.push(point);
    174                 }
    175                 let injected = state.armed == Some(point) && !state.fired;
    176                 let process_gate = if !state.fired {
    177                     state.process_barrier.as_mut().and_then(|barrier| {
    178                         if barrier.point != point {
    179                             return None;
    180                         }
    181                         barrier.seen = barrier.seen.saturating_add(1);
    182                         (barrier.seen == barrier.occurrence).then(|| Arc::clone(&barrier.gate))
    183                     })
    184                 } else {
    185                     None
    186                 };
    187                 if injected || process_gate.is_some() {
    188                     state.fired = true;
    189                 }
    190                 (injected, process_gate)
    191             };
    192             if let Some(process_gate) = process_gate {
    193                 process_gate.notify_and_wait()?;
    194                 return Err(DurabilityFailpointError);
    195             }
    196             if injected {
    197                 return Err(DurabilityFailpointError);
    198             }
    199         }
    200         #[cfg(not(test))]
    201         let _ = point;
    202         Ok(())
    203     }
    204 
    205     #[cfg(test)]
    206     pub(crate) fn armed(point: DurabilityFailpoint) -> Self {
    207         Self {
    208             state: Arc::new(Mutex::new(TestState {
    209                 armed: Some(point),
    210                 fired: false,
    211                 reached: Vec::new(),
    212                 observations: Vec::new(),
    213                 process_barrier: None,
    214             })),
    215         }
    216     }
    217 
    218     #[cfg(test)]
    219     pub(crate) fn process_barrier(point: DurabilityFailpoint, occurrence: u8) -> Self {
    220         assert!(
    221             matches!(occurrence, 1 | 2),
    222             "process occurrence must be 1 or 2"
    223         );
    224         Self {
    225             state: Arc::new(Mutex::new(TestState {
    226                 armed: None,
    227                 fired: false,
    228                 reached: Vec::new(),
    229                 observations: Vec::new(),
    230                 process_barrier: Some(TestProcessBarrier {
    231                     point,
    232                     occurrence,
    233                     seen: 0,
    234                     gate: Arc::new(TestProcessBarrierGate::default()),
    235                 }),
    236             })),
    237         }
    238     }
    239 
    240     #[cfg(test)]
    241     fn wait_for_process_barrier(&self) {
    242         let gate = self
    243             .state
    244             .lock()
    245             .expect("durability failpoint state")
    246             .process_barrier
    247             .as_ref()
    248             .map(|barrier| Arc::clone(&barrier.gate))
    249             .expect("process barrier");
    250         gate.wait_until_ready();
    251     }
    252 
    253     #[cfg(test)]
    254     fn release_process_barrier(&self) {
    255         let gate = self
    256             .state
    257             .lock()
    258             .expect("durability failpoint state")
    259             .process_barrier
    260             .as_ref()
    261             .map(|barrier| Arc::clone(&barrier.gate))
    262             .expect("process barrier");
    263         gate.release();
    264     }
    265 
    266     #[cfg(test)]
    267     pub(crate) fn arm(&self, point: DurabilityFailpoint) {
    268         let mut state = self.state.lock().expect("durability failpoint state");
    269         state.armed = Some(point);
    270         state.fired = false;
    271         state.reached.clear();
    272         state.observations.clear();
    273     }
    274 
    275     #[cfg(test)]
    276     pub(crate) fn disarm(&self) {
    277         let mut state = self.state.lock().expect("durability failpoint state");
    278         state.armed = None;
    279     }
    280 
    281     #[cfg(test)]
    282     pub(crate) fn fired(&self) -> bool {
    283         self.state.lock().is_ok_and(|state| state.fired)
    284     }
    285 
    286     #[cfg(test)]
    287     pub(crate) fn reached(&self) -> Vec<DurabilityFailpoint> {
    288         self.state
    289             .lock()
    290             .map_or_else(|_| Vec::new(), |state| state.reached.clone())
    291     }
    292 
    293     #[cfg(test)]
    294     pub(crate) fn observe(&self, point: DurabilityFailpoint, state_value: u8) {
    295         let mut state = self.state.lock().expect("durability failpoint state");
    296         state.observations.push((point, state_value));
    297     }
    298 
    299     #[cfg(test)]
    300     pub(crate) fn observation(&self, point: DurabilityFailpoint) -> Option<u8> {
    301         self.state.lock().ok().and_then(|state| {
    302             state
    303                 .observations
    304                 .iter()
    305                 .find_map(|(observed, value)| (*observed == point).then_some(*value))
    306         })
    307     }
    308 }
    309 
    310 #[cfg(test)]
    311 impl TestProcessBarrierGate {
    312     fn notify_and_wait(&self) -> Result<(), DurabilityFailpointError> {
    313         {
    314             let mut state = self.state.lock().map_err(|_| DurabilityFailpointError)?;
    315             state.ready = true;
    316             self.changed.notify_all();
    317         }
    318         let mut stdout = io::stdout().lock();
    319         stdout
    320             .write_all(PROCESS_BARRIER_READY)
    321             .and_then(|()| stdout.flush())
    322             .map_err(|_| DurabilityFailpointError)?;
    323         let state = self.state.lock().map_err(|_| DurabilityFailpointError)?;
    324         drop(
    325             self.changed
    326                 .wait_while(state, |state| !state.released)
    327                 .map_err(|_| DurabilityFailpointError)?,
    328         );
    329         Ok(())
    330     }
    331 
    332     fn wait_until_ready(&self) {
    333         let state = self.state.lock().expect("process barrier gate");
    334         drop(
    335             self.changed
    336                 .wait_while(state, |state| !state.ready)
    337                 .expect("process barrier ready"),
    338         );
    339     }
    340 
    341     fn release(&self) {
    342         let mut state = self.state.lock().expect("process barrier gate");
    343         state.released = true;
    344         self.changed.notify_all();
    345     }
    346 }
    347 
    348 impl fmt::Debug for DurabilityFailpoints {
    349     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    350         formatter.write_str("DurabilityFailpoints([redacted])")
    351     }
    352 }
    353 
    354 /// Source-free injected failure; subsystem adapters retain their stable error kind.
    355 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    356 pub(crate) struct DurabilityFailpointError;
    357 
    358 impl fmt::Display for DurabilityFailpointError {
    359     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    360         formatter.write_str("injected durability boundary failure")
    361     }
    362 }
    363 
    364 impl std::error::Error for DurabilityFailpointError {}
    365 
    366 #[cfg(test)]
    367 mod tests {
    368     use super::*;
    369     use std::thread;
    370 
    371     #[test]
    372     fn storage_full_injection_uses_the_semantic_io_kind() {
    373         assert_eq!(storage_full_error().kind(), io::ErrorKind::StorageFull);
    374     }
    375 
    376     #[test]
    377     fn every_closed_point_fires_once_on_its_owned_plan() {
    378         for point in DurabilityFailpoint::ALL {
    379             let plan = DurabilityFailpoints::armed(point);
    380             assert_eq!(plan.hit(point), Err(DurabilityFailpointError));
    381             assert_eq!(plan.hit(point), Ok(()));
    382             assert!(plan.fired());
    383             assert_eq!(plan.reached(), [point, point]);
    384         }
    385     }
    386 
    387     #[test]
    388     fn plans_are_instance_local_and_disabled_plan_never_fails() {
    389         let first = DurabilityFailpoints::armed(DurabilityFailpoint::TransactionBeforeCommit);
    390         let second = DurabilityFailpoints::armed(DurabilityFailpoint::BackupBeforeCopy);
    391         assert_eq!(first.hit(DurabilityFailpoint::BackupBeforeCopy), Ok(()));
    392         assert!(!first.fired());
    393         assert_eq!(
    394             second.hit(DurabilityFailpoint::BackupBeforeCopy),
    395             Err(DurabilityFailpointError)
    396         );
    397         assert!(!first.fired());
    398         assert!(second.fired());
    399         assert_eq!(
    400             DurabilityFailpoints::default().hit(DurabilityFailpoint::CloseBeforeDrain),
    401             Ok(())
    402         );
    403     }
    404 
    405     #[test]
    406     fn process_barrier_waits_for_the_selected_occurrence_and_releases_once() {
    407         let point = DurabilityFailpoint::MarkerAdvanceAfterDirectorySync;
    408         let plan = DurabilityFailpoints::process_barrier(point, 2);
    409         let worker_plan = plan.clone();
    410         let worker = thread::spawn(move || {
    411             assert_eq!(worker_plan.hit(point), Ok(()));
    412             worker_plan.hit(point)
    413         });
    414 
    415         plan.wait_for_process_barrier();
    416         assert!(plan.fired());
    417         plan.release_process_barrier();
    418         assert_eq!(
    419             worker.join().expect("barrier worker"),
    420             Err(DurabilityFailpointError)
    421         );
    422         assert_eq!(plan.reached(), [point, point]);
    423     }
    424 }