use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use reliar_core::{Classify, FailureKind};
use reliar_outbox::{
AcquireRequest, AcquiredBatch, CompletedRecord, FailedRecord, OutboxStats, OutboxStore,
PurgeReport, PurgeRequest, RecordRef, WorkerId,
};
#[derive(Clone, Copy, Debug)]
pub(crate) enum StoreStep {
Pass,
Transient,
Permanent,
DelayBefore(Duration),
}
#[derive(Debug)]
pub(crate) enum FaultyStoreError<E> {
Injected(FailureKind),
Inner(E),
}
impl<E: std::fmt::Display> std::fmt::Display for FaultyStoreError<E> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Injected(kind) => write!(f, "faulty store: injected {kind:?} failure"),
Self::Inner(inner) => write!(f, "faulty store: {inner}"),
}
}
}
impl<E: std::error::Error + 'static> std::error::Error for FaultyStoreError<E> {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Self::Injected(_) => None,
Self::Inner(inner) => Some(inner),
}
}
}
impl<E: Classify> Classify for FaultyStoreError<E> {
fn kind(&self) -> FailureKind {
match self {
Self::Injected(kind) => *kind,
Self::Inner(inner) => inner.kind(),
}
}
}
#[derive(Default)]
struct Script {
acquire: VecDeque<StoreStep>,
complete: VecDeque<StoreStep>,
fail: VecDeque<StoreStep>,
release: VecDeque<StoreStep>,
extend_lease: VecDeque<StoreStep>,
purge: VecDeque<StoreStep>,
stats: VecDeque<StoreStep>,
}
impl Script {
fn take(queue: &mut VecDeque<StoreStep>) -> StoreStep {
queue.pop_front().unwrap_or(StoreStep::Pass)
}
}
#[derive(Default)]
struct Counters {
acquire: AtomicU64,
complete: AtomicU64,
fail: AtomicU64,
release: AtomicU64,
extend_lease: AtomicU64,
purge: AtomicU64,
stats: AtomicU64,
}
pub(crate) struct FaultyStore<S> {
inner: S,
script: Arc<Mutex<Script>>,
counters: Arc<Counters>,
}
impl<S> FaultyStore<S> {
pub(crate) fn new(inner: S) -> Self {
Self {
inner,
script: Arc::new(Mutex::new(Script::default())),
counters: Arc::new(Counters::default()),
}
}
pub(crate) fn on_acquire(&self, steps: impl IntoIterator<Item = StoreStep>) {
self.script.lock().unwrap().acquire.extend(steps);
}
pub(crate) fn on_complete(&self, steps: impl IntoIterator<Item = StoreStep>) {
self.script.lock().unwrap().complete.extend(steps);
}
pub(crate) fn acquire_calls(&self) -> u64 {
self.counters.acquire.load(Ordering::Relaxed)
}
pub(crate) fn complete_calls(&self) -> u64 {
self.counters.complete.load(Ordering::Relaxed)
}
}
impl<S: Clone> Clone for FaultyStore<S> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
script: Arc::clone(&self.script),
counters: Arc::clone(&self.counters),
}
}
}
impl<S: OutboxStore> OutboxStore for FaultyStore<S> {
type Error = FaultyStoreError<S::Error>;
async fn acquire(&self, request: AcquireRequest) -> Result<AcquiredBatch, Self::Error> {
self.counters.acquire.fetch_add(1, Ordering::Relaxed);
let step = Script::take(&mut self.script.lock().unwrap().acquire);
match step {
StoreStep::Pass => self
.inner
.acquire(request)
.await
.map_err(FaultyStoreError::Inner),
StoreStep::DelayBefore(delay) => {
tokio::time::sleep(delay).await;
self.inner
.acquire(request)
.await
.map_err(FaultyStoreError::Inner)
}
StoreStep::Transient => Err(FaultyStoreError::Injected(FailureKind::Transient)),
StoreStep::Permanent => Err(FaultyStoreError::Injected(FailureKind::Permanent)),
}
}
async fn complete(
&self,
worker: &WorkerId,
items: &[CompletedRecord],
) -> Result<u64, Self::Error> {
self.counters.complete.fetch_add(1, Ordering::Relaxed);
let step = Script::take(&mut self.script.lock().unwrap().complete);
match step {
StoreStep::Pass => self
.inner
.complete(worker, items)
.await
.map_err(FaultyStoreError::Inner),
StoreStep::DelayBefore(delay) => {
tokio::time::sleep(delay).await;
self.inner
.complete(worker, items)
.await
.map_err(FaultyStoreError::Inner)
}
StoreStep::Transient => Err(FaultyStoreError::Injected(FailureKind::Transient)),
StoreStep::Permanent => Err(FaultyStoreError::Injected(FailureKind::Permanent)),
}
}
async fn fail(&self, worker: &WorkerId, items: &[FailedRecord]) -> Result<u64, Self::Error> {
self.counters.fail.fetch_add(1, Ordering::Relaxed);
let step = Script::take(&mut self.script.lock().unwrap().fail);
match step {
StoreStep::Pass => self
.inner
.fail(worker, items)
.await
.map_err(FaultyStoreError::Inner),
StoreStep::DelayBefore(delay) => {
tokio::time::sleep(delay).await;
self.inner
.fail(worker, items)
.await
.map_err(FaultyStoreError::Inner)
}
StoreStep::Transient => Err(FaultyStoreError::Injected(FailureKind::Transient)),
StoreStep::Permanent => Err(FaultyStoreError::Injected(FailureKind::Permanent)),
}
}
async fn release(&self, worker: &WorkerId, items: &[RecordRef]) -> Result<u64, Self::Error> {
self.counters.release.fetch_add(1, Ordering::Relaxed);
let step = Script::take(&mut self.script.lock().unwrap().release);
match step {
StoreStep::Pass => self
.inner
.release(worker, items)
.await
.map_err(FaultyStoreError::Inner),
StoreStep::DelayBefore(delay) => {
tokio::time::sleep(delay).await;
self.inner
.release(worker, items)
.await
.map_err(FaultyStoreError::Inner)
}
StoreStep::Transient => Err(FaultyStoreError::Injected(FailureKind::Transient)),
StoreStep::Permanent => Err(FaultyStoreError::Injected(FailureKind::Permanent)),
}
}
async fn extend_lease(
&self,
worker: &WorkerId,
items: &[RecordRef],
lease: Duration,
) -> Result<u64, Self::Error> {
self.counters.extend_lease.fetch_add(1, Ordering::Relaxed);
let step = Script::take(&mut self.script.lock().unwrap().extend_lease);
match step {
StoreStep::Pass => self
.inner
.extend_lease(worker, items, lease)
.await
.map_err(FaultyStoreError::Inner),
StoreStep::DelayBefore(delay) => {
tokio::time::sleep(delay).await;
self.inner
.extend_lease(worker, items, lease)
.await
.map_err(FaultyStoreError::Inner)
}
StoreStep::Transient => Err(FaultyStoreError::Injected(FailureKind::Transient)),
StoreStep::Permanent => Err(FaultyStoreError::Injected(FailureKind::Permanent)),
}
}
async fn purge(&self, request: PurgeRequest) -> Result<PurgeReport, Self::Error> {
self.counters.purge.fetch_add(1, Ordering::Relaxed);
let step = Script::take(&mut self.script.lock().unwrap().purge);
match step {
StoreStep::Pass => self
.inner
.purge(request)
.await
.map_err(FaultyStoreError::Inner),
StoreStep::DelayBefore(delay) => {
tokio::time::sleep(delay).await;
self.inner
.purge(request)
.await
.map_err(FaultyStoreError::Inner)
}
StoreStep::Transient => Err(FaultyStoreError::Injected(FailureKind::Transient)),
StoreStep::Permanent => Err(FaultyStoreError::Injected(FailureKind::Permanent)),
}
}
async fn stats(&self) -> Result<OutboxStats, Self::Error> {
self.counters.stats.fetch_add(1, Ordering::Relaxed);
let step = Script::take(&mut self.script.lock().unwrap().stats);
match step {
StoreStep::Pass => self.inner.stats().await.map_err(FaultyStoreError::Inner),
StoreStep::DelayBefore(delay) => {
tokio::time::sleep(delay).await;
self.inner.stats().await.map_err(FaultyStoreError::Inner)
}
StoreStep::Transient => Err(FaultyStoreError::Injected(FailureKind::Transient)),
StoreStep::Permanent => Err(FaultyStoreError::Injected(FailureKind::Permanent)),
}
}
}