use std::collections::VecDeque;
use std::sync::Mutex;
use std::time::{Duration, Instant};
pub type FeatureVector = Vec<f64>;
#[derive(Clone, Copy, Debug)]
pub enum WindowMode {
Count(usize),
Duration(Duration),
}
struct Entry {
features: FeatureVector,
at: Instant,
}
pub struct LiveWindow {
buffer: Mutex<VecDeque<Entry>>,
mode: WindowMode,
}
impl LiveWindow {
pub fn new(mode: WindowMode) -> Self {
Self {
buffer: Mutex::new(VecDeque::new()),
mode,
}
}
pub fn push(&self, features: FeatureVector) {
self.push_at(features, Instant::now());
}
pub(crate) fn push_at(&self, features: FeatureVector, at: Instant) {
let mut buf = self.buffer.lock().expect("LiveWindow mutex poisoned");
buf.push_back(Entry { features, at });
evict(&mut buf, self.mode, at);
}
pub fn snapshot(&self) -> Vec<FeatureVector> {
let mut buf = self.buffer.lock().expect("LiveWindow mutex poisoned");
evict(&mut buf, self.mode, Instant::now());
buf.iter().map(|e| e.features.clone()).collect()
}
pub fn column(&self, index: usize) -> Vec<f64> {
let mut buf = self.buffer.lock().expect("LiveWindow mutex poisoned");
evict(&mut buf, self.mode, Instant::now());
buf.iter()
.filter_map(|e| e.features.get(index).copied())
.collect()
}
pub fn len(&self) -> usize {
let mut buf = self.buffer.lock().expect("LiveWindow mutex poisoned");
evict(&mut buf, self.mode, Instant::now());
buf.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn clear(&self) {
self.buffer
.lock()
.expect("LiveWindow mutex poisoned")
.clear();
}
}
fn evict(buf: &mut VecDeque<Entry>, mode: WindowMode, now: Instant) {
match mode {
WindowMode::Count(n) => {
while buf.len() > n {
buf.pop_front();
}
}
WindowMode::Duration(d) => {
while let Some(front) = buf.front() {
if now.duration_since(front.at) > d {
buf.pop_front();
} else {
break;
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn count_mode_evicts_oldest() {
let window = LiveWindow::new(WindowMode::Count(3));
for i in 0..5 {
window.push(vec![i as f64]);
}
assert_eq!(window.snapshot(), vec![vec![2.0], vec![3.0], vec![4.0]]);
assert_eq!(window.len(), 3);
}
#[test]
fn column_extracts_feature() {
let window = LiveWindow::new(WindowMode::Count(10));
window.push(vec![1.0, 10.0]);
window.push(vec![2.0, 20.0]);
assert_eq!(window.column(0), vec![1.0, 2.0]);
assert_eq!(window.column(1), vec![10.0, 20.0]);
}
#[test]
fn duration_mode_evicts_expired() {
let t0 = Instant::now();
let mut buf = VecDeque::new();
buf.push_back(Entry {
features: vec![1.0],
at: t0,
});
buf.push_back(Entry {
features: vec![2.0],
at: t0 + Duration::from_secs(10),
});
evict(
&mut buf,
WindowMode::Duration(Duration::from_secs(5)),
t0 + Duration::from_secs(12),
);
assert_eq!(buf.len(), 1);
assert_eq!(buf[0].features, vec![2.0]);
}
}