use crate::window::OwnedWindow;
use crate::Error;
use crate::*;
use core::sync::atomic::*;
use histogram::{Bucket, Histogram};
use parking_lot::Mutex;
pub struct Heatmap {
slices: Vec<OwnedWindow>,
current: AtomicUsize,
lock: Mutex<()>,
start: AtomicInstant,
stop: AtomicInstant,
span: Duration,
resolution: Duration,
next_tick: AtomicInstant,
summary: Histogram,
}
pub struct Builder {
m: u32,
r: u32,
n: u32,
span: Duration,
resolution: Duration,
}
impl Builder {
pub fn build(self) -> Result<Heatmap, Error> {
Heatmap::new(self.m, self.r, self.n, self.span, self.resolution)
}
pub fn min_resolution(mut self, width: u64) -> Self {
self.m = 64 - width.leading_zeros();
self
}
pub fn min_resolution_range(mut self, value: u64) -> Self {
self.r = 64 - value.next_power_of_two().leading_zeros();
self
}
pub fn maximum_value(mut self, value: u64) -> Self {
self.n = 64 - value.next_power_of_two().leading_zeros();
self
}
pub fn span(mut self, duration: Duration) -> Self {
self.span = duration;
self
}
pub fn resolution(mut self, duration: Duration) -> Self {
self.resolution = duration;
self
}
}
impl Heatmap {
pub fn new(
m: u32,
r: u32,
n: u32,
span: Duration,
resolution: Duration,
) -> Result<Self, Error> {
let now = Instant::now();
let mut slices = Vec::new();
let mut true_span = Duration::from_nanos(0);
while true_span < span {
slices.push(OwnedWindow {
start: AtomicInstant::new(now + true_span),
stop: AtomicInstant::new(now + true_span + resolution),
histogram: Histogram::new(m, r, n)?,
});
true_span += resolution;
}
slices.push(OwnedWindow {
start: AtomicInstant::new(now + true_span),
stop: AtomicInstant::new(now + true_span + resolution),
histogram: Histogram::new(m, r, n)?,
});
slices.shrink_to_fit();
let start = AtomicInstant::new(now);
let stop = AtomicInstant::new(now + true_span);
let next_tick = AtomicInstant::new(now + resolution);
Ok(Self {
slices,
current: AtomicUsize::new(0),
lock: Mutex::new(()),
start,
stop,
span: true_span,
resolution,
next_tick,
summary: Histogram::new(m, r, n)?,
})
}
pub fn builder() -> Builder {
Builder {
m: 0,
r: 10,
n: 30,
span: Duration::from_secs(60),
resolution: Duration::from_secs(1),
}
}
pub fn windows(&self) -> usize {
self.slices.len()
}
pub fn buckets(&self) -> usize {
self.summary.buckets()
}
pub fn increment(&self, time: Instant, value: u64, count: u32) {
self.tick(time);
let current = self.current.load(Ordering::Relaxed);
if time >= self.slices[current].start.load(Ordering::Relaxed) {
let _ = self.summary.increment(value, count);
let _ = self.slices[current].histogram.increment(value, count);
}
let start = self.start.load(Ordering::Relaxed);
let stop = self.stop.load(Ordering::Relaxed);
if time < start {
return;
}
let offset = ((stop - time).as_nanos() / self.resolution.as_nanos()) as usize;
let index = if offset > current {
current + self.slices.len() - offset
} else {
current - offset
};
let _ = self.summary.increment(value, count);
let _ = self.slices[index].histogram.increment(value, count);
}
pub fn percentile(&self, percentile: f64) -> Result<Bucket, Error> {
self.tick(Instant::now());
self.summary.percentile(percentile).map_err(Error::from)
}
pub fn iter(&self) -> Iter {
self.into_iter()
}
pub fn summary(&self) -> &Histogram {
&self.summary
}
fn tick(&self, time: Instant) {
loop {
let next_tick = self.next_tick.load(Ordering::Relaxed);
if time < next_tick {
return;
} else {
if let Some(_lock) = self.lock.try_lock() {
if time < self.next_tick.load(Ordering::Relaxed) {
return;
}
let current = self.current.load(Ordering::Relaxed);
let mut next = current + 1;
if next >= self.slices.len() {
next -= self.slices.len();
}
self.current.store(next, Ordering::Relaxed);
self.next_tick.fetch_add(self.resolution, Ordering::Relaxed);
let mut to_clear = next + 1;
if to_clear >= self.slices.len() {
to_clear -= self.slices.len();
}
self.start.fetch_add(self.resolution, Ordering::Relaxed);
self.stop.fetch_add(self.resolution, Ordering::Relaxed);
if self.slices[to_clear].start.load(Ordering::Relaxed)
< self.start.load(Ordering::Relaxed)
{
let _ = self
.summary
.subtract_and_clear(&self.slices[to_clear].histogram);
self.slices[to_clear]
.start
.fetch_add(self.span, Ordering::Relaxed);
self.slices[to_clear]
.stop
.fetch_add(self.span, Ordering::Relaxed);
}
}
}
}
}
fn get_slice(&self, index: usize) -> Option<Window> {
if let Some(histogram) = self.slices.get(index) {
let shift = if index > self.current.load(Ordering::Relaxed) {
self.resolution.mul_f64(
(self.slices.len() + self.current.load(Ordering::Relaxed) - index) as f64,
)
} else {
self.resolution
.mul_f64((self.current.load(Ordering::Relaxed) - index) as f64)
};
Some(Window {
start: self.next_tick.load(Ordering::Relaxed) - shift - self.resolution,
stop: self.next_tick.load(Ordering::Relaxed) - shift,
histogram: &histogram.histogram,
})
} else {
None
}
}
}
impl Clone for Heatmap {
fn clone(&self) -> Self {
let slices = self.slices.clone();
let summary = self.summary.clone();
let start = AtomicInstant::new(self.start.load(Ordering::Relaxed));
let stop = AtomicInstant::new(self.stop.load(Ordering::Relaxed));
let span = self.span;
let resolution = self.resolution;
let current = AtomicUsize::new(self.current.load(Ordering::Relaxed));
let next_tick = AtomicInstant::new(self.next_tick.load(Ordering::Relaxed));
Heatmap {
slices,
current,
lock: Mutex::new(()),
start,
stop,
span,
resolution,
next_tick,
summary,
}
}
}
pub struct Iter<'a> {
inner: &'a Heatmap,
index: usize,
visited: usize,
}
impl<'a> Iter<'a> {
fn new(inner: &'a Heatmap) -> Iter<'a> {
let index = if inner.current.load(Ordering::Relaxed) < (inner.slices.len() - 1) {
inner.current.load(Ordering::Relaxed) + 1
} else {
0
};
Iter {
inner,
index,
visited: 0,
}
}
}
impl<'a> Iterator for Iter<'a> {
type Item = Window<'a>;
fn next(&mut self) -> Option<Window<'a>> {
if self.visited >= self.inner.slices.len() {
None
} else {
let bucket = self.inner.get_slice(self.index);
self.index += 1;
if self.index >= self.inner.slices.len() {
self.index = 0;
}
self.visited += 1;
bucket
}
}
}
impl<'a> IntoIterator for &'a Heatmap {
type Item = Window<'a>;
type IntoIter = Iter<'a>;
fn into_iter(self) -> Self::IntoIter {
Iter::new(self)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn age_out() {
let heatmap =
Heatmap::new(0, 4, 20, Duration::from_secs(1), Duration::from_millis(1)).unwrap();
assert_eq!(heatmap.percentile(0.0).map(|v| v.high()), Err(Error::Empty));
heatmap.increment(Instant::now(), 1, 1);
assert_eq!(heatmap.percentile(0.0).map(|v| v.high()), Ok(1));
std::thread::sleep(std::time::Duration::from_millis(100));
assert_eq!(heatmap.percentile(0.0).map(|v| v.high()), Ok(1));
std::thread::sleep(std::time::Duration::from_millis(2000));
assert_eq!(heatmap.percentile(0.0).map(|v| v.high()), Err(Error::Empty));
}
}