use crate::datatypes::{Crdt, DeltaBuffer, DeltaOrSet, OrSetDelta};
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum Shipment {
DeltaInterval {
delta: OrSetDelta,
ack_seq: u64,
},
FullState(DeltaOrSet),
UpToDate,
}
impl Shipment {
#[must_use]
pub fn wire_len(&self) -> usize {
match self {
Shipment::DeltaInterval { delta, .. } => delta.wire_len(),
Shipment::FullState(state) => state.wire_len(),
Shipment::UpToDate => 0,
}
}
}
#[must_use]
pub fn plan_shipment(state: &DeltaOrSet, buf: &DeltaBuffer, peer: &str) -> Shipment {
let Some(since) = buf.next_needed(peer) else {
return Shipment::FullState(state.clone());
};
match (buf.interval_since(since), buf.high_water()) {
(Some(delta), Some(hw)) => Shipment::DeltaInterval { delta, ack_seq: hw },
_ => Shipment::UpToDate,
}
}
pub fn apply_shipment(local: &mut DeltaOrSet, shipment: &Shipment) -> Option<u64> {
match shipment {
Shipment::DeltaInterval { delta, ack_seq } => {
local.merge_delta(delta);
Some(*ack_seq)
}
Shipment::FullState(state) => {
local.merge(state);
None
}
Shipment::UpToDate => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::datatypes::ActorId;
fn aid(name: &str) -> ActorId {
ActorId::new("dc1", name)
}
#[test]
fn first_contact_ships_full_state() {
let a = aid("a");
let mut src = DeltaOrSet::new();
let mut buf = DeltaBuffer::new();
buf.record(src.add(&a, b"x".to_vec()));
let ship = plan_shipment(&src, &buf, "peer-b");
assert!(matches!(ship, Shipment::FullState(_)));
let mut dst = DeltaOrSet::new();
apply_shipment(&mut dst, &ship);
assert_eq!(dst.value(), src.value());
}
#[test]
fn known_peer_ships_delta_interval() {
let a = aid("a");
let mut src = DeltaOrSet::new();
let mut buf = DeltaBuffer::new();
buf.record(src.add(&a, b"x".to_vec())); buf.ack("peer-b", 0);
buf.record(src.add(&a, b"y".to_vec()));
let ship = plan_shipment(&src, &buf, "peer-b");
let Shipment::DeltaInterval { ack_seq, .. } = &ship else {
panic!("expected delta interval, got {ship:?}");
};
assert_eq!(*ack_seq, 1);
let mut dst = DeltaOrSet::new();
dst.add(&a, b"x".to_vec());
let ack = apply_shipment(&mut dst, &ship);
assert_eq!(ack, Some(1));
assert!(dst.contains(b"y"));
}
#[test]
fn current_peer_ships_nothing() {
let a = aid("a");
let mut src = DeltaOrSet::new();
let mut buf = DeltaBuffer::new();
buf.record(src.add(&a, b"x".to_vec())); buf.ack("peer-b", 0);
let ship = plan_shipment(&src, &buf, "peer-b");
assert_eq!(ship, Shipment::UpToDate);
}
#[test]
fn delta_interval_is_smaller_than_full_state() {
let a = aid("a");
let mut src = DeltaOrSet::new();
let mut buf = DeltaBuffer::new();
for i in 0..100u32 {
let k = format!("key-{i:04}");
buf.record(src.add(&a, k.into_bytes()));
}
buf.ack("peer-b", 98);
let ship = plan_shipment(&src, &buf, "peer-b");
let full = src.wire_len();
assert!(
ship.wire_len() * 10 < full,
"delta {} not much smaller than full {}",
ship.wire_len(),
full
);
}
}