use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::Duration;
use reliar_core::{FailureKind, MessageType};
use crate::metrics::OutboxMetrics;
use crate::policy::RouteKind;
use crate::store::DeadReason;
#[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>,
purged: Option<(u64, u64)>,
publish_duration: Option<(Duration, MessageType)>,
routed: Vec<(RouteKind, MessageType)>,
}
impl RecordingMetrics {
fn lock(&self) -> MutexGuard<'_, Inner> {
self.inner.lock().unwrap_or_else(PoisonError::into_inner)
}
#[must_use]
pub fn claimed(&self) -> usize {
self.lock().claimed
}
#[must_use]
pub fn published(&self) -> Vec<MessageType> {
self.lock().published.clone()
}
#[must_use]
pub fn retried(&self) -> Vec<FailureKind> {
self.lock().retried.clone()
}
#[must_use]
pub fn dead(&self) -> Vec<DeadReason> {
self.lock().dead.clone()
}
#[must_use]
pub fn pending(&self) -> Option<u64> {
self.lock().pending
}
#[must_use]
pub fn expired_pending(&self) -> Option<u64> {
self.lock().expired_pending
}
#[must_use]
pub fn oldest_pending_age(&self) -> Option<Duration> {
self.lock().oldest_pending_age
}
#[must_use]
pub fn purged(&self) -> Option<(u64, u64)> {
self.lock().purged
}
#[must_use]
pub fn publish_duration(&self) -> Option<(Duration, MessageType)> {
self.lock().publish_duration.clone()
}
#[must_use]
pub fn routed(&self) -> Vec<(RouteKind, MessageType)> {
self.lock().routed.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);
}
fn purged(&self, published: u64, dead: u64) {
self.lock().purged = Some((published, dead));
}
fn routed(&self, route: RouteKind, message_type: &MessageType) {
self.lock().routed.push((route, message_type.clone()));
}
}