use std::sync::atomic::{AtomicI64, Ordering};
use crate::proto::tero::policy::v1::VolumeStats;
#[derive(Debug, Default)]
pub struct VolumeTracker {
log_records: AtomicI64,
log_bytes: AtomicI64,
metric_data_points: AtomicI64,
metric_bytes: AtomicI64,
spans: AtomicI64,
span_bytes: AtomicI64,
}
impl VolumeTracker {
pub fn new() -> Self {
Self::default()
}
pub(crate) fn record_log(&self) {
self.log_records.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_metric(&self) {
self.metric_data_points.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_span(&self) {
self.spans.fetch_add(1, Ordering::Relaxed);
}
pub fn add_log_bytes(&self, bytes: i64) {
self.log_bytes.fetch_add(bytes, Ordering::Relaxed);
}
pub fn add_metric_bytes(&self, bytes: i64) {
self.metric_bytes.fetch_add(bytes, Ordering::Relaxed);
}
pub fn add_span_bytes(&self, bytes: i64) {
self.span_bytes.fetch_add(bytes, Ordering::Relaxed);
}
pub fn collect(&self) -> Option<VolumeStats> {
let stats = VolumeStats {
log_records: self.log_records.swap(0, Ordering::Relaxed),
log_bytes: self.log_bytes.swap(0, Ordering::Relaxed),
metric_data_points: self.metric_data_points.swap(0, Ordering::Relaxed),
metric_bytes: self.metric_bytes.swap(0, Ordering::Relaxed),
spans: self.spans.swap(0, Ordering::Relaxed),
span_bytes: self.span_bytes.swap(0, Ordering::Relaxed),
};
(stats != VolumeStats::default()).then_some(stats)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn empty_tracker_reports_nothing() {
assert!(VolumeTracker::new().collect().is_none());
}
#[test]
fn counts_records_per_signal() {
let t = VolumeTracker::new();
t.record_log();
t.record_log();
t.record_metric();
t.record_span();
let stats = t.collect().unwrap();
assert_eq!(stats.log_records, 2);
assert_eq!(stats.metric_data_points, 1);
assert_eq!(stats.spans, 1);
assert_eq!(stats.log_bytes, 0, "bytes are opt-in");
}
#[test]
fn bytes_are_opt_in() {
let t = VolumeTracker::new();
t.add_log_bytes(120);
t.add_log_bytes(30);
t.add_metric_bytes(7);
t.add_span_bytes(9);
let stats = t.collect().unwrap();
assert_eq!(stats.log_bytes, 150);
assert_eq!(stats.metric_bytes, 7);
assert_eq!(stats.span_bytes, 9);
assert_eq!(stats.log_records, 0);
}
#[test]
fn collect_drains_so_deltas_never_repeat() {
let t = VolumeTracker::new();
t.record_log();
assert_eq!(t.collect().unwrap().log_records, 1);
assert!(t.collect().is_none());
}
#[test]
fn concurrent_collects_do_not_double_count() {
use std::sync::Arc;
use std::thread;
const RECORDS: i64 = 500;
let tracker = Arc::new(VolumeTracker::new());
for _ in 0..RECORDS {
tracker.record_log();
}
let drainers: Vec<_> = (0..8)
.map(|_| {
let tracker = Arc::clone(&tracker);
thread::spawn(move || tracker.collect().map(|s| s.log_records).unwrap_or_default())
})
.collect();
let reported: i64 = drainers.into_iter().map(|h| h.join().unwrap()).sum();
let left = tracker.collect().map(|s| s.log_records).unwrap_or_default();
assert_eq!(reported + left, RECORDS);
}
}