use std::collections::BTreeMap;
use std::collections::BTreeSet;
use std::hash::Hash;
use crate::Context;
use crate::cell::Source;
use crate::distributed::PeerId;
use crate::seq_crdt::SeqCrdt;
use crate::text_crdt::{OpId, TextCrdt};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
pub enum MergeMechanism {
Crdt,
Lww,
Ot,
Lease,
Custom,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct UnsupportedMechanism(pub MergeMechanism);
impl std::fmt::Display for UnsupportedMechanism {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"merge mechanism {:?} is reserved but not implemented; only `crdt` is supported",
self.0
)
}
}
impl std::error::Error for UnsupportedMechanism {}
impl MergeMechanism {
pub fn is_implemented(self) -> bool {
matches!(self, MergeMechanism::Crdt | MergeMechanism::Lww)
}
pub fn resolve(self) -> Result<Self, UnsupportedMechanism> {
if self.is_implemented() {
Ok(self)
} else {
Err(UnsupportedMechanism(self))
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct HlcStamp {
pub wall_time: u64,
pub logical: u64,
pub peer: PeerId,
}
impl HlcStamp {
fn new(wall_time: u64, logical: u64, peer: PeerId) -> Self {
Self {
wall_time,
logical,
peer,
}
}
}
#[cfg(feature = "ipc")]
impl From<HlcStamp> for crate::ipc::WireStamp {
fn from(stamp: HlcStamp) -> Self {
Self {
wall_time: stamp.wall_time,
logical: stamp.logical,
peer: stamp.peer.0,
}
}
}
#[cfg(feature = "ipc")]
impl From<crate::ipc::WireStamp> for HlcStamp {
fn from(wire: crate::ipc::WireStamp) -> Self {
Self {
wall_time: wire.wall_time,
logical: wire.logical,
peer: PeerId(wire.peer),
}
}
}
#[derive(Debug, Clone)]
pub struct Hlc {
peer: PeerId,
last_wall: u64,
last_logical: u64,
}
impl Hlc {
pub fn new(peer: PeerId) -> Self {
Self {
peer,
last_wall: 0,
last_logical: 0,
}
}
pub fn send(&mut self, now_micros: u64) -> HlcStamp {
if now_micros > self.last_wall {
self.last_wall = now_micros;
self.last_logical = 0;
} else {
self.last_logical += 1;
}
HlcStamp::new(self.last_wall, self.last_logical, self.peer)
}
pub fn recv(&mut self, remote: HlcStamp, now_micros: u64) -> HlcStamp {
let wall = self.last_wall.max(remote.wall_time).max(now_micros);
if wall == self.last_wall && wall == remote.wall_time {
self.last_logical = self.last_logical.max(remote.logical) + 1;
} else if wall == self.last_wall {
self.last_logical += 1;
} else if wall == remote.wall_time {
self.last_logical = remote.logical + 1;
} else {
self.last_logical = 0;
}
self.last_wall = wall;
HlcStamp::new(self.last_wall, self.last_logical, self.peer)
}
}
pub trait CellCrdt {
type Value;
fn merge_from(&mut self, other: &Self) -> bool;
fn value(&self) -> Self::Value;
}
pub trait RegisterCrdt: CellCrdt {
const MECHANISM: MergeMechanism;
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct LwwRegister<T> {
value: T,
stamp: HlcStamp,
}
impl<T: Clone> LwwRegister<T> {
pub fn new(value: T, stamp: HlcStamp) -> Self {
Self { value, stamp }
}
pub fn set(&mut self, value: T, stamp: HlcStamp) -> bool {
if stamp > self.stamp {
self.value = value;
self.stamp = stamp;
true
} else {
false
}
}
pub fn stamp(&self) -> HlcStamp {
self.stamp
}
}
impl<T: Clone + PartialEq> CellCrdt for LwwRegister<T> {
type Value = T;
fn merge_from(&mut self, other: &Self) -> bool {
if other.stamp > self.stamp {
let changed = self.value != other.value;
self.value = other.value.clone();
self.stamp = other.stamp;
changed
} else {
false
}
}
fn value(&self) -> T {
self.value.clone()
}
}
impl<T: Clone + PartialEq> RegisterCrdt for LwwRegister<T> {
const MECHANISM: MergeMechanism = MergeMechanism::Lww;
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct VersionVector(BTreeMap<u64, u64>);
impl VersionVector {
fn get(&self, peer: PeerId) -> u64 {
self.0.get(&peer.0).copied().unwrap_or(0)
}
fn bump(&mut self, peer: PeerId, floor: &VersionVector) {
let next = self.get(peer).max(floor.get(peer)) + 1;
self.0.insert(peer.0, next);
for (&p, &c) in &floor.0 {
let e = self.0.entry(p).or_insert(0);
*e = (*e).max(c);
}
}
fn dominates(&self, other: &VersionVector) -> bool {
other
.0
.iter()
.all(|(&p, &c)| self.0.get(&p).copied().unwrap_or(0) >= c)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct StampFrontier(BTreeMap<PeerId, HlcStamp>);
impl StampFrontier {
pub fn new() -> Self {
Self(BTreeMap::new())
}
pub fn observe(&mut self, peer: PeerId, stamp: HlcStamp) -> bool {
match self.0.get(&peer) {
Some(&cur) if cur >= stamp => false,
_ => {
self.0.insert(peer, stamp);
true
}
}
}
pub fn get(&self, peer: PeerId) -> Option<HlcStamp> {
self.0.get(&peer).copied()
}
pub fn merge(&mut self, other: &StampFrontier) -> bool {
let mut changed = false;
for (&peer, &stamp) in &other.0 {
changed |= self.observe(peer, stamp);
}
changed
}
pub fn frontier<I>(&self, membership: I) -> Option<HlcStamp>
where
I: IntoIterator<Item = PeerId>,
{
let mut min: Option<HlcStamp> = None;
for peer in membership {
let stamp = self.get(peer)?;
min = Some(match min {
Some(m) => m.min(stamp),
None => stamp,
});
}
min
}
pub fn dominates(&self, other: &StampFrontier) -> bool {
other
.0
.iter()
.all(|(peer, stamp)| self.0.get(peer).is_some_and(|cur| cur >= stamp))
}
pub fn iter(&self) -> impl Iterator<Item = (PeerId, HlcStamp)> + '_ {
self.0.iter().map(|(&peer, &stamp)| (peer, stamp))
}
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct OpIdFrontier(BTreeMap<u64, OpId>);
impl OpIdFrontier {
pub fn new() -> Self {
Self(BTreeMap::new())
}
pub fn observe(&mut self, op: OpId) -> bool {
match self.0.get(&op.peer()) {
Some(&cur) if cur >= op => false,
_ => {
self.0.insert(op.peer(), op);
true
}
}
}
pub fn get(&self, peer: u64) -> Option<OpId> {
self.0.get(&peer).copied()
}
pub fn merge(&mut self, other: &OpIdFrontier) -> bool {
let mut changed = false;
for &op in other.0.values() {
changed |= self.observe(op);
}
changed
}
pub fn frontier<I>(&self, membership: I) -> Option<OpId>
where
I: IntoIterator<Item = u64>,
{
let mut min: Option<OpId> = None;
for peer in membership {
let op = self.get(peer)?;
min = Some(match min {
Some(m) => m.min(op),
None => op,
});
}
min
}
pub fn dominates(&self, other: &OpIdFrontier) -> bool {
other
.0
.iter()
.all(|(peer, op)| self.0.get(peer).is_some_and(|cur| cur >= op))
}
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct MvRegister<T> {
entries: Vec<(VersionVector, T)>,
}
impl<T: Clone + PartialEq> MvRegister<T> {
pub fn new() -> Self {
Self {
entries: Vec::new(),
}
}
pub fn set(&mut self, value: T, peer: PeerId) -> bool {
let mut vv = VersionVector::default();
for (e, _) in &self.entries {
for (&p, &c) in &e.0 {
let slot = vv.0.entry(p).or_insert(0);
*slot = (*slot).max(c);
}
}
let mut next = VersionVector::default();
next.bump(peer, &vv);
let changed = !(self.entries.len() == 1 && self.entries[0].1 == value);
self.entries = vec![(next, value)];
changed
}
pub fn values(&self) -> Vec<T> {
self.entries.iter().map(|(_, v)| v.clone()).collect()
}
fn normalize(&mut self) {
let mut kept: Vec<(VersionVector, T)> = Vec::new();
for (vv, v) in self.entries.drain(..) {
if kept.iter().any(|(k, _)| k.dominates(&vv) && k != &vv) {
continue;
}
kept.retain(|(k, _)| !(vv.dominates(k) && k != &vv));
if !kept.iter().any(|(k, kv)| k == &vv && kv == &v) {
kept.push((vv, v));
}
}
self.entries = kept;
}
}
impl<T: Clone + PartialEq> Default for MvRegister<T> {
fn default() -> Self {
Self::new()
}
}
impl<T: Clone + PartialEq> CellCrdt for MvRegister<T> {
type Value = Vec<T>;
fn merge_from(&mut self, other: &Self) -> bool {
let before = self.values();
self.entries.extend(other.entries.iter().cloned());
self.normalize();
self.values() != before
}
fn value(&self) -> Vec<T> {
self.values()
}
}
impl<T: Clone + PartialEq> RegisterCrdt for MvRegister<T> {
const MECHANISM: MergeMechanism = MergeMechanism::Crdt;
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct PnCounter {
incr: BTreeMap<u64, u64>,
decr: BTreeMap<u64, u64>,
}
impl PnCounter {
pub fn new() -> Self {
Self::default()
}
pub fn increment(&mut self, peer: PeerId, amount: u64) {
*self.incr.entry(peer.0).or_insert(0) += amount;
}
pub fn decrement(&mut self, peer: PeerId, amount: u64) {
*self.decr.entry(peer.0).or_insert(0) += amount;
}
}
fn merge_max(into: &mut BTreeMap<u64, u64>, from: &BTreeMap<u64, u64>) {
for (&p, &c) in from {
let e = into.entry(p).or_insert(0);
*e = (*e).max(c);
}
}
impl CellCrdt for PnCounter {
type Value = i64;
fn merge_from(&mut self, other: &Self) -> bool {
let before = self.value();
merge_max(&mut self.incr, &other.incr);
merge_max(&mut self.decr, &other.decr);
self.value() != before
}
fn value(&self) -> i64 {
let p: u64 = self.incr.values().sum();
let n: u64 = self.decr.values().sum();
p as i64 - n as i64
}
}
impl RegisterCrdt for PnCounter {
const MECHANISM: MergeMechanism = MergeMechanism::Crdt;
}
pub struct ReplicatedCell<C: CellCrdt> {
crdt: C,
handle: Source<C::Value>,
}
impl<C> ReplicatedCell<C>
where
C: CellCrdt,
C::Value: PartialEq + Clone + 'static,
{
pub fn bind(ctx: &Context, crdt: C) -> Self {
let handle = ctx.source(crdt.value());
Self { crdt, handle }
}
pub fn handle(&self) -> Source<C::Value> {
self.handle
}
pub fn value(&self) -> C::Value {
self.crdt.value()
}
pub fn crdt(&self) -> &C {
&self.crdt
}
pub fn merge_remote(&mut self, ctx: &Context, remote: &C) -> bool {
if self.crdt.merge_from(remote) {
ctx.set(&self.handle, self.crdt.value());
true
} else {
false
}
}
pub fn update<F>(&mut self, ctx: &Context, mutate: F) -> bool
where
F: FnOnce(&mut C),
{
let before = self.crdt.value();
mutate(&mut self.crdt);
let after = self.crdt.value();
if after != before {
ctx.set(&self.handle, after);
true
} else {
false
}
}
}
impl<C> ReplicatedCell<C>
where
C: RegisterCrdt,
C::Value: PartialEq + Clone + 'static,
{
pub const MERGE: MergeMechanism = C::MECHANISM;
pub fn mechanism(&self) -> MergeMechanism {
C::MECHANISM
}
}
impl<T> ReplicatedCell<LwwRegister<T>>
where
T: Clone + PartialEq + 'static,
{
pub fn lww(ctx: &Context, value: T, stamp: HlcStamp) -> Self {
Self::bind(ctx, LwwRegister::new(value, stamp))
}
}
impl<T> ReplicatedCell<MvRegister<T>>
where
T: Clone + PartialEq + 'static,
{
pub fn multi_value(ctx: &Context) -> Self {
Self::bind(ctx, MvRegister::new())
}
}
impl ReplicatedCell<PnCounter> {
pub fn counter(ctx: &Context) -> Self {
Self::bind(ctx, PnCounter::new())
}
}
#[derive(Debug, Clone)]
pub struct CrdtPlane {
peer: PeerId,
clock: Hlc,
membership: BTreeSet<PeerId>,
frontier: StampFrontier,
op_frontier: OpIdFrontier,
}
impl CrdtPlane {
pub fn new(peer: PeerId) -> Self {
let mut membership = BTreeSet::new();
membership.insert(peer);
Self {
peer,
clock: Hlc::new(peer),
membership,
frontier: StampFrontier::new(),
op_frontier: OpIdFrontier::new(),
}
}
pub fn peer(&self) -> PeerId {
self.peer
}
pub fn add_peer(&mut self, peer: PeerId) {
self.membership.insert(peer);
}
pub fn membership(&self) -> impl Iterator<Item = PeerId> + '_ {
self.membership.iter().copied()
}
pub fn tick(&mut self, now_micros: u64) -> HlcStamp {
let stamp = self.clock.send(now_micros);
self.frontier.observe(self.peer, stamp);
stamp
}
pub fn observe_remote(&mut self, remote: HlcStamp, now_micros: u64) -> HlcStamp {
self.membership.insert(remote.peer);
self.frontier.observe(remote.peer, remote);
self.clock.recv(remote, now_micros)
}
pub fn stability_frontier(&self) -> Option<HlcStamp> {
self.frontier.frontier(self.membership.iter().copied())
}
pub fn frontier(&self) -> &StampFrontier {
&self.frontier
}
pub fn is_collectable(&self, stamp: HlcStamp) -> bool {
self.stability_frontier()
.is_some_and(|frontier| stamp <= frontier)
}
pub fn gc_seq<Id, V>(&self, seq: &mut SeqCrdt<Id, V>) -> usize
where
Id: Eq + Hash + Clone,
V: Clone + PartialEq,
{
match self.stability_frontier() {
Some(frontier) => seq.gc(frontier),
None => 0,
}
}
pub fn observe_op(&mut self, op: OpId) {
self.membership.insert(PeerId(op.peer()));
self.op_frontier.observe(op);
}
pub fn op_frontier(&self) -> &OpIdFrontier {
&self.op_frontier
}
pub fn op_stability_frontier(&self) -> Option<OpId> {
self.op_frontier
.frontier(self.membership.iter().map(|p| p.0))
}
pub fn is_op_collectable(&self, op: OpId) -> bool {
self.op_stability_frontier()
.is_some_and(|frontier| op <= frontier)
}
pub fn gc_text(&self, text: &mut TextCrdt) -> usize {
match self.op_stability_frontier() {
Some(frontier) => text.gc_with(|op| op <= frontier),
None => 0,
}
}
}
#[derive(Debug, Clone)]
pub struct OpLog<Op> {
ops: BTreeMap<HlcStamp, Op>,
frontier: StampFrontier,
}
impl<Op> Default for OpLog<Op> {
fn default() -> Self {
Self {
ops: BTreeMap::new(),
frontier: StampFrontier::new(),
}
}
}
impl<Op: Clone> OpLog<Op> {
pub fn new() -> Self {
Self::default()
}
pub fn record(&mut self, stamp: HlcStamp, op: Op) -> bool {
if self.ops.contains_key(&stamp) {
return false;
}
self.frontier.observe(stamp.peer, stamp);
self.ops.insert(stamp, op);
true
}
pub fn frontier(&self) -> &StampFrontier {
&self.frontier
}
pub fn missing_since(&self, since: &StampFrontier) -> Vec<(HlcStamp, Op)> {
self.ops
.iter()
.filter(|(stamp, _)| match since.get(stamp.peer) {
Some(seen) => **stamp > seen,
None => true,
})
.map(|(stamp, op)| (*stamp, op.clone()))
.collect()
}
pub fn apply_remote<I, F>(&mut self, ops: I, mut apply: F) -> usize
where
I: IntoIterator<Item = (HlcStamp, Op)>,
F: FnMut(&HlcStamp, &Op),
{
let mut incoming: Vec<(HlcStamp, Op)> = ops.into_iter().collect();
incoming.sort_by_key(|(stamp, _)| *stamp);
let mut applied = 0;
for (stamp, op) in incoming {
if self.ops.contains_key(&stamp) {
continue;
}
apply(&stamp, &op);
self.frontier.observe(stamp.peer, stamp);
self.ops.insert(stamp, op);
applied += 1;
}
applied
}
pub fn len(&self) -> usize {
self.ops.len()
}
pub fn is_empty(&self) -> bool {
self.ops.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
use proptest::prelude::*;
fn peer(n: u64) -> PeerId {
PeerId(n)
}
#[test]
fn merge_mechanism_crdt_and_lww_are_implemented() {
for m in [MergeMechanism::Crdt, MergeMechanism::Lww] {
assert!(m.is_implemented());
assert_eq!(m.resolve(), Ok(m));
}
for m in [
MergeMechanism::Ot,
MergeMechanism::Lease,
MergeMechanism::Custom,
] {
assert!(!m.is_implemented());
assert_eq!(m.resolve(), Err(UnsupportedMechanism(m)));
}
}
#[test]
fn hlc_send_is_monotonic_and_recv_observes_remote() {
let mut a = Hlc::new(peer(1));
let s1 = a.send(100);
let s2 = a.send(100); assert!(s2 > s1);
let s3 = a.send(50); assert!(s3 > s2);
let mut b = Hlc::new(peer(2));
let remote = a.send(200);
let got = b.recv(remote, 10);
assert!(
got > remote,
"recv must move past the observed remote stamp"
);
}
#[test]
fn lww_register_keeps_highest_stamp_and_merge_is_commutative_idempotent() {
let s_lo = HlcStamp::new(10, 0, peer(1));
let s_hi = HlcStamp::new(20, 0, peer(2));
let lo = LwwRegister::new("lo", s_lo);
let hi = LwwRegister::new("hi", s_hi);
let mut a = lo.clone();
a.merge_from(&hi);
let mut b = hi.clone();
b.merge_from(&lo);
assert_eq!(a.value(), "hi");
assert_eq!(b.value(), "hi", "merge is commutative");
assert!(!a.merge_from(&hi));
assert_eq!(a.value(), "hi");
}
#[test]
fn mv_register_surfaces_concurrent_writes_and_collapses_on_causal_write() {
let mut r1: MvRegister<&str> = MvRegister::new();
r1.set("from-1", peer(1));
let mut r2: MvRegister<&str> = MvRegister::new();
r2.set("from-2", peer(2));
let mut merged = r1.clone();
merged.merge_from(&r2);
let mut vals = merged.values();
vals.sort();
assert_eq!(
vals,
vec!["from-1", "from-2"],
"concurrent writes both survive"
);
let mut other = r2.clone();
other.merge_from(&r1);
let mut ov = other.values();
ov.sort();
assert_eq!(ov, vals);
assert!(!merged.merge_from(&r2));
merged.set("resolved", peer(1));
assert_eq!(merged.values(), vec!["resolved"]);
}
#[test]
fn pn_counter_merges_by_per_peer_max() {
let mut a = PnCounter::new();
a.increment(peer(1), 5);
a.decrement(peer(1), 2);
let mut b = PnCounter::new();
b.increment(peer(2), 3);
b.increment(peer(1), 5);
let mut m1 = a.clone();
m1.merge_from(&b);
let mut m2 = b.clone();
m2.merge_from(&a);
assert_eq!(m1.value(), 6);
assert_eq!(m2.value(), 6, "commutative");
assert!(!m1.merge_from(&b), "idempotent");
}
#[test]
fn replicated_cell_ingress_merge_recomputes_derived_and_suppresses_equal() {
use std::cell::Cell as StdCell;
use std::rc::Rc;
let ctx = Context::new();
let mut replica =
ReplicatedCell::bind(&ctx, LwwRegister::new(1i32, HlcStamp::new(1, 0, peer(1))));
assert_eq!(
ReplicatedCell::<LwwRegister<i32>>::MERGE,
MergeMechanism::Lww,
"an LwwRegister-backed cell declares the lww mechanism"
);
let handle = replica.handle();
let recomputes = Rc::new(StdCell::new(0usize));
let rc = recomputes.clone();
let doubled = ctx.computed(move |ctx| {
rc.set(rc.get() + 1);
ctx.get(&handle) * 2
});
assert_eq!(ctx.get(&doubled), 2);
assert_eq!(recomputes.get(), 1);
let remote = LwwRegister::new(10i32, HlcStamp::new(5, 0, peer(2)));
assert!(replica.merge_remote(&ctx, &remote));
assert_eq!(ctx.get(&doubled), 20);
assert_eq!(recomputes.get(), 2);
let stale = LwwRegister::new(99i32, HlcStamp::new(2, 0, peer(3)));
assert!(!replica.merge_remote(&ctx, &stale));
assert_eq!(ctx.get(&doubled), 20);
assert_eq!(recomputes.get(), 2, "losing merge invalidates nothing");
assert!(!replica.merge_remote(&ctx, &remote));
assert_eq!(recomputes.get(), 2);
}
#[test]
fn replicated_cells_converge_regardless_of_merge_order() {
let ctx_a = Context::new();
let ctx_b = Context::new();
let mut a =
ReplicatedCell::bind(&ctx_a, LwwRegister::new(0i32, HlcStamp::new(0, 0, peer(1))));
let mut b =
ReplicatedCell::bind(&ctx_b, LwwRegister::new(0i32, HlcStamp::new(0, 0, peer(2))));
a.update(&ctx_a, |c| {
c.set(7, HlcStamp::new(10, 0, peer(1)));
});
b.update(&ctx_b, |c| {
c.set(9, HlcStamp::new(20, 0, peer(2)));
});
let a_state = a.crdt().clone();
let b_state = b.crdt().clone();
a.merge_remote(&ctx_a, &b_state);
b.merge_remote(&ctx_b, &a_state);
assert_eq!(a.value(), b.value());
assert_eq!(a.value(), 9, "highest HLC stamp wins on both replicas");
assert_eq!(ctx_a.get(&a.handle()), ctx_b.get(&b.handle()));
}
#[test]
fn stamp_frontier_keeps_per_peer_max() {
let mut f = StampFrontier::new();
assert!(f.is_empty());
assert!(f.observe(peer(1), HlcStamp::new(10, 0, peer(1))));
assert!(f.observe(peer(1), HlcStamp::new(20, 0, peer(1))));
assert!(!f.observe(peer(1), HlcStamp::new(15, 0, peer(1))));
assert!(!f.observe(peer(1), HlcStamp::new(20, 0, peer(1))));
assert_eq!(f.get(peer(1)), Some(HlcStamp::new(20, 0, peer(1))));
assert_eq!(f.get(peer(2)), None);
assert_eq!(f.len(), 1);
}
#[test]
fn stamp_frontier_merge_is_commutative_and_idempotent() {
let a_stamp = HlcStamp::new(30, 0, peer(1));
let b_stamp = HlcStamp::new(40, 1, peer(2));
let mut left = StampFrontier::new();
left.observe(peer(1), a_stamp);
let mut right = StampFrontier::new();
right.observe(peer(2), b_stamp);
let mut lr = left.clone();
lr.merge(&right);
let mut rl = right.clone();
rl.merge(&left);
assert_eq!(lr, rl, "merge is commutative");
assert!(!lr.merge(&right));
assert!(!lr.merge(&left));
assert_eq!(lr.get(peer(1)), Some(a_stamp));
assert_eq!(lr.get(peer(2)), Some(b_stamp));
}
#[test]
fn stamp_frontier_is_min_over_membership_and_none_until_all_seen() {
let mut f = StampFrontier::new();
let s1 = HlcStamp::new(50, 0, peer(1));
let s2 = HlcStamp::new(40, 0, peer(2));
let members = [peer(1), peer(2)];
assert_eq!(f.frontier(std::iter::empty()), None);
f.observe(peer(1), s1);
assert_eq!(f.frontier(members), None);
f.observe(peer(2), s2);
assert_eq!(f.frontier(members), Some(s2));
}
#[test]
fn stamp_frontier_dominates() {
let mut bigger = StampFrontier::new();
bigger.observe(peer(1), HlcStamp::new(20, 0, peer(1)));
bigger.observe(peer(2), HlcStamp::new(30, 0, peer(2)));
let mut smaller = StampFrontier::new();
smaller.observe(peer(1), HlcStamp::new(10, 0, peer(1)));
assert!(bigger.dominates(&smaller));
assert!(!smaller.dominates(&bigger));
assert!(bigger.dominates(&bigger), "dominance is reflexive");
}
#[test]
fn crdt_plane_tick_advances_self_frontier() {
let mut plane = CrdtPlane::new(peer(1));
assert_eq!(plane.stability_frontier(), None);
let s1 = plane.tick(100);
let s2 = plane.tick(200);
assert!(s2 > s1, "ticks produce monotonically increasing stamps");
assert_eq!(plane.stability_frontier(), Some(s2));
assert_eq!(plane.frontier().get(peer(1)), Some(s2));
}
#[test]
fn crdt_plane_frontier_withheld_until_every_member_seen() {
let mut plane = CrdtPlane::new(peer(1));
plane.add_peer(peer(2));
plane.tick(100);
assert_eq!(plane.stability_frontier(), None);
plane.observe_remote(HlcStamp::new(50, 0, peer(2)), 110);
assert_eq!(
plane.stability_frontier(),
Some(HlcStamp::new(50, 0, peer(2))),
"frontier is the minimum across both members"
);
assert_eq!(
plane.membership().collect::<Vec<_>>(),
vec![peer(1), peer(2)]
);
}
#[test]
fn crdt_plane_observe_remote_advances_local_clock() {
let mut plane = CrdtPlane::new(peer(1));
plane.observe_remote(HlcStamp::new(1_000, 5, peer(2)), 100);
let local = plane.tick(200);
assert!(
local > HlcStamp::new(1_000, 5, peer(2)),
"local clock dominates the observed remote causal past"
);
}
#[test]
fn registers_declare_their_merge_mechanism() {
assert_eq!(
<LwwRegister<i32> as RegisterCrdt>::MECHANISM,
MergeMechanism::Lww
);
assert_eq!(
<MvRegister<i32> as RegisterCrdt>::MECHANISM,
MergeMechanism::Crdt
);
assert_eq!(<PnCounter as RegisterCrdt>::MECHANISM, MergeMechanism::Crdt);
}
#[test]
fn replicated_cell_constructors_pick_the_right_mechanism() {
let ctx = Context::new();
let lww = ReplicatedCell::lww(&ctx, 7i32, HlcStamp::new(1, 0, peer(1)));
assert_eq!(lww.mechanism(), MergeMechanism::Lww);
assert_eq!(lww.value(), 7);
let mv: ReplicatedCell<MvRegister<i32>> = ReplicatedCell::multi_value(&ctx);
assert_eq!(mv.mechanism(), MergeMechanism::Crdt);
assert!(mv.value().is_empty());
let pn = ReplicatedCell::counter(&ctx);
assert_eq!(pn.mechanism(), MergeMechanism::Crdt);
assert_eq!(pn.value(), 0);
}
#[test]
fn lww_constructor_drives_the_reactive_graph() {
let ctx = Context::new();
let mut cell = ReplicatedCell::lww(&ctx, 1i32, HlcStamp::new(1, 0, peer(1)));
let doubled = {
let h = cell.handle();
ctx.computed(move |ctx| ctx.get(&h) * 2)
};
assert_eq!(ctx.get(&doubled), 2);
cell.merge_remote(&ctx, &LwwRegister::new(10i32, HlcStamp::new(5, 0, peer(2))));
assert_eq!(ctx.get(&doubled), 20);
}
#[test]
fn pn_counter_constructor_increments_through_update() {
let ctx = Context::new();
let mut cell = ReplicatedCell::counter(&ctx);
cell.update(&ctx, |c| c.increment(peer(1), 3));
cell.update(&ctx, |c| c.decrement(peer(1), 1));
assert_eq!(cell.value(), 2);
assert_eq!(ctx.get(&cell.handle()), 2);
}
fn lww_from(parts: (u64, u64, u64)) -> LwwRegister<i32> {
let (wall, logical, p) = parts;
let v = (wall * 100 + logical * 10 + p) as i32;
LwwRegister::new(v, HlcStamp::new(wall, logical, peer(p)))
}
fn merged_lww(a: &LwwRegister<i32>, b: &LwwRegister<i32>) -> LwwRegister<i32> {
let mut out = a.clone();
out.merge_from(b);
out
}
fn pn_from(tallies: [(u64, u64); 3]) -> PnCounter {
let mut c = PnCounter::new();
for (i, (inc, dec)) in tallies.iter().enumerate() {
let p = peer(i as u64 + 1);
c.increment(p, *inc);
c.decrement(p, *dec);
}
c
}
fn merged_pn(a: &PnCounter, b: &PnCounter) -> PnCounter {
let mut out = a.clone();
out.merge_from(b);
out
}
proptest! {
#[test]
fn lww_merge_is_a_semilattice(
x in (0u64..4, 0u64..4, 1u64..4),
y in (0u64..4, 0u64..4, 1u64..4),
z in (0u64..4, 0u64..4, 1u64..4),
) {
let (a, b, c) = (lww_from(x), lww_from(y), lww_from(z));
prop_assert_eq!(merged_lww(&a, &b).value(), merged_lww(&b, &a).value());
let left = merged_lww(&merged_lww(&a, &b), &c);
let right = merged_lww(&a, &merged_lww(&b, &c));
prop_assert_eq!(left.value(), right.value());
let once = merged_lww(&a, &b);
let twice = merged_lww(&once, &b);
prop_assert_eq!(once.value(), twice.value());
prop_assert!(!once.clone().merge_from(&b), "re-merge is a no-op");
}
#[test]
fn pn_counter_merge_is_a_semilattice(
x in prop::array::uniform3((0u64..50, 0u64..50)),
y in prop::array::uniform3((0u64..50, 0u64..50)),
z in prop::array::uniform3((0u64..50, 0u64..50)),
) {
let (a, b, c) = (pn_from(x), pn_from(y), pn_from(z));
prop_assert_eq!(merged_pn(&a, &b).value(), merged_pn(&b, &a).value());
let left = merged_pn(&merged_pn(&a, &b), &c);
let right = merged_pn(&a, &merged_pn(&b, &c));
prop_assert_eq!(left.value(), right.value());
let once = merged_pn(&a, &b);
prop_assert!(!once.clone().merge_from(&b), "re-merge is a no-op");
}
#[test]
fn stamp_frontier_merge_is_a_semilattice(
xs in prop::collection::vec((1u64..5, 0u64..8), 0..6),
ys in prop::collection::vec((1u64..5, 0u64..8), 0..6),
zs in prop::collection::vec((1u64..5, 0u64..8), 0..6),
) {
let build = |obs: &[(u64, u64)]| {
let mut f = StampFrontier::new();
for &(p, w) in obs {
f.observe(peer(p), HlcStamp::new(w, 0, peer(p)));
}
f
};
let (a, b, c) = (build(&xs), build(&ys), build(&zs));
let mut ab = a.clone();
ab.merge(&b);
let mut ba = b.clone();
ba.merge(&a);
prop_assert_eq!(&ab, &ba, "commutative");
let mut left = ab.clone();
left.merge(&c);
let mut bc = b.clone();
bc.merge(&c);
let mut right = a.clone();
right.merge(&bc);
prop_assert_eq!(&left, &right, "associative");
prop_assert!(!ab.clone().merge(&b), "idempotent: re-merge changes nothing");
}
}
#[derive(Debug, Clone, PartialEq)]
struct LwwWrite(i32);
#[test]
fn op_log_record_is_idempotent_and_tracks_frontier() {
let mut log: OpLog<LwwWrite> = OpLog::new();
assert!(log.is_empty());
let s1 = HlcStamp::new(10, 0, peer(1));
assert!(log.record(s1, LwwWrite(1)));
assert!(!log.record(s1, LwwWrite(1)));
assert_eq!(log.len(), 1);
assert_eq!(log.frontier().get(peer(1)), Some(s1));
}
#[test]
fn op_log_missing_since_returns_only_unseen_ops_in_order() {
let mut log: OpLog<LwwWrite> = OpLog::new();
let a1 = HlcStamp::new(10, 0, peer(1));
let a2 = HlcStamp::new(20, 0, peer(1));
let b1 = HlcStamp::new(15, 0, peer(2));
log.record(a2, LwwWrite(2));
log.record(a1, LwwWrite(1));
log.record(b1, LwwWrite(3));
let mut since = StampFrontier::new();
since.observe(peer(1), a1);
let missing: Vec<HlcStamp> = log
.missing_since(&since)
.into_iter()
.map(|(s, _)| s)
.collect();
assert_eq!(missing, vec![b1, a2]);
let full = log.frontier().clone();
assert!(log.missing_since(&full).is_empty());
}
#[test]
fn op_log_apply_remote_is_idempotent_and_causally_ordered() {
let mut log: OpLog<LwwWrite> = OpLog::new();
let s1 = HlcStamp::new(10, 0, peer(2));
let s2 = HlcStamp::new(20, 0, peer(2));
let mut order = Vec::new();
let applied = log.apply_remote(
vec![(s2, LwwWrite(2)), (s1, LwwWrite(1)), (s2, LwwWrite(2))],
|stamp, op| order.push((*stamp, op.clone())),
);
assert_eq!(applied, 2, "the duplicate s2 is skipped");
assert_eq!(
order,
vec![(s1, LwwWrite(1)), (s2, LwwWrite(2))],
"applied in causal order"
);
let again = log.apply_remote(vec![(s1, LwwWrite(1)), (s2, LwwWrite(2))], |_, _| {
panic!("must not re-apply a stored op")
});
assert_eq!(again, 0);
}
#[test]
fn two_replicas_converge_through_anti_entropy_exchange() {
let mut hlc_a = Hlc::new(peer(1));
let mut hlc_b = Hlc::new(peer(2));
let mut log_a: OpLog<LwwWrite> = OpLog::new();
let mut log_b: OpLog<LwwWrite> = OpLog::new();
let mut reg_a = LwwRegister::new(0i32, HlcStamp::new(0, 0, peer(1)));
let mut reg_b = LwwRegister::new(0i32, HlcStamp::new(0, 0, peer(2)));
let sa = hlc_a.send(100);
reg_a.set(7, sa);
log_a.record(sa, LwwWrite(7));
let sb = hlc_b.send(200);
reg_b.set(9, sb);
log_b.record(sb, LwwWrite(9));
let to_a = log_b.missing_since(log_a.frontier());
log_a.apply_remote(to_a, |stamp, op| {
reg_a.set(op.0, *stamp);
});
let to_b = log_a.missing_since(log_b.frontier());
log_b.apply_remote(to_b, |stamp, op| {
reg_b.set(op.0, *stamp);
});
assert_eq!(reg_a.value(), reg_b.value(), "replicas converge");
assert_eq!(
reg_a.value(),
9,
"highest HLC stamp (peer 2 @ wall 200) wins"
);
assert!(log_b.missing_since(log_a.frontier()).is_empty());
assert!(log_a.missing_since(log_b.frontier()).is_empty());
}
#[test]
fn op_log_drives_a_pn_counter_to_convergence() {
let mut log_a: OpLog<u64> = OpLog::new();
let mut log_b: OpLog<u64> = OpLog::new();
let mut ctr_a = PnCounter::new();
let mut ctr_b = PnCounter::new();
let sa = HlcStamp::new(100, 0, peer(1));
ctr_a.increment(peer(1), 5);
log_a.record(sa, 5);
let sb = HlcStamp::new(110, 0, peer(2));
ctr_b.increment(peer(2), 3);
log_b.record(sb, 3);
for (stamp, amt) in log_b.missing_since(log_a.frontier()) {
log_a.apply_remote([(stamp, amt)], |s, a| ctr_a.increment(s.peer, *a));
}
for (stamp, amt) in log_a.missing_since(log_b.frontier()) {
log_b.apply_remote([(stamp, amt)], |s, a| ctr_b.increment(s.peer, *a));
}
assert_eq!(ctr_a.value(), 8);
assert_eq!(ctr_b.value(), 8, "both replicas reach the summed value");
}
#[test]
fn is_collectable_tracks_the_stability_frontier() {
let mut plane = CrdtPlane::new(peer(1));
plane.add_peer(peer(2));
let del = HlcStamp::new(200, 0, peer(1));
plane.tick(250); assert!(
!plane.is_collectable(del),
"unseen member withholds collectability"
);
plane.observe_remote(HlcStamp::new(150, 0, peer(2)), 260);
assert!(
!plane.is_collectable(del),
"a member behind the delete still withholds GC"
);
plane.observe_remote(HlcStamp::new(300, 0, peer(2)), 310);
assert!(plane.is_collectable(del));
}
#[test]
fn gc_seq_collects_a_tombstone_only_once_every_replica_has_observed_it() {
let mut plane = CrdtPlane::new(peer(1));
plane.add_peer(peer(2));
let mut seq: SeqCrdt<u64, &str> = SeqCrdt::new(peer(1));
seq.insert_back(1, "x", 100);
seq.remove(&1, 200); assert!(
!seq.contains(&1),
"removed element is tombstoned (still stored)"
);
plane.tick(250);
assert_eq!(
plane.gc_seq(&mut seq),
0,
"no GC until every member is observed"
);
plane.observe_remote(HlcStamp::new(150, 0, peer(2)), 260);
assert_eq!(
plane.gc_seq(&mut seq),
0,
"a lagging member blocks collection"
);
plane.observe_remote(HlcStamp::new(300, 0, peer(2)), 310);
assert_eq!(
plane.gc_seq(&mut seq),
1,
"tombstone collected once every replica has observed it"
);
assert_eq!(plane.gc_seq(&mut seq), 0);
}
#[test]
fn gc_seq_is_inert_for_a_single_writer_session() {
let mut plane = CrdtPlane::new(peer(1));
let mut seq: SeqCrdt<u64, &str> = SeqCrdt::new(peer(1));
seq.insert_back(1, "x", 100);
seq.remove(&1, 200);
assert_eq!(plane.gc_seq(&mut seq), 0);
plane.tick(250);
assert_eq!(plane.gc_seq(&mut seq), 1);
}
#[test]
fn op_id_frontier_takes_the_per_peer_max_and_membership_min() {
let a = TextCrdt::from_str(1, "abc");
let mut b = TextCrdt::from_str(2, "z");
let mut f = OpIdFrontier::new();
assert!(f.is_empty());
assert!(f.observe(a.clock()), "first observe advances");
assert!(
!f.observe(a.clock()),
"re-observing the same id is idempotent"
);
f.observe(b.clock());
assert_eq!(f.len(), 2);
assert_eq!(f.frontier([1, 2]), Some(b.clock().min(a.clock())));
assert_eq!(f.frontier([1, 2, 3]), None);
b.delete(0); let high = b.clock();
f.observe(high);
assert_eq!(f.get(2), Some(high));
f.observe(TextCrdt::from_str(2, "y").clock()); assert_eq!(f.get(2), Some(high), "max is sticky");
}
#[test]
fn is_op_collectable_tracks_the_op_stability_frontier() {
let mut plane = CrdtPlane::new(peer(1));
plane.add_peer(peer(2));
let mut a = TextCrdt::from_str(1, "abc");
a.delete(2); let del = a.clock();
plane.observe_op(a.clock()); assert!(
!plane.is_op_collectable(del),
"unseen member withholds collectability"
);
let behind = TextCrdt::from_str(2, "z"); plane.observe_op(behind.clock());
assert!(
!plane.is_op_collectable(del),
"a member behind the delete still withholds GC"
);
let mut ahead = behind.clone();
ahead.merge(&a); plane.observe_op(ahead.clock());
assert!(plane.is_op_collectable(del));
}
#[test]
fn gc_text_collects_a_tombstone_only_once_every_replica_has_observed_it() {
let mut plane = CrdtPlane::new(peer(1));
plane.add_peer(peer(2));
let mut a = TextCrdt::from_str(1, "abc");
let mut behind = a.fork(2); a.delete(2); assert_eq!(a.text(), "ab");
assert_eq!(a.tombstone_count(), 1);
plane.observe_op(a.clock());
assert_eq!(
plane.gc_text(&mut a),
0,
"no GC until every member is observed"
);
plane.observe_op(behind.clock());
assert_eq!(
plane.gc_text(&mut a),
0,
"a lagging member blocks collection"
);
behind.merge(&a);
plane.observe_op(behind.clock());
assert_eq!(
plane.gc_text(&mut a),
1,
"tombstone collected once every replica has observed it"
);
assert_eq!(a.tombstone_count(), 0);
assert_eq!(a.text(), "ab", "visible text is unchanged by GC");
assert_eq!(plane.gc_text(&mut a), 0);
}
#[test]
fn gc_text_is_inert_for_a_single_writer_session() {
let mut plane = CrdtPlane::new(peer(1));
let mut a = TextCrdt::from_str(1, "abc");
a.delete(2);
assert_eq!(plane.gc_text(&mut a), 0);
plane.observe_op(a.clock());
assert_eq!(plane.gc_text(&mut a), 1);
assert_eq!(a.text(), "ab");
}
}