use std::sync::{
Arc,
atomic::{AtomicU64, Ordering},
};
use parking_lot::RwLock;
use super::{
latency_metrics_entry_session::LatencyMetricsEntrySession,
latency_metrics_type::LatencyMetricsType,
};
pub struct GarnetLatencyMetricsSession {
monitor_iterations: Arc<AtomicU64>,
metrics: RwLock<Option<Vec<LatencyMetricsEntrySession>>>,
}
impl GarnetLatencyMetricsSession {
pub const DEFAULT_LATENCY_TYPES: &[LatencyMetricsType] = &LatencyMetricsType::ALL;
pub fn new(
monitor_iterations: Arc<AtomicU64>,
latency_types: &'static [LatencyMetricsType],
) -> Self {
Self {
monitor_iterations,
metrics: RwLock::new(Some(
latency_types
.iter()
.map(|_| LatencyMetricsEntrySession::new())
.collect(),
)),
}
}
#[inline]
pub fn version(&self) -> usize {
(self.monitor_iterations.load(Ordering::Relaxed) % 2) as usize
}
#[inline]
pub fn prior_version(&self) -> usize {
1 - self.version()
}
pub fn metrics_snapshot(&self) -> Option<Vec<LatencyMetricsEntrySession>> {
self.metrics.read().clone()
}
pub fn return_to_pool(&self) {
if let Some(mut entries) = self.metrics.write().take() {
for entry in &mut entries {
entry.return_to_pool();
}
}
}
#[inline]
pub fn start(&self, cmd: LatencyMetricsType, now_ticks: u64) {
if let Some(entry) = self
.metrics
.write()
.as_deref_mut()
.and_then(|m| m.get_mut(cmd.idx()))
{
entry.start(now_ticks);
}
}
#[inline]
pub fn get(&self, cmd: LatencyMetricsType) -> u64 {
self
.metrics
.read()
.as_deref()
.and_then(|m| m.get(cmd.idx()))
.map_or(0, |entry| entry.start_timestamp)
}
#[inline]
pub fn stop_and_switch(
&self,
old_cmd: LatencyMetricsType,
new_cmd: LatencyMetricsType,
now_ticks: u64,
) {
let mut metrics = self.metrics.write();
let Some(entries) = metrics.as_deref_mut() else {
return;
};
let (old_idx, new_idx) = (old_cmd.idx(), new_cmd.idx());
if old_idx == new_idx {
if let Some(entry) = entries.get_mut(old_idx) {
entry.start_timestamp = 0;
entry.record_value(self.version(), now_ticks);
}
return;
}
let pair = if old_idx < new_idx {
let (left, right) = entries.split_at_mut(new_idx);
left.get_mut(old_idx).zip(right.first_mut())
} else {
let (left, right) = entries.split_at_mut(old_idx);
right.first_mut().zip(left.get_mut(new_idx))
};
let Some((old_entry, new_entry)) = pair else {
return;
};
new_entry.start_timestamp = old_entry.start_timestamp;
old_entry.start_timestamp = 0;
new_entry.record_value(self.version(), now_ticks);
}
#[inline]
pub fn stop(&self, cmd: LatencyMetricsType, now_ticks: u64) {
let ver = self.version();
if let Some(entry) = self
.metrics
.write()
.as_deref_mut()
.and_then(|m| m.get_mut(cmd.idx()))
{
entry.record_value(ver, now_ticks);
}
}
#[inline]
pub fn record_value(&self, cmd: LatencyMetricsType, elapsed: i64) {
let ver = self.version();
if let Some(entry) = self
.metrics
.write()
.as_deref_mut()
.and_then(|m| m.get_mut(cmd.idx()))
{
entry.record_elapsed(ver, elapsed);
}
}
pub fn reset_all(&self) {
for cmd in Self::DEFAULT_LATENCY_TYPES {
self.reset(*cmd);
}
}
pub fn reset(&self, cmd: LatencyMetricsType) {
let ver = self.prior_version();
if let Some(entry) = self
.metrics
.write()
.as_deref_mut()
.and_then(|m| m.get_mut(cmd.idx()))
{
entry.latency[ver].reset();
}
}
}