sinusoidal 0.10.0

The official SDK to write rust apps for the Sinusoidal Systems Digital Measurement Platform
Documentation
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,
    }
  }

  /// Pushes a new event into the pulser.
  ///
  /// If the new event completes a set for a given timestamp, it returns `Some((timestamp, events))`.
  /// Otherwise, it returns `None`.
  ///
  /// This method also clears out old entries that have timed out. An entry is considered
  /// timed out if its timestamp `ts_old` is such that `ts_high - ts_old > timeout`, where
  /// `ts_high` is the highest timestamp this `pulser` has seen.
  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(&timestamp)
      .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); // age is 10, not > 10, so not timed out
    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 {
            // 10% chance of invalid key
            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) {
            // Property 1: Emitted timestamp is the same as the event's timestamp that completed the set.
            assert_eq!(emitted_ts, event.ts);

            // Property 2: Completeness - an event for every key.
            let emitted_keys: BTreeSet<_> = emitted_events.keys().cloned().collect();
            assert_eq!(emitted_keys, valid_keys);

            // Property 3: Correct values.
            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(); // Should not fail if pulser is correct
                assert_eq!(&source_event.value, emitted_value);
            }
          }

          // Property 4: Timeout check.
          let buffer_timestamps = pulser.get_buffer_timestamps();
          assert!(buffer_timestamps.iter().all(|&t| t >= test_cutoff));
        }
      }
    }
  }
}