use std::sync::{
Arc, Mutex,
atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering},
};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct ClientStats {
pub queued_commands: usize,
pub queued_bytes: usize,
pub queued_bytes_high_water: usize,
pub shed_commands: u64,
pub reconnections: u64,
pub connected: bool,
}
#[derive(Debug, Default)]
pub(crate) struct StatsRecorder {
queued_commands: AtomicUsize,
queued_bytes: AtomicUsize,
queued_bytes_high_water: AtomicUsize,
shed_commands: AtomicU64,
reconnections: AtomicU64,
connected: AtomicBool,
server_version: Mutex<Option<Arc<str>>>,
}
impl StatsRecorder {
pub(crate) fn new() -> Arc<Self> {
Arc::new(Self::default())
}
pub(crate) fn set_queued(&self, commands: usize, bytes: usize) {
self.queued_commands.store(commands, Ordering::Relaxed);
self.queued_bytes.store(bytes, Ordering::Relaxed);
self.queued_bytes_high_water
.fetch_max(bytes, Ordering::Relaxed);
}
pub(crate) fn record_shed(&self) {
self.shed_commands.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_reconnection(&self) {
self.reconnections.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn set_connected(&self, connected: bool) {
self.connected.store(connected, Ordering::Relaxed);
}
pub(crate) fn set_server_version(&self, version: Option<Arc<str>>) {
if let Ok(mut guard) = self.server_version.lock() {
*guard = version;
}
}
pub(crate) fn server_version(&self) -> Option<Arc<str>> {
self.server_version.lock().ok()?.clone()
}
pub(crate) fn connected(&self) -> bool {
self.connected.load(Ordering::Relaxed)
}
pub(crate) fn snapshot(&self) -> ClientStats {
ClientStats {
queued_commands: self.queued_commands.load(Ordering::Relaxed),
queued_bytes: self.queued_bytes.load(Ordering::Relaxed),
queued_bytes_high_water: self.queued_bytes_high_water.load(Ordering::Relaxed),
shed_commands: self.shed_commands.load(Ordering::Relaxed),
reconnections: self.reconnections.load(Ordering::Relaxed),
connected: self.connected.load(Ordering::Relaxed),
}
}
}