minerva 0.2.0

Causal ordering for distributed systems
extern crate alloc;

use alloc::vec::Vec;

use crate::metis::{Dot, DotFun, DotSet, DotStore, Dotted};

/// The multi-value register over the payload store: dots are the writes,
/// values are the payloads, concurrent writes are siblings.
type Register = Dotted<DotFun<&'static str>>;

/// The write-that-supersedes flow, via the documented path: a fresh dot in the
/// store, the fresh dot plus the observed old dots in the context
/// (`try_new`), merged in. Returns the fresh dot.
fn write(reg: &mut Register, station: u32, value: &'static str) -> Dot {
    let dot = reg.next_dot(station);
    // Supersede everything the register currently carries: a register holds
    // one logical value, and its surviving dots are the siblings this write
    // replaces.
    let mut context = DotSet::new();
    let _ = context.insert(dot);
    for observed in reg.store().dots() {
        let _ = context.insert(observed);
    }
    let store = DotFun::singleton(dot, value);
    let delta = Dotted::try_new(store, context).expect("the fresh pair is covered");
    *reg = reg.merge(&delta);
    dot
}

/// The observed-clear flow: a pure-context delta of exactly the dots the
/// register currently carries, superseding them without resurrection.
fn clear_observed(reg: &mut Register) {
    let mut context = DotSet::new();
    for observed in reg.store().dots() {
        let _ = context.insert(observed);
    }
    let delta = Dotted::from_context(context);
    *reg = reg.merge(&delta);
}

#[test]
fn test_register_write_supersedes_observed() {
    // A replica writes "a", then writes "b" superseding "a"; after the merge
    // only "b" survives (the single-value read with no concurrency).
    let mut reg = Register::new();
    let d_a = write(&mut reg, 1, "a");
    assert_eq!(reg.store().get(d_a), Some(&"a"));

    let d_b = write(&mut reg, 1, "b");
    assert_eq!(reg.store().get(d_a), None); // "a" superseded
    assert_eq!(reg.store().get(d_b), Some(&"b"));
    let live: Vec<&&str> = reg.store().values().collect();
    assert_eq!(live, [&"b"]);
    // The context never forgets the superseded dot.
    assert!(reg.context().contains(d_a));
}

#[test]
fn test_register_concurrent_writes_are_siblings() {
    // Two replicas write concurrently from the same (empty) ancestor. After
    // the exchange BOTH values survive, readable via `values()`: the
    // multi-value read a collaborative editor shows as a conflict.
    let mut a = Register::new();
    let mut b = Register::new();
    let d_a = write(&mut a, 1, "alpha"); // station 1
    let d_b = write(&mut b, 2, "beta"); // station 2, concurrent

    let merged = a.merge(&b);
    assert_eq!(merged.store().get(d_a), Some(&"alpha"));
    assert_eq!(merged.store().get(d_b), Some(&"beta"));
    let mut siblings: Vec<&str> = merged.store().values().copied().collect();
    siblings.sort_unstable();
    assert_eq!(siblings, ["alpha", "beta"]);

    // A superseding write collapses the siblings to one value.
    let mut joined = merged;
    let d_c = write(&mut joined, 1, "gamma");
    assert_eq!(joined.store().get(d_a), None);
    assert_eq!(joined.store().get(d_b), None);
    let live: Vec<&&str> = joined.store().values().collect();
    assert_eq!(live, [&"gamma"]);
    assert_eq!(joined.store().get(d_c), Some(&"gamma"));
}

#[test]
fn test_register_observed_clear_empties_without_resurrect() {
    // An observed clear empties the register at every replica: a stale peer's
    // re-merge cannot resurrect the cleared value, because the clearing
    // context covers the dot with no survivor.
    let mut a = Register::new();
    let d = write(&mut a, 1, "held");

    // B learns the write, holding a stale copy.
    let b = Register::new().merge(&a);
    assert_eq!(b.store().get(d), Some(&"held"));

    // A clears what it observed; its store empties.
    clear_observed(&mut a);
    assert!(a.store().is_bottom());

    // The stale copy at B cannot bring it back on re-merge, symmetrically.
    let rejoined = a.merge(&b);
    assert!(rejoined.store().is_bottom());
    assert_eq!(rejoined, b.merge(&a));

    // A replica that never saw the write learns the clear the same way.
    let fresh = Register::new().merge(&a);
    assert!(fresh.store().is_bottom());
    assert!(fresh.context().contains(d));
}

#[test]
fn test_register_folds_converge_in_both_orders() {
    // Convergence: a history of concurrent writes and a clear over three
    // replicas folds to the same register in both orders (state and delta are
    // one type, merged in any order).
    let mut r0 = Register::new();
    let mut r1 = Register::new();
    let mut r2 = Register::new();

    let _ = write(&mut r0, 0, "zero");
    let _ = write(&mut r1, 1, "one");
    r1 = r1.merge(&r0); // r1 learns r0
    let _ = write(&mut r2, 2, "two");
    clear_observed(&mut r1); // r1 clears what it observed (zero, one)

    let replicas = [r0, r1, r2];
    let forward = replicas.iter().fold(Register::new(), |acc, r| acc.merge(r));
    let backward = replicas
        .iter()
        .rev()
        .fold(Register::new(), |acc, r| acc.merge(r));
    assert_eq!(forward, backward);

    // And the fold absorbs every replica (idempotent at the top).
    for r in &replicas {
        assert_eq!(&forward.merge(r), &forward);
    }
}