minerva 0.2.0

Causal ordering for distributed systems
extern crate alloc;

use super::super::support::drain;
use super::support::producer;
use crate::kairos::{Clock, TickCounter, VirtualTimeSource};
use crate::metis::{CausalIdeal, Producer};
use alloc::sync::Arc;
use alloc::vec::Vec;

#[test]
fn test_successive_produces_are_monotonic() {
    let mut p = producer(1);
    let first = p.produce(0u16, ());
    assert_eq!(first.deps.get(1), 1);
    let second = p.produce(0u16, ());
    assert!(second.stamp > first.stamp);
    assert_eq!(second.deps.get(1), 2);
}

#[test]
fn test_observe_advances_knowledge_and_stamp() {
    let mut a = producer(1);
    let ea = a.produce(0u16, ());
    let mut b = producer(2);
    b.observe(&ea);
    assert_eq!(b.knowledge().get(1), 1);
    let eb = b.produce(0u16, ());
    assert!(eb.stamp > ea.stamp);
    assert_eq!(eb.deps.get(1), 1);
    assert_eq!(eb.deps.get(2), 1);
}

#[test]
fn test_concurrent_first_events() {
    let mut a = producer(1);
    let mut b = producer(2);
    let ea = a.produce(0u16, 0u32);
    let eb = b.produce(0u16, 1u32);
    assert!(ea.deps.concurrent(&eb.deps));
    assert_eq!(ea.deps.get(1), 1);
    assert_eq!(eb.deps.get(2), 1);

    let mut buf = CausalIdeal::new();
    buf.insert(ea);
    buf.insert(eb);
    assert_eq!(drain(&mut buf).len(), 2);
}

#[test]
fn test_producer_feeds_buffer_end_to_end() {
    let mut a = producer(1);
    let mut b = producer(2);

    let a0 = a.produce(0u16, 0u32);
    b.observe(&a0);
    let b0 = b.produce(0u16, 1u32);
    a.observe(&b0);
    let a1 = a.produce(0u16, 2u32);

    assert_eq!(b0.deps.get(1), 1);
    assert_eq!(a1.deps.get(2), 1);
    assert!(a0.stamp < b0.stamp && b0.stamp < a1.stamp);

    let mut buf = CausalIdeal::new();
    buf.insert(a1);
    buf.insert(b0);
    buf.insert(a0);
    assert_eq!(drain(&mut buf), (0u32..3).collect::<Vec<u32>>());
}

#[test]
fn test_self_echo_ticks_clock() {
    let frozen = || {
        let clock = Clock::with_default_config(VirtualTimeSource::new(100), 1).unwrap();
        Producer::new(clock)
    };

    let mut with_echo = frozen();
    let e0 = with_echo.produce(0u16, ());
    let knowledge_before = with_echo.knowledge().clone();
    with_echo.observe(&e0);
    assert_eq!(with_echo.knowledge(), &knowledge_before);
    let echoed = with_echo.produce(0u16, ());

    let mut no_echo = frozen();
    let _ = no_echo.produce(0u16, ());
    let plain = no_echo.produce(0u16, ());

    assert!(echoed.stamp > plain.stamp);
    assert_eq!(echoed.stamp.logical(), plain.stamp.logical() + 1);
}

#[test]
fn test_shared_clock_is_the_single_high_water() {
    let clock = Arc::new(Clock::with_default_config(VirtualTimeSource::new(100), 3).unwrap());
    let mut p = Producer::new(Arc::clone(&clock));
    assert_eq!(p.station_id(), 3);

    // A sibling stamping site (a seal, a ledger append) and the producer share
    // one high-water word: stamps interleave strictly under frozen time.
    let seal = clock.now(0u16);
    let event = p.produce(0u16, ());
    let later_seal = clock.now(0u16);
    assert!(seal < event.stamp && event.stamp < later_seal);
}

#[test]
fn test_borrowed_clock_producer_composes() {
    let clock = Clock::with_default_config(TickCounter::new(), 4).unwrap();
    let mut p = Producer::new(&clock);
    let event = p.produce(0u16, ());
    assert_eq!(event.deps.get(4), 1);
    drop(p);
    // The clock outlives the producer and still dominates its output.
    assert!(clock.now(0u16) > event.stamp);
}

#[test]
fn test_try_observe_within_bound_is_observe() {
    let mut a = producer(1);
    let ea = a.produce(0u16, ());
    let mut b = producer(2);
    b.try_observe(&ea, u64::MAX).unwrap();
    assert_eq!(b.knowledge().get(1), 1);
    let eb = b.produce(0u16, ());
    assert!(eb.stamp > ea.stamp);
}

#[test]
fn test_try_observe_rejection_leaves_no_partial_trace() {
    let mut far =
        Producer::new(Clock::with_default_config(VirtualTimeSource::new(10_000), 1).unwrap());
    let remote = far.produce(0u16, ());

    let mut b = Producer::new(Clock::with_default_config(VirtualTimeSource::new(100), 2).unwrap());
    let err = b.try_observe(&remote, 1_000).unwrap_err();
    assert_eq!(err.observed_forward_skew, 9_900);
    // All-or-nothing: neither knowledge nor clock carries the rejected event.
    assert_eq!(b.knowledge().get(1), 0);
    let minted = b.produce(0u16, ());
    assert_eq!(minted.stamp.physical(), 100);
    assert!(minted.stamp < remote.stamp);
}