agent_effects/fault.rs
1//! Crash simulation for testing recovery. The injector is enabled by the
2//! `fault-injection` feature.
3//!
4//! A `FaultInjector` armed at a [`FaultPoint`] stops the runtime there:
5//!
6//! - `Fault::Crash` panics inside the runtime's task. Nothing after the
7//! point runs: no further writes, no lease release, no heartbeat. The
8//! caller gets a `RuntimeError::Internal`, and a new runtime over the same
9//! store plays the restarted process.
10//! - `Fault::Abort` calls [`std::process::abort`], killing the process at
11//! that exact point, for tests that run the runtime in a subprocess.
12//!
13//! An action already spawned when the crash hits keeps running, like a
14//! request already on the wire.
15//!
16//! ```ignore
17//! let injector = Arc::new(FaultInjector::new().at(FaultPoint::AfterActionReturned).crash());
18//! let runtime = Runtime::builder(store).fault_injector(Arc::clone(&injector)).build();
19//! ```
20
21/// Where in an effect's execution a crash can be injected, in the order
22/// they are reached.
23#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
24#[non_exhaustive]
25pub enum FaultPoint {
26 /// The call started; nothing is recorded yet.
27 BeforeInsert,
28 /// The record exists (`Pending`); no lease is held yet.
29 AfterInsert,
30 /// `RequestApproval` is persisted (`AwaitingApproval`); the approval
31 /// provider was not asked.
32 AfterApprovalRequested,
33 /// `StartAttempt` is persisted (`Executing`); the action was not called.
34 AfterAttemptPersisted,
35 /// The action is running: its request may or may not reach the remote
36 /// system.
37 AfterActionStarted,
38 /// The action returned; its result is not persisted.
39 AfterActionReturned,
40 /// `StartVerification` is persisted (`Verifying`); no check has run.
41 AfterVerificationStarted,
42 /// A compensation attempt is recorded (`Compensating`); the compensation
43 /// was not called.
44 AfterCompensationStarted,
45}
46
47#[cfg(feature = "fault-injection")]
48pub use injector::{Arming, Fault, FaultInjector};
49
50#[cfg(feature = "fault-injection")]
51mod injector {
52 use std::collections::HashMap;
53 use std::sync::{Mutex, PoisonError};
54
55 use super::FaultPoint;
56
57 /// What happens at an armed point.
58 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
59 pub enum Fault {
60 /// Panic in the runtime's task: an in-process crash.
61 Crash,
62 /// Abort the process.
63 Abort,
64 }
65
66 /// Stops the runtime at chosen points. Each armed fault fires once.
67 #[derive(Debug, Default)]
68 pub struct FaultInjector {
69 armed: Mutex<HashMap<FaultPoint, Fault>>,
70 reached: Mutex<Vec<FaultPoint>>,
71 }
72
73 /// A point being armed; finish with [`Arming::crash`] or
74 /// [`Arming::abort`].
75 #[must_use]
76 pub struct Arming {
77 injector: FaultInjector,
78 point: FaultPoint,
79 }
80
81 impl Arming {
82 /// Crash in-process at the point.
83 pub fn crash(self) -> FaultInjector {
84 self.arm(Fault::Crash)
85 }
86
87 /// Abort the process at the point.
88 pub fn abort(self) -> FaultInjector {
89 self.arm(Fault::Abort)
90 }
91
92 fn arm(self, fault: Fault) -> FaultInjector {
93 self.injector
94 .armed
95 .lock()
96 .unwrap_or_else(PoisonError::into_inner)
97 .insert(self.point, fault);
98 self.injector
99 }
100 }
101
102 impl FaultInjector {
103 /// An injector with nothing armed.
104 pub fn new() -> Self {
105 Self::default()
106 }
107
108 /// Arms `point`.
109 pub fn at(self, point: FaultPoint) -> Arming {
110 Arming {
111 injector: self,
112 point,
113 }
114 }
115
116 /// The points reached so far, in order, including any that crashed.
117 pub fn reached(&self) -> Vec<FaultPoint> {
118 self.reached
119 .lock()
120 .unwrap_or_else(PoisonError::into_inner)
121 .clone()
122 }
123
124 pub(crate) fn reach(&self, point: FaultPoint) {
125 self.reached
126 .lock()
127 .unwrap_or_else(PoisonError::into_inner)
128 .push(point);
129 let fault = self
130 .armed
131 .lock()
132 .unwrap_or_else(PoisonError::into_inner)
133 .remove(&point);
134 match fault {
135 Some(Fault::Crash) => panic!("fault injected at {point:?}"),
136 Some(Fault::Abort) => std::process::abort(),
137 None => {}
138 }
139 }
140 }
141}