Skip to main content

turnframe_store/
stores.rs

1//! [`Stores`]: one value carrying an implementation of every persistence trait.
2//!
3//! The runtime needs all seven stores and should not care whether they come
4//! from one backend or seven. `Stores` holds them behind `Arc<dyn …>`, so it is
5//! cheap to clone, shareable across tasks, and mixable: a deployment may keep
6//! conversations in PostgreSQL and the outbox somewhere else without any type
7//! elsewhere in the library learning about it.
8//!
9//! # Building one
10//!
11//! * [`Stores::in_memory`] for tests, examples and single-process runs.
12//! * [`Stores::from_memory`] when the test needs to keep the
13//!   [`MemoryStores`] handle too, to drive its clock or arm a failure.
14//! * [`Stores::builder`] when the implementations differ.
15//!
16//! ```rust
17//! use std::sync::Arc;
18//!
19//! use turnframe_store::memory::MemoryStores;
20//! use turnframe_store::stores::Stores;
21//!
22//! # fn main() -> Result<(), turnframe_store::stores::StoresBuilderError> {
23//! let backend = Arc::new(MemoryStores::new());
24//! let stores = Stores::builder()
25//!     .conversations(backend.clone())
26//!     .interactions(backend.clone())
27//!     .journal(backend.clone())
28//!     .events(backend.clone())
29//!     .outbox(backend.clone())
30//!     .replay(backend.clone())
31//!     .commit(backend)
32//!     .build()?;
33//! assert!(format!("{stores:?}").contains("conversations"));
34//! # Ok(())
35//! # }
36//! ```
37
38use std::fmt;
39use std::sync::Arc;
40
41use crate::commit::CommitStore;
42use crate::conversation::{ConversationReader, ConversationStore};
43use crate::events::{EventJournal, EventJournalReader};
44use crate::interaction::{InteractionReader, InteractionStore};
45use crate::journal::{CommandJournal, CommandJournalReader};
46use crate::memory::MemoryStores;
47use crate::outbox::{OutboxReader, OutboxStore};
48use crate::replay::{ReplayReader, ReplayStore};
49
50/// A complete set of stores.
51///
52/// Cloning shares the same underlying implementations; it never copies data.
53#[derive(Clone)]
54pub struct Stores {
55    conversations: Arc<dyn ConversationStore>,
56    interactions: Arc<dyn InteractionStore>,
57    journal: Arc<dyn CommandJournal>,
58    events: Arc<dyn EventJournal>,
59    outbox: Arc<dyn OutboxStore>,
60    replay: Arc<dyn ReplayStore>,
61    commit: Arc<dyn CommitStore>,
62}
63
64impl Stores {
65    /// An empty [`MemoryStores`] behind every trait.
66    #[must_use]
67    pub fn in_memory() -> Self {
68        Self::from_memory(Arc::new(MemoryStores::new()))
69    }
70
71    /// Every trait served by one [`MemoryStores`] the caller keeps a handle to.
72    ///
73    /// Use it when the test also has to drive the store's clock or arm a
74    /// [`FailurePoint`](crate::memory::FailurePoint).
75    #[must_use]
76    pub fn from_memory(backend: Arc<MemoryStores>) -> Self {
77        Self {
78            conversations: backend.clone(),
79            interactions: backend.clone(),
80            journal: backend.clone(),
81            events: backend.clone(),
82            outbox: backend.clone(),
83            replay: backend.clone(),
84            commit: backend,
85        }
86    }
87
88    /// A builder that requires every store to be supplied.
89    #[must_use]
90    pub fn builder() -> StoresBuilder {
91        StoresBuilder::default()
92    }
93
94    /// Conversations, turns and phase markers.
95    #[must_use]
96    pub fn conversations(&self) -> &Arc<dyn ConversationStore> {
97        &self.conversations
98    }
99
100    /// Persisted interactions.
101    #[must_use]
102    pub fn interactions(&self) -> &Arc<dyn InteractionStore> {
103        &self.interactions
104    }
105
106    /// The command journal.
107    #[must_use]
108    pub fn journal(&self) -> &Arc<dyn CommandJournal> {
109        &self.journal
110    }
111
112    /// The claim ledger.
113    #[must_use]
114    pub fn events(&self) -> &Arc<dyn EventJournal> {
115        &self.events
116    }
117
118    /// The external-effect outbox.
119    #[must_use]
120    pub fn outbox(&self) -> &Arc<dyn OutboxStore> {
121        &self.outbox
122    }
123
124    /// Replay records.
125    #[must_use]
126    pub fn replay(&self) -> &Arc<dyn ReplayStore> {
127        &self.replay
128    }
129
130    /// The all-or-nothing bundle writer.
131    #[must_use]
132    pub fn commit(&self) -> &Arc<dyn CommitStore> {
133        &self.commit
134    }
135
136    /// The same persistence layer with the write side removed
137    /// (spec §23, plan-only path).
138    ///
139    /// Each store is upcast to its read half, so the returned value has no
140    /// method that changes anything and no way to recover one: there is no
141    /// downcast back to [`Stores`], and [`CommitStore`] — which is write-only —
142    /// is simply absent. It is what a caller hands to a path that must be
143    /// unable to write, such as
144    /// [`Orchestrator::plan_turn`](https://docs.rs/turnframe-runtime), so that
145    /// "this path does not persist anything" is checked by the compiler rather
146    /// than by review.
147    #[must_use]
148    pub fn read_only(&self) -> ReadOnlyStores {
149        ReadOnlyStores {
150            conversations: self.conversations.clone(),
151            interactions: self.interactions.clone(),
152            journal: self.journal.clone(),
153            events: self.events.clone(),
154            outbox: self.outbox.clone(),
155            replay: self.replay.clone(),
156        }
157    }
158}
159
160/// The read half of a persistence layer: six readers, no writer of any kind.
161///
162/// Build one with [`Stores::read_only`], or with [`ReadOnlyStores::builder`]
163/// when the reading implementations are not backed by a full [`Stores`] at all.
164/// Cloning shares the same underlying implementations.
165///
166/// [`CommitStore`] has no read half — a commit is
167/// the write — so it has no counterpart here.
168#[derive(Clone)]
169pub struct ReadOnlyStores {
170    conversations: Arc<dyn ConversationReader>,
171    interactions: Arc<dyn InteractionReader>,
172    journal: Arc<dyn CommandJournalReader>,
173    events: Arc<dyn EventJournalReader>,
174    outbox: Arc<dyn OutboxReader>,
175    replay: Arc<dyn ReplayReader>,
176}
177
178impl ReadOnlyStores {
179    /// An empty [`MemoryStores`] behind every reader.
180    #[must_use]
181    pub fn in_memory() -> Self {
182        Stores::in_memory().read_only()
183    }
184
185    /// A builder that requires every reader to be supplied.
186    #[must_use]
187    pub fn builder() -> ReadOnlyStoresBuilder {
188        ReadOnlyStoresBuilder::default()
189    }
190
191    /// Conversations, turns and phase markers.
192    #[must_use]
193    pub fn conversations(&self) -> &Arc<dyn ConversationReader> {
194        &self.conversations
195    }
196
197    /// Persisted interactions.
198    #[must_use]
199    pub fn interactions(&self) -> &Arc<dyn InteractionReader> {
200        &self.interactions
201    }
202
203    /// The command journal.
204    #[must_use]
205    pub fn journal(&self) -> &Arc<dyn CommandJournalReader> {
206        &self.journal
207    }
208
209    /// The claim ledger.
210    #[must_use]
211    pub fn events(&self) -> &Arc<dyn EventJournalReader> {
212        &self.events
213    }
214
215    /// The external-effect outbox.
216    #[must_use]
217    pub fn outbox(&self) -> &Arc<dyn OutboxReader> {
218        &self.outbox
219    }
220
221    /// Replay records.
222    #[must_use]
223    pub fn replay(&self) -> &Arc<dyn ReplayReader> {
224        &self.replay
225    }
226}
227
228impl From<&Stores> for ReadOnlyStores {
229    fn from(stores: &Stores) -> Self {
230        stores.read_only()
231    }
232}
233
234impl fmt::Debug for ReadOnlyStores {
235    /// Names the roles, never the implementations.
236    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
237        f.debug_struct("ReadOnlyStores")
238            .field("roles", &ReadOnlyStoresBuilder::ROLES)
239            .finish()
240    }
241}
242
243/// Collects one reader per role.
244#[derive(Clone, Default)]
245pub struct ReadOnlyStoresBuilder {
246    conversations: Option<Arc<dyn ConversationReader>>,
247    interactions: Option<Arc<dyn InteractionReader>>,
248    journal: Option<Arc<dyn CommandJournalReader>>,
249    events: Option<Arc<dyn EventJournalReader>>,
250    outbox: Option<Arc<dyn OutboxReader>>,
251    replay: Option<Arc<dyn ReplayReader>>,
252}
253
254impl ReadOnlyStoresBuilder {
255    /// The six role names, in the order [`ReadOnlyStoresBuilder::build`] checks
256    /// them.
257    pub const ROLES: [&'static str; 6] = [
258        "conversations",
259        "interactions",
260        "journal",
261        "events",
262        "outbox",
263        "replay",
264    ];
265
266    /// An empty builder.
267    #[must_use]
268    pub fn new() -> Self {
269        Self::default()
270    }
271
272    /// Sets the conversation reader.
273    #[must_use]
274    pub fn conversations(mut self, store: Arc<dyn ConversationReader>) -> Self {
275        self.conversations = Some(store);
276        self
277    }
278
279    /// Sets the interaction reader.
280    #[must_use]
281    pub fn interactions(mut self, store: Arc<dyn InteractionReader>) -> Self {
282        self.interactions = Some(store);
283        self
284    }
285
286    /// Sets the command journal reader.
287    #[must_use]
288    pub fn journal(mut self, store: Arc<dyn CommandJournalReader>) -> Self {
289        self.journal = Some(store);
290        self
291    }
292
293    /// Sets the event journal reader.
294    #[must_use]
295    pub fn events(mut self, store: Arc<dyn EventJournalReader>) -> Self {
296        self.events = Some(store);
297        self
298    }
299
300    /// Sets the outbox reader.
301    #[must_use]
302    pub fn outbox(mut self, store: Arc<dyn OutboxReader>) -> Self {
303        self.outbox = Some(store);
304        self
305    }
306
307    /// Sets the replay reader.
308    #[must_use]
309    pub fn replay(mut self, store: Arc<dyn ReplayReader>) -> Self {
310        self.replay = Some(store);
311        self
312    }
313
314    /// Builds the set.
315    ///
316    /// # Errors
317    /// * [`StoresBuilderError::MissingStore`] naming the first role that was
318    ///   never supplied.
319    pub fn build(self) -> Result<ReadOnlyStores, StoresBuilderError> {
320        fn required<T: ?Sized>(
321            store: Option<Arc<T>>,
322            role: &'static str,
323        ) -> Result<Arc<T>, StoresBuilderError> {
324            store.ok_or(StoresBuilderError::MissingStore { role })
325        }
326
327        Ok(ReadOnlyStores {
328            conversations: required(self.conversations, Self::ROLES[0])?,
329            interactions: required(self.interactions, Self::ROLES[1])?,
330            journal: required(self.journal, Self::ROLES[2])?,
331            events: required(self.events, Self::ROLES[3])?,
332            outbox: required(self.outbox, Self::ROLES[4])?,
333            replay: required(self.replay, Self::ROLES[5])?,
334        })
335    }
336}
337
338impl fmt::Debug for ReadOnlyStoresBuilder {
339    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
340        let supplied: Vec<&'static str> = Self::ROLES
341            .iter()
342            .zip([
343                self.conversations.is_some(),
344                self.interactions.is_some(),
345                self.journal.is_some(),
346                self.events.is_some(),
347                self.outbox.is_some(),
348                self.replay.is_some(),
349            ])
350            .filter_map(|(role, present)| present.then_some(*role))
351            .collect();
352        f.debug_struct("ReadOnlyStoresBuilder")
353            .field("supplied", &supplied)
354            .finish()
355    }
356}
357
358impl fmt::Debug for Stores {
359    /// Names the roles, never the implementations: a store may hold a
360    /// connection string and this output is safe to log.
361    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
362        f.debug_struct("Stores")
363            .field("roles", &StoresBuilder::ROLES)
364            .finish()
365    }
366}
367
368/// Why a [`Stores`] could not be built.
369#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
370#[non_exhaustive]
371pub enum StoresBuilderError {
372    /// One of the seven roles was not supplied.
373    #[error("no implementation supplied for the {role} store")]
374    MissingStore {
375        /// The role, one of [`StoresBuilder::ROLES`].
376        role: &'static str,
377    },
378}
379
380/// Collects one implementation per role.
381///
382/// Every role is mandatory: a half-configured persistence layer is a bug that
383/// should surface at start-up, not at the first turn that needs the missing
384/// store.
385#[derive(Clone, Default)]
386pub struct StoresBuilder {
387    conversations: Option<Arc<dyn ConversationStore>>,
388    interactions: Option<Arc<dyn InteractionStore>>,
389    journal: Option<Arc<dyn CommandJournal>>,
390    events: Option<Arc<dyn EventJournal>>,
391    outbox: Option<Arc<dyn OutboxStore>>,
392    replay: Option<Arc<dyn ReplayStore>>,
393    commit: Option<Arc<dyn CommitStore>>,
394}
395
396impl StoresBuilder {
397    /// The seven role names, in the order [`StoresBuilder::build`] checks them.
398    pub const ROLES: [&'static str; 7] = [
399        "conversations",
400        "interactions",
401        "journal",
402        "events",
403        "outbox",
404        "replay",
405        "commit",
406    ];
407
408    /// An empty builder.
409    #[must_use]
410    pub fn new() -> Self {
411        Self::default()
412    }
413
414    /// Sets the conversation store.
415    #[must_use]
416    pub fn conversations(mut self, store: Arc<dyn ConversationStore>) -> Self {
417        self.conversations = Some(store);
418        self
419    }
420
421    /// Sets the interaction store.
422    #[must_use]
423    pub fn interactions(mut self, store: Arc<dyn InteractionStore>) -> Self {
424        self.interactions = Some(store);
425        self
426    }
427
428    /// Sets the command journal.
429    #[must_use]
430    pub fn journal(mut self, store: Arc<dyn CommandJournal>) -> Self {
431        self.journal = Some(store);
432        self
433    }
434
435    /// Sets the event journal.
436    #[must_use]
437    pub fn events(mut self, store: Arc<dyn EventJournal>) -> Self {
438        self.events = Some(store);
439        self
440    }
441
442    /// Sets the outbox.
443    #[must_use]
444    pub fn outbox(mut self, store: Arc<dyn OutboxStore>) -> Self {
445        self.outbox = Some(store);
446        self
447    }
448
449    /// Sets the replay store.
450    #[must_use]
451    pub fn replay(mut self, store: Arc<dyn ReplayStore>) -> Self {
452        self.replay = Some(store);
453        self
454    }
455
456    /// Sets the commit store.
457    #[must_use]
458    pub fn commit(mut self, store: Arc<dyn CommitStore>) -> Self {
459        self.commit = Some(store);
460        self
461    }
462
463    /// Builds the set.
464    ///
465    /// # Errors
466    /// * [`StoresBuilderError::MissingStore`] naming the first role that was
467    ///   never supplied.
468    pub fn build(self) -> Result<Stores, StoresBuilderError> {
469        fn required<T: ?Sized>(
470            store: Option<Arc<T>>,
471            role: &'static str,
472        ) -> Result<Arc<T>, StoresBuilderError> {
473            store.ok_or(StoresBuilderError::MissingStore { role })
474        }
475
476        Ok(Stores {
477            conversations: required(self.conversations, Self::ROLES[0])?,
478            interactions: required(self.interactions, Self::ROLES[1])?,
479            journal: required(self.journal, Self::ROLES[2])?,
480            events: required(self.events, Self::ROLES[3])?,
481            outbox: required(self.outbox, Self::ROLES[4])?,
482            replay: required(self.replay, Self::ROLES[5])?,
483            commit: required(self.commit, Self::ROLES[6])?,
484        })
485    }
486}
487
488impl fmt::Debug for StoresBuilder {
489    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
490        let supplied: Vec<&'static str> = Self::ROLES
491            .iter()
492            .zip([
493                self.conversations.is_some(),
494                self.interactions.is_some(),
495                self.journal.is_some(),
496                self.events.is_some(),
497                self.outbox.is_some(),
498                self.replay.is_some(),
499                self.commit.is_some(),
500            ])
501            .filter_map(|(role, present)| present.then_some(*role))
502            .collect();
503        f.debug_struct("StoresBuilder")
504            .field("supplied", &supplied)
505            .finish()
506    }
507}
508
509#[cfg(test)]
510mod tests {
511    use super::*;
512    use turnframe_core::ids::{AccountId, ConversationId};
513
514    #[test]
515    fn a_store_set_is_shareable_across_tasks() {
516        // The runtime holds one `Stores` and hands it to every turn, on any
517        // executor thread. If this stops compiling, a trait lost `Send + Sync`.
518        fn assert_send_sync<T: Send + Sync + 'static>() {}
519        assert_send_sync::<Stores>();
520        assert_send_sync::<StoresBuilder>();
521        assert_send_sync::<MemoryStores>();
522    }
523
524    #[test]
525    fn builder_names_the_first_missing_role() {
526        assert_eq!(
527            Stores::builder().build().unwrap_err(),
528            StoresBuilderError::MissingStore {
529                role: "conversations"
530            }
531        );
532        let backend = Arc::new(MemoryStores::new());
533        let err = Stores::builder()
534            .conversations(backend.clone())
535            .interactions(backend.clone())
536            .build()
537            .unwrap_err();
538        assert_eq!(err, StoresBuilderError::MissingStore { role: "journal" });
539        assert!(
540            err.to_string()
541                .contains("no implementation supplied for the journal store")
542        );
543    }
544
545    #[test]
546    fn builder_accepts_one_backend_for_every_role() {
547        let backend = Arc::new(MemoryStores::new());
548        let stores = Stores::builder()
549            .conversations(backend.clone())
550            .interactions(backend.clone())
551            .journal(backend.clone())
552            .events(backend.clone())
553            .outbox(backend.clone())
554            .replay(backend.clone())
555            .commit(backend)
556            .build()
557            .unwrap();
558        assert!(format!("{stores:?}").contains("conversations"));
559    }
560
561    #[tokio::test]
562    async fn the_read_only_view_answers_what_the_full_set_answers() {
563        let backend = Arc::new(MemoryStores::new());
564        let stores = Stores::from_memory(backend);
565        let account = AccountId::from("a");
566        let conversation = ConversationId::new();
567        stores
568            .conversations()
569            .create_conversation(crate::conversation::ConversationRecord::new(
570                conversation,
571                account.clone(),
572                chrono::Utc::now(),
573            ))
574            .await
575            .unwrap();
576
577        let reading = stores.read_only();
578        assert_eq!(
579            reading
580                .conversations()
581                .load_conversation(&account, &conversation)
582                .await
583                .unwrap()
584                .account_id,
585            account
586        );
587        // There is nothing to call here that writes: the view carries the six
588        // readers and no commit store at all.
589        assert!(format!("{reading:?}").contains("conversations"));
590        assert!(!format!("{reading:?}").contains("commit"));
591    }
592
593    #[test]
594    fn the_read_only_builder_names_the_first_missing_role() {
595        assert_eq!(
596            ReadOnlyStores::builder().build().unwrap_err(),
597            StoresBuilderError::MissingStore {
598                role: "conversations"
599            }
600        );
601        let backend = Arc::new(MemoryStores::new());
602        let reading = ReadOnlyStores::builder()
603            .conversations(backend.clone())
604            .interactions(backend.clone())
605            .journal(backend.clone())
606            .events(backend.clone())
607            .outbox(backend.clone())
608            .replay(backend)
609            .build()
610            .unwrap();
611        assert!(format!("{reading:?}").contains("replay"));
612    }
613
614    #[test]
615    fn a_read_only_view_is_shareable_across_tasks() {
616        fn assert_send_sync<T: Send + Sync + 'static>() {}
617        assert_send_sync::<ReadOnlyStores>();
618        assert_send_sync::<ReadOnlyStoresBuilder>();
619    }
620
621    #[tokio::test]
622    async fn in_memory_shares_one_state_across_roles() {
623        let backend = Arc::new(MemoryStores::new());
624        let stores = Stores::from_memory(backend.clone());
625        let account = AccountId::from("a");
626        let conversation = ConversationId::new();
627        stores
628            .conversations()
629            .create_conversation(crate::conversation::ConversationRecord::new(
630                conversation,
631                account.clone(),
632                backend.now(),
633            ))
634            .await
635            .unwrap();
636        // The same state answers through the aggregate and through the handle.
637        assert!(
638            crate::conversation::ConversationReader::load_conversation(
639                backend.as_ref(),
640                &account,
641                &conversation
642            )
643            .await
644            .is_ok()
645        );
646        let cloned = stores.clone();
647        assert!(
648            cloned
649                .conversations()
650                .load_conversation(&account, &conversation)
651                .await
652                .is_ok()
653        );
654    }
655}