mod clock;
mod fault;
mod impls;
mod state;
use std::sync::{Arc, Mutex, MutexGuard};
use chrono::{DateTime, Utc};
pub use self::clock::{Clock, ManualClock, SystemClock};
use self::fault::FaultQueue;
pub use self::fault::{FailurePoint, UnknownFailurePoint};
use self::state::Inner;
use crate::error::StoreError;
#[derive(Debug)]
pub struct MemoryStores {
inner: Mutex<Inner>,
faults: Mutex<FaultQueue>,
clock: Arc<dyn Clock>,
}
impl MemoryStores {
#[must_use]
pub fn new() -> Self {
Self::with_clock(Arc::new(SystemClock))
}
#[must_use]
pub fn with_clock(clock: Arc<dyn Clock>) -> Self {
Self {
inner: Mutex::new(Inner::default()),
faults: Mutex::new(FaultQueue::default()),
clock,
}
}
#[must_use]
pub fn now(&self) -> DateTime<Utc> {
self.clock.now()
}
pub fn fail_next(&self, point: FailurePoint, error: StoreError) -> Result<(), StoreError> {
self.faults()?.arm(point, error);
Ok(())
}
pub fn armed_failures(&self) -> Result<Vec<FailurePoint>, StoreError> {
Ok(self.faults()?.armed_points())
}
pub fn clear_failures(&self) -> Result<(), StoreError> {
self.faults()?.clear();
Ok(())
}
fn inner(&self) -> Result<MutexGuard<'_, Inner>, StoreError> {
self.inner.lock().map_err(|_| StoreError::Corrupt)
}
fn faults(&self) -> Result<MutexGuard<'_, FaultQueue>, StoreError> {
self.faults.lock().map_err(|_| StoreError::Corrupt)
}
fn fire(&self, point: FailurePoint) -> Result<(), StoreError> {
match self.faults()?.take(point) {
Some(error) => Err(error),
None => Ok(()),
}
}
}
impl Default for MemoryStores {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::conversation::{ConversationRecord, ConversationWriter};
use turnframe_core::ids::{AccountId, ConversationId};
#[test]
fn failures_are_armed_consumed_and_cleared() {
let stores = MemoryStores::new();
assert!(stores.armed_failures().unwrap().is_empty());
stores
.fail_next(FailurePoint::BeforeOutboxDispatch, StoreError::Unavailable)
.unwrap();
stores
.fail_next(FailurePoint::BeforeJournalInsert, StoreError::Timeout)
.unwrap();
assert_eq!(
stores.armed_failures().unwrap(),
vec![
FailurePoint::BeforeOutboxDispatch,
FailurePoint::BeforeJournalInsert
]
);
assert_eq!(
stores.fire(FailurePoint::BeforeJournalInsert),
Err(StoreError::Timeout)
);
assert_eq!(
stores.armed_failures().unwrap(),
vec![FailurePoint::BeforeOutboxDispatch]
);
stores.clear_failures().unwrap();
assert!(stores.armed_failures().unwrap().is_empty());
assert_eq!(stores.fire(FailurePoint::BeforeOutboxDispatch), Ok(()));
}
#[tokio::test]
async fn clock_drives_the_stamps_the_store_owns() {
let clock = Arc::new(ManualClock::epoch());
let stores = MemoryStores::with_clock(clock.clone());
let account = AccountId::from("a");
let conversation = ConversationId::new();
stores
.create_conversation(ConversationRecord::new(
conversation,
account.clone(),
clock.now(),
))
.await
.unwrap();
assert_eq!(stores.now(), DateTime::<Utc>::UNIX_EPOCH);
clock.advance(chrono::TimeDelta::seconds(90));
assert_eq!(
stores.now(),
DateTime::<Utc>::UNIX_EPOCH + chrono::TimeDelta::seconds(90)
);
}
}