use std::fmt;
use std::sync::Arc;
use crate::commit::CommitStore;
use crate::conversation::{ConversationReader, ConversationStore};
use crate::events::{EventJournal, EventJournalReader};
use crate::interaction::{InteractionReader, InteractionStore};
use crate::journal::{CommandJournal, CommandJournalReader};
use crate::memory::MemoryStores;
use crate::outbox::{OutboxReader, OutboxStore};
use crate::replay::{ReplayReader, ReplayStore};
#[derive(Clone)]
pub struct Stores {
conversations: Arc<dyn ConversationStore>,
interactions: Arc<dyn InteractionStore>,
journal: Arc<dyn CommandJournal>,
events: Arc<dyn EventJournal>,
outbox: Arc<dyn OutboxStore>,
replay: Arc<dyn ReplayStore>,
commit: Arc<dyn CommitStore>,
}
impl Stores {
#[must_use]
pub fn in_memory() -> Self {
Self::from_memory(Arc::new(MemoryStores::new()))
}
#[must_use]
pub fn from_memory(backend: Arc<MemoryStores>) -> Self {
Self {
conversations: backend.clone(),
interactions: backend.clone(),
journal: backend.clone(),
events: backend.clone(),
outbox: backend.clone(),
replay: backend.clone(),
commit: backend,
}
}
#[must_use]
pub fn builder() -> StoresBuilder {
StoresBuilder::default()
}
#[must_use]
pub fn conversations(&self) -> &Arc<dyn ConversationStore> {
&self.conversations
}
#[must_use]
pub fn interactions(&self) -> &Arc<dyn InteractionStore> {
&self.interactions
}
#[must_use]
pub fn journal(&self) -> &Arc<dyn CommandJournal> {
&self.journal
}
#[must_use]
pub fn events(&self) -> &Arc<dyn EventJournal> {
&self.events
}
#[must_use]
pub fn outbox(&self) -> &Arc<dyn OutboxStore> {
&self.outbox
}
#[must_use]
pub fn replay(&self) -> &Arc<dyn ReplayStore> {
&self.replay
}
#[must_use]
pub fn commit(&self) -> &Arc<dyn CommitStore> {
&self.commit
}
#[must_use]
pub fn read_only(&self) -> ReadOnlyStores {
ReadOnlyStores {
conversations: self.conversations.clone(),
interactions: self.interactions.clone(),
journal: self.journal.clone(),
events: self.events.clone(),
outbox: self.outbox.clone(),
replay: self.replay.clone(),
}
}
}
#[derive(Clone)]
pub struct ReadOnlyStores {
conversations: Arc<dyn ConversationReader>,
interactions: Arc<dyn InteractionReader>,
journal: Arc<dyn CommandJournalReader>,
events: Arc<dyn EventJournalReader>,
outbox: Arc<dyn OutboxReader>,
replay: Arc<dyn ReplayReader>,
}
impl ReadOnlyStores {
#[must_use]
pub fn in_memory() -> Self {
Stores::in_memory().read_only()
}
#[must_use]
pub fn builder() -> ReadOnlyStoresBuilder {
ReadOnlyStoresBuilder::default()
}
#[must_use]
pub fn conversations(&self) -> &Arc<dyn ConversationReader> {
&self.conversations
}
#[must_use]
pub fn interactions(&self) -> &Arc<dyn InteractionReader> {
&self.interactions
}
#[must_use]
pub fn journal(&self) -> &Arc<dyn CommandJournalReader> {
&self.journal
}
#[must_use]
pub fn events(&self) -> &Arc<dyn EventJournalReader> {
&self.events
}
#[must_use]
pub fn outbox(&self) -> &Arc<dyn OutboxReader> {
&self.outbox
}
#[must_use]
pub fn replay(&self) -> &Arc<dyn ReplayReader> {
&self.replay
}
}
impl From<&Stores> for ReadOnlyStores {
fn from(stores: &Stores) -> Self {
stores.read_only()
}
}
impl fmt::Debug for ReadOnlyStores {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ReadOnlyStores")
.field("roles", &ReadOnlyStoresBuilder::ROLES)
.finish()
}
}
#[derive(Clone, Default)]
pub struct ReadOnlyStoresBuilder {
conversations: Option<Arc<dyn ConversationReader>>,
interactions: Option<Arc<dyn InteractionReader>>,
journal: Option<Arc<dyn CommandJournalReader>>,
events: Option<Arc<dyn EventJournalReader>>,
outbox: Option<Arc<dyn OutboxReader>>,
replay: Option<Arc<dyn ReplayReader>>,
}
impl ReadOnlyStoresBuilder {
pub const ROLES: [&'static str; 6] = [
"conversations",
"interactions",
"journal",
"events",
"outbox",
"replay",
];
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn conversations(mut self, store: Arc<dyn ConversationReader>) -> Self {
self.conversations = Some(store);
self
}
#[must_use]
pub fn interactions(mut self, store: Arc<dyn InteractionReader>) -> Self {
self.interactions = Some(store);
self
}
#[must_use]
pub fn journal(mut self, store: Arc<dyn CommandJournalReader>) -> Self {
self.journal = Some(store);
self
}
#[must_use]
pub fn events(mut self, store: Arc<dyn EventJournalReader>) -> Self {
self.events = Some(store);
self
}
#[must_use]
pub fn outbox(mut self, store: Arc<dyn OutboxReader>) -> Self {
self.outbox = Some(store);
self
}
#[must_use]
pub fn replay(mut self, store: Arc<dyn ReplayReader>) -> Self {
self.replay = Some(store);
self
}
pub fn build(self) -> Result<ReadOnlyStores, StoresBuilderError> {
fn required<T: ?Sized>(
store: Option<Arc<T>>,
role: &'static str,
) -> Result<Arc<T>, StoresBuilderError> {
store.ok_or(StoresBuilderError::MissingStore { role })
}
Ok(ReadOnlyStores {
conversations: required(self.conversations, Self::ROLES[0])?,
interactions: required(self.interactions, Self::ROLES[1])?,
journal: required(self.journal, Self::ROLES[2])?,
events: required(self.events, Self::ROLES[3])?,
outbox: required(self.outbox, Self::ROLES[4])?,
replay: required(self.replay, Self::ROLES[5])?,
})
}
}
impl fmt::Debug for ReadOnlyStoresBuilder {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let supplied: Vec<&'static str> = Self::ROLES
.iter()
.zip([
self.conversations.is_some(),
self.interactions.is_some(),
self.journal.is_some(),
self.events.is_some(),
self.outbox.is_some(),
self.replay.is_some(),
])
.filter_map(|(role, present)| present.then_some(*role))
.collect();
f.debug_struct("ReadOnlyStoresBuilder")
.field("supplied", &supplied)
.finish()
}
}
impl fmt::Debug for Stores {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Stores")
.field("roles", &StoresBuilder::ROLES)
.finish()
}
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum StoresBuilderError {
#[error("no implementation supplied for the {role} store")]
MissingStore {
role: &'static str,
},
}
#[derive(Clone, Default)]
pub struct StoresBuilder {
conversations: Option<Arc<dyn ConversationStore>>,
interactions: Option<Arc<dyn InteractionStore>>,
journal: Option<Arc<dyn CommandJournal>>,
events: Option<Arc<dyn EventJournal>>,
outbox: Option<Arc<dyn OutboxStore>>,
replay: Option<Arc<dyn ReplayStore>>,
commit: Option<Arc<dyn CommitStore>>,
}
impl StoresBuilder {
pub const ROLES: [&'static str; 7] = [
"conversations",
"interactions",
"journal",
"events",
"outbox",
"replay",
"commit",
];
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn conversations(mut self, store: Arc<dyn ConversationStore>) -> Self {
self.conversations = Some(store);
self
}
#[must_use]
pub fn interactions(mut self, store: Arc<dyn InteractionStore>) -> Self {
self.interactions = Some(store);
self
}
#[must_use]
pub fn journal(mut self, store: Arc<dyn CommandJournal>) -> Self {
self.journal = Some(store);
self
}
#[must_use]
pub fn events(mut self, store: Arc<dyn EventJournal>) -> Self {
self.events = Some(store);
self
}
#[must_use]
pub fn outbox(mut self, store: Arc<dyn OutboxStore>) -> Self {
self.outbox = Some(store);
self
}
#[must_use]
pub fn replay(mut self, store: Arc<dyn ReplayStore>) -> Self {
self.replay = Some(store);
self
}
#[must_use]
pub fn commit(mut self, store: Arc<dyn CommitStore>) -> Self {
self.commit = Some(store);
self
}
pub fn build(self) -> Result<Stores, StoresBuilderError> {
fn required<T: ?Sized>(
store: Option<Arc<T>>,
role: &'static str,
) -> Result<Arc<T>, StoresBuilderError> {
store.ok_or(StoresBuilderError::MissingStore { role })
}
Ok(Stores {
conversations: required(self.conversations, Self::ROLES[0])?,
interactions: required(self.interactions, Self::ROLES[1])?,
journal: required(self.journal, Self::ROLES[2])?,
events: required(self.events, Self::ROLES[3])?,
outbox: required(self.outbox, Self::ROLES[4])?,
replay: required(self.replay, Self::ROLES[5])?,
commit: required(self.commit, Self::ROLES[6])?,
})
}
}
impl fmt::Debug for StoresBuilder {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let supplied: Vec<&'static str> = Self::ROLES
.iter()
.zip([
self.conversations.is_some(),
self.interactions.is_some(),
self.journal.is_some(),
self.events.is_some(),
self.outbox.is_some(),
self.replay.is_some(),
self.commit.is_some(),
])
.filter_map(|(role, present)| present.then_some(*role))
.collect();
f.debug_struct("StoresBuilder")
.field("supplied", &supplied)
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
use turnframe_core::ids::{AccountId, ConversationId};
#[test]
fn a_store_set_is_shareable_across_tasks() {
fn assert_send_sync<T: Send + Sync + 'static>() {}
assert_send_sync::<Stores>();
assert_send_sync::<StoresBuilder>();
assert_send_sync::<MemoryStores>();
}
#[test]
fn builder_names_the_first_missing_role() {
assert_eq!(
Stores::builder().build().unwrap_err(),
StoresBuilderError::MissingStore {
role: "conversations"
}
);
let backend = Arc::new(MemoryStores::new());
let err = Stores::builder()
.conversations(backend.clone())
.interactions(backend.clone())
.build()
.unwrap_err();
assert_eq!(err, StoresBuilderError::MissingStore { role: "journal" });
assert!(
err.to_string()
.contains("no implementation supplied for the journal store")
);
}
#[test]
fn builder_accepts_one_backend_for_every_role() {
let backend = Arc::new(MemoryStores::new());
let stores = Stores::builder()
.conversations(backend.clone())
.interactions(backend.clone())
.journal(backend.clone())
.events(backend.clone())
.outbox(backend.clone())
.replay(backend.clone())
.commit(backend)
.build()
.unwrap();
assert!(format!("{stores:?}").contains("conversations"));
}
#[tokio::test]
async fn the_read_only_view_answers_what_the_full_set_answers() {
let backend = Arc::new(MemoryStores::new());
let stores = Stores::from_memory(backend);
let account = AccountId::from("a");
let conversation = ConversationId::new();
stores
.conversations()
.create_conversation(crate::conversation::ConversationRecord::new(
conversation,
account.clone(),
chrono::Utc::now(),
))
.await
.unwrap();
let reading = stores.read_only();
assert_eq!(
reading
.conversations()
.load_conversation(&account, &conversation)
.await
.unwrap()
.account_id,
account
);
assert!(format!("{reading:?}").contains("conversations"));
assert!(!format!("{reading:?}").contains("commit"));
}
#[test]
fn the_read_only_builder_names_the_first_missing_role() {
assert_eq!(
ReadOnlyStores::builder().build().unwrap_err(),
StoresBuilderError::MissingStore {
role: "conversations"
}
);
let backend = Arc::new(MemoryStores::new());
let reading = ReadOnlyStores::builder()
.conversations(backend.clone())
.interactions(backend.clone())
.journal(backend.clone())
.events(backend.clone())
.outbox(backend.clone())
.replay(backend)
.build()
.unwrap();
assert!(format!("{reading:?}").contains("replay"));
}
#[test]
fn a_read_only_view_is_shareable_across_tasks() {
fn assert_send_sync<T: Send + Sync + 'static>() {}
assert_send_sync::<ReadOnlyStores>();
assert_send_sync::<ReadOnlyStoresBuilder>();
}
#[tokio::test]
async fn in_memory_shares_one_state_across_roles() {
let backend = Arc::new(MemoryStores::new());
let stores = Stores::from_memory(backend.clone());
let account = AccountId::from("a");
let conversation = ConversationId::new();
stores
.conversations()
.create_conversation(crate::conversation::ConversationRecord::new(
conversation,
account.clone(),
backend.now(),
))
.await
.unwrap();
assert!(
crate::conversation::ConversationReader::load_conversation(
backend.as_ref(),
&account,
&conversation
)
.await
.is_ok()
);
let cloned = stores.clone();
assert!(
cloned
.conversations()
.load_conversation(&account, &conversation)
.await
.is_ok()
);
}
}