Skip to main content

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}