haematite 0.6.1

Content-addressed, branchable, actor-native storage engine
Documentation
use crate::carrier::{
    AttemptGeneration, AttemptSlot, AuthenticatedPeerBinding, CarrierAction,
    CarrierConnectFailureCause, CarrierConnectFailureDisposition, CarrierConnectFailureStage,
    CarrierCounterSnapshot, CarrierCounters, CarrierEstablishedCause, CarrierEvent,
    CarrierLostCause, CarrierMachine, CarrierMachineInput, CarrierParkedCause, CarrierState,
    ConnectionId, ConnectionKey, EpisodeId, Incarnation, step,
};

#[derive(Debug, Default)]
struct FakeClock {
    value: u128,
}

impl FakeClock {
    fn advance(&mut self, amount: u128) {
        self.value += amount;
    }
}

fn key(value: u8, incarnation: u64) -> ConnectionKey {
    ConnectionKey {
        connection_id: ConnectionId::new([value; 16]),
        incarnation: Incarnation::new(incarnation),
    }
}

fn apply(
    machine: CarrierMachine,
    counters: &CarrierCounters,
    input: CarrierMachineInput,
) -> (CarrierMachine, Vec<CarrierAction>) {
    let transition = step(machine, input);
    counters.record_actions(&transition.actions);
    (transition.machine, transition.actions)
}

fn current_generation(machine: CarrierMachine) -> Option<AttemptGeneration> {
    match machine.state() {
        CarrierState::Parked {
            episode: _,
            cause: _,
            attempt: AttemptSlot::InFlight { generation },
        } => Some(generation),
        CarrierState::Parked {
            episode: _,
            cause: _,
            attempt: AttemptSlot::Idle,
        }
        | CarrierState::Connected {
            episode: _,
            generation: _,
            key: _,
        }
        | CarrierState::Stopped {
            episode: _,
            cause: _,
        } => None,
    }
}

#[test]
fn disconnected_carrier_has_zero_timer_wakes_and_connect_attempts_without_qualifying_event() {
    let counters = CarrierCounters::default();
    let mut clock = FakeClock::default();
    let machine = CarrierMachine::new(EpisodeId::new(40));
    let (machine, _) = apply(machine, &counters, CarrierMachineInput::ExplicitAction);
    let dial_generation = current_generation(machine).unwrap_or(AttemptGeneration::new(0));
    let dial_failure = CarrierConnectFailureCause::new(
        CarrierConnectFailureStage::Dial,
        CarrierConnectFailureDisposition::Park,
    );
    let (machine, failure_actions) = apply(
        machine,
        &counters,
        CarrierMachineInput::AttemptFailed {
            generation: dial_generation,
            cause: dial_failure,
        },
    );
    assert!(matches!(
        failure_actions.as_slice(),
        [
            CarrierAction::EmitEvent(CarrierEvent::ConnectFailed { .. }),
            CarrierAction::EmitEvent(CarrierEvent::Parked { .. })
        ]
    ));

    let after_dial_failure = counters.snapshot();
    assert_eq!(after_dial_failure.armed_timer_gauge, 0);
    clock.advance(u128::MAX / 2);
    assert_eq!(counters.snapshot(), after_dial_failure);

    let (machine, duplicate_actions) = apply(
        machine,
        &counters,
        CarrierMachineInput::AttemptFailed {
            generation: dial_generation,
            cause: dial_failure,
        },
    );
    assert!(duplicate_actions.is_empty());
    let (machine, stale_actions) = apply(
        machine,
        &counters,
        CarrierMachineInput::AttemptEstablished {
            generation: AttemptGeneration::new(0),
            key: key(1, 1),
            authenticated_peer_binding: AuthenticatedPeerBinding::new(vec![9]),
            cause: CarrierEstablishedCause::ExplicitAction,
        },
    );
    assert!(stale_actions.is_empty());
    clock.advance(1_000_000_000_000);
    assert_eq!(counters.snapshot(), after_dial_failure);

    let (machine, _) = apply(machine, &counters, CarrierMachineInput::ExplicitAction);
    let handshake_generation = current_generation(machine).unwrap_or(AttemptGeneration::new(0));
    let handshake_failure = CarrierConnectFailureCause::new(
        CarrierConnectFailureStage::Handshake,
        CarrierConnectFailureDisposition::Park,
    );
    let (machine, _) = apply(
        machine,
        &counters,
        CarrierMachineInput::AttemptFailed {
            generation: handshake_generation,
            cause: handshake_failure,
        },
    );
    let after_handshake_failure = counters.snapshot();
    assert_eq!(
        after_handshake_failure,
        CarrierCounterSnapshot {
            armed_timer_gauge: 0,
            timer_wakes: 0,
            connect_attempts: after_dial_failure.connect_attempts + 1,
            frames_sent: after_dial_failure.frames_sent,
        }
    );
    clock.advance(u128::MAX / 3);
    assert_eq!(counters.snapshot(), after_handshake_failure);

    let (machine, duplicate_actions) = apply(
        machine,
        &counters,
        CarrierMachineInput::AttemptFailed {
            generation: handshake_generation,
            cause: handshake_failure,
        },
    );
    assert!(duplicate_actions.is_empty());
    let (_, old_fate_actions) = apply(
        machine,
        &counters,
        CarrierMachineInput::EstablishedFate {
            key: key(2, 2),
            cause: CarrierLostCause::TransportError,
        },
    );
    assert!(old_fate_actions.is_empty());
    clock.advance(9_999_999_999_999);
    assert_eq!(counters.snapshot(), after_handshake_failure);
    assert!(clock.value > 0);
}

