use std::{
future::Future,
sync::{
Arc,
atomic::{AtomicU64, Ordering},
},
time::Duration,
};
use parking_lot::Mutex;
use super::{
command_stats::CommandStats,
garnet_server_metrics::GarnetServerMetrics,
garnet_session_metrics::GarnetSessionMetrics,
info_metrics_type::InfoMetricsType,
latency::{
garnet_latency_metrics_session::GarnetLatencyMetricsSession,
latency_metrics_type::LatencyMetricsType,
},
};
pub struct GarnetServerMonitor {
pub reset_event_flags: [bool; InfoMetricsType::ALL.len()],
pub reset_latency_metrics: [bool; LatencyMetricsType::ALL.len()],
monitor_sampling_frequency: Duration,
pub monitor_iterations: Arc<AtomicU64>,
state: Mutex<MonitorState>,
}
struct MonitorState {
global_metrics: GarnetServerMetrics,
acc_session_metrics: GarnetSessionMetrics,
acc_command_stats: Option<CommandStats>,
instant_input_net_bytes: u64,
instant_output_net_bytes: u64,
instant_commands_processed: u64,
}
pub struct SessionSample<'a> {
pub metrics: &'a GarnetSessionMetrics,
pub command_stats: Option<&'a CommandStats>,
pub latency: Option<&'a GarnetLatencyMetricsSession>,
}
pub struct ServerSample<'a> {
pub total_connections_received: i64,
pub total_connections_disposed: i64,
pub total_connections_active: i64,
pub sessions: Vec<SessionSample<'a>>,
}
pub struct MonitorIterationInputs<'a> {
pub servers: &'a [ServerSample<'a>],
pub reset_all_session_latency: &'a mut dyn FnMut(),
pub reset_active_sessions: &'a mut dyn FnMut(),
pub reset_active_command_stats: &'a mut dyn FnMut(),
pub reset_session_latency: &'a mut dyn FnMut(LatencyMetricsType),
}
impl GarnetServerMonitor {
pub fn new(
metrics_sampling_frequency_secs: u64,
track_stats: bool,
track_latency: bool,
track_command_stats: bool,
) -> Self {
Self {
reset_event_flags: [false; InfoMetricsType::ALL.len()],
reset_latency_metrics: [false; LatencyMetricsType::ALL.len()],
monitor_sampling_frequency: Duration::from_secs(metrics_sampling_frequency_secs),
monitor_iterations: Arc::new(AtomicU64::new(0)),
state: Mutex::new(MonitorState {
global_metrics: GarnetServerMetrics::new(track_stats, track_latency, track_command_stats),
acc_session_metrics: GarnetSessionMetrics::default(),
acc_command_stats: track_command_stats.then(CommandStats::new),
instant_input_net_bytes: 0,
instant_output_net_bytes: 0,
instant_commands_processed: 0,
}),
}
}
pub fn global_latency_metrics(
&self,
) -> Option<Arc<parking_lot::Mutex<super::latency::garnet_latency_metrics::GarnetLatencyMetrics>>>
{
self
.state
.lock()
.global_metrics
.global_latency_metrics
.clone()
}
pub fn global_metrics_snapshot(&self) -> (i64, i64, i64, f64, f64, f64) {
let state = self.state.lock();
let g = &state.global_metrics;
(
g.total_connections_received,
g.total_connections_disposed,
g.total_connections_active,
g.instantaneous_cmd_per_sec,
g.instantaneous_net_input_tpt,
g.instantaneous_net_output_tpt,
)
}
pub fn shared_iterations(&self) -> Arc<AtomicU64> {
self.monitor_iterations.clone()
}
pub fn add_metrics_history_session_dispose(
&self,
curr_session_metrics: Option<&GarnetSessionMetrics>,
curr_latency_metrics: Option<&GarnetLatencyMetricsSession>,
curr_command_stats: Option<&CommandStats>,
) {
let mut state = self.state.lock();
if let Some(metrics) = curr_session_metrics
&& let Some(history) = &mut state.global_metrics.history_session_metrics
{
history.add(metrics);
}
if let Some(latency) = curr_latency_metrics {
if let Some(global_latency) = &mut state.global_metrics.global_latency_metrics {
global_latency.lock().merge(latency);
}
latency.return_to_pool();
}
if let Some(stats) = curr_command_stats
&& let Some(history) = &mut state.global_metrics.history_command_stats
{
history.add(stats);
}
}
pub fn get_all_locksets(locksets: impl Iterator<Item = (i64, String)>) -> String {
let mut result = String::new();
for (session_id, lockset) in locksets {
if !lockset.is_empty() {
result += &format!("{session_id}: {lockset}\n");
}
}
result
}
fn track_latency(state: &MonitorState) -> bool {
state.global_metrics.global_latency_metrics.is_some()
}
fn update_instantaneous_metrics(state: &mut MonitorState, elapsed_sec: f64) {
let elapsed_units = elapsed_sec * GarnetServerMetrics::BYTE_UNIT as f64;
let g = &mut state.global_metrics;
g.instantaneous_net_input_tpt = (g
.global_session_metrics
.as_ref()
.map_or(0, GarnetSessionMetrics::get_total_net_input_bytes)
- state.instant_input_net_bytes) as f64
/ elapsed_units;
g.instantaneous_net_output_tpt = (g
.global_session_metrics
.as_ref()
.map_or(0, GarnetSessionMetrics::get_total_net_output_bytes)
- state.instant_output_net_bytes) as f64
/ elapsed_units;
g.instantaneous_cmd_per_sec = (g
.global_session_metrics
.as_ref()
.map_or(0, GarnetSessionMetrics::get_total_commands_processed)
- state.instant_commands_processed) as f64
/ elapsed_sec;
g.instantaneous_net_input_tpt = round2(g.instantaneous_net_input_tpt);
g.instantaneous_net_output_tpt = round2(g.instantaneous_net_output_tpt);
g.instantaneous_cmd_per_sec = g.instantaneous_cmd_per_sec.round();
state.instant_input_net_bytes = g
.global_session_metrics
.as_ref()
.map_or(0, GarnetSessionMetrics::get_total_net_input_bytes);
state.instant_output_net_bytes = g
.global_session_metrics
.as_ref()
.map_or(0, GarnetSessionMetrics::get_total_net_output_bytes);
state.instant_commands_processed = g
.global_session_metrics
.as_ref()
.map_or(0, GarnetSessionMetrics::get_total_commands_processed);
}
fn add_current_server_stats(state: &mut MonitorState, server: &ServerSample<'_>) {
for session in &server.sessions {
state.acc_session_metrics.add(session.metrics);
if let (Some(acc), Some(stats)) = (&mut state.acc_command_stats, session.command_stats) {
acc.add(stats);
}
if let (Some(global_latency), Some(latency)) = (
&state.global_metrics.global_latency_metrics,
session.latency,
) {
global_latency.lock().merge(latency);
}
}
if let Some(global_session) = &mut state.global_metrics.global_session_metrics {
global_session.reset();
global_session.add(&state.acc_session_metrics);
}
if let Some(global_stats) = &mut state.global_metrics.global_command_stats {
global_stats.reset();
if let Some(acc) = &state.acc_command_stats {
global_stats.add(acc);
}
}
}
fn reset_and_add_global_history(state: &mut MonitorState) {
state.acc_session_metrics.reset();
if let Some(history) = &state.global_metrics.history_session_metrics {
state.acc_session_metrics.add(history);
}
if let Some(acc) = &mut state.acc_command_stats {
acc.reset();
if let Some(history) = &state.global_metrics.history_command_stats {
acc.add(history);
}
}
}
fn cleanup_global_stats(
state: &mut MonitorState,
flags: &mut [bool],
reset_active_sessions: &mut dyn FnMut(),
reset_active_command_stats: &mut dyn FnMut(),
) {
if flags[InfoMetricsType::Stats as usize] {
log::info!("Resetting latency metrics for commands");
state.global_metrics.instantaneous_net_input_tpt = 0.0;
state.global_metrics.instantaneous_net_output_tpt = 0.0;
state.global_metrics.instantaneous_cmd_per_sec = 0.0;
state.global_metrics.total_connections_received = 0;
state.global_metrics.total_connections_disposed = 0;
if let Some(global_session) = &mut state.global_metrics.global_session_metrics {
global_session.reset();
}
if let Some(history) = &mut state.global_metrics.history_session_metrics {
history.reset();
}
reset_active_sessions();
flags[InfoMetricsType::Stats as usize] = false;
}
if flags[InfoMetricsType::CommandStats as usize] {
log::info!("Resetting command stats");
if let Some(global_stats) = &mut state.global_metrics.global_command_stats {
global_stats.reset();
}
if let Some(history) = &mut state.global_metrics.history_command_stats {
history.reset();
}
reset_active_command_stats();
flags[InfoMetricsType::CommandStats as usize] = false;
}
}
fn cleanup_global_latency_metrics(
state: &mut MonitorState,
flags: &mut [bool],
reset_session_latency: &mut dyn FnMut(LatencyMetricsType),
) {
if !Self::track_latency(state) {
return;
}
for (idx, flagged) in flags.iter_mut().enumerate() {
if !*flagged {
continue;
}
let Some(event_type) = LatencyMetricsType::ALL.get(idx).copied() else {
continue;
};
log::info!("Resetting server-side stats {event_type:?}");
reset_session_latency(event_type);
if let Some(global_latency) = &state.global_metrics.global_latency_metrics {
global_latency.lock().reset(event_type);
}
*flagged = false;
}
}
fn reset_latency_session_metrics(
state: &MonitorState,
reset_all_session_latency: &mut dyn FnMut(),
) {
if Self::track_latency(state) {
reset_all_session_latency();
}
}
fn monitor_iteration(&mut self, inputs: &mut MonitorIterationInputs<'_>) {
let mut state = self.state.lock();
Self::reset_latency_session_metrics(&state, inputs.reset_all_session_latency);
self.monitor_iterations.fetch_add(1, Ordering::Relaxed);
Self::reset_and_add_global_history(&mut state);
let (mut total_received, mut total_disposed, mut total_active) = (0i64, 0i64, 0i64);
for server in inputs.servers {
total_received += server.total_connections_received;
total_disposed += server.total_connections_disposed;
total_active += server.total_connections_active;
Self::add_current_server_stats(&mut state, server);
}
Self::update_instantaneous_metrics(&mut state, self.monitor_sampling_frequency.as_secs_f64());
state.global_metrics.total_connections_received = total_received;
state.global_metrics.total_connections_disposed = total_disposed;
state.global_metrics.total_connections_active = total_active;
Self::cleanup_global_stats(
&mut state,
&mut self.reset_event_flags,
inputs.reset_active_sessions,
inputs.reset_active_command_stats,
);
Self::cleanup_global_latency_metrics(
&mut state,
&mut self.reset_latency_metrics,
inputs.reset_session_latency,
);
}
pub async fn main_monitor_task_async<S, Fut>(
&mut self,
mut sleep: S,
cancelled: impl Fn() -> bool,
mut resolve: impl FnMut() -> MonitorIterationInputs<'static>,
) where
S: FnMut(Duration) -> Fut,
Fut: Future<Output = ()>,
{
while !cancelled() {
sleep(self.monitor_sampling_frequency).await;
self.monitor_iteration(&mut resolve());
}
}
}
#[inline]
fn round2(v: f64) -> f64 {
(v * 100.0).round() / 100.0
}