use std::{
sync::{
Arc,
atomic::{AtomicU64, Ordering},
},
time::Instant,
};
#[derive(Clone, Debug)]
pub struct QueueMetrics {
claimed_runs: Arc<AtomicU64>,
claim_errors: Arc<AtomicU64>,
executions: Arc<AtomicU64>,
completed_executions: Arc<AtomicU64>,
failed_executions: Arc<AtomicU64>,
lease_lost_executions: Arc<AtomicU64>,
cancelled_executions: Arc<AtomicU64>,
suspended_executions: Arc<AtomicU64>,
unhandled_executions: Arc<AtomicU64>,
execution_duration_nanoseconds: Arc<AtomicU64>,
}
impl QueueMetrics {
pub(crate) fn new() -> Self {
Self {
claimed_runs: Arc::new(AtomicU64::new(0)),
claim_errors: Arc::new(AtomicU64::new(0)),
executions: Arc::new(AtomicU64::new(0)),
completed_executions: Arc::new(AtomicU64::new(0)),
failed_executions: Arc::new(AtomicU64::new(0)),
lease_lost_executions: Arc::new(AtomicU64::new(0)),
cancelled_executions: Arc::new(AtomicU64::new(0)),
suspended_executions: Arc::new(AtomicU64::new(0)),
unhandled_executions: Arc::new(AtomicU64::new(0)),
execution_duration_nanoseconds: Arc::new(AtomicU64::new(0)),
}
}
pub fn claimed_runs(&self) -> u64 {
self.claimed_runs.load(Ordering::Relaxed)
}
pub fn claim_errors(&self) -> u64 {
self.claim_errors.load(Ordering::Relaxed)
}
pub fn executions(&self) -> u64 {
self.executions.load(Ordering::Relaxed)
}
pub fn completed_executions(&self) -> u64 {
self.completed_executions.load(Ordering::Relaxed)
}
pub fn failed_executions(&self) -> u64 {
self.failed_executions.load(Ordering::Relaxed)
}
pub fn lease_lost_executions(&self) -> u64 {
self.lease_lost_executions.load(Ordering::Relaxed)
}
pub fn cancelled_executions(&self) -> u64 {
self.cancelled_executions.load(Ordering::Relaxed)
}
pub fn suspended_executions(&self) -> u64 {
self.suspended_executions.load(Ordering::Relaxed)
}
pub fn unhandled_executions(&self) -> u64 {
self.unhandled_executions.load(Ordering::Relaxed)
}
pub fn execution_duration_nanoseconds(&self) -> u64 {
self.execution_duration_nanoseconds.load(Ordering::Relaxed)
}
fn increment(counter: &AtomicU64, value: u64) {
let _ = counter.try_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
Some(current.saturating_add(value))
});
}
pub(crate) fn record_claimed(&self, count: usize) {
Self::increment(&self.claimed_runs, u64::try_from(count).unwrap_or(u64::MAX));
}
pub(crate) fn record_claim_error(&self) {
Self::increment(&self.claim_errors, 1);
}
}
#[derive(Clone, Copy, Debug)]
pub(crate) enum ExecutionOutcome {
Completed,
Failed,
LeaseLost,
Cancelled,
Suspended,
}
#[derive(Debug)]
pub(crate) struct TaskExecution {
metrics: QueueMetrics,
start: Instant,
finished: bool,
}
impl TaskExecution {
pub(crate) fn start(metrics: QueueMetrics) -> Self {
Self { metrics, start: Instant::now(), finished: false }
}
pub(crate) fn finish(mut self, outcome: ExecutionOutcome) {
self.record(outcome);
self.finished = true;
}
fn record(&self, outcome: ExecutionOutcome) {
QueueMetrics::increment(&self.metrics.executions, 1);
let outcome_counter = match outcome {
ExecutionOutcome::Completed => &self.metrics.completed_executions,
ExecutionOutcome::Failed => &self.metrics.failed_executions,
ExecutionOutcome::LeaseLost => &self.metrics.lease_lost_executions,
ExecutionOutcome::Cancelled => &self.metrics.cancelled_executions,
ExecutionOutcome::Suspended => &self.metrics.suspended_executions,
};
QueueMetrics::increment(outcome_counter, 1);
let elapsed = u64::try_from(self.start.elapsed().as_nanos()).unwrap_or(u64::MAX);
QueueMetrics::increment(&self.metrics.execution_duration_nanoseconds, elapsed);
}
}
impl Drop for TaskExecution {
fn drop(&mut self) {
if !self.finished {
QueueMetrics::increment(&self.metrics.executions, 1);
QueueMetrics::increment(&self.metrics.unhandled_executions, 1);
let elapsed = u64::try_from(self.start.elapsed().as_nanos()).unwrap_or(u64::MAX);
QueueMetrics::increment(&self.metrics.execution_duration_nanoseconds, elapsed);
}
}
}
#[cfg(test)]
mod tests {
use super::{ExecutionOutcome, QueueMetrics, TaskExecution};
#[test]
fn cloned_metrics_observe_same_counters() {
let metrics = QueueMetrics::new();
let exporter = metrics.clone();
metrics.record_claimed(3);
metrics.record_claim_error();
assert_eq!(exporter.claimed_runs(), 3);
assert_eq!(exporter.claim_errors(), 1);
}
#[test]
fn executions_record_one_bounded_outcome() {
let metrics = QueueMetrics::new();
TaskExecution::start(metrics.clone()).finish(ExecutionOutcome::Completed);
assert_eq!(metrics.executions(), 1);
assert_eq!(metrics.completed_executions(), 1);
assert_eq!(metrics.unhandled_executions(), 0);
}
#[test]
fn lease_loss_is_a_bounded_execution_outcome() {
let metrics = QueueMetrics::new();
TaskExecution::start(metrics.clone()).finish(ExecutionOutcome::LeaseLost);
assert_eq!(metrics.executions(), 1);
assert_eq!(metrics.lease_lost_executions(), 1);
assert_eq!(metrics.unhandled_executions(), 0);
}
#[test]
fn dropped_executions_are_counted_as_unhandled() {
let metrics = QueueMetrics::new();
drop(TaskExecution::start(metrics.clone()));
assert_eq!(metrics.executions(), 1);
assert_eq!(metrics.unhandled_executions(), 1);
}
}