reliar-store-postgres 0.7.0

PostgreSQL provider for the Reliar transactional outbox and inbox: migrations, migrate(), enqueue and the SKIP LOCKED claim.
Documentation
//! `FaultyStore<S>` — a delegating decorator over a real [`OutboxStore`], scripted per method
//! (ADR 0043 §4). Postgres cannot be told to fail its next `complete`, and coupling a dispatcher
//! **control-flow** trial to whichever real database error happens to classify transient would
//! test error classification instead of the control flow this crate's dispatcher trials actually
//! want to prove. `FaultyStore` closes that gap without becoming a second implementation:
//!
//! - **It never fabricates success.** [`StoreStep::Pass`] and [`StoreStep::DelayBefore`] delegate
//!   to the wrapped store; [`StoreStep::Transient`]/[`StoreStep::Permanent`] return
//!   [`FaultyStoreError::Injected`] **without calling through** — there is no path that returns
//!   `Ok` with invented rows, counts or ids.
//! - **It owns no domain state** — only a per-method script and call counter. Every assertion in
//!   a faulted trial still reads Postgres through the wrapped store or a `sqlx` query, never
//!   through this decorator (ADR 0043 §2's stimulus/oracle line).
//! - **`Classify` forwards**: [`FaultyStoreError::Inner`] delegates to the wrapped error's own
//!   `kind()`; [`FaultyStoreError::Injected`] reports the scripted kind directly.
//! - **It is scripted per method** (`on_acquire`, `on_complete`), so a trial names exactly which
//!   call it is perturbing.
//! - **It lives here, in the test harness, only** — never in a published crate, never behind a
//!   feature, so nothing outside these trials can mistake it for an implementation of the
//!   contract (ADR 0043 Amendment A.1).

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,
};

/// One scripted response for a single call to a `FaultyStore`-wrapped method.
#[derive(Clone, Copy, Debug)]
pub(crate) enum StoreStep {
    /// Delegate to the wrapped store, unmodified.
    Pass,

    /// Return an injected transient error; the wrapped store is **not** called.
    Transient,

    /// Return an injected permanent error; the wrapped store is **not** called.
    Permanent,

    /// Sleep for the given duration, then delegate.
    DelayBefore(Duration),
}

/// [`OutboxStore::Error`] for a [`FaultyStore`]: either the wrapped store's own error, or a
/// failure this decorator injected instead of ever calling through.
#[derive(Debug)]
pub(crate) enum FaultyStoreError<E> {
    /// A failure this decorator invented, never the wrapped store's.
    Injected(FailureKind),

    /// The wrapped store's own error, passed through unchanged.
    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(),
        }
    }
}

/// Every method [`FaultyStore`] can script, one queue each. An empty queue always answers
/// [`StoreStep::Pass`] — a trial only ever scripts the calls it cares about perturbing.
#[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)
    }
}

/// Per-method call counts — stimulus bookkeeping only (ADR 0043 §2), never read by an assertion
/// as a stand-in for store state.
#[derive(Default)]
struct Counters {
    acquire: AtomicU64,

    complete: AtomicU64,

    fail: AtomicU64,

    release: AtomicU64,

    extend_lease: AtomicU64,

    purge: AtomicU64,

    stats: AtomicU64,
}

/// A delegating decorator over a real [`OutboxStore`] `S`, scripted per method (module docs).
/// Cloning a `FaultyStore` shares its script and counters — the same scripted answers apply
/// through every clone, exactly as a real store's state is shared through its own clones.
pub(crate) struct FaultyStore<S> {
    inner: S,

    script: Arc<Mutex<Script>>,

    counters: Arc<Counters>,
}

impl<S> FaultyStore<S> {
    /// Wraps `inner` with an empty script — every call passes through until a trial scripts one.
    pub(crate) fn new(inner: S) -> Self {
        Self {
            inner,
            script: Arc::new(Mutex::new(Script::default())),
            counters: Arc::new(Counters::default()),
        }
    }

    /// Scripts `steps` for the next calls to [`OutboxStore::acquire`], in order; calls past the
    /// end of `steps` fall back to [`StoreStep::Pass`].
    pub(crate) fn on_acquire(&self, steps: impl IntoIterator<Item = StoreStep>) {
        self.script.lock().unwrap().acquire.extend(steps);
    }

    /// Scripts `steps` for [`OutboxStore::complete`], as [`Self::on_acquire`].
    pub(crate) fn on_complete(&self, steps: impl IntoIterator<Item = StoreStep>) {
        self.script.lock().unwrap().complete.extend(steps);
    }

    /// How many times [`OutboxStore::acquire`] has been called on this store (or any clone of
    /// it) so far — stimulus bookkeeping, never a stand-in for a Postgres-read assertion.
    pub(crate) fn acquire_calls(&self) -> u64 {
        self.counters.acquire.load(Ordering::Relaxed)
    }

    /// As [`Self::acquire_calls`], for [`OutboxStore::complete`].
    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)),
        }
    }
}