minerva 0.2.0

Causal ordering for distributed systems
extern crate alloc;

use crate::kairos::Kairos;
use crate::metis::{CausalIdeal, Event, Gate, Ideal, VersionVector};
use alloc::rc::Rc;
use alloc::vec::Vec;
use core::cell::Cell;

#[test]
fn test_frontier_is_the_ready_antichain() {
    let mut s1 = VersionVector::new();
    let _ = s1.increment(1);
    let mut s2 = VersionVector::new();
    let _ = s2.increment(2);
    let high = Event {
        stamp: Kairos::new(50, 0, 1, 0u16),
        deps: s1,
        payload: 1u32,
    };
    let low = Event {
        stamp: Kairos::new(40, 0, 2, 0u16),
        deps: s2,
        payload: 2u32,
    };

    let mut buf = CausalIdeal::new();
    buf.insert(high);
    buf.insert(low);

    let mut front: Vec<u32> = buf.frontier().map(|e| e.payload).collect();
    front.sort_unstable();
    assert_eq!(front, [1u32, 2]);
    let members: Vec<&Event<u32>> = buf.frontier().collect();
    assert!(members[0].deps.concurrent(&members[1].deps));

    assert_eq!(
        buf.frontier().min_by_key(|e| e.stamp).map(|e| e.payload),
        Some(2)
    );
    assert_eq!(buf.pop_ready(), Some(2));

    let rest: Vec<u32> = buf.frontier().map(|e| e.payload).collect();
    assert_eq!(rest, [1u32]);
}

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

    let mut buf = CausalIdeal::new();
    buf.insert(second);
    assert_eq!(buf.pending_len(), 1);
    assert_eq!(buf.frontier().count(), 0);
    assert_eq!(buf.pop_ready(), None);
}

#[test]
fn test_pop_ready_event_carries_stamp_and_deps() {
    let mut s1 = VersionVector::new();
    let _ = s1.increment(1);
    let a_stamp = Kairos::new(10, 0, 1, 0u16);
    let a = Event {
        stamp: a_stamp,
        deps: s1.clone(),
        payload: 0u32,
    };
    let mut s2 = s1;
    let _ = s2.increment(2);
    let b = Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: s2,
        payload: 1u32,
    };

    let mut buf = CausalIdeal::new();
    buf.insert(b);
    buf.insert(a);

    let first = buf.pop_ready_event().expect("a is deliverable");
    assert_eq!(*first.payload(), 0);
    assert_eq!(first.stamp(), a_stamp);
    assert_eq!(first.deps().get(1), 1);
    let second = buf.pop_ready_event().expect("b is now deliverable");
    assert_eq!(*second.payload(), 1);
    assert!(first.deps().happens_before(second.deps()));
    assert!(buf.pop_ready_event().is_none());
    assert_eq!(buf.pending_len(), 0);
}

#[test]
fn selective_release_skips_a_ready_event_without_advancing_it() {
    let mut early_dep = VersionVector::new();
    let _ = early_dep.increment(1);
    let early = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: early_dep,
        payload: 1u32,
    };
    let mut chosen_dep = VersionVector::new();
    let _ = chosen_dep.increment(2);
    let chosen = Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: chosen_dep,
        payload: 2u32,
    };

    let mut buf = CausalIdeal::new();
    buf.insert(chosen);
    buf.insert(early);

    let released = buf
        .pop_ready_event_where(|event| event.payload == 2)
        .expect("the eligible concurrent event is ready");
    assert_eq!(*released.payload(), 2);
    assert_eq!(buf.delivered().get(1), 0);
    assert_eq!(buf.delivered().get(2), 1);
    assert_eq!(
        buf.frontier()
            .map(|event| event.payload)
            .collect::<Vec<_>>(),
        [1]
    );
}

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

    let mut buf = CausalIdeal::new();
    buf.insert(event);
    assert!(buf.pop_ready_event_where(|_| false).is_none());
    assert_eq!(buf.pending_len(), 1);
    assert_eq!(buf.delivered(), &VersionVector::new());
    assert_eq!(buf.pop_ready(), Some(1));
}

struct MutableDependency;

impl Gate for MutableDependency {
    type Dep = Rc<Cell<bool>>;
    type Progress = u64;

