Skip to main content

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}