use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Limits {
pub max_depth: usize,
pub max_key_bytes: usize,
pub max_scalar_bytes: usize,
pub max_metadata_bytes: usize,
pub max_columns: usize,
pub max_record_bytes: usize,
pub max_capture_bytes: usize,
pub max_output_bytes: Option<u64>,
}
impl Default for Limits {
fn default() -> Self {
Limits {
max_depth: 256,
max_key_bytes: 64 * 1024,
max_scalar_bytes: 16 * 1024 * 1024,
max_metadata_bytes: 16 * 1024 * 1024,
max_columns: 10_000,
max_record_bytes: 64 * 1024 * 1024,
max_capture_bytes: 64 * 1024 * 1024,
max_output_bytes: None,
}
}
}
impl Limits {
pub fn unlimited() -> Self {
Limits {
max_depth: usize::MAX,
max_key_bytes: usize::MAX,
max_scalar_bytes: usize::MAX,
max_metadata_bytes: usize::MAX,
max_columns: usize::MAX,
max_record_bytes: usize::MAX,
max_capture_bytes: usize::MAX,
max_output_bytes: None,
}
}
}
pub const NODE_BYTES: usize = 16;
#[derive(Debug, Default)]
pub struct Metrics {
pub events: AtomicU64,
pub keys: AtomicU64,
pub scalars: AtomicU64,
pub rows: AtomicU64,
pub captured_bytes: AtomicU64,
pub captured_bytes_high: AtomicU64,
pub retained_bytes_high: AtomicU64,
pub output_bytes: AtomicU64,
}
impl Metrics {
pub fn new() -> Arc<Metrics> {
Arc::new(Metrics::default())
}
#[inline]
pub fn add(counter: &AtomicU64, n: u64) {
counter.fetch_add(n, Ordering::Relaxed);
}
#[inline]
pub fn raise(high: &AtomicU64, value: u64) {
high.fetch_max(value, Ordering::Relaxed);
}
pub fn capture(&self, bytes: u64) {
let now = self.captured_bytes.fetch_add(bytes, Ordering::Relaxed) + bytes;
self.captured_bytes_high.fetch_max(now, Ordering::Relaxed);
self.retained_bytes_high.fetch_max(now, Ordering::Relaxed);
}
pub fn release(&self, bytes: u64) {
self.captured_bytes.fetch_sub(bytes, Ordering::Relaxed);
}
pub fn get(counter: &AtomicU64) -> u64 {
counter.load(Ordering::Relaxed)
}
pub fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"events": Metrics::get(&self.events),
"keys": Metrics::get(&self.keys),
"scalars": Metrics::get(&self.scalars),
"rows": Metrics::get(&self.rows),
"captured_bytes": Metrics::get(&self.captured_bytes),
"captured_bytes_high": Metrics::get(&self.captured_bytes_high),
"retained_bytes_high": Metrics::get(&self.retained_bytes_high),
"output_bytes": Metrics::get(&self.output_bytes),
})
}
}
#[derive(Clone, Debug, Default)]
pub struct AbortFlag(Arc<AtomicBool>);
impl AbortFlag {
pub fn new() -> Self {
AbortFlag::default()
}
pub fn abort(&self) {
self.0.store(true, Ordering::Relaxed);
}
pub fn is_aborted(&self) -> bool {
self.0.load(Ordering::Relaxed)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn capture_tracks_high_water() {
let m = Metrics::default();
m.capture(10);
m.capture(20);
m.release(10);
m.capture(5);
assert_eq!(Metrics::get(&m.captured_bytes), 25);
assert_eq!(Metrics::get(&m.captured_bytes_high), 30);
assert_eq!(Metrics::get(&m.retained_bytes_high), 30);
assert_eq!(m.to_json()["captured_bytes_high"], 30);
}
#[test]
fn abort_is_shared() {
let a = AbortFlag::new();
let b = a.clone();
assert!(!b.is_aborted());
a.abort();
assert!(b.is_aborted());
}
}