reliar-outbox 0.6.0

Storage-agnostic transactional outbox: OutboxStore/Publisher contracts, retry policy, settings and dispatcher (no storage or transport dependency).
Documentation
//! [`RecordingMetrics`]: an [`OutboxMetrics`] fake that remembers every call.

use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::Duration;

use reliar_core::{FailureKind, MessageType};

use crate::metrics::OutboxMetrics;
use crate::store::DeadReason;

/// Records every [`OutboxMetrics`] call for assertion. Each getter below has the same name as
/// the hook it observes but a different arity — an inherent method always shadows a trait method
/// of the same name for a concrete `RecordingMetrics` value, so `metrics.claimed()` (the getter)
/// and `OutboxMetrics::claimed(&metrics, n)` (the hook, reached through the trait bound a
/// generic dispatcher holds it behind) never collide in practice.
///
/// ```
/// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
///
/// let metrics = RecordingMetrics::default();
/// OutboxMetrics::claimed(&metrics, 2);
/// OutboxMetrics::claimed(&metrics, 1);
/// assert_eq!(metrics.claimed(), 3);
/// ```
#[derive(Clone, Debug, Default)]
pub struct RecordingMetrics {
    inner: Arc<Mutex<Inner>>,
}

#[derive(Debug, Default)]
struct Inner {
    claimed: usize,

    published: Vec<MessageType>,

    retried: Vec<FailureKind>,

    dead: Vec<DeadReason>,

    pending: Option<u64>,

    expired_pending: Option<u64>,

    oldest_pending_age: Option<Duration>,

    publish_duration: Option<(Duration, MessageType)>,
}

impl RecordingMetrics {
    fn lock(&self) -> MutexGuard<'_, Inner> {
        self.inner.lock().unwrap_or_else(PoisonError::into_inner)
    }

    /// The total claimed across every [`OutboxMetrics::claimed`] call.
    ///
    /// ```
    /// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
    /// let metrics = RecordingMetrics::default();
    /// OutboxMetrics::claimed(&metrics, 5);
    /// assert_eq!(metrics.claimed(), 5);
    /// ```
    #[must_use]
    pub fn claimed(&self) -> usize {
        self.lock().claimed
    }

    /// One entry per published message, in call order.
    ///
    /// ```
    /// use reliar_core::MessageType;
    /// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
    /// let metrics = RecordingMetrics::default();
    /// let message_type = MessageType::new("orders.created", 1);
    /// OutboxMetrics::published(&metrics, 1, &message_type);
    /// assert_eq!(metrics.published(), vec![message_type]);
    /// ```
    #[must_use]
    pub fn published(&self) -> Vec<MessageType> {
        self.lock().published.clone()
    }

    /// One entry per retried message, in call order.
    ///
    /// ```
    /// use reliar_core::FailureKind;
    /// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
    /// let metrics = RecordingMetrics::default();
    /// OutboxMetrics::retried(&metrics, 1, FailureKind::Transient);
    /// assert_eq!(metrics.retried(), vec![FailureKind::Transient]);
    /// ```
    #[must_use]
    pub fn retried(&self) -> Vec<FailureKind> {
        self.lock().retried.clone()
    }

    /// One entry per message moved to dead, in call order.
    ///
    /// ```
    /// use reliar_outbox::{DeadReason, OutboxMetrics, RecordingMetrics};
    /// let metrics = RecordingMetrics::default();
    /// OutboxMetrics::dead(&metrics, 1, DeadReason::AttemptsExhausted);
    /// assert_eq!(metrics.dead(), vec![DeadReason::AttemptsExhausted]);
    /// ```
    #[must_use]
    pub fn dead(&self) -> Vec<DeadReason> {
        self.lock().dead.clone()
    }

    /// The last observed pending count, if [`OutboxMetrics::pending`] was ever called.
    ///
    /// ```
    /// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
    /// let metrics = RecordingMetrics::default();
    /// assert!(metrics.pending().is_none());
    /// OutboxMetrics::pending(&metrics, 7);
    /// assert_eq!(metrics.pending(), Some(7));
    /// ```
    #[must_use]
    pub fn pending(&self) -> Option<u64> {
        self.lock().pending
    }

    /// The last observed expired-pending count, if [`OutboxMetrics::expired_pending`] was ever
    /// called.
    ///
    /// ```
    /// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
    /// let metrics = RecordingMetrics::default();
    /// OutboxMetrics::expired_pending(&metrics, 2);
    /// assert_eq!(metrics.expired_pending(), Some(2));
    /// ```
    #[must_use]
    pub fn expired_pending(&self) -> Option<u64> {
        self.lock().expired_pending
    }

    /// The last observed outbox lag, if [`OutboxMetrics::oldest_pending_age`] was ever called.
    /// `None` also when the dispatcher correctly skipped the call on an empty backlog.
    ///
    /// ```
    /// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
    /// use std::time::Duration;
    /// let metrics = RecordingMetrics::default();
    /// assert!(metrics.oldest_pending_age().is_none());
    /// OutboxMetrics::oldest_pending_age(&metrics, Duration::from_secs(3));
    /// assert_eq!(metrics.oldest_pending_age(), Some(Duration::from_secs(3)));
    /// ```
    #[must_use]
    pub fn oldest_pending_age(&self) -> Option<Duration> {
        self.lock().oldest_pending_age
    }

    /// The last `(duration, message_type)` pair observed via
    /// [`OutboxMetrics::publish_duration`], if it was ever called.
    ///
    /// ```
    /// use reliar_core::MessageType;
    /// use reliar_outbox::{OutboxMetrics, RecordingMetrics};
    /// use std::time::Duration;
    /// let metrics = RecordingMetrics::default();
    /// let message_type = MessageType::new("orders.created", 1);
    /// OutboxMetrics::publish_duration(&metrics, Duration::from_millis(5), &message_type);
    /// assert_eq!(
    ///     metrics.publish_duration(),
    ///     Some((Duration::from_millis(5), message_type))
    /// );
    /// ```
    #[must_use]
    pub fn publish_duration(&self) -> Option<(Duration, MessageType)> {
        self.lock().publish_duration.clone()
    }
}

impl OutboxMetrics for RecordingMetrics {
    fn claimed(&self, n: usize) {
        self.lock().claimed += n;
    }

    fn published(&self, n: usize, message_type: &MessageType) {
        let mut guard = self.lock();

        guard
            .published
            .extend(std::iter::repeat_n(message_type.clone(), n));
    }

    fn retried(&self, n: usize, kind: FailureKind) {
        let mut guard = self.lock();

        guard.retried.extend(std::iter::repeat_n(kind, n));
    }

    fn dead(&self, n: usize, reason: DeadReason) {
        let mut guard = self.lock();

        guard.dead.extend(std::iter::repeat_n(reason, n));
    }

    fn publish_duration(&self, d: Duration, message_type: &MessageType) {
        self.lock().publish_duration = Some((d, message_type.clone()));
    }

    fn pending(&self, n: u64) {
        self.lock().pending = Some(n);
    }

    fn expired_pending(&self, n: u64) {
        self.lock().expired_pending = Some(n);
    }

    fn oldest_pending_age(&self, age: Duration) {
        self.lock().oldest_pending_age = Some(age);
    }
}