use std::time::Instant;
use super::schema::FeatureValue;
#[derive(Debug, Clone)]
pub struct BufferedEntry {
pub group: String,
pub entity_key: String,
pub feature: String,
pub value: FeatureValue,
pub event_ts: u64,
pub ingestion_ts: u64,
}
pub struct WriteAheadBuffer {
buffer: Vec<BufferedEntry>,
max_entries: usize,
max_age_micros: u64,
last_flush: Instant,
}
impl WriteAheadBuffer {
pub fn new(max_entries: usize, max_age_micros: u64) -> Self {
Self {
buffer: Vec::with_capacity(max_entries),
max_entries,
max_age_micros,
last_flush: Instant::now(),
}
}
pub fn push(&mut self, entry: BufferedEntry) -> Option<Vec<BufferedEntry>> {
self.buffer.push(entry);
if self.buffer.len() >= self.max_entries {
Some(self.flush())
} else {
None
}
}
pub fn flush(&mut self) -> Vec<BufferedEntry> {
self.last_flush = Instant::now();
std::mem::take(&mut self.buffer)
}
pub fn should_flush(&self) -> bool {
if self.buffer.is_empty() {
return false;
}
let elapsed_micros = self.last_flush.elapsed().as_micros() as u64;
elapsed_micros >= self.max_age_micros
}
pub fn len(&self) -> usize {
self.buffer.len()
}
pub fn is_empty(&self) -> bool {
self.buffer.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn test_entry(group: &str, entity: &str, feature: &str) -> BufferedEntry {
BufferedEntry {
group: group.to_string(),
entity_key: entity.to_string(),
feature: feature.to_string(),
value: FeatureValue::Float64(1.0),
event_ts: 1_000_000,
ingestion_ts: 1_000_001,
}
}
#[test]
fn test_buffer_push_no_flush() {
let mut buffer = WriteAheadBuffer::new(10, 10_000);
let result = buffer.push(test_entry("g", "e", "f"));
assert!(result.is_none());
assert_eq!(buffer.len(), 1);
}
#[test]
fn test_buffer_auto_flush_on_threshold() {
let mut buffer = WriteAheadBuffer::new(3, 10_000);
assert!(buffer.push(test_entry("g", "e", "f1")).is_none());
assert!(buffer.push(test_entry("g", "e", "f2")).is_none());
let result = buffer.push(test_entry("g", "e", "f3"));
assert!(result.is_some());
let entries = result.unwrap();
assert_eq!(entries.len(), 3);
assert!(buffer.is_empty());
}
#[test]
fn test_buffer_manual_flush() {
let mut buffer = WriteAheadBuffer::new(100, 10_000);
buffer.push(test_entry("g", "e", "f1"));
buffer.push(test_entry("g", "e", "f2"));
let entries = buffer.flush();
assert_eq!(entries.len(), 2);
assert!(buffer.is_empty());
}
#[test]
fn test_buffer_should_flush_empty() {
let buffer = WriteAheadBuffer::new(100, 10_000);
assert!(!buffer.should_flush());
}
#[test]
fn test_buffer_should_flush_age() {
let mut buffer = WriteAheadBuffer::new(100, 0); buffer.push(test_entry("g", "e", "f"));
assert!(buffer.should_flush());
}
}