use std::fmt;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Default)]
pub struct ConnectionStats {
bytes_sent: AtomicU64,
bytes_received: AtomicU64,
requests_sent: AtomicU64,
responses_received: AtomicU64,
}
impl ConnectionStats {
pub fn new() -> Arc<Self> {
Arc::new(Self::default())
}
pub(crate) fn record_sent(&self, bytes: usize) {
metrics::counter!("kafka_bytes_sent_total").increment(u64::try_from(bytes).unwrap_or(0));
self.bytes_sent
.fetch_add(u64::try_from(bytes).unwrap_or(u64::MAX), Ordering::Relaxed);
self.requests_sent.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_received(&self, bytes: usize) {
metrics::counter!("kafka_bytes_received_total")
.increment(u64::try_from(bytes).unwrap_or(0));
self.bytes_received
.fetch_add(u64::try_from(bytes).unwrap_or(u64::MAX), Ordering::Relaxed);
self.responses_received.fetch_add(1, Ordering::Relaxed);
}
pub fn bytes_sent(&self) -> u64 {
self.bytes_sent.load(Ordering::Relaxed)
}
pub fn bytes_received(&self) -> u64 {
self.bytes_received.load(Ordering::Relaxed)
}
pub fn requests_sent(&self) -> u64 {
self.requests_sent.load(Ordering::Relaxed)
}
pub fn responses_received(&self) -> u64 {
self.responses_received.load(Ordering::Relaxed)
}
pub fn snapshot(&self) -> StatsSnapshot {
StatsSnapshot {
bytes_sent: self.bytes_sent(),
bytes_received: self.bytes_received(),
requests_sent: self.requests_sent(),
responses_received: self.responses_received(),
}
}
}
impl fmt::Debug for ConnectionStats {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.snapshot().fmt(f)
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct StatsSnapshot {
pub bytes_sent: u64,
pub bytes_received: u64,
pub requests_sent: u64,
pub responses_received: u64,
}
impl StatsSnapshot {
pub fn since(&self, earlier: &StatsSnapshot) -> StatsSnapshot {
StatsSnapshot {
bytes_sent: self.bytes_sent.saturating_sub(earlier.bytes_sent),
bytes_received: self.bytes_received.saturating_sub(earlier.bytes_received),
requests_sent: self.requests_sent.saturating_sub(earlier.requests_sent),
responses_received: self
.responses_received
.saturating_sub(earlier.responses_received),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn deltas_are_taken_against_a_snapshot_not_a_live_counter() {
let stats = ConnectionStats::new();
stats.record_sent(100);
stats.record_received(400);
let mark = stats.snapshot();
stats.record_received(50);
let delta = stats.snapshot().since(&mark);
assert_eq!(delta.bytes_received, 50);
assert_eq!(delta.bytes_sent, 0);
assert_eq!(delta.responses_received, 1);
}
#[test]
fn deltas_never_underflow() {
let later = StatsSnapshot::default();
let earlier = StatsSnapshot {
bytes_received: 10,
..Default::default()
};
assert_eq!(later.since(&earlier).bytes_received, 0);
}
}