minerva 0.2.0

Causal ordering for distributed systems
extern crate alloc;

use crate::kairos::Kairos;
use crate::metis::{BoundedInsert, CausalIdeal, Event, VersionVector};
use alloc::vec::Vec;
use core::num::NonZeroUsize;

#[test]
fn test_try_insert_buffers_below_capacity() {
    let mut v = VersionVector::new();
    let _ = v.increment(1);
    let event = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: v,
        payload: 7u32,
    };

    let mut buf = CausalIdeal::new();
    let cap = NonZeroUsize::new(4).unwrap();
    assert!(buf.try_insert(event, cap).is_ok());
    assert_eq!(buf.pending_len(), 1);
    assert_eq!(buf.pop_ready(), Some(7));
}

#[test]
fn test_try_insert_rejects_at_capacity_unchanged() {
    let cap = NonZeroUsize::new(1).unwrap();
    let mut v1 = VersionVector::new();
    let _ = v1.increment(1);
    let _ = v1.increment(1);
    let blocker = Event {
        stamp: Kairos::new(11, 0, 1, 0u16),
        deps: v1,
        payload: 1u32,
    };
    let mut v2 = VersionVector::new();
    let _ = v2.increment(2);
    let newcomer = Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: v2,
        payload: 2u32,
    };

    let mut buf = CausalIdeal::new();
    assert!(buf.try_insert(blocker, cap).is_ok());
    let rejected = buf.try_insert(newcomer, cap).unwrap_err();
    assert_eq!(rejected.event.payload, 2);
    assert_eq!(rejected.capacity, cap);
    assert_eq!(buf.pending_len(), 1);
}

#[test]
fn test_try_insert_drops_stale_without_consuming_capacity() {
    let cap = NonZeroUsize::new(1).unwrap();
    let mut v = VersionVector::new();
    let _ = v.increment(1);
    let event = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: v,
        payload: 0u32,
    };

    let mut buf = CausalIdeal::new();
    buf.insert(event.clone());
    assert_eq!(buf.pop_ready(), Some(0));
    assert!(buf.try_insert(event, cap).is_ok());
    assert_eq!(buf.pending_len(), 0);
}

#[test]
fn test_try_insert_reclaims_dead_twin_when_full() {
    let mut v = VersionVector::new();
    let _ = v.increment(1);
    let original = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: v,
        payload: 0u32,
    };
    let duplicate = original.clone();

    let mut buf = CausalIdeal::new();
    buf.insert(original);
    buf.insert(duplicate);
    assert_eq!(buf.pop_ready(), Some(0));
    assert_eq!(buf.pending_len(), 1);

    let cap = NonZeroUsize::new(1).unwrap();
    let mut v2 = VersionVector::new();
    let _ = v2.increment(2);
    let fresh = Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: v2,
        payload: 9u32,
    };
    assert!(buf.try_insert(fresh, cap).is_ok());
    assert_eq!(buf.pending_len(), 1);
    assert_eq!(buf.pop_ready(), Some(9));
}

#[test]
fn test_try_insert_evicting_where_replaces_without_advancing_progress() {
    let cap = NonZeroUsize::new(2).unwrap();
    let mut buf = CausalIdeal::new();
    for (physical, station, payload) in [(20, 2, 2u32), (10, 1, 1)] {
        let mut deps = VersionVector::new();
        let _ = deps.increment(station);
        let _ = deps.increment(station);
        buf.try_insert(
            Event {
                stamp: Kairos::new(physical, 0, station, 0u16),
                deps,
                payload,
            },
            cap,
        )
        .unwrap();
    }

    let mut deps = VersionVector::new();
    let _ = deps.increment(3);
    let outcome = buf
        .try_insert_evicting_where(
            Event {
                stamp: Kairos::new(30, 0, 3, 0u16),
                deps,
                payload: 3,
            },
            cap,
            |event| event.payload <= 2,
        )
        .unwrap();

    let BoundedInsert::Evicted { evicted } = outcome else {
        panic!("the full buffer must replace one approved victim");
    };
    assert_eq!(evicted.payload, 1, "the earliest approved victim wins");
    assert_eq!(buf.pending_len(), 2);
    assert_eq!(buf.delivered(), &VersionVector::new());
}

#[test]
fn test_try_insert_evicting_where_refuses_without_an_approved_victim() {
    let cap = NonZeroUsize::new(1).unwrap();
    let mut blocked = VersionVector::new();
    let _ = blocked.increment(1);
    let _ = blocked.increment(1);
    let mut buf = CausalIdeal::new();
    buf.try_insert(
        Event {
            stamp: Kairos::new(10, 0, 1, 0u16),
            deps: blocked,
            payload: 1u32,
        },
        cap,
    )
    .unwrap();

    let mut fresh = VersionVector::new();
    let _ = fresh.increment(2);
    let refused = buf
        .try_insert_evicting_where(
            Event {
                stamp: Kairos::new(20, 0, 2, 0u16),
                deps: fresh,
                payload: 2,
            },
            cap,
            |_| false,
        )
        .unwrap_err();

    assert_eq!(refused.event.payload, 2);
    assert_eq!(buf.pending_len(), 1);
    assert_eq!(buf.delivered(), &VersionVector::new());
}

