#![expect(
clippy::let_underscore_must_use,
reason = "best effort, with the real outcome reported elsewhere"
)]
use std::fs::{File, OpenOptions};
use std::io::Write;
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering::Relaxed};
use std::sync::{Arc, Mutex};
use std::time::{Instant, SystemTime, UNIX_EPOCH};
#[derive(Debug, Default)]
pub(crate) struct ShardStats {
pub used_memory: AtomicU64,
pub used_memory_peak: AtomicU64,
pub keys: AtomicU64,
pub expires: AtomicU64,
pub expired_keys: AtomicU64,
pub evicted_keys: AtomicU64,
pub commands_processed: AtomicU64,
pub connections_received: AtomicU64,
pub clients_connected: AtomicU64,
pub blocked_clients: AtomicU64,
pub tick_gap_max_us: AtomicU64,
pub query_buffer_disconnections: AtomicU64,
pub ticks_total: AtomicU64,
pub tier: TierGauges,
pub alloc: AllocGauges,
}
#[derive(Debug, Default)]
pub(crate) struct TierGauges {
pub enabled: AtomicU64,
pub budget: AtomicU64,
pub effective_target: AtomicU64,
pub reserved_bytes: AtomicU64,
pub stub_bytes: AtomicU64,
pub cold_keys: AtomicU64,
pub cold_bytes: AtomicU64,
pub demotions_total: AtomicU64,
pub promotions_total: AtomicU64,
pub peek_preads_total: AtomicU64,
pub batch_submissions_total: AtomicU64,
pub vlog_files: AtomicU64,
pub vlog_bytes: AtomicU64,
pub vlog_live_bytes: AtomicU64,
pub vlog_epoch: AtomicU64,
}
#[derive(Debug, Default)]
pub(crate) struct AllocGauges {
pub reporting: AtomicU64,
pub mapped: AtomicU64,
pub live: AtomicU64,
pub rounding: AtomicU64,
pub cache: AtomicU64,
pub span_free: AtomicU64,
pub returned: AtomicU64,
pub virgin: AtomicU64,
pub hysteresis: AtomicU64,
pub segment_overhead: AtomicU64,
pub large_count: AtomicU64,
pub spans_assigned: AtomicU64,
}
#[derive(Default)]
pub(crate) struct Totals {
pub used_memory: u64,
pub used_memory_peak: u64,
pub keys: u64,
pub expires: u64,
pub expired_keys: u64,
pub evicted_keys: u64,
pub commands_processed: u64,
pub connections_received: u64,
pub clients_connected: u64,
pub blocked_clients: u64,
pub tick_gap_max_us: u64,
pub query_buffer_disconnections: u64,
pub ticks_total: u64,
pub tier_enabled: bool,
pub tier: TierTotals,
pub alloc_shards: u64,
pub alloc: AllocTotals,
}
#[derive(Default)]
pub(crate) struct TierTotals {
pub budget: u64,
pub effective_target: u64,
pub reserved_bytes: u64,
pub stub_bytes: u64,
pub cold_keys: u64,
pub cold_bytes: u64,
pub demotions_total: u64,
pub promotions_total: u64,
pub peek_preads_total: u64,
pub batch_submissions_total: u64,
pub vlog_files: u64,
pub vlog_bytes: u64,
pub vlog_live_bytes: u64,
pub vlog_epoch: u64,
}
#[derive(Default)]
pub(crate) struct AllocTotals {
pub mapped: u64,
pub live: u64,
pub rounding: u64,
pub cache: u64,
pub span_free: u64,
pub returned: u64,
pub virgin: u64,
pub hysteresis: u64,
pub segment_overhead: u64,
pub large_count: u64,
pub spans_assigned: u64,
}
impl AllocTotals {
fn add(&mut self, g: &AllocGauges) {
self.mapped += g.mapped.load(Relaxed);
self.live += g.live.load(Relaxed);
self.rounding += g.rounding.load(Relaxed);
self.cache += g.cache.load(Relaxed);
self.span_free += g.span_free.load(Relaxed);
self.returned += g.returned.load(Relaxed);
self.virgin += g.virgin.load(Relaxed);
self.hysteresis += g.hysteresis.load(Relaxed);
self.segment_overhead += g.segment_overhead.load(Relaxed);
self.large_count += g.large_count.load(Relaxed);
self.spans_assigned += g.spans_assigned.load(Relaxed);
}
pub fn accounted(&self) -> u64 {
self.live
+ self.rounding
+ self.cache
+ self.span_free
+ self.returned
+ self.virgin
+ self.hysteresis
+ self.segment_overhead
}
}
const OPS_WINDOW: usize = 16;
#[derive(Debug)]
pub(crate) struct ObsState {
audit: Option<Mutex<File>>,
shard_stats: Box<[Arc<ShardStats>]>,
ops_ring: Mutex<Vec<(u128, u64)>>,
repl_views: Box<[Mutex<ReplShardView>]>,
start: Instant,
}
#[derive(Debug, Clone, Default)]
pub(crate) struct ReplShardView {
pub(crate) offset: u64,
pub(crate) replicas: Vec<kevy_rt::ReplicaViewRow>,
}
impl ObsState {
pub(crate) fn new(audit_log_path: &Path, nshards: usize) -> Self {
Self {
audit: open_audit_log(audit_log_path),
shard_stats: (0..nshards.max(1)).map(|_| Arc::new(ShardStats::default())).collect(),
ops_ring: Mutex::new(Vec::new()),
repl_views: (0..nshards.max(1)).map(|_| Mutex::new(ReplShardView::default())).collect(),
start: Instant::now(),
}
}
pub(crate) fn publish_repl_view(&self, shard: usize, view: ReplShardView) {
if let Some(slot) = self.repl_views.get(shard) {
*slot.lock().expect("repl_views poisoned") = view;
}
}
pub(crate) fn repl_views(&self) -> Vec<ReplShardView> {
self.repl_views.iter().map(|s| s.lock().expect("repl_views poisoned").clone()).collect()
}
pub(crate) fn slot(&self, shard: usize) -> Option<Arc<ShardStats>> {
self.shard_stats.get(shard).cloned()
}
pub(crate) fn aggregate(&self) -> Totals {
let mut t = Totals::default();
for s in &self.shard_stats {
t.used_memory += s.used_memory.load(Relaxed);
t.used_memory_peak += s.used_memory_peak.load(Relaxed);
t.keys += s.keys.load(Relaxed);
t.expires += s.expires.load(Relaxed);
t.expired_keys += s.expired_keys.load(Relaxed);
t.evicted_keys += s.evicted_keys.load(Relaxed);
t.commands_processed += s.commands_processed.load(Relaxed);
t.connections_received += s.connections_received.load(Relaxed);
t.clients_connected += s.clients_connected.load(Relaxed);
t.blocked_clients += s.blocked_clients.load(Relaxed);
t.tick_gap_max_us = t.tick_gap_max_us.max(s.tick_gap_max_us.load(Relaxed));
t.ticks_total += s.ticks_total.load(Relaxed);
t.query_buffer_disconnections += s.query_buffer_disconnections.load(Relaxed);
t.tier_enabled |= s.tier.enabled.load(Relaxed) != 0;
t.tier.budget += s.tier.budget.load(Relaxed);
t.tier.effective_target += s.tier.effective_target.load(Relaxed);
t.tier.reserved_bytes += s.tier.reserved_bytes.load(Relaxed);
t.tier.stub_bytes += s.tier.stub_bytes.load(Relaxed);
t.tier.cold_keys += s.tier.cold_keys.load(Relaxed);
t.tier.cold_bytes += s.tier.cold_bytes.load(Relaxed);
t.tier.demotions_total += s.tier.demotions_total.load(Relaxed);
t.tier.promotions_total += s.tier.promotions_total.load(Relaxed);
t.tier.peek_preads_total += s.tier.peek_preads_total.load(Relaxed);
t.tier.batch_submissions_total += s.tier.batch_submissions_total.load(Relaxed);
t.tier.vlog_files += s.tier.vlog_files.load(Relaxed);
t.tier.vlog_bytes += s.tier.vlog_bytes.load(Relaxed);
t.tier.vlog_live_bytes += s.tier.vlog_live_bytes.load(Relaxed);
t.tier.vlog_epoch += s.tier.vlog_epoch.load(Relaxed);
if s.alloc.reporting.load(Relaxed) != 0 {
t.alloc_shards += 1;
t.alloc.add(&s.alloc);
}
}
t
}
pub(crate) fn push_ops_sample(&self, total_commands: u64) {
let mut ring = self.ops_ring.lock().expect("ops_ring poisoned");
ring.push((self.elapsed_ms(), total_commands));
if ring.len() > OPS_WINDOW {
let drop = ring.len() - OPS_WINDOW;
ring.drain(0..drop);
}
}
pub(crate) fn instantaneous_ops_per_sec(&self, current: u64) -> u64 {
let ring = self.ops_ring.lock().expect("ops_ring poisoned");
let Some(&(oldest_ms, oldest_cmds)) = ring.first() else {
return 0;
};
let dt_ms = self.elapsed_ms().saturating_sub(oldest_ms);
if dt_ms == 0 {
return 0;
}
let dc = current.saturating_sub(oldest_cmds);
((u128::from(dc) * 1000) / dt_ms) as u64
}
fn elapsed_ms(&self) -> u128 {
self.start.elapsed().as_millis()
}
pub(crate) fn audit_record(&self, args: &[&[u8]]) {
let Some(mu) = &self.audit else { return };
let mut line = String::with_capacity(128);
let micros =
SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_micros()).unwrap_or(0);
line.push_str(µs.to_string());
for arg in args {
line.push('\t');
let s = String::from_utf8_lossy(&arg[..arg.len().min(256)]);
for c in s.chars() {
match c {
'\t' | '\n' | '\r' => line.push(' '),
_ => line.push(c),
}
}
if arg.len() > 256 {
line.push('…');
}
}
line.push('\n');
if let Ok(mut f) = mu.lock() {
let _ = f.write_all(line.as_bytes());
let _ = f.flush();
}
}
}
fn open_audit_log(path: &Path) -> Option<Mutex<File>> {
if path.as_os_str().is_empty() {
return None;
}
match OpenOptions::new().create(true).append(true).open(path) {
Ok(f) => Some(Mutex::new(f)),
Err(e) => {
eprintln!("kevy: audit log {} could not open: {e}", path.display());
None
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
#[test]
fn audit_off_when_path_empty() {
let obs = ObsState::new(Path::new(""), 1);
obs.audit_record(&[b"CONFIG", b"SET", b"maxmemory", b"1g"]);
assert!(obs.audit.is_none());
}
#[test]
fn audit_records_one_sanitised_line() {
let nanos = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
let path: PathBuf = std::env::temp_dir().join(format!("kevy-audit-{nanos}"));
let obs = ObsState::new(&path, 1);
obs.audit_record(&[b"DEBUG", b"tab\there"]);
let text = std::fs::read_to_string(&path).unwrap();
let _ = std::fs::remove_file(&path);
assert_eq!(text.lines().count(), 1);
assert!(text.contains("DEBUG\ttab here"), "got: {text:?}");
}
#[test]
fn aggregate_sums_across_slots() {
let obs = ObsState::new(Path::new(""), 3);
obs.shard_stats[0].keys.store(2, Relaxed);
obs.shard_stats[2].keys.store(5, Relaxed);
assert_eq!(obs.aggregate().keys, 7);
}
#[test]
fn ops_ring_caps_at_window() {
let obs = ObsState::new(Path::new(""), 1);
for i in 0..(OPS_WINDOW as u64 + 10) {
obs.push_ops_sample(i);
}
assert_eq!(obs.ops_ring.lock().unwrap().len(), OPS_WINDOW);
}
#[test]
fn only_shards_that_reported_are_folded_in() {
let obs = ObsState::new(Path::new(""), 3);
let a = obs.slot(0).expect("slot 0");
a.alloc.mapped.store(8_388_608, Relaxed);
a.alloc.live.store(700_000, Relaxed);
a.alloc.hysteresis.store(7_688_608, Relaxed);
a.alloc.reporting.store(1, Relaxed);
let b = obs.slot(1).expect("slot 1");
b.alloc.mapped.store(4_194_304, Relaxed);
b.alloc.live.store(100_000, Relaxed);
b.alloc.hysteresis.store(4_094_304, Relaxed);
b.alloc.reporting.store(1, Relaxed);
let t = obs.aggregate();
assert_eq!(t.alloc_shards, 2, "a silent shard was counted as a reporting one");
assert_eq!(t.alloc.mapped, 12_582_912);
assert_eq!(t.alloc.live, 800_000);
assert_eq!(t.alloc.accounted(), t.alloc.mapped, "the summed terms stopped partitioning");
}
}