use std::sync::{
Arc,
atomic::{AtomicU64, Ordering},
};
use crate::wire::BufferPoolStats;
#[derive(Debug, Clone, Default)]
pub(super) struct ConsumerMetrics {
inner: Arc<ConsumerMetricsInner>,
}
#[derive(Debug, Default)]
#[expect(
clippy::struct_field_names,
reason = "Every counter is a running total."
)]
struct ConsumerMetricsInner {
poll_total: AtomicU64,
records_consumed_total: AtomicU64,
fetch_total: AtomicU64,
commit_total: AtomicU64,
heartbeat_total: AtomicU64,
rebalance_total: AtomicU64,
}
impl ConsumerMetrics {
pub(super) fn record_poll(&self) {
let _previous = self.inner.poll_total.fetch_add(1, Ordering::Relaxed);
}
pub(super) fn record_records(&self, count: usize) {
let count = u64::try_from(count).unwrap_or(u64::MAX);
let _previous = self
.inner
.records_consumed_total
.fetch_add(count, Ordering::Relaxed);
}
pub(super) fn record_fetch(&self) {
let _previous = self.inner.fetch_total.fetch_add(1, Ordering::Relaxed);
}
pub(super) fn record_commit(&self) {
let _previous = self.inner.commit_total.fetch_add(1, Ordering::Relaxed);
}
pub(super) fn record_heartbeat(&self) {
let _previous = self.inner.heartbeat_total.fetch_add(1, Ordering::Relaxed);
}
pub(super) fn record_rebalance(&self) {
let _previous = self.inner.rebalance_total.fetch_add(1, Ordering::Relaxed);
}
pub(super) fn snapshot(&self, buffer_pool: BufferPoolStats) -> ConsumerMetricsSnapshot {
ConsumerMetricsSnapshot {
poll_total: self.inner.poll_total.load(Ordering::Relaxed),
records_consumed_total: self.inner.records_consumed_total.load(Ordering::Relaxed),
fetch_total: self.inner.fetch_total.load(Ordering::Relaxed),
commit_total: self.inner.commit_total.load(Ordering::Relaxed),
heartbeat_total: self.inner.heartbeat_total.load(Ordering::Relaxed),
rebalance_total: self.inner.rebalance_total.load(Ordering::Relaxed),
buffer_pool,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerMetricsSnapshot {
pub poll_total: u64,
pub records_consumed_total: u64,
pub fetch_total: u64,
pub commit_total: u64,
pub heartbeat_total: u64,
pub rebalance_total: u64,
pub buffer_pool: BufferPoolStats,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn snapshot_reflects_recorded_counters() {
let metrics = ConsumerMetrics::default();
metrics.record_poll();
metrics.record_poll();
metrics.record_records(7);
metrics.record_fetch();
metrics.record_commit();
metrics.record_heartbeat();
metrics.record_rebalance();
let snapshot = metrics.snapshot(BufferPoolStats::default());
assert_eq!(snapshot.poll_total, 2);
assert_eq!(snapshot.records_consumed_total, 7);
assert_eq!(snapshot.fetch_total, 1);
assert_eq!(snapshot.commit_total, 1);
assert_eq!(snapshot.heartbeat_total, 1);
assert_eq!(snapshot.rebalance_total, 1);
}
#[test]
fn clones_share_the_same_aggregate() {
let metrics = ConsumerMetrics::default();
let clone = metrics.clone();
clone.record_poll();
assert_eq!(metrics.snapshot(BufferPoolStats::default()).poll_total, 1);
}
}