1use 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#[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 #[must_use]
67 pub fn in_memory() -> Self {
68 Self::from_memory(Arc::new(MemoryStores::new()))
69 }
70
71 #[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 #[must_use]
90 pub fn builder() -> StoresBuilder {
91 StoresBuilder::default()
92 }
93
94 #[must_use]
96 pub fn conversations(&self) -> &Arc<dyn ConversationStore> {
97 &self.conversations
98 }
99
100 #[must_use]
102 pub fn interactions(&self) -> &Arc<dyn InteractionStore> {
103 &self.interactions
104 }
105
106 #[must_use]
108 pub fn journal(&self) -> &Arc<dyn CommandJournal> {
109 &self.journal
110 }
111
112 #[must_use]
114 pub fn events(&self) -> &Arc<dyn EventJournal> {
115 &self.events
116 }
117
118 #[must_use]
120 pub fn outbox(&self) -> &Arc<dyn OutboxStore> {
121 &self.outbox
122 }
123
124 #[must_use]
126 pub fn replay(&self) -> &Arc<dyn ReplayStore> {
127 &self.replay
128 }
129
130 #[must_use]
132 pub fn commit(&self) -> &Arc<dyn CommitStore> {
133 &self.commit
134 }
135
136 #[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#[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 #[must_use]
181 pub fn in_memory() -> Self {
182 Stores::in_memory().read_only()
183 }
184
185 #[must_use]
187 pub fn builder() -> ReadOnlyStoresBuilder {
188 ReadOnlyStoresBuilder::default()
189 }
190
191 #[must_use]
193 pub fn conversations(&self) -> &Arc<dyn ConversationReader> {
194 &self.conversations
195 }
196
197 #[must_use]
199 pub fn interactions(&self) -> &Arc<dyn InteractionReader> {
200 &self.interactions
201 }
202
203 #[must_use]
205 pub fn journal(&self) -> &Arc<dyn CommandJournalReader> {
206 &self.journal
207 }
208
209 #[must_use]
211 pub fn events(&self) -> &Arc<dyn EventJournalReader> {
212 &self.events
213 }
214
215 #[must_use]
217 pub fn outbox(&self) -> &Arc<dyn OutboxReader> {
218 &self.outbox
219 }
220
221 #[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 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#[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 pub const ROLES: [&'static str; 6] = [
258 "conversations",
259 "interactions",
260 "journal",
261 "events",
262 "outbox",
263 "replay",
264 ];
265
266 #[must_use]
268 pub fn new() -> Self {
269 Self::default()
270 }
271
272 #[must_use]
274 pub fn conversations(mut self, store: Arc<dyn ConversationReader>) -> Self {
275 self.conversations = Some(store);
276 self
277 }
278
279 #[must_use]
281 pub fn interactions(mut self, store: Arc<dyn InteractionReader>) -> Self {
282 self.interactions = Some(store);
283 self
284 }
285
286 #[must_use]
288 pub fn journal(mut self, store: Arc<dyn CommandJournalReader>) -> Self {
289 self.journal = Some(store);
290 self
291 }
292
293 #[must_use]
295 pub fn events(mut self, store: Arc<dyn EventJournalReader>) -> Self {
296 self.events = Some(store);
297 self
298 }
299
300 #[must_use]
302 pub fn outbox(mut self, store: Arc<dyn OutboxReader>) -> Self {
303 self.outbox = Some(store);
304 self
305 }
306
307 #[must_use]
309 pub fn replay(mut self, store: Arc<dyn ReplayReader>) -> Self {
310 self.replay = Some(store);
311 self
312 }
313
314 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 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#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
370#[non_exhaustive]
371pub enum StoresBuilderError {
372 #[error("no implementation supplied for the {role} store")]
374 MissingStore {
375 role: &'static str,
377 },
378}
379
380#[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 pub const ROLES: [&'static str; 7] = [
399 "conversations",
400 "interactions",
401 "journal",
402 "events",
403 "outbox",
404 "replay",
405 "commit",
406 ];
407
408 #[must_use]
410 pub fn new() -> Self {
411 Self::default()
412 }
413
414 #[must_use]
416 pub fn conversations(mut self, store: Arc<dyn ConversationStore>) -> Self {
417 self.conversations = Some(store);
418 self
419 }
420
421 #[must_use]
423 pub fn interactions(mut self, store: Arc<dyn InteractionStore>) -> Self {
424 self.interactions = Some(store);
425 self
426 }
427
428 #[must_use]
430 pub fn journal(mut self, store: Arc<dyn CommandJournal>) -> Self {
431 self.journal = Some(store);
432 self
433 }
434
435 #[must_use]
437 pub fn events(mut self, store: Arc<dyn EventJournal>) -> Self {
438 self.events = Some(store);
439 self
440 }
441
442 #[must_use]
444 pub fn outbox(mut self, store: Arc<dyn OutboxStore>) -> Self {
445 self.outbox = Some(store);
446 self
447 }
448
449 #[must_use]
451 pub fn replay(mut self, store: Arc<dyn ReplayStore>) -> Self {
452 self.replay = Some(store);
453 self
454 }
455
456 #[must_use]
458 pub fn commit(mut self, store: Arc<dyn CommitStore>) -> Self {
459 self.commit = Some(store);
460 self
461 }
462
463 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 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 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 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}