use std::collections::HashSet;
use calimero_crypto::Nonce;
use calimero_primitives::hash::Hash;
use calimero_primitives::identity::PublicKey;
pub const DEFAULT_BUFFER_CAPACITY: usize = 10_000;
pub const MIN_RECOMMENDED_CAPACITY: usize = 100;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PushResult {
Added,
Duplicate,
Evicted([u8; 32]),
DroppedZeroCapacity([u8; 32]),
}
impl PushResult {
#[must_use]
pub fn was_added(&self) -> bool {
matches!(self, Self::Added | Self::Evicted(_))
}
#[must_use]
pub fn had_data_loss(&self) -> bool {
matches!(self, Self::Evicted(_) | Self::DroppedZeroCapacity(_))
}
#[must_use]
pub fn lost_delta_id(&self) -> Option<[u8; 32]> {
match self {
Self::Evicted(id) | Self::DroppedZeroCapacity(id) => Some(*id),
Self::Added | Self::Duplicate => None,
}
}
}
#[derive(Debug, Clone)]
pub struct BufferedDelta {
pub id: [u8; 32],
pub parents: Vec<[u8; 32]>,
pub hlc: u64,
pub payload: Vec<u8>,
pub nonce: Nonce,
pub author_id: PublicKey,
pub root_hash: Hash,
pub events: Option<Vec<u8>>,
pub source_peer: libp2p::PeerId,
}
#[derive(Debug)]
pub struct DeltaBuffer {
deltas: std::collections::VecDeque<BufferedDelta>,
seen_ids: HashSet<[u8; 32]>,
sync_start_hlc: u64,
capacity: usize,
drops: u64,
}
impl DeltaBuffer {
#[must_use]
pub fn new(capacity: usize, sync_start_hlc: u64) -> Self {
Self {
deltas: std::collections::VecDeque::with_capacity(capacity.min(1000)),
seen_ids: HashSet::with_capacity(capacity.min(1000)),
sync_start_hlc,
capacity,
drops: 0,
}
}
#[must_use]
pub fn is_capacity_below_recommended(&self) -> bool {
self.capacity < MIN_RECOMMENDED_CAPACITY
}
pub fn push(&mut self, delta: BufferedDelta) -> PushResult {
let delta_id = delta.id;
if self.capacity == 0 {
self.drops += 1;
return PushResult::DroppedZeroCapacity(delta_id);
}
if self.seen_ids.contains(&delta_id) {
return PushResult::Duplicate;
}
if self.deltas.len() >= self.capacity {
if let Some(evicted) = self.deltas.pop_front() {
self.seen_ids.remove(&evicted.id);
let evicted_id = evicted.id;
self.drops += 1;
self.seen_ids.insert(delta_id);
self.deltas.push_back(delta);
PushResult::Evicted(evicted_id)
} else {
self.seen_ids.insert(delta_id);
self.deltas.push_back(delta);
PushResult::Added
}
} else {
self.seen_ids.insert(delta_id);
self.deltas.push_back(delta);
PushResult::Added
}
}
#[must_use]
pub fn drain(&mut self) -> Vec<BufferedDelta> {
self.seen_ids.clear();
self.deltas.drain(..).collect()
}
#[must_use]
pub fn contains(&self, id: &[u8; 32]) -> bool {
self.seen_ids.contains(id)
}
#[must_use]
pub fn len(&self) -> usize {
self.deltas.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.deltas.is_empty()
}
#[must_use]
pub fn sync_start_hlc(&self) -> u64 {
self.sync_start_hlc
}
#[must_use]
pub fn drops(&self) -> u64 {
self.drops
}
#[must_use]
pub fn capacity(&self) -> usize {
self.capacity
}
}
#[cfg(test)]
mod tests {
use super::*;
fn make_test_delta(id: u8) -> BufferedDelta {
BufferedDelta {
id: [id; 32],
parents: vec![[0; 32]],
hlc: 12345,
payload: vec![1, 2, 3],
nonce: [0; 12],
author_id: PublicKey::from([0; 32]),
root_hash: Hash::from([0; 32]),
events: None,
source_peer: libp2p::PeerId::random(),
}
}
#[test]
fn test_buffer_basic() {
let mut buffer = DeltaBuffer::new(100, 12345);
assert!(buffer.is_empty());
assert_eq!(buffer.sync_start_hlc(), 12345);
assert_eq!(buffer.capacity(), 100);
assert_eq!(buffer.drops(), 0);
assert!(!buffer.is_capacity_below_recommended());
let result = buffer.push(make_test_delta(1));
assert_eq!(result, PushResult::Added, "Should add without eviction");
assert!(result.was_added());
assert!(!result.had_data_loss());
assert_eq!(buffer.len(), 1);
let drained = buffer.drain();
assert_eq!(drained.len(), 1);
assert!(buffer.is_empty());
}
#[test]
fn test_buffer_only_during_sync() {
let mut buffer = DeltaBuffer::new(10, 12345);
assert!(buffer.is_empty());
assert_eq!(buffer.push(make_test_delta(1)), PushResult::Added);
assert_eq!(buffer.push(make_test_delta(2)), PushResult::Added);
assert_eq!(buffer.len(), 2);
let drained = buffer.drain();
assert_eq!(drained.len(), 2);
assert_eq!(drained[0].id[0], 1);
assert_eq!(drained[1].id[0], 2);
}
#[test]
fn test_buffer_overflow_drops_oldest() {
let mut buffer = DeltaBuffer::new(2, 0);
assert_eq!(buffer.push(make_test_delta(1)), PushResult::Added);
assert_eq!(buffer.push(make_test_delta(2)), PushResult::Added);
assert_eq!(buffer.drops(), 0);
let result = buffer.push(make_test_delta(3));
assert_eq!(result, PushResult::Evicted([1; 32]), "Should evict delta 1");
assert!(result.had_data_loss());
assert_eq!(result.lost_delta_id(), Some([1; 32]));
assert_eq!(buffer.drops(), 1);
assert_eq!(buffer.len(), 2);
let result = buffer.push(make_test_delta(4));
assert_eq!(result, PushResult::Evicted([2; 32]), "Should evict delta 2");
assert_eq!(buffer.drops(), 2);
assert_eq!(buffer.len(), 2);
let drained = buffer.drain();
assert_eq!(drained.len(), 2);
assert_eq!(drained[0].id[0], 3);
assert_eq!(drained[1].id[0], 4);
}
#[test]
fn test_zero_capacity_drops_immediately() {
let mut buffer = DeltaBuffer::new(0, 0);
assert!(buffer.is_empty());
assert_eq!(buffer.capacity(), 0);
assert_eq!(buffer.drops(), 0);
assert!(buffer.is_capacity_below_recommended());
let result = buffer.push(make_test_delta(1));
assert_eq!(
result,
PushResult::DroppedZeroCapacity([1; 32]),
"Zero capacity should drop incoming delta"
);
assert!(result.had_data_loss());
assert!(!result.was_added());
assert_eq!(result.lost_delta_id(), Some([1; 32]));
assert_eq!(buffer.drops(), 1);
assert!(buffer.is_empty(), "Buffer should remain empty");
assert_eq!(buffer.len(), 0);
let result = buffer.push(make_test_delta(2));
assert_eq!(result, PushResult::DroppedZeroCapacity([2; 32]));
assert_eq!(buffer.drops(), 2);
assert!(buffer.is_empty());
}
#[test]
fn test_finish_sync_returns_fifo() {
let mut buffer = DeltaBuffer::new(100, 0);
buffer.push(make_test_delta(1));
buffer.push(make_test_delta(2));
buffer.push(make_test_delta(3));
let drained = buffer.drain();
assert_eq!(drained.len(), 3);
assert_eq!(drained[0].id[0], 1);
assert_eq!(drained[1].id[0], 2);
assert_eq!(drained[2].id[0], 3);
}
#[test]
fn test_cancel_sync_clears_buffer() {
let mut buffer = DeltaBuffer::new(100, 0);
buffer.push(make_test_delta(1));
buffer.push(make_test_delta(2));
assert_eq!(buffer.len(), 2);
let _ = buffer.drain();
assert!(buffer.is_empty());
assert_eq!(buffer.len(), 0);
}
#[test]
fn test_drops_counter_observable() {
let mut buffer = DeltaBuffer::new(1, 0);
assert_eq!(buffer.drops(), 0);
buffer.push(make_test_delta(1));
assert_eq!(buffer.drops(), 0);
buffer.push(make_test_delta(2));
assert_eq!(buffer.drops(), 1);
buffer.push(make_test_delta(3));
assert_eq!(buffer.drops(), 2);
buffer.push(make_test_delta(4));
assert_eq!(buffer.drops(), 3);
}
#[test]
fn test_deduplication_prevents_double_buffering() {
let mut buffer = DeltaBuffer::new(10, 0);
assert_eq!(buffer.push(make_test_delta(1)), PushResult::Added);
assert_eq!(buffer.len(), 1);
let result = buffer.push(make_test_delta(1));
assert_eq!(result, PushResult::Duplicate);
assert!(!result.had_data_loss());
assert!(!result.was_added()); assert_eq!(buffer.len(), 1);
assert_eq!(buffer.push(make_test_delta(2)), PushResult::Added);
assert_eq!(buffer.len(), 2);
}
#[test]
fn test_deduplication_cleared_on_drain() {
let mut buffer = DeltaBuffer::new(10, 0);
assert_eq!(buffer.push(make_test_delta(1)), PushResult::Added);
assert!(buffer.contains(&[1; 32]));
let _ = buffer.drain();
assert!(!buffer.contains(&[1; 32]));
assert_eq!(buffer.push(make_test_delta(1)), PushResult::Added);
assert_eq!(buffer.len(), 1);
}
#[test]
fn test_deduplication_cleared_on_eviction() {
let mut buffer = DeltaBuffer::new(2, 0);
buffer.push(make_test_delta(1));
buffer.push(make_test_delta(2));
assert!(buffer.contains(&[1; 32]));
buffer.push(make_test_delta(3));
assert!(!buffer.contains(&[1; 32])); assert!(buffer.contains(&[2; 32]));
assert!(buffer.contains(&[3; 32]));
let result = buffer.push(make_test_delta(1));
assert_eq!(result, PushResult::Evicted([2; 32])); }
#[test]
fn test_capacity_below_recommended() {
let buffer = DeltaBuffer::new(50, 0);
assert!(buffer.is_capacity_below_recommended());
let buffer = DeltaBuffer::new(MIN_RECOMMENDED_CAPACITY, 0);
assert!(!buffer.is_capacity_below_recommended());
let buffer = DeltaBuffer::new(MIN_RECOMMENDED_CAPACITY + 1, 0);
assert!(!buffer.is_capacity_below_recommended());
}
}