use std::{env, num::NonZeroUsize, thread};
use crate::metrics::{
garnet_session_metrics::GarnetSessionMetrics, info_metrics_type::InfoMetricsType,
latency::garnet_latency_metrics::fmt_n2, metrics_item::MetricsItem,
system_metrics::SystemMetrics,
};
pub const DEFAULT_INFO: &[InfoMetricsType] = &[
InfoMetricsType::Server,
InfoMetricsType::Memory,
InfoMetricsType::Cluster,
InfoMetricsType::Replication,
InfoMetricsType::Stats,
InfoMetricsType::Store,
InfoMetricsType::Persistence,
InfoMetricsType::Clients,
InfoMetricsType::Modules,
InfoMetricsType::BpStats,
InfoMetricsType::CInfo,
];
pub const ALL_INFO_SET: &[InfoMetricsType] = &[
InfoMetricsType::Server,
InfoMetricsType::Memory,
InfoMetricsType::Cluster,
InfoMetricsType::Replication,
InfoMetricsType::Stats,
InfoMetricsType::Store,
InfoMetricsType::Persistence,
InfoMetricsType::Clients,
InfoMetricsType::Keyspace,
InfoMetricsType::BpStats,
InfoMetricsType::CInfo,
];
#[derive(Debug, Default, Clone)]
pub struct DbSnapshot {
pub id: i32,
pub current_version: i64,
pub last_checkpointed_version: i64,
pub system_state: String,
pub index_bucket_count: i64,
pub index_bucket_size_bytes: i64,
pub index_memory_size_bytes: i64,
pub index_overflow_bucket_count: i64,
pub index_overflow_memory_size_bytes: i64,
pub index_total_memory_size_bytes: i64,
pub log_page_size_bytes: i64,
pub log_max_allocated_page_count: i64,
pub log_allocated_page_count: i64,
pub log_max_memory_size_bytes: i64,
pub log_memory_size_bytes: i64,
pub log_heap_size_bytes: i64,
pub log_begin_address: i64,
pub log_head_address: i64,
pub log_safe_readonly_address: i64,
pub log_flushed_until_address: i64,
pub log_tail_address: i64,
pub read_cache: Option<ReadCacheSnapshot>,
pub mainlog_target_size: Option<i64>,
pub readcache_target_size: Option<i64>,
pub aof_memory_size_bytes: i64,
pub aof: Option<AofSnapshot>,
pub hash_distribution_dump: String,
pub revivification_dump: String,
}
#[derive(Debug, Default, Clone)]
pub struct ReadCacheSnapshot {
pub page_size_bytes: i64,
pub max_allocated_page_count: i64,
pub allocated_page_count: i64,
pub max_memory_size_bytes: i64,
pub memory_size_bytes: i64,
pub heap_size_bytes: i64,
pub begin_address: i64,
pub head_address: i64,
pub tail_address: i64,
}
#[derive(Debug, Default, Clone)]
pub struct AofSnapshot {
pub committed_begin_address: i64,
pub committed_until_address: i64,
pub flushed_until_address: i64,
pub begin_address: i64,
pub tail_address: i64,
}
#[derive(Debug, Default, Clone)]
pub struct GlobalMetricsSnapshot {
pub total_connections_active: i64,
pub total_connections_received: i64,
pub total_connections_disposed: i64,
pub instantaneous_cmd_per_sec: f64,
pub instantaneous_net_input_tpt: f64,
pub instantaneous_net_output_tpt: f64,
pub global_session_metrics: GarnetSessionMetrics,
}
#[derive(Debug, Clone)]
pub struct ServerFacts {
pub version: String,
pub run_id: String,
pub redis_protocol_version: String,
pub enable_cluster: bool,
pub enable_aof: bool,
pub metrics_sampling_frequency: i32,
pub latency_monitor: bool,
pub command_stats_monitor: bool,
pub startup_timestamp_unix_secs: i64,
pub log_dir: String,
}
pub trait InfoProvider {
fn server_facts(&self) -> ServerFacts;
fn databases(&self) -> Vec<DbSnapshot>;
fn max_database_id(&self) -> i32;
fn global_metrics(&self) -> Option<GlobalMetricsSnapshot>;
fn command_stats(&self) -> Vec<(String, u64, u64)>;
fn keyspace_stats(&self, db_id: i32) -> (u64, u64);
fn replication_info(&self) -> Option<Vec<MetricsItem>>;
fn gossip_stats(&self, metrics_disabled: bool) -> Vec<MetricsItem>;
fn buffer_pool_stats(&self) -> Vec<(String, String)>;
fn checkpoint_info(&self) -> Option<Vec<MetricsItem>>;
fn hlog_scan_dump(&self) -> Vec<(String, String)>;
fn safe_aof_address(&self) -> i64;
}
pub struct GarnetInfoMetrics {
server_info: Option<Vec<MetricsItem>>,
memory_info: Option<Vec<MetricsItem>>,
cluster_info: Option<Vec<MetricsItem>>,
replication_info: Option<Vec<MetricsItem>>,
stats_info: Option<Vec<MetricsItem>>,
store_info: Option<Vec<Vec<MetricsItem>>>,
store_hash_distr_info: Option<Vec<Vec<MetricsItem>>>,
store_reviv_info: Option<Vec<Vec<MetricsItem>>>,
persistence_info: Option<Vec<Vec<MetricsItem>>>,
clients_info: Option<Vec<MetricsItem>>,
keyspace_info: Option<Vec<MetricsItem>>,
buffer_pool_stats: Option<Vec<MetricsItem>>,
checkpoint_stats: Option<Vec<MetricsItem>>,
hlog_scan_stats: Option<Vec<Vec<MetricsItem>>>,
command_stats_info: Option<Vec<MetricsItem>>,
}
impl Default for GarnetInfoMetrics {
fn default() -> Self {
Self::new()
}
}
impl GarnetInfoMetrics {
pub fn new() -> Self {
Self {
server_info: None,
memory_info: None,
cluster_info: None,
replication_info: None,
stats_info: None,
store_info: None,
store_hash_distr_info: None,
store_reviv_info: None,
persistence_info: None,
clients_info: None,
keyspace_info: None,
buffer_pool_stats: None,
checkpoint_stats: None,
hlog_scan_stats: None,
command_stats_info: None,
}
}
fn populate_server_info(&mut self, provider: &dyn InfoProvider) {
let facts = provider.server_facts();
let uptime_secs =
coarsetime::Clock::now_since_epoch().as_secs() as i64 - facts.startup_timestamp_unix_secs;
let uptime_secs = uptime_secs.max(0);
self.server_info = Some(vec![
MetricsItem::new("garnet_version", facts.version),
MetricsItem::new("server_name", "garnet"),
MetricsItem::new("os", env::consts::OS),
MetricsItem::new(
"processor_count",
thread::available_parallelism()
.map_or(1, NonZeroUsize::get)
.to_string(),
),
MetricsItem::new(
"arch_bits",
if cfg!(target_pointer_width = "64") {
"64"
} else {
"32"
},
),
MetricsItem::new("uptime_in_seconds", uptime_secs.to_string()),
MetricsItem::new("uptime_in_days", (uptime_secs / 86_400).to_string()),
MetricsItem::new(
"monitor_task",
if facts.metrics_sampling_frequency > 0 {
"enabled"
} else {
"disabled"
},
),
MetricsItem::new("monitor_freq", facts.metrics_sampling_frequency.to_string()),
MetricsItem::new(
"latency_monitor",
if facts.latency_monitor {
"enabled"
} else {
"disabled"
},
),
MetricsItem::new(
"commandstats_monitor",
if facts.command_stats_monitor {
"enabled"
} else {
"disabled"
},
),
MetricsItem::new("run_id", facts.run_id),
MetricsItem::new("redis_version", facts.redis_protocol_version),
MetricsItem::new(
"redis_mode",
if facts.enable_cluster {
"cluster"
} else {
"standalone"
},
),
]);
}
fn populate_memory_info(&mut self, provider: &dyn InfoProvider) {
let facts = provider.server_facts();
let mut store_index_size = 0i64;
let mut store_mainlog_memory_target_size = 0i64;
let mut store_mainlog_memory_size = 0i64;
let mut store_readcache_memory_size = 0i64;
let mut aof_log_memory_size = if facts.enable_aof { 0 } else { -1 };
for db in provider.databases() {
store_index_size += db.index_total_memory_size_bytes;
aof_log_memory_size += db.aof_memory_size_bytes;
match db.mainlog_target_size {
None => store_mainlog_memory_size += db.log_memory_size_bytes,
Some(target) => {
store_mainlog_memory_target_size += target;
store_mainlog_memory_size += db.log_memory_size_bytes;
}
}
match db.readcache_target_size {
None => {
store_readcache_memory_size += db.read_cache.as_ref().map_or(0, |rc| rc.memory_size_bytes)
}
Some(_) => {
store_readcache_memory_size +=
db.read_cache.as_ref().map_or(0, |rc| rc.memory_size_bytes);
}
}
}
let total_store_size =
store_index_size + store_mainlog_memory_size + store_readcache_memory_size;
let m = |name: &str, value: i64| MetricsItem::new(name, value.to_string());
self.memory_info = Some(vec![
MetricsItem::new("system_page_size", page_size().to_string()),
m("total_system_memory", SystemMetrics::get_total_memory(1)),
m(
"total_system_memory(MB)",
SystemMetrics::get_total_memory(1 << 20),
),
m(
"available_system_memory",
SystemMetrics::get_physical_available_memory(1),
),
m(
"available_system_memory(MB)",
SystemMetrics::get_physical_available_memory(1 << 20),
),
m(
"proc_paged_memory_size",
SystemMetrics::get_paged_memory_size(1),
),
m(
"proc_paged_memory_size(MB)",
SystemMetrics::get_paged_memory_size(1 << 20),
),
m(
"proc_peak_paged_memory_size",
SystemMetrics::get_peak_paged_memory_size(1),
),
m(
"proc_peak_paged_memory_size(MB)",
SystemMetrics::get_peak_paged_memory_size(1 << 20),
),
m(
"proc_pageable_memory_size",
SystemMetrics::get_paged_system_memory_size(1),
),
m(
"proc_pageable_memory_size(MB)",
SystemMetrics::get_paged_system_memory_size(1 << 20),
),
m(
"proc_private_memory_size",
SystemMetrics::get_private_memory_size64(1),
),
m(
"proc_private_memory_size(MB)",
SystemMetrics::get_private_memory_size64(1 << 20),
),
m(
"proc_virtual_memory_size",
SystemMetrics::get_virtual_memory_size64(1),
),
m(
"proc_virtual_memory_size(MB)",
SystemMetrics::get_virtual_memory_size64(1 << 20),
),
m(
"proc_peak_virtual_memory_size",
SystemMetrics::get_peak_virtual_memory_size64(1),
),
m(
"proc_peak_virtual_memory_size(MB)",
SystemMetrics::get_peak_virtual_memory_size64(1 << 20),
),
m(
"proc_physical_memory_size",
SystemMetrics::get_physical_memory_usage(1),
),
m(
"proc_physical_memory_size(MB)",
SystemMetrics::get_physical_memory_usage(1 << 20),
),
m(
"proc_peak_physical_memory_size",
SystemMetrics::get_peak_physical_memory_usage(1),
),
m(
"proc_peak_physical_memory_size(MB)",
SystemMetrics::get_peak_physical_memory_usage(1 << 20),
),
m("gc_committed_bytes", 0),
m("gc_heap_bytes", 0),
m("gc_managed_memory_bytes_excluding_heap", 0),
m("gc_fragmented_bytes", 0),
m("native_allocator_bytes", 0),
m("store_index_size", store_index_size),
m("store_mainlog_memory_size", store_mainlog_memory_size),
m("store_readcache_memory_size", store_readcache_memory_size),
m("total_main_store_size", total_store_size),
m(
"store_heap_memory_target_size",
store_mainlog_memory_target_size,
),
m("aof_memory_size", aof_log_memory_size),
]);
}
fn populate_cluster_info(&mut self, provider: &dyn InfoProvider) {
let facts = provider.server_facts();
self.cluster_info = Some(vec![MetricsItem::new(
"cluster_enabled",
if facts.enable_cluster { "1" } else { "0" },
)]);
}
fn populate_replication_info(&mut self, provider: &dyn InfoProvider) {
self.replication_info = Some(provider.replication_info().unwrap_or_else(|| {
vec![
MetricsItem::new("role", "master"),
MetricsItem::new("connected_slaves", "0"),
MetricsItem::new("master_failover_state", "no-failover"),
MetricsItem::new("master_replid", generate_default_hex_id()),
MetricsItem::new("master_replid2", generate_default_hex_id()),
MetricsItem::new("master_repl_offset", "N/A"),
MetricsItem::new("second_repl_offset", "N/A"),
MetricsItem::new("store_current_safe_aof_address", "N/A"),
MetricsItem::new("store_recovered_safe_aof_address", "N/A"),
]
}));
}
fn populate_stats_info(&mut self, provider: &dyn InfoProvider) {
let facts = provider.server_facts();
let Some(global) = provider.global_metrics() else {
self.stats_info = Some(vec![
MetricsItem::new("total_connections_active", "0"),
MetricsItem::new("total_connections_received", "0"),
MetricsItem::new("total_connections_disposed", "0"),
MetricsItem::new("total_commands_processed", "0"),
MetricsItem::new("instantaneous_ops_per_sec", "0"),
MetricsItem::new("total_net_input_bytes", "0"),
MetricsItem::new("total_net_output_bytes", "0"),
MetricsItem::new("instantaneous_net_input_KBps", "0"),
MetricsItem::new("instantaneous_net_output_KBps", "0"),
MetricsItem::new("total_pending", "0"),
MetricsItem::new("total_found", "0"),
MetricsItem::new("total_notfound", "0"),
MetricsItem::new("garnet_hit_rate", fmt_n2(0.0)),
MetricsItem::new("total_cluster_commands_processed", "0"),
MetricsItem::new("total_write_commands_processed", "0"),
MetricsItem::new("total_read_commands_processed", "0"),
MetricsItem::new("total_number_resp_server_session_exceptions", "0"),
MetricsItem::new("total_transaction_commands_received", "0"),
MetricsItem::new("total_transaction_commands_execution_failed", "0"),
]);
return;
};
let session = &global.global_session_metrics;
let tt = session.get_total_found() + session.get_total_notfound();
let garnet_hit_rate = if tt > 0 {
session.get_total_found() as f64 / tt as f64
} else {
0.0
} * 100.0;
let mut items = vec![
MetricsItem::new(
"total_connections_active",
global.total_connections_active.to_string(),
),
MetricsItem::new(
"total_connections_received",
global.total_connections_received.to_string(),
),
MetricsItem::new(
"total_connections_disposed",
global.total_connections_disposed.to_string(),
),
MetricsItem::new(
"total_commands_processed",
session.get_total_commands_processed().to_string(),
),
MetricsItem::new(
"instantaneous_ops_per_sec",
format!("{}", global.instantaneous_cmd_per_sec),
),
MetricsItem::new(
"total_net_input_bytes",
session.get_total_net_input_bytes().to_string(),
),
MetricsItem::new(
"total_net_output_bytes",
session.get_total_net_output_bytes().to_string(),
),
MetricsItem::new(
"instantaneous_net_input_KBps",
format!("{}", global.instantaneous_net_input_tpt),
),
MetricsItem::new(
"instantaneous_net_output_KBps",
format!("{}", global.instantaneous_net_output_tpt),
),
MetricsItem::new("total_pending", session.get_total_pending().to_string()),
MetricsItem::new("total_found", session.get_total_found().to_string()),
MetricsItem::new("total_notfound", session.get_total_notfound().to_string()),
MetricsItem::new("garnet_hit_rate", fmt_n2(garnet_hit_rate)),
MetricsItem::new(
"total_cluster_commands_processed",
session.get_total_cluster_commands_processed().to_string(),
),
MetricsItem::new(
"total_write_commands_processed",
session.get_total_write_commands_processed().to_string(),
),
MetricsItem::new(
"total_read_commands_processed",
session.get_total_read_commands_processed().to_string(),
),
MetricsItem::new(
"total_number_resp_server_session_exceptions",
session
.get_total_number_resp_server_session_exceptions()
.to_string(),
),
MetricsItem::new(
"total_transaction_commands_received",
session
.get_total_transaction_commands_received()
.to_string(),
),
MetricsItem::new(
"total_transaction_commands_execution_failed",
session
.get_total_transaction_commands_execution_failed()
.to_string(),
),
];
if facts.enable_cluster {
items.extend(provider.gossip_stats(false));
}
self.stats_info = Some(items);
}
fn populate_command_stats_info(&mut self, provider: &dyn InfoProvider) {
if !provider.server_facts().command_stats_monitor {
self.command_stats_info = Some(vec![MetricsItem::new(
"",
"Command stats monitoring is disabled. Enable with --commandstats-monitor flag.",
)]);
return;
}
let stats = provider.command_stats();
self.command_stats_info = if stats.is_empty() {
None
} else {
Some(
stats
.into_iter()
.map(|(name, calls, rejected)| {
MetricsItem::new(
format!("cmdstat_{name}"),
format!(
"calls={calls},usec=0,usec_per_call=0.00,rejected_calls={rejected},failed_calls=0"
),
)
})
.collect(),
)
};
}
fn populate_store_stats(&mut self, provider: &dyn InfoProvider) {
let mut store_info = vec![Vec::new(); (provider.max_database_id() + 1).max(0) as usize];
for db in provider.databases() {
let stats = Self::get_database_store_stats(provider, &db);
store_info[db.id as usize] = stats;
}
self.store_info = Some(store_info);
}
fn get_database_store_stats(provider: &dyn InfoProvider, db: &DbSnapshot) -> Vec<MetricsItem> {
let n = |v: i64| v.to_string();
let mut items = vec![
MetricsItem::new("CurrentVersion", n(db.current_version)),
MetricsItem::new("LastCheckpointedVersion", n(db.last_checkpointed_version)),
MetricsItem::new("SystemState", db.system_state.clone()),
MetricsItem::new("IndexBucketCount", n(db.index_bucket_count)),
MetricsItem::new("IndexBucketSizeBytes", n(db.index_bucket_size_bytes)),
MetricsItem::new("IndexMemorySizeBytes", n(db.index_memory_size_bytes)),
MetricsItem::new(
"IndexOverflowBucketCount",
n(db.index_overflow_bucket_count),
),
MetricsItem::new(
"IndexOverflowMemorySizeBytes",
n(db.index_overflow_memory_size_bytes),
),
MetricsItem::new(
"IndexTotalMemorySizeBytes",
n(db.index_total_memory_size_bytes),
),
MetricsItem::new("LogDir", provider.server_facts().log_dir),
MetricsItem::new("Log.PageSizeBytes", n(db.log_page_size_bytes)),
MetricsItem::new("Log.MaxPageCount", n(db.log_max_allocated_page_count)),
MetricsItem::new("Log.AllocatedPageCount", n(db.log_allocated_page_count)),
MetricsItem::new("Log.MaxMemorySizeBytes", n(db.log_max_memory_size_bytes)),
MetricsItem::new("Log.CurrentMemorySizeBytes", n(db.log_memory_size_bytes)),
MetricsItem::new("Log.CurrentHeapSizeBytes", n(db.log_heap_size_bytes)),
MetricsItem::new("Log.BeginAddress", n(db.log_begin_address)),
MetricsItem::new("Log.HeadAddress", n(db.log_head_address)),
MetricsItem::new("Log.SafeReadOnlyAddress", n(db.log_safe_readonly_address)),
MetricsItem::new("Log.FlushedUntilAddress", n(db.log_flushed_until_address)),
MetricsItem::new("Log.TailAddress", n(db.log_tail_address)),
];
let read_cache = db.read_cache.as_ref();
let na = |v: Option<i64>| v.map_or_else(|| "N/A".to_string(), |v| v.to_string());
items.push(MetricsItem::new(
"ReadCache.PageSizeBytes",
na(read_cache.map(|rc| rc.page_size_bytes)),
));
items.push(MetricsItem::new(
"ReadCache.MaxPageCount",
na(read_cache.map(|rc| rc.max_allocated_page_count)),
));
items.push(MetricsItem::new(
"ReadCache.AllocatedPageCount",
na(read_cache.map(|rc| rc.allocated_page_count)),
));
items.push(MetricsItem::new(
"ReadCache.MaxMemorySizeBytes",
na(read_cache.map(|rc| rc.max_memory_size_bytes)),
));
items.push(MetricsItem::new(
"ReadCache.CurrentMemorySizeBytes",
na(read_cache.map(|rc| rc.memory_size_bytes)),
));
items.push(MetricsItem::new(
"ReadCache.CurrentHeapSizeBytes",
na(read_cache.map(|rc| rc.heap_size_bytes)),
));
items.push(MetricsItem::new(
"ReadCache.BeginAddress",
na(read_cache.map(|rc| rc.begin_address)),
));
items.push(MetricsItem::new(
"ReadCache.HeadAddress",
na(read_cache.map(|rc| rc.head_address)),
));
items.push(MetricsItem::new(
"ReadCache.TailAddress",
na(read_cache.map(|rc| rc.tail_address)),
));
items
}
fn populate_store_hash_distribution(&mut self, provider: &dyn InfoProvider) {
let mut info = vec![Vec::new(); (provider.max_database_id() + 1).max(0) as usize];
for db in provider.databases() {
info[db.id as usize] = vec![MetricsItem::new("", db.hash_distribution_dump.clone())];
}
self.store_hash_distr_info = Some(info);
}
fn populate_store_reviv_info(&mut self, provider: &dyn InfoProvider) {
let mut info = vec![Vec::new(); (provider.max_database_id() + 1).max(0) as usize];
for db in provider.databases() {
info[db.id as usize] = vec![MetricsItem::new("", db.revivification_dump.clone())];
}
self.store_reviv_info = Some(info);
}
fn populate_persistence_info(&mut self, provider: &dyn InfoProvider) {
let mut info = vec![Vec::new(); (provider.max_database_id() + 1).max(0) as usize];
for db in provider.databases() {
info[db.id as usize] = Self::get_database_persistence_stats(provider, &db);
}
self.persistence_info = Some(info);
}
fn get_database_persistence_stats(
provider: &dyn InfoProvider,
db: &DbSnapshot,
) -> Vec<MetricsItem> {
let aof_enabled = provider.server_facts().enable_aof;
let na = |v: Option<i64>| v.map_or_else(|| "N/A".to_string(), |v| v.to_string());
let aof = db.aof.as_ref();
vec![
MetricsItem::new(
"CommittedBeginAddress",
if !aof_enabled {
"N/A".into()
} else {
na(aof.map(|a| a.committed_begin_address))
},
),
MetricsItem::new(
"CommittedUntilAddress",
if !aof_enabled {
"N/A".into()
} else {
na(aof.map(|a| a.committed_until_address))
},
),
MetricsItem::new(
"FlushedUntilAddress",
if !aof_enabled {
"N/A".into()
} else {
na(aof.map(|a| a.flushed_until_address))
},
),
MetricsItem::new(
"BeginAddress",
if !aof_enabled {
"N/A".into()
} else {
na(aof.map(|a| a.begin_address))
},
),
MetricsItem::new(
"TailAddress",
if !aof_enabled {
"N/A".into()
} else {
na(aof.map(|a| a.tail_address))
},
),
MetricsItem::new(
"SafeAofAddress",
if !aof_enabled {
"N/A".into()
} else {
provider.safe_aof_address().to_string()
},
),
]
}
fn populate_clients_info(&mut self, provider: &dyn InfoProvider) {
let connected = provider.global_metrics().map_or(0, |g| {
g.total_connections_received - g.total_connections_disposed
});
self.clients_info = Some(vec![MetricsItem::new(
"connected_clients",
if provider.global_metrics().is_some() {
connected.to_string()
} else {
"0".to_string()
},
)]);
}
fn populate_keyspace_info(&mut self, provider: &dyn InfoProvider) {
let mut items = None;
for db in provider.databases() {
let (key_count, expire_count) = provider.keyspace_stats(db.id);
if key_count == 0 {
continue;
}
items.get_or_insert_with(Vec::new).push(MetricsItem::new(
format!("db{}", db.id),
format!("keys={key_count},expires={expire_count},avg_ttl=0"),
));
}
self.keyspace_info = items;
}
fn populate_cluster_buffer_pool_stats(&mut self, provider: &dyn InfoProvider) {
self.buffer_pool_stats = Some(
provider
.buffer_pool_stats()
.into_iter()
.map(|(name, stats)| MetricsItem::new(name, stats))
.collect(),
);
}
fn populate_checkpoint_info(&mut self, provider: &dyn InfoProvider) {
self.checkpoint_stats = provider.checkpoint_info();
}
fn populate_hlog_scan_info(&mut self, provider: &dyn InfoProvider) {
let mut result = Vec::new();
for (i, (main, object)) in provider.hlog_scan_dump().into_iter().enumerate() {
let main = if main.is_empty() {
"Empty".to_string()
} else {
main
};
let object = if object.is_empty() {
"Empty".to_string()
} else {
object
};
result.push(vec![
MetricsItem::new(format!("MainStore_HLog_{i}"), main),
MetricsItem::new(format!("ObjectStore_HLog_{i}"), object),
]);
}
self.hlog_scan_stats = Some(result);
}
pub fn get_section_header(info_type: InfoMetricsType, db_id: i32) -> String {
match info_type {
InfoMetricsType::Server => "Server".into(),
InfoMetricsType::Memory => "Memory".into(),
InfoMetricsType::Cluster => "Cluster".into(),
InfoMetricsType::Replication => "Replication".into(),
InfoMetricsType::Stats => "Stats".into(),
InfoMetricsType::Store => format!("Store_DB_{db_id}"),
InfoMetricsType::StoreHashtable => format!("StoreHashTableDistribution_DB_{db_id}"),
InfoMetricsType::StoreReviv => format!("StoreDeletedRecordRevivification_DB_{db_id}"),
InfoMetricsType::Persistence => format!("Persistence_DB_{db_id}"),
InfoMetricsType::Clients => "Clients".into(),
InfoMetricsType::Keyspace => "Keyspace".into(),
InfoMetricsType::Modules => "Modules".into(),
InfoMetricsType::BpStats => "BufferPoolStats".into(),
InfoMetricsType::CInfo => "CheckpointInfo".into(),
InfoMetricsType::HlogScan => format!("MainStoreHLogScan_DB_{db_id}"),
InfoMetricsType::CommandStats => "Commandstats".into(),
}
}
fn get_section_resp_info(
section_header: &str,
info: Option<&[MetricsItem]>,
sb_response: &mut String,
) {
sb_response.push_str(&format!("# {section_header}\r\n"));
let Some(info) = info else {
return;
};
if info.first().is_some_and(|item| item.name.is_empty()) {
sb_response.push_str(&format!("{}\r\n", info[0].value));
} else {
for item in info {
sb_response.push_str(&format!("{}:{}\r\n", item.name, item.value));
}
}
}
fn get_resp_info_single(
&mut self,
section: InfoMetricsType,
db_id: i32,
provider: &dyn InfoProvider,
sb_response: &mut String,
) {
let header = Self::get_section_header(section, db_id);
match section {
InfoMetricsType::Server => {
self.populate_server_info(provider);
Self::get_section_resp_info(&header, self.server_info.as_deref(), sb_response);
}
InfoMetricsType::Memory => {
self.populate_memory_info(provider);
Self::get_section_resp_info(&header, self.memory_info.as_deref(), sb_response);
}
InfoMetricsType::Cluster => {
self.populate_cluster_info(provider);
Self::get_section_resp_info(&header, self.cluster_info.as_deref(), sb_response);
}
InfoMetricsType::Replication => {
self.populate_replication_info(provider);
Self::get_section_resp_info(&header, self.replication_info.as_deref(), sb_response);
}
InfoMetricsType::Stats => {
self.populate_stats_info(provider);
Self::get_section_resp_info(&header, self.stats_info.as_deref(), sb_response);
}
InfoMetricsType::Store => {
self.populate_store_stats(provider);
Self::get_section_resp_info(
&header,
self
.store_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.map(|v| &v[..]),
sb_response,
);
}
InfoMetricsType::StoreHashtable => {
self.populate_store_hash_distribution(provider);
Self::get_section_resp_info(
&header,
self
.store_hash_distr_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.map(|v| &v[..]),
sb_response,
);
}
InfoMetricsType::StoreReviv => {
self.populate_store_reviv_info(provider);
Self::get_section_resp_info(
&header,
self
.store_reviv_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.map(|v| &v[..]),
sb_response,
);
}
InfoMetricsType::Persistence => {
if !provider.server_facts().enable_aof {
return;
}
self.populate_persistence_info(provider);
Self::get_section_resp_info(
&header,
self
.persistence_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.map(|v| &v[..]),
sb_response,
);
}
InfoMetricsType::Clients => {
self.populate_clients_info(provider);
Self::get_section_resp_info(&header, self.clients_info.as_deref(), sb_response);
}
InfoMetricsType::Keyspace => {
self.populate_keyspace_info(provider);
Self::get_section_resp_info(&header, self.keyspace_info.as_deref(), sb_response);
}
InfoMetricsType::Modules => {
Self::get_section_resp_info(&header, None, sb_response);
}
InfoMetricsType::BpStats => {
self.populate_cluster_buffer_pool_stats(provider);
Self::get_section_resp_info(&header, self.buffer_pool_stats.as_deref(), sb_response);
}
InfoMetricsType::CInfo => {
self.populate_checkpoint_info(provider);
Self::get_section_resp_info(&header, self.checkpoint_stats.as_deref(), sb_response);
}
InfoMetricsType::HlogScan => {
self.populate_hlog_scan_info(provider);
Self::get_section_resp_info(
&header,
self
.hlog_scan_stats
.as_deref()
.and_then(|v| v.get(db_id as usize))
.map(|v| &v[..]),
sb_response,
);
}
InfoMetricsType::CommandStats => {
self.populate_command_stats_info(provider);
Self::get_section_resp_info(&header, self.command_stats_info.as_deref(), sb_response);
}
}
}
pub fn get_resp_info(
&mut self,
sections: &[InfoMetricsType],
db_id: i32,
provider: &dyn InfoProvider,
) -> String {
let mut sb_response = String::new();
for (i, section) in sections.iter().enumerate() {
self.get_resp_info_single(*section, db_id, provider, &mut sb_response);
if i != sections.len() - 1 {
sb_response.push_str("\r\n");
}
}
sb_response
}
fn get_metric_internal(
&mut self,
section: InfoMetricsType,
db_id: i32,
provider: &dyn InfoProvider,
) -> Option<Vec<MetricsItem>> {
match section {
InfoMetricsType::Server => {
self.populate_server_info(provider);
self.server_info.clone()
}
InfoMetricsType::Memory => {
self.populate_memory_info(provider);
self.memory_info.clone()
}
InfoMetricsType::Cluster => {
self.populate_cluster_info(provider);
self.cluster_info.clone()
}
InfoMetricsType::Replication => {
self.populate_replication_info(provider);
self.replication_info.clone()
}
InfoMetricsType::Stats => {
self.populate_stats_info(provider);
self.stats_info.clone()
}
InfoMetricsType::Store => {
self.populate_store_stats(provider);
self
.store_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.cloned()
}
InfoMetricsType::StoreHashtable => {
self.populate_store_hash_distribution(provider);
self
.store_hash_distr_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.cloned()
}
InfoMetricsType::StoreReviv => {
self.populate_store_reviv_info(provider);
self
.store_reviv_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.cloned()
}
InfoMetricsType::Persistence => {
if !provider.server_facts().enable_aof {
return None;
}
self.populate_persistence_info(provider);
self
.persistence_info
.as_deref()
.and_then(|v| v.get(db_id as usize))
.cloned()
}
InfoMetricsType::Clients => {
self.populate_clients_info(provider);
self.clients_info.clone()
}
InfoMetricsType::Keyspace => {
self.populate_keyspace_info(provider);
self.keyspace_info.clone()
}
InfoMetricsType::Modules => None,
InfoMetricsType::CommandStats => {
self.populate_command_stats_info(provider);
self.command_stats_info.clone()
}
_ => None,
}
}
pub fn get_metric(
&mut self,
section: InfoMetricsType,
db_id: i32,
provider: &dyn InfoProvider,
) -> Option<Vec<MetricsItem>> {
self.get_metric_internal(section, db_id, provider)
}
pub fn get_info_metrics(
&mut self,
sections: &[InfoMetricsType],
db_id: i32,
provider: &dyn InfoProvider,
) -> Vec<(InfoMetricsType, Vec<MetricsItem>)> {
sections
.iter()
.filter_map(|§ion| {
self
.get_metric_internal(section, db_id, provider)
.map(|items| (section, items))
})
.collect()
}
}
fn page_size() -> usize {
4096
}
pub fn generate_default_hex_id() -> String {
use std::sync::atomic::{AtomicU64, Ordering};
static STATE: AtomicU64 = AtomicU64::new(0);
let mut state = STATE.fetch_add(1, Ordering::Relaxed);
if state == 0 {
state = coarsetime::Clock::now_since_epoch().as_nanos() | 1;
}
let mut out = String::with_capacity(40);
while out.len() < 40 {
state ^= state >> 12;
state ^= state << 25;
state ^= state >> 27;
let x = state.wrapping_mul(0x2545_F491_4F6C_DD1D);
out.push_str(&format!("{x:016x}"));
}
out.truncate(40);
out
}