#[test]
fn qualifying_events_start_at_most_one_fresh_generation() {
    let counters = CarrierCounters::default();
    let machine = CarrierMachine::new(EpisodeId::new(1));
    let (machine, first) = apply(machine, &counters, CarrierMachineInput::ExplicitAction);
    assert_eq!(first.len(), 1);
    let (machine, overlapping) = apply(machine, &counters, CarrierMachineInput::ProvedOnline);
    assert!(overlapping.is_empty());

    let generation = current_generation(machine).unwrap_or(AttemptGeneration::new(0));
    let established_key = key(3, 7);
    let (machine, _) = apply(
        machine,
        &counters,
        CarrierMachineInput::AttemptEstablished {
            generation,
            key: established_key,
            authenticated_peer_binding: AuthenticatedPeerBinding::new(vec![1, 2]),
            cause: CarrierEstablishedCause::ExplicitAction,
        },
    );
    let (machine, offline) = apply(machine, &counters, CarrierMachineInput::Offline);
    assert!(offline.is_empty());
    let (machine, fate) = apply(
        machine,
        &counters,
        CarrierMachineInput::EstablishedFate {
            key: established_key,
            cause: CarrierLostCause::CleanClose,
        },
    );
    assert_eq!(
        fate.iter()
            .filter(|action| matches!(action, CarrierAction::StartAttempt { .. }))
            .count(),
        1
    );
    let (_, duplicate) = apply(
        machine,
        &counters,
        CarrierMachineInput::EstablishedFate {
            key: established_key,
            cause: CarrierLostCause::TransportError,
        },
    );
    assert!(duplicate.is_empty());
    assert_eq!(counters.snapshot().connect_attempts, 2);
}

#[test]
fn current_failed_attempt_parks_with_its_typed_cause_and_no_start_action() {
    let counters = CarrierCounters::default();
    let machine = CarrierMachine::new(EpisodeId::new(9));
    let (machine, _) = apply(machine, &counters, CarrierMachineInput::ExplicitAction);
    let generation = current_generation(machine).unwrap_or(AttemptGeneration::new(0));
    let cause = CarrierConnectFailureCause::new(
        CarrierConnectFailureStage::Handshake,
        CarrierConnectFailureDisposition::Park,
    );
    let (machine, actions) = apply(
        machine,
        &counters,
        CarrierMachineInput::AttemptFailed { generation, cause },
    );
    assert_eq!(
        machine.state(),
        CarrierState::Parked {
            episode: EpisodeId::new(9),
            cause: CarrierParkedCause::ConnectFailed(cause),
            attempt: AttemptSlot::Idle,
        }
    );
    assert!(
        !actions
            .iter()
            .any(|action| matches!(action, CarrierAction::StartAttempt { .. }))
    );
}