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 }