use std::collections::HashMap;
use std::path::Path;
use std::time::SystemTime;
use std::time::UNIX_EPOCH;
use parking_lot::Mutex;
use crate::aae::persist::PersistError;
use crate::aae::tictac::Tree;
#[derive(Debug, Default)]
pub struct AaeMetrics {
inner: Mutex<AaeInner>,
}
#[derive(Debug, Default)]
struct AaeInner {
exchange_attempts: HashMap<PeerKey, u64>,
exchange_success: HashMap<PeerKey, u64>,
divergent_keys: HashMap<PeerKey, u64>,
repair_dispatched: HashMap<PeerKey, u64>,
segments_dirty: HashMap<u32, u64>,
full_sweep_last_completed_unix: HashMap<u32, u64>,
snapshot_save_total: u64,
snapshot_load_total: u64,
snapshot_corruption_total: u64,
}
#[derive(Debug, Eq, PartialEq, Hash, Clone)]
struct PeerKey {
peer_idx: u32,
dc: String,
rack: String,
}
impl AaeMetrics {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn record_exchange_attempt(&self, peer_idx: u32, dc: &str, rack: &str) {
let key = peer_key(peer_idx, dc, rack);
let mut inner = self.inner.lock();
*inner.exchange_attempts.entry(key).or_insert(0) += 1;
}
pub fn record_exchange_success(&self, peer_idx: u32, dc: &str, rack: &str) {
let key = peer_key(peer_idx, dc, rack);
let mut inner = self.inner.lock();
*inner.exchange_success.entry(key).or_insert(0) += 1;
}
pub fn record_divergent_keys(&self, peer_idx: u32, dc: &str, rack: &str, count: u64) {
if count == 0 {
return;
}
let key = peer_key(peer_idx, dc, rack);
let mut inner = self.inner.lock();
*inner.divergent_keys.entry(key).or_insert(0) += count;
}
pub fn record_repair_dispatched(&self, peer_idx: u32, dc: &str, rack: &str, count: u64) {
if count == 0 {
return;
}
let key = peer_key(peer_idx, dc, rack);
let mut inner = self.inner.lock();
*inner.repair_dispatched.entry(key).or_insert(0) += count;
}
pub fn set_segments_dirty(&self, peer_idx: u32, count: u64) {
let mut inner = self.inner.lock();
inner.segments_dirty.insert(peer_idx, count);
}
pub fn mark_full_sweep_completed(&self, peer_idx: u32, unix_seconds: u64) {
let mut inner = self.inner.lock();
inner
.full_sweep_last_completed_unix
.insert(peer_idx, unix_seconds);
}
pub fn mark_full_sweep_completed_now(&self, peer_idx: u32) {
let secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
self.mark_full_sweep_completed(peer_idx, secs);
}
pub fn record_snapshot_save(&self) {
self.inner.lock().snapshot_save_total += 1;
}
pub fn record_snapshot_load(&self) {
self.inner.lock().snapshot_load_total += 1;
}
pub fn record_snapshot_corruption(&self) {
self.inner.lock().snapshot_corruption_total += 1;
}
#[must_use]
pub fn snapshot(&self) -> AaeMetricsSnapshot {
let inner = self.inner.lock();
AaeMetricsSnapshot {
exchange_attempts: collect_peer_entries(&inner.exchange_attempts),
exchange_success: collect_peer_entries(&inner.exchange_success),
divergent_keys: collect_peer_entries(&inner.divergent_keys),
repair_dispatched: collect_peer_entries(&inner.repair_dispatched),
segments_dirty: collect_simple_peer_entries(&inner.segments_dirty),
full_sweep_last_completed_unix: collect_simple_peer_entries(
&inner.full_sweep_last_completed_unix,
),
snapshot_save_total: inner.snapshot_save_total,
snapshot_load_total: inner.snapshot_load_total,
snapshot_corruption_total: inner.snapshot_corruption_total,
}
}
}
pub fn save_snapshot_with_metrics(
tree: &Tree,
path: &Path,
metrics: &AaeMetrics,
) -> Result<(), PersistError> {
tree.save_snapshot(path)?;
metrics.record_snapshot_save();
Ok(())
}
pub fn load_snapshot_with_metrics(path: &Path, metrics: &AaeMetrics) -> Result<Tree, PersistError> {
match Tree::load_snapshot(path) {
Ok(t) => {
metrics.record_snapshot_load();
Ok(t)
}
Err(e) => {
match &e {
PersistError::Corrupted(_)
| PersistError::VersionSkew { .. }
| PersistError::BadShape(_) => metrics.record_snapshot_corruption(),
PersistError::Io(_) => {}
}
Err(e)
}
}
}
fn peer_key(peer_idx: u32, dc: &str, rack: &str) -> PeerKey {
PeerKey {
peer_idx,
dc: dc.to_owned(),
rack: rack.to_owned(),
}
}
fn collect_peer_entries(map: &HashMap<PeerKey, u64>) -> Vec<PeerEntry> {
let mut out: Vec<PeerEntry> = map
.iter()
.map(|(k, v)| PeerEntry {
peer_idx: k.peer_idx,
dc: k.dc.clone(),
rack: k.rack.clone(),
count: *v,
})
.collect();
out.sort_by(|a, b| {
a.peer_idx
.cmp(&b.peer_idx)
.then(a.dc.cmp(&b.dc))
.then(a.rack.cmp(&b.rack))
});
out
}
fn collect_simple_peer_entries(map: &HashMap<u32, u64>) -> Vec<SimplePeerEntry> {
let mut out: Vec<SimplePeerEntry> = map
.iter()
.map(|(k, v)| SimplePeerEntry {
peer_idx: *k,
value: *v,
})
.collect();
out.sort_by_key(|e| e.peer_idx);
out
}
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub struct PeerEntry {
pub peer_idx: u32,
pub dc: String,
pub rack: String,
pub count: u64,
}
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub struct SimplePeerEntry {
pub peer_idx: u32,
pub value: u64,
}
#[derive(Clone, Debug, Default)]
pub struct AaeMetricsSnapshot {
pub exchange_attempts: Vec<PeerEntry>,
pub exchange_success: Vec<PeerEntry>,
pub divergent_keys: Vec<PeerEntry>,
pub repair_dispatched: Vec<PeerEntry>,
pub segments_dirty: Vec<SimplePeerEntry>,
pub full_sweep_last_completed_unix: Vec<SimplePeerEntry>,
pub snapshot_save_total: u64,
pub snapshot_load_total: u64,
pub snapshot_corruption_total: u64,
}
impl AaeMetricsSnapshot {
#[must_use]
pub fn is_empty(&self) -> bool {
self.exchange_attempts.is_empty()
&& self.exchange_success.is_empty()
&& self.divergent_keys.is_empty()
&& self.repair_dispatched.is_empty()
&& self.segments_dirty.is_empty()
&& self.full_sweep_last_completed_unix.is_empty()
&& self.snapshot_save_total == 0
&& self.snapshot_load_total == 0
&& self.snapshot_corruption_total == 0
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::aae::tictac::TreeShape;
use tempfile::TempDir;
fn shape() -> TreeShape {
TreeShape {
n_time_buckets: 2,
n_segments: 16,
time_window_seconds: 60,
}
}
#[test]
fn fresh_metrics_snapshot_is_empty() {
let m = AaeMetrics::new();
assert!(m.snapshot().is_empty());
}
#[test]
fn exchange_attempts_keyed_by_peer_dc_rack() {
let m = AaeMetrics::new();
m.record_exchange_attempt(1, "dc1", "rA");
m.record_exchange_attempt(1, "dc1", "rA");
m.record_exchange_attempt(2, "dc1", "rA");
let s = m.snapshot();
assert_eq!(s.exchange_attempts.len(), 2);
assert_eq!(s.exchange_attempts[0].peer_idx, 1);
assert_eq!(s.exchange_attempts[0].count, 2);
assert_eq!(s.exchange_attempts[1].peer_idx, 2);
assert_eq!(s.exchange_attempts[1].count, 1);
}
#[test]
fn divergent_and_repair_zero_values_are_dropped() {
let m = AaeMetrics::new();
m.record_divergent_keys(1, "dc1", "rA", 0);
m.record_repair_dispatched(1, "dc1", "rA", 0);
let s = m.snapshot();
assert!(s.divergent_keys.is_empty());
assert!(s.repair_dispatched.is_empty());
}
#[test]
fn segments_dirty_gauge_is_set_by_peer() {
let m = AaeMetrics::new();
m.set_segments_dirty(1, 17);
m.set_segments_dirty(2, 4);
m.set_segments_dirty(1, 5); let s = m.snapshot();
assert_eq!(s.segments_dirty.len(), 2);
let one = s.segments_dirty.iter().find(|e| e.peer_idx == 1).unwrap();
assert_eq!(one.value, 5);
}
#[test]
fn full_sweep_completed_records_unix_seconds() {
let m = AaeMetrics::new();
m.mark_full_sweep_completed(7, 1_700_000_000);
let s = m.snapshot();
assert_eq!(s.full_sweep_last_completed_unix.len(), 1);
assert_eq!(s.full_sweep_last_completed_unix[0].peer_idx, 7);
assert_eq!(s.full_sweep_last_completed_unix[0].value, 1_700_000_000);
}
#[test]
fn snapshot_counters_increment() {
let m = AaeMetrics::new();
m.record_snapshot_save();
m.record_snapshot_save();
m.record_snapshot_load();
m.record_snapshot_corruption();
let s = m.snapshot();
assert_eq!(s.snapshot_save_total, 2);
assert_eq!(s.snapshot_load_total, 1);
assert_eq!(s.snapshot_corruption_total, 1);
}
#[test]
fn save_with_metrics_bumps_save_total() {
let m = AaeMetrics::new();
let dir = TempDir::new().unwrap();
let path = dir.path().join("aae.snap");
let tree = Tree::new(shape());
save_snapshot_with_metrics(&tree, &path, &m).unwrap();
assert_eq!(m.snapshot().snapshot_save_total, 1);
}
#[test]
fn load_with_metrics_bumps_load_total() {
let m = AaeMetrics::new();
let dir = TempDir::new().unwrap();
let path = dir.path().join("aae.snap");
let tree = Tree::new(shape());
tree.save_snapshot(&path).unwrap();
let _ = load_snapshot_with_metrics(&path, &m).unwrap();
assert_eq!(m.snapshot().snapshot_load_total, 1);
assert_eq!(m.snapshot().snapshot_corruption_total, 0);
}
#[test]
fn load_with_metrics_corrupted_bumps_corruption_total() {
use std::fs;
let m = AaeMetrics::new();
let dir = TempDir::new().unwrap();
let path = dir.path().join("aae.snap");
fs::write(&path, b"not a snapshot").unwrap();
let err = load_snapshot_with_metrics(&path, &m).unwrap_err();
assert!(matches!(err, PersistError::Corrupted(_)));
assert_eq!(m.snapshot().snapshot_corruption_total, 1);
assert_eq!(m.snapshot().snapshot_load_total, 0);
}
#[test]
fn load_with_metrics_missing_file_does_not_count_as_corruption() {
let m = AaeMetrics::new();
let dir = TempDir::new().unwrap();
let path = dir.path().join("nope.snap");
let err = load_snapshot_with_metrics(&path, &m).unwrap_err();
assert!(matches!(err, PersistError::Io(_)));
assert_eq!(m.snapshot().snapshot_corruption_total, 0);
assert_eq!(m.snapshot().snapshot_load_total, 0);
}
}