use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::fmt::Debug;
use std::hash::Hash;
#[derive(Debug)]
pub struct Pulser<K, V>
where
K: Hash + Eq + Clone + Ord,
{
keys: BTreeSet<K>,
timeout: u64,
buffer: BTreeMap<u64, HashMap<K, V>>,
cutoff: u64,
}
impl<K, V> Pulser<K, V>
where
K: Hash + Eq + Clone + Ord + Debug,
{
pub fn new<I>(keys: I, timeout: u64) -> Self
where
I: IntoIterator<Item = K>,
{
Pulser {
keys: keys.into_iter().collect(),
timeout,
buffer: BTreeMap::new(),
cutoff: 0,
}
}
pub fn push(&mut self, timestamp: u64, key: K, value: V) -> Option<(u64, HashMap<K, V>)> {
self.advance_clock(timestamp);
if !self.is_valid_key(&key) {
return None;
}
if timestamp < self.cutoff {
return None;
}
if !self.push_event(timestamp, key, value) {
return None;
}
self.pop_events(timestamp)
}
fn is_valid_key(&mut self, key: &K) -> bool {
self.keys.contains(key)
}
fn advance_clock(&mut self, timestamp: u64) {
let new_cutoff = timestamp.saturating_sub(self.timeout);
self.cutoff = self.cutoff.max(new_cutoff);
self.buffer = self.buffer.split_off(&self.cutoff);
}
fn push_event(&mut self, timestamp: u64, key: K, value: V) -> bool {
let events_at_ts = self.buffer.entry(timestamp).or_default();
events_at_ts.insert(key, value);
events_at_ts.len() == self.keys.len()
}
fn pop_events(&mut self, timestamp: u64) -> Option<(u64, HashMap<K, V>)> {
self
.buffer
.remove(×tamp)
.map(|events| (timestamp, events))
}
#[cfg(test)]
fn get_buffer_timestamps(&self) -> Vec<u64> {
self.buffer.keys().cloned().collect()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_pulser_new() {
let keys = vec!["a", "b", "c"];
let pulser: Pulser<&str, i32> = Pulser::new(keys, 10);
assert_eq!(pulser.keys.len(), 3);
assert!(pulser.keys.contains("a"));
assert!(pulser.keys.contains("b"));
assert!(pulser.keys.contains("c"));
assert_eq!(pulser.timeout, 10);
assert!(pulser.buffer.is_empty());
}
#[test]
fn test_pulser_push_and_emit() {
let keys = vec!["a", "b"];
let mut pulser = Pulser::new(keys, 10);
let r1 = pulser.push(100, "a", 1);
assert!(r1.is_none());
let r2 = pulser.push(100, "b", 2);
assert!(r2.is_some());
if let Some((ts, events)) = r2 {
assert_eq!(ts, 100);
assert_eq!(events.len(), 2);
assert_eq!(events.get("a"), Some(&1));
assert_eq!(events.get("b"), Some(&2));
}
assert!(pulser.buffer.is_empty());
}
#[test]
fn test_pulser_different_timestamps() {
let keys = vec!["a", "b"];
let mut pulser = Pulser::new(keys, 10);
assert!(pulser.push(100, "a", 1).is_none());
assert!(pulser.push(101, "a", 2).is_none());
assert_eq!(pulser.buffer.len(), 2);
let result = pulser.push(100, "b", 3);
assert!(result.is_some());
if let Some((ts, events)) = result {
assert_eq!(ts, 100);
assert_eq!(events.get("a"), Some(&1));
assert_eq!(events.get("b"), Some(&3));
}
assert_eq!(pulser.buffer.len(), 1);
assert!(pulser.buffer.contains_key(&101));
let result2 = pulser.push(101, "b", 4);
assert!(result2.is_some());
if let Some((ts, events)) = result2 {
assert_eq!(ts, 101);
assert_eq!(events.get("a"), Some(&2));
assert_eq!(events.get("b"), Some(&4));
}
assert!(pulser.buffer.is_empty());
}
#[test]
fn test_pulser_timeout() {
let keys = vec!["a", "b"];
let mut pulser = Pulser::new(keys.clone(), 10);
pulser.push(100, "a", 1);
assert_eq!(pulser.buffer.len(), 1);
pulser.push(111, "a", 2);
assert_eq!(pulser.buffer.len(), 1);
assert!(!pulser.buffer.contains_key(&100));
assert!(pulser.buffer.contains_key(&111));
pulser.push(120, "b", 3);
assert_eq!(pulser.buffer.len(), 2);
assert!(pulser.buffer.contains_key(&111));
assert!(pulser.buffer.contains_key(&120));
let result = pulser.push(122, "a", 4);
assert!(result.is_none());
assert_eq!(pulser.buffer.len(), 2);
assert!(!pulser.buffer.contains_key(&111));
assert!(pulser.buffer.contains_key(&120));
assert!(pulser.buffer.contains_key(&122));
}
#[test]
fn test_unknown_key() {
let keys = vec!["a", "b"];
let mut pulser = Pulser::new(keys, 10);
let result = pulser.push(100, "c", 1);
assert!(result.is_none());
assert!(pulser.buffer.is_empty());
}
#[test]
fn test_duplicate_push_overwrites() {
let keys = vec!["a", "b"];
let mut pulser = Pulser::new(keys, 10);
assert!(pulser.push(100, "a", 1).is_none());
assert_eq!(pulser.buffer.get(&100).unwrap().get("a"), Some(&1));
assert!(pulser.push(100, "a", 2).is_none());
assert_eq!(pulser.buffer.get(&100).unwrap().get("a"), Some(&2));
let result = pulser.push(100, "b", 3);
assert!(result.is_some());
if let Some((_ts, events)) = result {
assert_eq!(events.get("a"), Some(&2));
}
}
#[test]
fn test_timeout_with_exact_boundary() {
let keys = vec!["a", "b"];
let mut pulser = Pulser::new(keys, 10);
pulser.push(100, "a", 1);
pulser.push(110, "a", 2); assert_eq!(pulser.buffer.len(), 2, "Buffer should contain two entries");
}
mod properties {
use super::*;
use quickcheck::{Arbitrary, Gen, quickcheck};
#[derive(Clone, Debug)]
struct Event {
ts: u64,
key: String,
value: u32,
}
impl Arbitrary for Event {
fn arbitrary(g: &mut Gen) -> Self {
Event {
ts: u64::arbitrary(g),
key: String::arbitrary(g),
value: u32::arbitrary(g),
}
}
fn shrink(&self) -> Box<dyn Iterator<Item = Self>> {
let ts = self.ts;
let key = self.key.clone();
let value = self.value;
Box::new(
(ts, key, value)
.shrink()
.map(|(ts, key, value)| Event { ts, key, value }),
)
}
}
#[derive(Clone, Debug)]
struct TestScenario {
keys: BTreeSet<String>,
events: Vec<Event>,
timeout: u64,
}
impl Arbitrary for TestScenario {
fn arbitrary(g: &mut Gen) -> Self {
let mut keys: BTreeSet<String> = Arbitrary::arbitrary(g);
if keys.is_empty() {
keys.insert("default_key".to_string());
}
let num_events = u64::arbitrary(g) % 100;
let mut events = Vec::with_capacity(num_events as usize);
let valid_keys: Vec<String> = keys.iter().cloned().collect();
for _ in 0..num_events {
let key = if (u64::arbitrary(g) % 10) == 0 {
let mut invalid_key: String = Arbitrary::arbitrary(g);
while keys.contains(&invalid_key) {
invalid_key = Arbitrary::arbitrary(g);
}
invalid_key
} else {
g.choose(&valid_keys).unwrap().clone()
};
let ts = u64::arbitrary(g) % 1000;
let value = u32::arbitrary(g);
events.push(Event { ts, key, value });
}
let timeout = 1 + (u64::arbitrary(g) % 199);
TestScenario {
keys,
events,
timeout,
}
}
fn shrink(&self) -> Box<dyn Iterator<Item = Self>> {
let keys = self.keys.clone();
let events = self.events.clone();
let timeout = self.timeout;
Box::new(
(keys, events, timeout)
.shrink()
.map(|(keys, events, timeout)| TestScenario {
keys,
events,
timeout,
}),
)
}
}
quickcheck! {
fn prop_pulser_scenario_test(scenario: TestScenario) -> () {
let mut pulser = Pulser::new(scenario.keys.clone(), scenario.timeout);
let valid_keys = scenario.keys;
let all_events = scenario.events;
let mut test_cutoff = 0;
for (event_idx, event) in all_events.iter().enumerate() {
let new_cutoff = event.ts.saturating_sub(scenario.timeout);
test_cutoff = test_cutoff.max(new_cutoff);
if let Some((emitted_ts, emitted_events)) = pulser.push(event.ts, event.key.clone(), event.value) {
assert_eq!(emitted_ts, event.ts);
let emitted_keys: BTreeSet<_> = emitted_events.keys().cloned().collect();
assert_eq!(emitted_keys, valid_keys);
for (emitted_key, emitted_value) in &emitted_events {
let source_event = all_events[..=event_idx]
.iter()
.rev()
.find(|e| e.ts == emitted_ts && &e.key == emitted_key)
.unwrap(); assert_eq!(&source_event.value, emitted_value);
}
}
let buffer_timestamps = pulser.get_buffer_timestamps();
assert!(buffer_timestamps.iter().all(|&t| t >= test_cutoff));
}
}
}
}
}