use std::sync::OnceLock;
pub use histogram::{Bucket, Config, Error, Histogram};
use metriken_core::Window;
use parking_lot::RwLock;
use crate::{HistogramMetric, Metric, Value};
pub struct AtomicHistogram {
inner: OnceLock<histogram::AtomicHistogram>,
config: Config,
}
impl AtomicHistogram {
pub const fn new(grouping_power: u8, max_value_power: u8) -> Self {
let config = match ::histogram::Config::new(grouping_power, max_value_power) {
Ok(c) => c,
Err(_) => panic!("invalid histogram config"),
};
Self {
inner: OnceLock::new(),
config,
}
}
pub fn increment(&self, value: u64) -> Result<(), Error> {
self.get_or_init().increment(value)
}
pub fn config(&self) -> Config {
self.config
}
pub fn load(&self) -> Option<Histogram> {
self.inner.get().map(|h| h.load())
}
fn get_or_init(&self) -> &::histogram::AtomicHistogram {
self.inner
.get_or_init(|| ::histogram::AtomicHistogram::with_config(&self.config))
}
}
impl HistogramMetric for AtomicHistogram {
fn config(&self) -> Config {
self.config
}
fn load(&self) -> Option<Histogram> {
self.load()
}
}
impl Metric for AtomicHistogram {
fn as_any(&self) -> Option<&dyn std::any::Any> {
Some(self)
}
fn value(&self) -> Option<Value<'_>> {
Some(Value::Histogram(self))
}
}
struct HistogramState {
histogram: histogram::Histogram,
window: Option<Window>,
}
pub struct RwLockHistogram {
inner: OnceLock<RwLock<HistogramState>>,
config: Config,
}
impl RwLockHistogram {
pub const fn new(grouping_power: u8, max_value_power: u8) -> Self {
let config = match ::histogram::Config::new(grouping_power, max_value_power) {
Ok(c) => c,
Err(_e) => panic!("invalid histogram config"),
};
Self {
inner: OnceLock::new(),
config,
}
}
pub fn update_from(&self, data: &[u64]) -> Result<(), Error> {
if data.len() != self.config.total_buckets() {
return Err(Error::IncompatibleParameters);
}
let mut state = self.get_or_init().write();
state.histogram.as_mut_slice().copy_from_slice(data);
Ok(())
}
pub fn config(&self) -> Config {
self.config
}
pub fn load(&self) -> Option<Histogram> {
self.inner.get().map(|h| h.read().histogram.clone())
}
pub fn set_with_window(&self, data: &[u64], window: Window) -> Result<(), Error> {
if data.len() != self.config.total_buckets() {
return Err(Error::IncompatibleParameters);
}
let mut state = self.get_or_init().write();
state.histogram.as_mut_slice().copy_from_slice(data);
state.window = Some(window);
Ok(())
}
pub fn load_with_window(&self) -> (Option<Histogram>, Option<Window>) {
match self.inner.get() {
Some(lock) => {
let state = lock.read();
(Some(state.histogram.clone()), state.window)
}
None => (None, None),
}
}
fn get_or_init(&self) -> &RwLock<HistogramState> {
self.inner.get_or_init(|| {
RwLock::new(HistogramState {
histogram: ::histogram::Histogram::with_config(&self.config),
window: None,
})
})
}
}
impl HistogramMetric for RwLockHistogram {
fn config(&self) -> Config {
self.config
}
fn load(&self) -> Option<Histogram> {
self.load()
}
}
impl Metric for RwLockHistogram {
fn as_any(&self) -> Option<&dyn std::any::Any> {
Some(self)
}
fn value(&self) -> Option<Value<'_>> {
Some(Value::Histogram(self))
}
fn load_window(&self) -> Option<Window> {
self.inner.get().and_then(|h| h.read().window)
}
fn value_with_window(&self) -> (Option<Value<'_>>, Option<Window>) {
let window = self.inner.get().and_then(|h| h.read().window);
(Some(Value::Histogram(self)), window)
}
}
#[cfg(test)]
mod tests {
use super::*;
use metriken_core::Window;
#[test]
fn round_trip() {
let h = RwLockHistogram::new(7, 64);
let buckets = vec![0u64; h.config().total_buckets()];
h.set_with_window(&buckets, Window::new(10, 20)).unwrap();
let (hist, window) = h.load_with_window();
assert!(hist.is_some());
assert_eq!(window, Some(Window::new(10, 20)));
assert_eq!(
<RwLockHistogram as Metric>::load_window(&h),
Some(Window::new(10, 20))
);
}
#[test]
fn unset_window_is_none() {
let h = RwLockHistogram::new(7, 64);
let (hist, window) = h.load_with_window();
assert!(hist.is_none());
assert!(window.is_none());
assert!(<RwLockHistogram as Metric>::load_window(&h).is_none());
}
#[test]
fn wrong_length_is_rejected_and_records_no_window() {
let h = RwLockHistogram::new(7, 64);
assert!(h.set_with_window(&[1, 2, 3], Window::new(1, 2)).is_err());
assert!(<RwLockHistogram as Metric>::load_window(&h).is_none());
}
#[test]
fn update_from_is_windowless() {
let h = RwLockHistogram::new(7, 64);
let buckets = vec![0u64; h.config().total_buckets()];
h.update_from(&buckets).unwrap();
assert!(h.load().is_some());
let (_hist, window) = h.load_with_window();
assert!(window.is_none());
}
#[test]
fn value_with_window_returns_window_with_histogram_ref() {
let h = RwLockHistogram::new(7, 64);
let buckets = vec![0u64; h.config().total_buckets()];
h.set_with_window(&buckets, Window::new(10, 20)).unwrap();
let (value, window) = <RwLockHistogram as Metric>::value_with_window(&h);
assert!(matches!(value, Some(Value::Histogram(_))));
assert_eq!(window, Some(Window::new(10, 20)));
}
#[test]
fn set_with_window_torn_read_stress() {
use std::sync::Arc;
use std::thread;
const ITERS: u64 = 50_000;
let h = Arc::new(RwLockHistogram::new(4, 16));
let total = h.config().total_buckets();
{
let mut buckets = vec![0u64; total];
buckets[0] = 0;
h.set_with_window(&buckets, Window::new(0, 1)).unwrap();
}
let writer = {
let h = h.clone();
thread::spawn(move || {
let mut buckets = vec![0u64; total];
for v in 1..ITERS {
buckets[0] = v;
h.set_with_window(&buckets, Window::new(v, v + 1)).unwrap();
}
})
};
let reader = {
let h = h.clone();
thread::spawn(move || {
for _ in 0..ITERS {
let (hist, w) = h.load_with_window();
if let (Some(hist), Some(w)) = (hist, w) {
let v = hist.as_slice()[0];
assert_eq!(w.begin_ns, v, "torn read: buckets {v} paired with {w:?}");
assert_eq!(w.end_ns, v + 1, "torn read: buckets {v} paired with {w:?}");
}
}
})
};
writer.join().unwrap();
reader.join().unwrap();
}
}