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, FailedRecord, OutboxStats, OutboxStore, PurgeReport,
PurgeRequest, RecordRef, WorkerId,
};
#[derive(Clone, Copy, Debug)]
pub(crate) enum StoreStep {
Pass,
Transient,
Permanent,
DelayBefore(Duration),
DetachOnSignal,
}
#[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)
}
fn take_detachable(queue: &mut VecDeque<StoreStep>) -> StoreStep {
if matches!(queue.front(), Some(StoreStep::DetachOnSignal)) {
StoreStep::DetachOnSignal
} else {
Self::take(queue)
}
}
fn disarm_detach(queue: &mut VecDeque<StoreStep>) {
if matches!(queue.front(), Some(StoreStep::DetachOnSignal)) {
queue.pop_front();
}
}
}
#[derive(Default)]
struct Counters {
acquire: AtomicU64,
complete: AtomicU64,
fail: AtomicU64,
release: AtomicU64,
extend_lease: AtomicU64,
purge: AtomicU64,
stats: AtomicU64,
}
struct DetachedCall<T> {
worker: WorkerId,
items: Vec<T>,
}
pub(crate) struct FaultyStore<S> {
inner: S,
script: Arc<Mutex<Script>>,
counters: Arc<Counters>,
detached_complete: Arc<Mutex<Option<DetachedCall<RecordRef>>>>,
detached_fail: Arc<Mutex<Option<DetachedCall<FailedRecord>>>>,
}
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()),
detached_complete: Arc::new(Mutex::new(None)),
detached_fail: Arc::new(Mutex::new(None)),
}
}
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 on_fail(&self, steps: impl IntoIterator<Item = StoreStep>) {
self.script.lock().unwrap().fail.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)
}
pub(crate) fn has_detached_complete(&self) -> bool {
self.detached_complete.lock().unwrap().is_some()
}
pub(crate) fn has_detached_fail(&self) -> bool {
self.detached_fail.lock().unwrap().is_some()
}
}
impl<S: OutboxStore> FaultyStore<S> {
pub(crate) async fn release_detached_complete(&self) -> Result<u64, S::Error> {
Script::disarm_detach(&mut self.script.lock().unwrap().complete);
let DetachedCall { worker, items } = self
.detached_complete
.lock()
.unwrap()
.take()
.expect("release_detached_complete called with no call captured");
self.inner.complete(&worker, &items).await
}
pub(crate) async fn release_detached_fail(&self) -> Result<u64, S::Error> {
Script::disarm_detach(&mut self.script.lock().unwrap().fail);
let DetachedCall { worker, items } = self
.detached_fail
.lock()
.unwrap()
.take()
.expect("release_detached_fail called with no call captured");
self.inner.fail(&worker, &items).await
}
}
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),
detached_complete: Arc::clone(&self.detached_complete),
detached_fail: Arc::clone(&self.detached_fail),
}
}
}
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)),
StoreStep::DetachOnSignal => {
unimplemented!("StoreStep::DetachOnSignal is only supported for complete/fail")
}
}
}
async fn complete(&self, worker: &WorkerId, items: &[RecordRef]) -> Result<u64, Self::Error> {
self.counters.complete.fetch_add(1, Ordering::Relaxed);
let step = Script::take_detachable(&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)),
StoreStep::DetachOnSignal => {
let mut slot = self.detached_complete.lock().unwrap();
if slot.is_none() {
*slot = Some(DetachedCall {
worker: worker.clone(),
items: items.to_vec(),
});
}
Err(FaultyStoreError::Injected(FailureKind::Transient))
}
}
}
async fn fail(&self, worker: &WorkerId, items: &[FailedRecord]) -> Result<u64, Self::Error> {
self.counters.fail.fetch_add(1, Ordering::Relaxed);
let step = Script::take_detachable(&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)),
StoreStep::DetachOnSignal => {
let mut slot = self.detached_fail.lock().unwrap();
if slot.is_none() {
*slot = Some(DetachedCall {
worker: worker.clone(),
items: items.to_vec(),
});
}
Err(FaultyStoreError::Injected(FailureKind::Transient))
}
}
}
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)),
StoreStep::DetachOnSignal => {
unimplemented!("StoreStep::DetachOnSignal is only supported for complete/fail")
}
}
}
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)),
StoreStep::DetachOnSignal => {
unimplemented!("StoreStep::DetachOnSignal is only supported for complete/fail")
}
}
}
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)),
StoreStep::DetachOnSignal => {
unimplemented!("StoreStep::DetachOnSignal is only supported for complete/fail")
}
}
}
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)),
StoreStep::DetachOnSignal => {
unimplemented!("StoreStep::DetachOnSignal is only supported for complete/fail")
}
}
}
}