use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::HashMap;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct LwwRegister {
pub value: Value,
pub version: u64,
pub replica: String,
}
impl LwwRegister {
pub fn new(value: Value, version: u64, replica: impl Into<String>) -> Self {
Self {
value,
version,
replica: replica.into(),
}
}
fn dominates(&self, other: &LwwRegister) -> bool {
(self.version, self.replica.as_str()) > (other.version, other.replica.as_str())
}
pub fn merge(a: &LwwRegister, b: &LwwRegister) -> LwwRegister {
if a.dominates(b) {
a.clone()
} else {
b.clone()
}
}
}
pub type LwwMap = HashMap<String, LwwRegister>;
pub fn merge_maps(a: &LwwMap, b: &LwwMap) -> LwwMap {
let mut out = a.clone();
for (k, rb) in b {
out.entry(k.clone())
.and_modify(|ra| *ra = LwwRegister::merge(ra, rb))
.or_insert_with(|| rb.clone());
}
out
}
pub fn merge_many(replicas: &[LwwMap]) -> LwwMap {
let mut iter = replicas.iter();
match iter.next() {
Some(first) => iter.fold(first.clone(), |acc, r| merge_maps(&acc, r)),
None => LwwMap::new(),
}
}
pub fn materialize(m: &LwwMap) -> HashMap<String, Value> {
m.iter()
.map(|(k, r)| (k.clone(), r.value.clone()))
.collect()
}
pub fn export_lww(
snapshot: &HashMap<String, Value>,
versions: &HashMap<String, u64>,
replica: impl Into<String>,
) -> LwwMap {
let replica = replica.into();
snapshot
.iter()
.map(|(k, v)| {
let version = versions.get(k).copied().unwrap_or(0);
(
k.clone(),
LwwRegister::new(v.clone(), version, replica.clone()),
)
})
.collect()
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Claim {
pub claimant: String,
pub version: u64,
pub replica: String,
}
impl Claim {
fn earlier_than(&self, other: &Claim) -> bool {
(self.version, self.replica.as_str()) < (other.version, other.replica.as_str())
}
fn merge(a: &Claim, b: &Claim) -> Claim {
if a.earlier_than(b) {
a.clone()
} else {
b.clone()
}
}
}
pub type ClaimRegistry = HashMap<String, Claim>;
pub fn merge_claims(a: &ClaimRegistry, b: &ClaimRegistry) -> ClaimRegistry {
let mut out = a.clone();
for (task, cb) in b {
out.entry(task.clone())
.and_modify(|ca| *ca = Claim::merge(ca, cb))
.or_insert_with(|| cb.clone());
}
out
}
pub fn merge_claims_many(registries: &[ClaimRegistry]) -> ClaimRegistry {
let mut iter = registries.iter();
match iter.next() {
Some(first) => iter.fold(first.clone(), |acc, r| merge_claims(&acc, r)),
None => ClaimRegistry::new(),
}
}
pub fn claim(
registry: &mut ClaimRegistry,
task: impl Into<String>,
claimant: impl Into<String>,
version: u64,
replica: impl Into<String>,
) -> bool {
let task = task.into();
let candidate = Claim {
claimant: claimant.into(),
version,
replica: replica.into(),
};
match registry.get(&task) {
Some(existing) if !candidate.earlier_than(existing) => {
existing.claimant == candidate.claimant
}
_ => {
let owns = candidate.claimant.clone();
registry.insert(task, candidate);
let _ = owns;
true
}
}
}
pub fn owner<'a>(registry: &'a ClaimRegistry, task: &str) -> Option<&'a str> {
registry.get(task).map(|c| c.claimant.as_str())
}
pub fn tasks_claimed_by<'a>(registry: &'a ClaimRegistry, claimant: &str) -> Vec<&'a str> {
let mut tasks: Vec<&str> = registry
.iter()
.filter(|(_, c)| c.claimant == claimant)
.map(|(t, _)| t.as_str())
.collect();
tasks.sort_unstable(); tasks
}
pub fn claim_owners(registry: &ClaimRegistry) -> HashMap<String, String> {
registry
.iter()
.map(|(t, c)| (t.clone(), c.claimant.clone()))
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn reg(v: Value, ver: u64, rep: &str) -> LwwRegister {
LwwRegister::new(v, ver, rep)
}
#[test]
fn higher_version_wins() {
let a = reg(json!("a"), 1, "r1");
let b = reg(json!("b"), 2, "r1");
assert_eq!(LwwRegister::merge(&a, &b).value, json!("b"));
assert_eq!(LwwRegister::merge(&b, &a).value, json!("b")); }
#[test]
fn replica_breaks_version_ties_deterministically() {
let a = reg(json!("a"), 5, "r1");
let b = reg(json!("b"), 5, "r2"); assert_eq!(LwwRegister::merge(&a, &b).value, json!("b"));
assert_eq!(LwwRegister::merge(&b, &a).value, json!("b"));
}
#[test]
fn merge_is_idempotent() {
let a = reg(json!(1), 3, "r1");
assert_eq!(LwwRegister::merge(&a, &a), a);
}
fn map(entries: &[(&str, Value, u64, &str)]) -> LwwMap {
entries
.iter()
.map(|(k, v, ver, rep)| (k.to_string(), reg(v.clone(), *ver, rep)))
.collect()
}
#[test]
fn divergent_replicas_converge() {
let a = map(&[("x", json!(1), 2, "r1"), ("shared", json!("a"), 1, "r1")]);
let b = map(&[("y", json!(2), 1, "r2"), ("shared", json!("b"), 3, "r2")]);
let ab = merge_maps(&a, &b);
let ba = merge_maps(&b, &a);
assert_eq!(ab, ba, "merge must be commutative");
assert_eq!(ab.get("shared").unwrap().value, json!("b"));
assert_eq!(ab.get("x").unwrap().value, json!(1));
assert_eq!(ab.get("y").unwrap().value, json!(2));
}
#[test]
fn merge_many_is_order_independent() {
let r1 = map(&[("k", json!("one"), 1, "r1")]);
let r2 = map(&[("k", json!("two"), 2, "r2")]);
let r3 = map(&[("k", json!("three"), 3, "r3")]);
let forward = merge_many(&[r1.clone(), r2.clone(), r3.clone()]);
let reverse = merge_many(&[r3, r2, r1]);
assert_eq!(forward, reverse);
assert_eq!(forward.get("k").unwrap().value, json!("three")); }
#[test]
fn associativity() {
let a = map(&[("k", json!("a"), 1, "r1")]);
let b = map(&[("k", json!("b"), 2, "r2")]);
let c = map(&[("k", json!("c"), 2, "r3")]);
let left = merge_maps(&merge_maps(&a, &b), &c);
let right = merge_maps(&a, &merge_maps(&b, &c));
assert_eq!(left, right);
}
#[test]
fn materialize_drops_tags() {
let m = map(&[("k", json!(42), 1, "r1")]);
let plain = materialize(&m);
assert_eq!(plain.get("k"), Some(&json!(42)));
}
#[test]
fn export_tags_snapshot_with_versions_and_replica() {
let snapshot: HashMap<String, Value> =
[("a".to_string(), json!(1)), ("b".to_string(), json!(2))].into();
let versions: HashMap<String, u64> = [("a".to_string(), 5)].into(); let m = export_lww(&snapshot, &versions, "dev1");
assert_eq!(m["a"].version, 5);
assert_eq!(m["a"].replica, "dev1");
assert_eq!(m["b"].version, 0);
assert_eq!(m["b"].value, json!(2));
}
#[test]
fn export_then_merge_round_trip() {
let snap_a: HashMap<String, Value> = [("k".to_string(), json!("a"))].into();
let snap_b: HashMap<String, Value> = [("k".to_string(), json!("b"))].into();
let dev_a = export_lww(&snap_a, &[("k".to_string(), 1)].into(), "A");
let dev_b = export_lww(&snap_b, &[("k".to_string(), 2)].into(), "B");
let merged = merge_many(&[dev_a, dev_b]);
assert_eq!(materialize(&merged).get("k"), Some(&json!("b")));
}
#[test]
fn claim_is_first_wins_and_idempotent() {
let mut reg = ClaimRegistry::new();
assert!(claim(&mut reg, "t", "agent-1", 5, "r1")); assert_eq!(owner(®, "t"), Some("agent-1"));
assert!(!claim(&mut reg, "t", "agent-2", 9, "r2"));
assert_eq!(owner(®, "t"), Some("agent-1"));
assert!(claim(&mut reg, "t", "agent-1", 5, "r1"));
assert_eq!(owner(®, "t"), Some("agent-1"));
}
#[test]
fn tasks_claimed_by_and_owners() {
let mut reg = ClaimRegistry::new();
claim(&mut reg, "t1", "a1", 1, "r1");
claim(&mut reg, "t2", "a1", 1, "r1");
claim(&mut reg, "t3", "a2", 1, "r2");
assert_eq!(tasks_claimed_by(®, "a1"), vec!["t1", "t2"]);
let owners = claim_owners(®);
assert_eq!(owners.get("t3"), Some(&"a2".to_string()));
}
#[test]
fn merge_claims_many_resolves_one_owner_per_task() {
let mut a = ClaimRegistry::new();
claim(&mut a, "t", "a1", 2, "r1");
let mut b = ClaimRegistry::new();
claim(&mut b, "t", "a2", 1, "r2"); let merged = merge_claims_many(&[a, b]);
assert_eq!(owner(&merged, "t"), Some("a2"));
}
#[test]
fn first_claim_wins_and_converges() {
let a: ClaimRegistry = [(
"t".to_string(),
Claim {
claimant: "agent-1".into(),
version: 1,
replica: "r1".into(),
},
)]
.into();
let b: ClaimRegistry = [(
"t".to_string(),
Claim {
claimant: "agent-2".into(),
version: 2,
replica: "r2".into(),
},
)]
.into();
let ab = merge_claims(&a, &b);
let ba = merge_claims(&b, &a);
assert_eq!(ab, ba, "claim merge must be commutative");
assert_eq!(ab.get("t").unwrap().claimant, "agent-1");
}
}