turnframe_store/memory/mod.rs
1//! [`MemoryStores`]: one deterministic implementation of every persistence
2//! trait, over a single shared state.
3//!
4//! It exists for three jobs. It is what tests, examples and the `turnframe`
5//! facade use when no database is configured. It is the reference an adapter
6//! author reads when the prose in a trait leaves a question open — if
7//! PostgreSQL and this store disagree, one of them is wrong. And it is the
8//! subject the [`conformance`](crate::conformance) suite is developed against,
9//! so a suite failure means the adapter is wrong rather than the suite.
10//!
11//! # What "deterministic" buys
12//!
13//! * **Order.** Every collection is a `BTreeMap` keyed by account first, and
14//! every list is sorted by the key its trait documents. Two runs over the
15//! same writes produce the same lists, so `assert_eq!` on a whole list is a
16//! fair test rather than a flaky one.
17//! * **Time.** The store never reads the wall clock directly; it asks a
18//! [`Clock`]. With a [`ManualClock`], "the card expired" is something a test
19//! makes true, not something it waits for.
20//! * **Failure.** [`MemoryStores::fail_next`] arms a [`FailurePoint`], so the
21//! recovery paths of spec §27.7 are reachable without killing a process.
22//! * **Cost.** A commit is all-or-nothing without copying the state: the bundle
23//! is applied in place and every mutation records how to undo itself, so a
24//! failure replays the reversals instead of discarding a clone. A commit
25//! costs what it writes, not what the store already holds, which is what
26//! keeps a suite of thousands of turns from slowing down as it goes.
27//!
28//! # Concurrency
29//!
30//! One `std::sync::Mutex` guards the whole state and a second, independent one
31//! guards the armed failures. Every trait method locks, runs one synchronous
32//! rule function and releases before returning: no guard is ever alive across
33//! an await, which is why these futures are `Send` at all. A poisoned lock is
34//! reported as [`StoreError::Corrupt`] rather than propagating a panic.
35//!
36//! # Example
37//!
38//! ```rust
39//! use std::sync::Arc;
40//!
41//! use chrono::TimeDelta;
42//! use turnframe_store::memory::{Clock, ManualClock, MemoryStores};
43//!
44//! let clock = Arc::new(ManualClock::epoch());
45//! let stores = MemoryStores::with_clock(clock.clone());
46//! assert_eq!(stores.now(), clock.now());
47//!
48//! clock.advance(TimeDelta::minutes(5));
49//! assert_eq!(stores.now(), clock.now());
50//! ```
51
52mod clock;
53mod fault;
54mod impls;
55mod state;
56
57use std::sync::{Arc, Mutex, MutexGuard};
58
59use chrono::{DateTime, Utc};
60
61pub use self::clock::{Clock, ManualClock, SystemClock};
62use self::fault::FaultQueue;
63pub use self::fault::{FailurePoint, UnknownFailurePoint};
64use self::state::Inner;
65use crate::error::StoreError;
66
67/// A complete, deterministic set of stores backed by one in-process state.
68///
69/// Implements [`ConversationStore`](crate::conversation::ConversationStore),
70/// [`InteractionStore`](crate::interaction::InteractionStore),
71/// [`CommandJournal`](crate::journal::CommandJournal),
72/// [`EventJournal`](crate::events::EventJournal),
73/// [`OutboxStore`](crate::outbox::OutboxStore),
74/// [`ReplayStore`](crate::replay::ReplayStore) and
75/// [`CommitStore`](crate::commit::CommitStore) over the same state, so a card
76/// written through the interaction trait is the card a commit bundle settles.
77///
78/// Wrap it in [`Stores::from_memory`](crate::stores::Stores::from_memory), or
79/// take the shortcut [`Stores::in_memory`](crate::stores::Stores::in_memory).
80#[derive(Debug)]
81pub struct MemoryStores {
82 inner: Mutex<Inner>,
83 faults: Mutex<FaultQueue>,
84 clock: Arc<dyn Clock>,
85}
86
87impl MemoryStores {
88 /// An empty store on the wall clock.
89 #[must_use]
90 pub fn new() -> Self {
91 Self::with_clock(Arc::new(SystemClock))
92 }
93
94 /// An empty store that stamps its own columns from `clock`.
95 ///
96 /// Keep a clone of the clock to drive it: with a [`ManualClock`] the test
97 /// decides when a claim goes stale or a phase marker moves.
98 #[must_use]
99 pub fn with_clock(clock: Arc<dyn Clock>) -> Self {
100 Self {
101 inner: Mutex::new(Inner::default()),
102 faults: Mutex::new(FaultQueue::default()),
103 clock,
104 }
105 }
106
107 /// The instant this store would stamp a write with right now.
108 #[must_use]
109 pub fn now(&self) -> DateTime<Utc> {
110 self.clock.now()
111 }
112
113 /// Arms one failure: the next call reaching `point` returns `error` instead
114 /// of its normal result.
115 ///
116 /// Whether the write that precedes the boundary survives is part of the
117 /// point's meaning; see [`FailurePoint`]. Arming several failures is
118 /// allowed, including several at the same point, and they fire in arming
119 /// order. Inside
120 /// [`CommitStore::commit`](crate::commit::CommitStore::commit) a fired
121 /// failure discards the whole bundle.
122 ///
123 /// # Errors
124 /// * `Corrupt` when the internal lock is poisoned.
125 pub fn fail_next(&self, point: FailurePoint, error: StoreError) -> Result<(), StoreError> {
126 self.faults()?.arm(point, error);
127 Ok(())
128 }
129
130 /// The points still armed, in arming order.
131 ///
132 /// # Errors
133 /// * `Corrupt` when the internal lock is poisoned.
134 pub fn armed_failures(&self) -> Result<Vec<FailurePoint>, StoreError> {
135 Ok(self.faults()?.armed_points())
136 }
137
138 /// Disarms every pending failure.
139 ///
140 /// # Errors
141 /// * `Corrupt` when the internal lock is poisoned.
142 pub fn clear_failures(&self) -> Result<(), StoreError> {
143 self.faults()?.clear();
144 Ok(())
145 }
146
147 /// The state guard. Held for one synchronous rule call and released; the
148 /// lock is never alive across an await.
149 fn inner(&self) -> Result<MutexGuard<'_, Inner>, StoreError> {
150 self.inner.lock().map_err(|_| StoreError::Corrupt)
151 }
152
153 fn faults(&self) -> Result<MutexGuard<'_, FaultQueue>, StoreError> {
154 self.faults.lock().map_err(|_| StoreError::Corrupt)
155 }
156
157 /// Returns the armed error for `point`, consuming the arming.
158 fn fire(&self, point: FailurePoint) -> Result<(), StoreError> {
159 match self.faults()?.take(point) {
160 Some(error) => Err(error),
161 None => Ok(()),
162 }
163 }
164}
165
166impl Default for MemoryStores {
167 fn default() -> Self {
168 Self::new()
169 }
170}
171
172#[cfg(test)]
173mod tests {
174 use super::*;
175 use crate::conversation::{ConversationRecord, ConversationWriter};
176 use turnframe_core::ids::{AccountId, ConversationId};
177
178 #[test]
179 fn failures_are_armed_consumed_and_cleared() {
180 let stores = MemoryStores::new();
181 assert!(stores.armed_failures().unwrap().is_empty());
182 stores
183 .fail_next(FailurePoint::BeforeOutboxDispatch, StoreError::Unavailable)
184 .unwrap();
185 stores
186 .fail_next(FailurePoint::BeforeJournalInsert, StoreError::Timeout)
187 .unwrap();
188 assert_eq!(
189 stores.armed_failures().unwrap(),
190 vec![
191 FailurePoint::BeforeOutboxDispatch,
192 FailurePoint::BeforeJournalInsert
193 ]
194 );
195 assert_eq!(
196 stores.fire(FailurePoint::BeforeJournalInsert),
197 Err(StoreError::Timeout)
198 );
199 assert_eq!(
200 stores.armed_failures().unwrap(),
201 vec![FailurePoint::BeforeOutboxDispatch]
202 );
203 stores.clear_failures().unwrap();
204 assert!(stores.armed_failures().unwrap().is_empty());
205 assert_eq!(stores.fire(FailurePoint::BeforeOutboxDispatch), Ok(()));
206 }
207
208 #[tokio::test]
209 async fn clock_drives_the_stamps_the_store_owns() {
210 let clock = Arc::new(ManualClock::epoch());
211 let stores = MemoryStores::with_clock(clock.clone());
212 let account = AccountId::from("a");
213 let conversation = ConversationId::new();
214 stores
215 .create_conversation(ConversationRecord::new(
216 conversation,
217 account.clone(),
218 clock.now(),
219 ))
220 .await
221 .unwrap();
222 assert_eq!(stores.now(), DateTime::<Utc>::UNIX_EPOCH);
223 clock.advance(chrono::TimeDelta::seconds(90));
224 assert_eq!(
225 stores.now(),
226 DateTime::<Utc>::UNIX_EPOCH + chrono::TimeDelta::seconds(90)
227 );
228 }
229}