minerva 0.2.0

Causal ordering for distributed systems
extern crate alloc;

use super::support::{physical_reading, remote_physical};
use crate::kairos::{Clock, Kairos, SkewExceeded, TickCounter, VirtualTimeSource};
use proptest::prelude::*;

proptest! {
    /// Successive `now` stamps strictly increase under arbitrary physical skew.
    #[test]
    fn prop_now_is_monotonic_under_skew(
        readings in prop::collection::vec(physical_reading(), 1..64),
    ) {
        let source = VirtualTimeSource::new(readings[0]);
        let driver = source.clone();
        let clock = Clock::with_default_config(source, 1).unwrap();
        let mut prev = clock.now(0u16);
        for &reading in &readings[1..] {
            driver.set(reading);
            let next = clock.now(0u16);
            prop_assert!(next > prev, "stamps must strictly increase: {next:?} !> {prev:?}");
            prev = next;
        }
    }

    /// `after(remote)` exceeds `remote`, and local output stays strictly increasing.
    #[test]
    fn prop_after_dominates_remote_and_stays_monotonic(
        remotes in prop::collection::vec(
            (remote_physical(), any::<u16>(), any::<u16>(), any::<u32>()),
            1..64,
        ),
    ) {
        let clock = Clock::with_default_config(
            TickCounter::new(), 7,
        ).unwrap();
        let mut prev = clock.now(0u16);
        for (physical, logical, kairotic, station_id) in remotes {
            let remote = Kairos::new(physical, logical, station_id, kairotic);
            let next = clock.after(remote, 0u16);
            prop_assert!(next > remote, "after(remote) must exceed remote: {next:?} !> {remote:?}");
            prop_assert!(next > prev, "clock output must stay strictly increasing");
            prev = next;
        }
    }

    /// After `observe(remote)`, the next `now` dominates `remote`.
    #[test]
    fn prop_observe_then_now_dominates_remote_and_stays_monotonic(
        remotes in prop::collection::vec(
            (remote_physical(), any::<u16>(), any::<u16>(), any::<u32>()),
            1..64,
        ),
    ) {
        let clock = Clock::with_default_config(
            TickCounter::new(), 9,
        ).unwrap();
        let mut prev = clock.now(0u16);
        for (physical, logical, kairotic, station_id) in remotes {
            let remote = Kairos::new(physical, logical, station_id, kairotic);
            clock.observe(remote);
            let next = clock.now(0u16);
            prop_assert!(next > remote, "now after observe must exceed remote: {next:?} !> {remote:?}");
            prop_assert!(next > prev, "clock output must stay strictly increasing");
            prev = next;
        }
    }

    /// `try_observe` admits within-bound remotes and leaves rejected remotes inert.
    #[test]
    fn prop_try_observe_admits_within_bound_and_rejection_is_inert(
        reading in physical_reading(),
        remote in (remote_physical(), any::<u16>(), any::<u16>(), any::<u32>()),
        max_forward_skew in any::<u64>(),
    ) {
        let (remote_phys, logical, kairotic, station_id) = remote;
        let clock = Clock::with_default_config(
            VirtualTimeSource::new(reading), 1,
        ).unwrap();
        let remote = Kairos::new(remote_phys, logical, station_id, kairotic);
        let expected_skew = remote_phys.saturating_sub(reading);

        let result = clock.try_observe(remote, max_forward_skew);

        if expected_skew > max_forward_skew {
            prop_assert_eq!(
                result,
                Err(SkewExceeded { observed_forward_skew: expected_skew, max_forward_skew })
            );
            prop_assert_eq!(
                clock.stats().max_forward_skew_ns, 0,
                "a rejected remote records no forward skew"
            );
            let next = clock.now(0u16);
            prop_assert_eq!(
                next.physical(), reading,
                "a rejected remote must not move physical off local time"
            );
        } else {
            prop_assert_eq!(result, Ok(()));
            let next = clock.now(0u16);
            prop_assert!(
                next > remote,
                "after an admitted try_observe, now dominates remote: {:?} !> {:?}", next, remote
            );
        }
    }
}

/// One step of a scripted clock tape: the reading both shells will see, and
/// the operation both will run against it.
#[derive(Debug, Clone)]
enum TapeOp {
    /// A bare mint.
    Now,
    /// A run mint of the carried length.
    NowRun(u32),
    /// Receive-and-send against the carried remote stamp.
    After(Kairos),
    /// A pure receive of the carried remote stamp.
    Observe(Kairos),
    /// A bounded receive of the carried remote under the carried bound.
    TryObserve(Kairos, u64),
}

/// An arbitrary remote stamp below the `after` ceiling.
fn tape_remote() -> impl Strategy<Value = Kairos> {
    (remote_physical(), any::<u16>(), any::<u16>(), 1u32..).prop_map(
        |(physical, logical, kairotic, station)| Kairos::new(physical, logical, station, kairotic),
    )
}

