turnframe_store/memory/fault.rs
1//! Failure injection at the boundaries spec §27.7 asks chaos tests to break.
2//!
3//! Recovery code is the part of a system that is never exercised by the happy
4//! path, so the in-memory store lets a test decide exactly where the process
5//! "dies". [`MemoryStores::fail_next`](super::MemoryStores::fail_next) arms one
6//! failure at a [`FailurePoint`]; the next call that reaches that boundary
7//! returns the armed [`StoreError`] and the arming is consumed. Several
8//! failures may be armed at once, including several at the same point: they
9//! fire in the order they were armed.
10//!
11//! Whether the write that precedes the boundary survives is the whole point,
12//! and it differs per point. [`FailurePoint`] documents it variant by variant;
13//! the rule of thumb is that a point named `After…` leaves the preceding write
14//! **visible** — that is the partial state recovery has to cope with — while a
15//! point named `Before…` writes nothing.
16//!
17//! Inside [`CommitStore::commit`](crate::commit::CommitStore::commit) the rule
18//! is different, and it has to be: the bundle is atomic. A failure armed at a
19//! point the bundle passes through aborts the whole bundle, and *nothing* of it
20//! becomes visible — not even the stages that had already been applied to the
21//! staged state.
22
23use std::collections::VecDeque;
24
25use crate::error::StoreError;
26
27/// A boundary at which the in-memory store can be made to fail (spec §27.7).
28///
29/// The names describe the moment in a turn, not the method: one point may be
30/// reached from more than one method, and the variant documentation names them.
31#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
32#[non_exhaustive]
33pub enum FailurePoint {
34 /// The card was written and the caller is told it was not.
35 ///
36 /// Fires in [`InteractionWriter::insert`](crate::interaction::InteractionWriter::insert)
37 /// and
38 /// [`insert_replacing_blocking`](crate::interaction::InteractionWriter::insert_replacing_blocking)
39 /// **after** the insert: the interaction (and any invalidation the insert
40 /// performed) stays visible. Spec §15.5 requires interaction persistence to
41 /// succeed before the response mentions the card, so the truthful reaction
42 /// is to omit the card, not to claim it.
43 ///
44 /// Inside a commit bundle it fires after stage 5 (interaction inserts) and
45 /// aborts the bundle.
46 AfterInteractionPersistence,
47 /// The command was never admitted.
48 ///
49 /// Fires in [`CommandJournalWriter::begin`](crate::journal::CommandJournalWriter::begin)
50 /// **before** anything is written, so the key stays free and the command
51 /// may be admitted again from scratch.
52 BeforeJournalInsert,
53 /// The command was admitted and then the process died before the domain
54 /// commit.
55 ///
56 /// Fires in [`CommandJournalWriter::begin`](crate::journal::CommandJournalWriter::begin)
57 /// **after** the entry is persisted: the entry survives in `Pending`, which
58 /// is exactly the state
59 /// [`pending_for_turn`](crate::journal::CommandJournalReader::pending_for_turn)
60 /// exists to find, and recovery must resume it by idempotency key rather
61 /// than admitting a second command (spec §23.1).
62 ///
63 /// Inside a commit bundle it fires after stage 1 (journal completions) and
64 /// aborts the bundle.
65 AfterJournalInsertBeforeCommit,
66 /// The events were appended and the caller could not read them back.
67 ///
68 /// Fires in [`EventJournalWriter::append`](crate::events::EventJournalWriter::append)
69 /// **after** the batch is appended: the events are in the ledger, so a
70 /// recovery that re-reads by command finds them and must not append them a
71 /// second time.
72 ///
73 /// Inside a commit bundle it fires after stage 2 (event batches) and aborts
74 /// the bundle.
75 AfterCommitBeforeEventReadback,
76 /// The dispatcher never got the work.
77 ///
78 /// Fires in [`OutboxWriter::claim_due`](crate::outbox::OutboxWriter::claim_due)
79 /// **before** anything is claimed: no row moves to `Dispatching` and the
80 /// next sweep picks the same rows up.
81 BeforeOutboxDispatch,
82 /// The external call was made and its result could not be recorded.
83 ///
84 /// Fires in [`mark_completed`](crate::outbox::OutboxWriter::mark_completed),
85 /// [`mark_failed`](crate::outbox::OutboxWriter::mark_failed) and
86 /// [`mark_outcome_unknown`](crate::outbox::OutboxWriter::mark_outcome_unknown)
87 /// **before** the write: the row stays `Dispatching` with its claim, which
88 /// is the state
89 /// [`release_expired_claims`](crate::outbox::OutboxWriter::release_expired_claims)
90 /// exists to reap. It is the worst case of spec §16.5 — an effect may exist
91 /// and nothing local says so — so a retry is only safe when the remote
92 /// guarantees idempotency.
93 AfterOutboxDispatch,
94 /// The answer was composed and never stored.
95 ///
96 /// Fires in
97 /// [`append_assistant_turn`](crate::conversation::ConversationWriter::append_assistant_turn)
98 /// **before** the write. The turn keeps its user side and its phase marker,
99 /// so recovery regenerates the response from committed events and stored
100 /// answer tasks instead of re-executing anything (spec §23.1).
101 BeforeResponsePersistence,
102}
103
104impl FailurePoint {
105 /// Every point, in declaration order.
106 pub const ALL: [Self; 7] = [
107 Self::AfterInteractionPersistence,
108 Self::BeforeJournalInsert,
109 Self::AfterJournalInsertBeforeCommit,
110 Self::AfterCommitBeforeEventReadback,
111 Self::BeforeOutboxDispatch,
112 Self::AfterOutboxDispatch,
113 Self::BeforeResponsePersistence,
114 ];
115
116 /// The boundary's stable name, matching the crash boundaries the
117 /// specification enumerates.
118 ///
119 /// Stable across releases: test kits and fixtures address a boundary by
120 /// this string, so renaming one is a breaking change.
121 #[must_use]
122 pub const fn as_str(self) -> &'static str {
123 match self {
124 Self::AfterInteractionPersistence => "after_interaction_persistence",
125 Self::BeforeJournalInsert => "before_journal_insert",
126 Self::AfterJournalInsertBeforeCommit => "after_journal_insert_before_commit",
127 Self::AfterCommitBeforeEventReadback => "after_commit_before_event_readback",
128 Self::BeforeOutboxDispatch => "before_outbox_dispatch",
129 Self::AfterOutboxDispatch => "after_outbox_dispatch",
130 Self::BeforeResponsePersistence => "before_response_persistence",
131 }
132 }
133
134 /// Parses a boundary from the name [`as_str`](Self::as_str) renders.
135 ///
136 /// # Errors
137 /// Returns the unknown name when it matches no boundary.
138 pub fn parse(name: &str) -> Result<Self, UnknownFailurePoint> {
139 Self::ALL
140 .into_iter()
141 .find(|point| point.as_str() == name)
142 .ok_or_else(|| UnknownFailurePoint {
143 name: name.to_owned(),
144 })
145 }
146
147 /// Returns `true` when the write that precedes the boundary stays visible
148 /// after the injected failure (outside a commit bundle, which is always
149 /// all-or-nothing).
150 #[must_use]
151 pub fn leaves_write_visible(self) -> bool {
152 matches!(
153 self,
154 Self::AfterInteractionPersistence
155 | Self::AfterJournalInsertBeforeCommit
156 | Self::AfterCommitBeforeEventReadback
157 )
158 }
159}
160
161impl core::fmt::Display for FailurePoint {
162 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
163 f.write_str(self.as_str())
164 }
165}
166
167/// A name that matches no crash boundary.
168#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
169#[error("unknown crash boundary: {name}")]
170pub struct UnknownFailurePoint {
171 /// The name that did not match.
172 pub name: String,
173}
174
175/// The queue of armed failures. Ordered, so a test can script a sequence.
176#[derive(Debug, Default)]
177pub(super) struct FaultQueue {
178 armed: VecDeque<(FailurePoint, StoreError)>,
179}
180
181impl FaultQueue {
182 /// Arms one failure at `point`.
183 pub(super) fn arm(&mut self, point: FailurePoint, error: StoreError) {
184 self.armed.push_back((point, error));
185 }
186
187 /// Consumes and returns the first failure armed at `point`, if any.
188 pub(super) fn take(&mut self, point: FailurePoint) -> Option<StoreError> {
189 let index = self.armed.iter().position(|(armed, _)| *armed == point)?;
190 self.armed.remove(index).map(|(_, error)| error)
191 }
192
193 /// The points still armed, in arming order.
194 pub(super) fn armed_points(&self) -> Vec<FailurePoint> {
195 self.armed.iter().map(|(point, _)| *point).collect()
196 }
197
198 /// Disarms everything.
199 pub(super) fn clear(&mut self) {
200 self.armed.clear();
201 }
202}
203
204#[cfg(test)]
205mod tests {
206 #[test]
207 fn every_boundary_has_a_stable_name_that_round_trips() {
208 let mut seen = std::collections::BTreeSet::new();
209 for point in FailurePoint::ALL {
210 let name = point.as_str();
211 assert!(seen.insert(name), "boundary names must be unique: {name}");
212 assert_eq!(FailurePoint::parse(name), Ok(point));
213 assert_eq!(point.to_string(), name);
214 }
215 assert!(FailurePoint::parse("nope").is_err());
216 }
217
218 use super::*;
219
220 #[test]
221 fn queue_fires_in_arming_order_per_point() {
222 let mut queue = FaultQueue::default();
223 queue.arm(FailurePoint::BeforeJournalInsert, StoreError::Unavailable);
224 queue.arm(FailurePoint::AfterOutboxDispatch, StoreError::Timeout);
225 queue.arm(FailurePoint::BeforeJournalInsert, StoreError::Timeout);
226 assert_eq!(
227 queue.armed_points(),
228 vec![
229 FailurePoint::BeforeJournalInsert,
230 FailurePoint::AfterOutboxDispatch,
231 FailurePoint::BeforeJournalInsert
232 ]
233 );
234 assert_eq!(
235 queue.take(FailurePoint::BeforeJournalInsert),
236 Some(StoreError::Unavailable)
237 );
238 assert_eq!(
239 queue.take(FailurePoint::BeforeJournalInsert),
240 Some(StoreError::Timeout)
241 );
242 assert_eq!(queue.take(FailurePoint::BeforeJournalInsert), None);
243 assert_eq!(
244 queue.armed_points(),
245 vec![FailurePoint::AfterOutboxDispatch]
246 );
247 queue.clear();
248 assert!(queue.armed_points().is_empty());
249 }
250
251 #[test]
252 fn write_visibility_is_declared_per_point() {
253 for point in FailurePoint::ALL {
254 let expected = matches!(
255 point,
256 FailurePoint::AfterInteractionPersistence
257 | FailurePoint::AfterJournalInsertBeforeCommit
258 | FailurePoint::AfterCommitBeforeEventReadback
259 );
260 assert_eq!(point.leaves_write_visible(), expected, "{point:?}");
261 }
262 }
263}