    fn deliverable(_progress: &Self::Progress, _sender: u32, dep: &Self::Dep) -> bool {
        dep.get()
    }

    fn advance(progress: &mut Self::Progress, _sender: u32, _dep: &Self::Dep) {
        *progress += 1;
    }

    fn stale(_progress: &Self::Progress, _sender: u32, _dep: &Self::Dep) -> bool {
        false
    }
}

#[test]
fn selective_release_rechecks_readiness_after_the_predicate() {
    let ready = Rc::new(Cell::new(true));
    let event = Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: Rc::clone(&ready),
        payload: Rc::clone(&ready),
    };
    let mut buf: Ideal<_, MutableDependency> = Ideal::default();
    buf.insert(event);

    assert!(
        buf.pop_ready_event_where(|event| {
            event.payload.set(false);
            true
        })
        .is_none()
    );
    assert_eq!(buf.pending_len(), 1);
    assert_eq!(*buf.progress(), 0);

    ready.set(true);
    assert!(buf.pop_ready_event().is_some());
    assert_eq!(*buf.progress(), 1);
}

#[test]
fn selective_release_skips_candidates_newly_blocked_by_the_predicate() {
    let first_ready = Rc::new(Cell::new(true));
    let second_ready = Rc::new(Cell::new(true));
    let third_ready = Rc::new(Cell::new(true));
    let mut buf: Ideal<_, MutableDependency> = Ideal::default();

    for (time, station, ready, payload) in [
        (30, 3, Rc::clone(&third_ready), 3u32),
        (10, 1, Rc::clone(&first_ready), 1),
        (20, 2, Rc::clone(&second_ready), 2),
    ] {
        buf.insert(Event {
            stamp: Kairos::new(time, 0, station, 0u16),
            deps: ready,
            payload,
        });
    }

    let mut visited = Vec::new();
    let released = buf
        .pop_ready_event_where(|event| {
            visited.push(event.payload);
            if event.payload == 1 {
                second_ready.set(false);
                return false;
            }
            event.payload == 3
        })
        .expect("the third candidate remains ready and eligible");

    assert_eq!(visited, [1, 3]);
    assert_eq!(*released.payload(), 3);
    assert_eq!(buf.pending_len(), 2);
    assert_eq!(*buf.progress(), 1);
}

#[test]
fn selective_release_restarts_for_a_newly_ready_lower_candidate() {
    let lower_ready = Rc::new(Cell::new(false));
    let higher_ready = Rc::new(Cell::new(true));
    let mut buf: Ideal<_, MutableDependency> = Ideal::default();
    buf.insert(Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: Rc::clone(&higher_ready),
        payload: 2u32,
    });
    buf.insert(Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: Rc::clone(&lower_ready),
        payload: 1u32,
    });

    let mut visited = Vec::new();
    let released = buf
        .pop_ready_event_where(|event| {
            visited.push(event.payload);
            if event.payload == 2 {
                lower_ready.set(true);
            }
            true
        })
        .expect("the newly ready lower candidate takes precedence");

    assert_eq!(visited, [2, 1]);
    assert_eq!(*released.payload(), 1);
    assert_eq!(*buf.progress(), 1);
}

#[test]
fn selective_release_reconsiders_an_admitted_candidate_when_reenabled() {
    let lower_ready = Rc::new(Cell::new(true));
    let higher_ready = Rc::new(Cell::new(true));
    let mut buf: Ideal<_, MutableDependency> = Ideal::default();
    buf.insert(Event {
        stamp: Kairos::new(20, 0, 2, 0u16),
        deps: Rc::clone(&higher_ready),
        payload: 2u32,
    });
    buf.insert(Event {
        stamp: Kairos::new(10, 0, 1, 0u16),
        deps: Rc::clone(&lower_ready),
        payload: 1u32,
    });

    let mut visited = Vec::new();
    let released = buf
        .pop_ready_event_where(|event| {
            visited.push(event.payload);
            match event.payload {
                1 => lower_ready.set(false),
                2 => lower_ready.set(true),
                _ => unreachable!(),
            }
            true
        })
        .expect("the admitted lower candidate becomes ready again");

    assert_eq!(visited, [1, 2]);
    assert_eq!(*released.payload(), 1);
    assert_eq!(*buf.progress(), 1);
}