/// One tape step: a reading plus an operation.
fn tape_step() -> impl Strategy<Value = (u64, TapeOp)> {
    (
        physical_reading(),
        prop_oneof![
            Just(TapeOp::Now),
            (1u32..200).prop_map(TapeOp::NowRun),
            tape_remote().prop_map(TapeOp::After),
            tape_remote().prop_map(TapeOp::Observe),
            (tape_remote(), any::<u64>()).prop_map(|(r, bound)| TapeOp::TryObserve(r, bound)),
        ],
    )
}

proptest! {
    /// The two minting shells tell one story: driven by identical readings
    /// through an identical operation tape, `LocalClock` and `Clock` mint
    /// identical stamps, return identical bounded-receive verdicts, and
    /// record identical skew maxima. The shells share the pure fold by
    /// construction (S188); this pins that neither shell's commitment story
    /// (CAS loop with per-attempt stat recording, `Cell` store) leaks into
    /// the time algebra.
    #[test]
    fn prop_local_clock_agrees_with_atomic_clock_on_any_tape(
        tape in prop::collection::vec(tape_step(), 1..48),
    ) {
        let atomic_source = VirtualTimeSource::new(0);
        let atomic_driver = atomic_source.clone();
        let atomic = Clock::with_default_config(atomic_source, 11).unwrap();

        let local_source = VirtualTimeSource::new(0);
        let local_driver = local_source.clone();
        let local = crate::kairos::LocalClock::new(local_source, 11).unwrap();

        for (reading, op) in tape {
            atomic_driver.set(reading);
            local_driver.set(reading);
            match op {
                TapeOp::Now => {
                    prop_assert_eq!(atomic.now(5u16), local.now(5u16));
                }
                TapeOp::NowRun(len) => {
                    let len = core::num::NonZeroU32::new(len).expect("tape lens are nonzero");
                    prop_assert_eq!(atomic.now_run(len, 5u16), local.now_run(len, 5u16));
                }
                TapeOp::After(remote) => {
                    prop_assert_eq!(atomic.after(remote, 5u16), local.after(remote, 5u16));
                }
                TapeOp::Observe(remote) => {
                    atomic.observe(remote);
                    local.observe(remote);
                    // The observed state is invisible until the next mint;
                    // compare through it.
                    prop_assert_eq!(atomic.now(5u16), local.now(5u16));
                }
                TapeOp::TryObserve(remote, bound) => {
                    prop_assert_eq!(
                        atomic.try_observe(remote, bound),
                        local.try_observe(remote, bound)
                    );
                }
            }
        }

        let atomic_stats = atomic.stats();
        let local_stats = local.stats();
        prop_assert_eq!(atomic_stats.max_forward_skew_ns, local_stats.max_forward_skew_ns);
        prop_assert_eq!(atomic_stats.max_backward_skew_ns, local_stats.max_backward_skew_ns);
    }

    /// The run mint is the frozen-reading mint succession, from any state a
    /// prior receive may have left: a run of `len` and `len` individual
    /// mints against the same frozen reading commit identical stamp
    /// successions and identical end states, and the members strictly
    /// ascend and chain by the codec's own successor rule.
    #[test]
    fn prop_a_run_equals_frozen_individual_mints(
        reading in physical_reading(),
        primer in prop::option::of((remote_physical(), any::<u16>(), any::<u16>(), 1u32..)),
        len in 1u32..2_048,
    ) {
        let bulk = Clock::with_default_config(
            VirtualTimeSource::new(reading), 13,
        ).unwrap();
        let single = Clock::with_default_config(
            VirtualTimeSource::new(reading), 13,
        ).unwrap();
        if let Some((physical, logical, kairotic, station)) = primer {
            let remote = Kairos::new(physical, logical, station, kairotic);
            bulk.observe(remote);
            single.observe(remote);
        }

        let run = bulk.now_run(core::num::NonZeroU32::new(len).unwrap(), 5u16);
        let mut prev: Option<Kairos> = None;
        for index in 0..len {
            let minted = single.now(5u16);
            let member = run.get(index).expect("members within the reservation");
            prop_assert_eq!(member, minted, "member {} diverged from the mint", index);
            if let Some(prev) = prev {
                if prev.physical() < u64::MAX {
                    prop_assert!(member > prev, "members must strictly ascend");
                }
                prop_assert_eq!(
                    (member.physical(), member.logical()),
                    crate::metis::rank_successor(prev.physical(), prev.logical()),
                    "members must chain by the codec successor rule"
                );
            }
            prev = Some(member);
        }
        // Both clocks committed the same position.
        prop_assert_eq!(bulk.now(5u16), single.now(5u16));
    }
}