#[test]
fn test_try_insert_evicting_where_preserves_a_resident_predecessor() {
    let cap = NonZeroUsize::new(2).unwrap();
    let predecessor = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: {
            let mut deps = VersionVector::new();
            let _ = deps.increment(1);
            deps
        },
        payload: 1u32,
    };
    let dependent = Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: {
            let mut deps = VersionVector::new();
            deps.observe(1, 1);
            let _ = deps.increment(2);
            deps
        },
        payload: 2u32,
    };
    let arrival = Event {
        stamp: Kairos::new(30, 0, 3, 0u16),
        deps: {
            let mut deps = VersionVector::new();
            let _ = deps.increment(3);
            deps
        },
        payload: 3u32,
    };

    let mut buf = CausalIdeal::new();
    buf.try_insert(predecessor, cap).unwrap();
    buf.try_insert(dependent, cap).unwrap();

    let mut inspected = Vec::new();
    let refused = buf
        .try_insert_evicting_where(arrival, cap, |event| {
            inspected.push(event.payload);
            event.payload == 1
        })
        .unwrap_err();

    assert_eq!(
        inspected,
        [2],
        "policy sees only causally maximal residents"
    );
    assert_eq!(refused.event.payload, 3);
    assert_eq!(buf.pop_ready(), Some(1));
    assert_eq!(buf.pop_ready(), Some(2));
}

#[test]
fn test_try_insert_evicting_where_preserves_an_arrival_predecessor() {
    let cap = NonZeroUsize::new(1).unwrap();
    let predecessor = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: {
            let mut deps = VersionVector::new();
            let _ = deps.increment(1);
            deps
        },
        payload: 1u32,
    };
    let arrival = Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: {
            let mut deps = VersionVector::new();
            deps.observe(1, 1);
            let _ = deps.increment(2);
            deps
        },
        payload: 2u32,
    };

    let mut buf = CausalIdeal::new();
    buf.try_insert(predecessor, cap).unwrap();
    let refused = buf
        .try_insert_evicting_where(arrival, cap, |_| {
            panic!("an arrival predecessor is not a candidate")
        })
        .unwrap_err();

    assert_eq!(refused.event.payload, 2);
    assert_eq!(buf.pop_ready(), Some(1));
}

#[test]
fn test_try_insert_evicting_where_does_not_consult_policy_below_capacity_or_for_stale_input() {
    let cap = NonZeroUsize::new(2).unwrap();
    let mut deps = VersionVector::new();
    let _ = deps.increment(1);
    let event = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps,
        payload: 1u32,
    };
    let mut buf = CausalIdeal::new();

    let outcome = buf
        .try_insert_evicting_where(event.clone(), cap, |_| panic!("buffer has capacity"))
        .unwrap();
    assert!(matches!(outcome, BoundedInsert::Buffered));
    assert_eq!(buf.pop_ready(), Some(1));

    let outcome = buf
        .try_insert_evicting_where(event, cap, |_| panic!("stale input has no victim policy"))
        .unwrap();
    assert!(matches!(outcome, BoundedInsert::DroppedStale));
    assert_eq!(buf.pending_len(), 0);
}

#[test]
fn test_try_insert_evicting_where_refuses_when_a_smaller_call_bound_is_already_exceeded() {
    let large = NonZeroUsize::new(2).unwrap();
    let small = NonZeroUsize::new(1).unwrap();
    let mut buf = CausalIdeal::new();
    for station in [1, 2] {
        let mut deps = VersionVector::new();
        let _ = deps.increment(station);
        let _ = deps.increment(station);
        buf.try_insert(
            Event {
                stamp: Kairos::new(u64::from(station), 0, station, 0u16),
                deps,
                payload: station,
            },
            large,
        )
        .unwrap();
    }

    let mut fresh = VersionVector::new();
    let _ = fresh.increment(3);
    let refused = buf
        .try_insert_evicting_where(
            Event {
                stamp: Kairos::new(3, 0, 3, 0u16),
                deps: fresh,
                payload: 3,
            },
            small,
            |_| panic!("one replacement cannot restore an already-exceeded bound"),
        )
        .unwrap_err();

    assert_eq!(refused.event.payload, 3);
    assert_eq!(buf.pending_len(), 2);
    assert_eq!(buf.delivered(), &VersionVector::new());
}