use std::time::SystemTime;
use kevy_config::Config;
use super::{appendfsync_str, eviction_str, memory, replication};
use crate::state::Ctx;
pub(super) fn info_server(cfg: &Config, b: &mut String) {
let now = SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).map_or(0, |d| d.as_secs());
b.push_str("# Server\r\n");
b.push_str("redis_version:7.4.0\r\n"); b.push_str(&format!("kevy_version:{}\r\n", env!("CARGO_PKG_VERSION")));
b.push_str("redis_mode:standalone\r\n");
b.push_str(&format!("process_id:{}\r\n", std::process::id()));
b.push_str(&format!("tcp_port:{}\r\n", cfg.server.port));
b.push_str(&format!("server_time_usec:{}\r\n", now * 1_000_000));
b.push_str("\r\n");
}
pub(super) fn info_clients(cfg: &Config, totals: &crate::state::Totals, b: &mut String) {
b.push_str("# Clients\r\n");
b.push_str(&format!("connected_clients:{}\r\n", totals.clients_connected));
b.push_str(&format!("blocked_clients:{}\r\n", totals.blocked_clients));
b.push_str(&format!("maxclients:{}\r\n", cfg.server.max_clients));
b.push_str("\r\n");
}
pub(super) fn info_memory(cfg: &Config, totals: &crate::state::Totals, b: &mut String) {
let used = totals.used_memory;
let peak = totals.used_memory_peak;
b.push_str("# Memory\r\n");
b.push_str(&format!("used_memory:{used}\r\n"));
b.push_str(&format!("used_memory_human:{}\r\n", memory::format_bytes_human(used)));
b.push_str(&format!("used_memory_peak:{peak}\r\n"));
b.push_str(&format!("used_memory_peak_human:{}\r\n", memory::format_bytes_human(peak)));
b.push_str(&format!("maxmemory:{}\r\n", cfg.memory.maxmemory));
b.push_str(&format!(
"maxmemory_human:{}\r\n",
memory::format_bytes_human(cfg.memory.maxmemory)
));
b.push_str(&format!("maxmemory_policy:{}\r\n", eviction_str(cfg.memory.maxmemory_policy)));
b.push_str(&format!("evicted_keys:{}\r\n", totals.evicted_keys));
b.push_str(&format!("process_rss_bytes:{}\r\n", kevy_sys::process_rss_bytes()));
b.push_str("\r\n");
}
pub(super) fn info_tiering(totals: &crate::state::Totals, b: &mut String) {
let t = &totals.tier;
b.push_str("# Tiering\r\n");
b.push_str("tiering_enabled:1\r\n");
b.push_str(&format!("tier_budget_bytes:{}\r\n", t.budget));
b.push_str(&format!("tier_effective_target:{}\r\n", t.effective_target));
b.push_str(&format!("cold_keys:{}\r\n", t.cold_keys));
b.push_str(&format!("cold_bytes:{}\r\n", t.cold_bytes));
b.push_str(&format!("stub_bytes:{}\r\n", t.stub_bytes));
b.push_str(&format!("index_reserved_bytes:{}\r\n", t.reserved_bytes));
b.push_str(&format!("vlog_size_bytes:{}\r\n", t.vlog_bytes));
b.push_str(&format!("vlog_live_bytes:{}\r\n", t.vlog_live_bytes));
b.push_str(&format!("vlog_files:{}\r\n", t.vlog_files));
b.push_str(&format!("vlog_epoch:{}\r\n", t.vlog_epoch));
b.push_str(&format!("demotions_total:{}\r\n", t.demotions_total));
b.push_str(&format!("promotions_total:{}\r\n", t.promotions_total));
b.push_str(&format!("peek_preads_total:{}\r\n", t.peek_preads_total));
b.push_str(&format!("batch_submissions_total:{}\r\n", t.batch_submissions_total));
b.push_str("\r\n");
}
pub(super) fn info_allocator(totals: &crate::state::Totals, b: &mut String) {
if totals.alloc_shards == 0 {
return;
}
let a = &totals.alloc;
b.push_str("# Allocator\r\n");
b.push_str("allocator_impl:kevy-alloc\r\n");
b.push_str(&format!("alloc_shards_reporting:{}\r\n", totals.alloc_shards));
b.push_str(&format!("alloc_mapped:{}\r\n", a.mapped));
b.push_str(&format!("alloc_accounted:{}\r\n", a.accounted()));
b.push_str(&format!("alloc_live:{}\r\n", a.live));
b.push_str(&format!("alloc_rounding:{}\r\n", a.rounding));
b.push_str(&format!("alloc_cache:{}\r\n", a.cache));
b.push_str(&format!("alloc_span_free:{}\r\n", a.span_free));
b.push_str(&format!("alloc_returned:{}\r\n", a.returned));
b.push_str(&format!("alloc_virgin:{}\r\n", a.virgin));
b.push_str(&format!("alloc_hysteresis:{}\r\n", a.hysteresis));
b.push_str(&format!("alloc_segment_overhead:{}\r\n", a.segment_overhead));
b.push_str(&format!("alloc_large_count:{}\r\n", a.large_count));
b.push_str(&format!("alloc_spans_assigned:{}\r\n", a.spans_assigned));
b.push_str("\r\n");
}
pub(super) fn info_persistence(ctx: &Ctx<'_>, cfg: &Config, b: &mut String) {
let (in_flight, rewrites) = ctx.shard.persist_stats();
b.push_str("# Persistence\r\n");
b.push_str(&format!("loading:{}\r\n", i32::from(ctx.state.replication.loading())));
b.push_str(&format!("aof_enabled:{}\r\n", i32::from(cfg.persistence.aof)));
b.push_str(&format!("appendfsync:{}\r\n", appendfsync_str(cfg.persistence.appendfsync)));
b.push_str(&format!("aof_rewrite_in_progress:{}\r\n", i32::from(in_flight)));
b.push_str(&format!("aof_rewrites_total:{rewrites}\r\n"));
b.push_str(&format!(
"aof_format:{}\r\n",
match ctx.shard.aof_format() {
1 => "v1",
2 => "v2",
_ => "off",
}
));
b.push_str("aof_last_rewrite_time_sec:-1\r\n");
info_replay_verdict(ctx, b);
b.push_str("\r\n");
}
fn info_replay_verdict(ctx: &Ctx<'_>, b: &mut String) {
let (dropped, corrupt) = ctx.shard.replay_report();
b.push_str(&format!("aof_last_open_dropped_bytes:{dropped}\r\n"));
b.push_str(&format!("aof_last_open_corrupt:{}\r\n", i32::from(corrupt)));
}
pub(super) fn info_stats(ctx: &Ctx<'_>, totals: &crate::state::Totals, b: &mut String) {
b.push_str("# Stats\r\n");
b.push_str(&format!("total_connections_received:{}\r\n", totals.connections_received));
b.push_str(&format!("total_commands_processed:{}\r\n", totals.commands_processed));
b.push_str(&format!(
"instantaneous_ops_per_sec:{}\r\n",
ctx.state.obs.instantaneous_ops_per_sec(totals.commands_processed)
));
b.push_str(&format!("expired_keys:{}\r\n", totals.expired_keys));
b.push_str(&format!("evicted_keys:{}\r\n", totals.evicted_keys));
b.push_str(&format!(
"client_query_buffer_limit_disconnections:{}\r\n",
totals.query_buffer_disconnections
));
b.push_str(&format!("reactor_tick_gap_max_us:{}\r\n", totals.tick_gap_max_us));
b.push_str(&format!("reactor_ticks_total:{}\r\n", totals.ticks_total));
b.push_str("\r\n");
}
pub(super) fn info_replication(ctx: &Ctx<'_>, b: &mut String) {
b.push_str("# Replication\r\n");
match ctx.state.replication.current_upstream() {
Some((host, port)) => info_repl_replica(ctx, b, host, port),
None => info_repl_master(ctx, b),
}
b.push_str("\r\n");
}
pub(super) fn info_repl_replica(ctx: &Ctx<'_>, b: &mut String, host: std::net::IpAddr, port: u16) {
b.push_str("role:slave\r\n");
b.push_str(&format!("master_host:{host}\r\n"));
b.push_str(&format!("master_port:{port}\r\n"));
let (up, applied, lag, last_io) = ctx.state.replication.replica_link_view();
b.push_str(if up { "master_link_status:up\r\n" } else { "master_link_status:down\r\n" });
b.push_str(&format!("master_last_io_seconds_ago:{last_io}\r\n"));
b.push_str("master_sync_in_progress:0\r\n");
b.push_str(if ctx.state.replication.read_only() {
"slave_read_only:1\r\n"
} else {
"slave_read_only:0\r\n"
});
b.push_str(&format!("slave_repl_offset:{applied}\r\n"));
b.push_str(&format!("slave_lag_frames:{lag}\r\n"));
}
pub(super) fn info_repl_master(ctx: &Ctx<'_>, b: &mut String) {
let views = ctx.state.obs.repl_views();
let (offset, replicas) = replication::aggregate_replication(&views);
b.push_str("role:master\r\n");
b.push_str(&format!("connected_slaves:{}\r\n", replicas.len()));
for (i, agg) in replicas.iter().enumerate() {
let acked_v = agg.acked.unwrap_or(0);
let lag = offset.saturating_sub(acked_v);
let state = if agg.acked.is_some() { "online" } else { "syncing" };
b.push_str(&format!(
"slave{i}:ip={},port={},state={state},offset={acked_v},sent={},lag={lag}\r\n",
agg.ip, agg.port, agg.sent,
));
}
b.push_str("master_replid:0000000000000000000000000000000000000000\r\n");
b.push_str(&format!("master_repl_offset:{offset}\r\n"));
}
pub(super) fn info_cluster(cfg: &Config, b: &mut String) {
b.push_str("# Cluster\r\n");
b.push_str(if cfg.cluster.enabled { "cluster_enabled:1\r\n" } else { "cluster_enabled:0\r\n" });
b.push_str("\r\n");
}
pub(super) fn info_modules(totals: &crate::state::Totals, b: &mut String) {
b.push_str("# Modules\r\n");
b.push_str(if cfg!(feature = "kevy-alloc") {
"module:name=alloc,impl=kevy-alloc\r\n"
} else {
"module:name=alloc,impl=system\r\n"
});
b.push_str(if totals.tier_enabled {
"module:name=tiering,status=on\r\n"
} else {
"module:name=tiering,status=off\r\n"
});
for name in ["indexes", "tables", "views", "text", "vector", "cdc", "pubsub", "lua"] {
b.push_str(&format!("module:name={name},status=on\r\n"));
}
b.push_str("\r\n");
}
pub(super) fn info_keyspace(totals: &crate::state::Totals, b: &mut String) {
b.push_str("# Keyspace\r\n");
if totals.keys > 0 {
b.push_str(&format!("db0:keys={},expires={},avg_ttl=0\r\n", totals.keys, totals.expires));
}
b.push_str("\r\n");
}
#[cfg(test)]
mod allocator_section_tests {
use super::info_allocator;
use crate::state::Totals;
#[test]
fn no_section_is_written_when_no_shard_reported() {
let mut quiet = String::new();
info_allocator(&Totals::default(), &mut quiet);
assert!(quiet.is_empty(), "wrote a section for zero heaps: {quiet:?}");
}
fn two_reporting_shards() -> Totals {
let mut t = Totals { alloc_shards: 2, ..Totals::default() };
t.alloc.mapped = 16_777_216;
t.alloc.live = 152_610;
t.alloc.rounding = 3_150;
t.alloc.span_free = 972_720;
t.alloc.virgin = 444_384;
t.alloc.hysteresis = 14_742_208;
t.alloc.segment_overhead = 262_144;
t.alloc.returned = 200_000;
assert_eq!(t.alloc.accounted(), t.alloc.mapped, "the fixture must balance first");
t
}
#[test]
fn the_section_names_every_term_it_summed() {
let t = two_reporting_shards();
let mut out = String::new();
info_allocator(&t, &mut out);
assert!(out.starts_with("# Allocator\r\n"), "{out}");
assert!(out.ends_with("\r\n\r\n"), "a section ends with a blank line: {out:?}");
for line in [
"alloc_shards_reporting:2\r\n",
"alloc_mapped:16777216\r\n",
"alloc_accounted:16777216\r\n",
"alloc_hysteresis:14742208\r\n",
"alloc_returned:200000\r\n",
] {
assert!(out.contains(line), "missing {line:?} in:\n{out}");
}
}
#[test]
fn the_printed_terms_sum_to_the_printed_map() {
let t = two_reporting_shards();
let mut out = String::new();
info_allocator(&t, &mut out);
let sum: u64 = [
"live",
"rounding",
"cache",
"span_free",
"returned",
"virgin",
"hysteresis",
"segment_overhead",
]
.iter()
.map(|k| {
out.lines()
.find_map(|l| l.trim_end().strip_prefix(&format!("alloc_{k}:")))
.unwrap_or_else(|| panic!("no alloc_{k} in:\n{out}"))
.parse::<u64>()
.unwrap()
})
.sum();
assert_eq!(sum, t.alloc.mapped, "the printed terms do not sum to the printed map");
}
}