haematite 0.7.0

Content-addressed, branchable, actor-native storage engine
Documentation
use crate::carrier::{
    BoundedFrameQueue, CarrierAction, CarrierCounters, CarrierEvent, CarrierRefusal,
    CarrierSendOutcome, CarrierSendRefusal, ConnectionId, ConnectionKey, Incarnation,
    MAX_CARRIER_FRAME_BYTES, MAX_QUEUED_CARRIER_BYTES, MAX_QUEUED_CARRIER_FRAMES, OpaqueBytes,
};

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

#[test]
fn frame_count_breach_is_typed_closes_and_transitions_to_ready() {
    let counters = CarrierCounters::default();
    let mut queue = BoundedFrameQueue::new(key());
    for _ in 0..MAX_QUEUED_CARRIER_FRAMES {
        assert_eq!(
            queue.enqueue(OpaqueBytes::default()).outcome,
            CarrierSendOutcome::Accepted
        );
    }
    assert_eq!(queue.queued_frames(), MAX_QUEUED_CARRIER_FRAMES);

    let breach = queue.enqueue(OpaqueBytes::default());
    assert_eq!(
        breach.outcome,
        CarrierSendOutcome::Refused(CarrierSendRefusal::QueueCapacityExceeded {
            queued_frames: MAX_QUEUED_CARRIER_FRAMES,
            queued_bytes: MAX_QUEUED_CARRIER_FRAMES * 4,
        })
    );
    assert!(matches!(
        breach.actions.as_slice(),
        [
            CarrierAction::EmitEvent(CarrierEvent::Refused {
                refusal: CarrierRefusal::QueueCapacityExceeded { .. },
                ..
            }),
            CarrierAction::CloseSocket { .. }
        ]
    ));

    let ready = queue.confirm_handoff(&counters);
    assert_eq!(
        ready,
        vec![CarrierAction::EmitEvent(CarrierEvent::SendReady {
            key: key(),
        })]
    );
    assert_eq!(counters.snapshot().frames_sent, 1);
}

#[test]
fn byte_bound_is_checked_before_taking_an_over_cap_frame() {
    let counters = CarrierCounters::default();
    let mut queue = BoundedFrameQueue::new(key());
    let payload_len = MAX_QUEUED_CARRIER_BYTES / 2 - 4;
    for _ in 0..2 {
        assert_eq!(
            queue
                .enqueue(OpaqueBytes::new(vec![7; payload_len]))
                .outcome,
            CarrierSendOutcome::Accepted
        );
    }
    assert_eq!(queue.queued_bytes(), MAX_QUEUED_CARRIER_BYTES);

    let breach = queue.enqueue(OpaqueBytes::default());
    assert_eq!(
        breach.outcome,
        CarrierSendOutcome::Refused(CarrierSendRefusal::QueueCapacityExceeded {
            queued_frames: 2,
            queued_bytes: MAX_QUEUED_CARRIER_BYTES,
        })
    );
    assert_eq!(queue.queued_bytes(), MAX_QUEUED_CARRIER_BYTES);
    assert!(matches!(
        queue.confirm_handoff(&counters).as_slice(),
        [CarrierAction::EmitEvent(CarrierEvent::SendReady { .. })]
    ));
}

#[test]
fn payload_over_frame_cap_is_refused_without_queue_growth() {
    let mut queue = BoundedFrameQueue::new(key());
    let too_large = OpaqueBytes::new(vec![0; MAX_CARRIER_FRAME_BYTES + 1]);
    assert_eq!(
        queue.enqueue(too_large).outcome,
        CarrierSendOutcome::Refused(CarrierSendRefusal::FrameTooLarge {
            announced: MAX_CARRIER_FRAME_BYTES + 1,
            maximum: MAX_CARRIER_FRAME_BYTES,
        })
    );
    assert!(queue.is_empty());
    assert_eq!(queue.queued_bytes(), 0);
}