haematite 0.6.1

Content-addressed, branchable, actor-native storage engine
Documentation
use std::collections::VecDeque;

use super::{
    CarrierAction, CarrierCloseCause, CarrierCounters, CarrierEvent, CarrierRefusal,
    CarrierSendOutcome, CarrierSendRefusal, ConnectionKey, MAX_CARRIER_FRAME_BYTES, OpaqueBytes,
    encode_frame,
};

/// Maximum number of encoded frames owned by one carrier queue.
pub const MAX_QUEUED_CARRIER_FRAMES: usize = 64;
/// Maximum encoded bytes owned by one carrier queue (32 MiB).
pub const MAX_QUEUED_CARRIER_BYTES: usize = 33_554_432;
const LENGTH_PREFIX_BYTES: usize = 4;

/// Result and actions produced by one bounded queue enqueue.
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct QueueEnqueueResult {
    /// Immediate result returned to the sender.
    pub outcome: CarrierSendOutcome,
    /// Ordered refusal and closure actions, if the bound was breached.
    pub actions: Vec<CarrierAction>,
}

/// Bounded FIFO of fully encoded carrier frames for one connection.
#[derive(Debug)]
pub struct BoundedFrameQueue {
    key: ConnectionKey,
    frames: VecDeque<Vec<u8>>,
    queued_bytes: usize,
    blocked_encoded_len: Option<usize>,
}

impl BoundedFrameQueue {
    /// Creates an empty queue bound to one established connection.
    pub const fn new(key: ConnectionKey) -> Self {
        Self {
            key,
            frames: VecDeque::new(),
            queued_bytes: 0,
            blocked_encoded_len: None,
        }
    }

    /// Checks both bounds before encoding or storing the supplied bytes.
    pub fn enqueue(&mut self, payload: OpaqueBytes) -> QueueEnqueueResult {
        if payload.len() > MAX_CARRIER_FRAME_BYTES {
            return QueueEnqueueResult {
                outcome: CarrierSendOutcome::Refused(CarrierSendRefusal::FrameTooLarge {
                    announced: payload.len(),
                    maximum: MAX_CARRIER_FRAME_BYTES,
                }),
                actions: Vec::new(),
            };
        }

        let encoded_len = payload.len() + LENGTH_PREFIX_BYTES;
        if !self.can_admit(encoded_len) {
            self.blocked_encoded_len = Some(encoded_len);
            let refusal = CarrierRefusal::QueueCapacityExceeded {
                queued_frames: self.frames.len(),
                queued_bytes: self.queued_bytes,
            };
            return QueueEnqueueResult {
                outcome: CarrierSendOutcome::Refused(CarrierSendRefusal::QueueCapacityExceeded {
                    queued_frames: self.frames.len(),
                    queued_bytes: self.queued_bytes,
                }),
                actions: vec![
                    CarrierAction::EmitEvent(CarrierEvent::Refused {
                        key: self.key,
                        refusal,
                    }),
                    CarrierAction::CloseSocket {
                        key: self.key,
                        cause: CarrierCloseCause::Refused(refusal),
                    },
                ],
            };
        }

        let payload_len = payload.len();
        let payload = payload.into_vec();
        let Ok(encoded) = encode_frame(&payload) else {
            return QueueEnqueueResult {
                outcome: CarrierSendOutcome::Refused(CarrierSendRefusal::FrameTooLarge {
                    announced: payload_len,
                    maximum: MAX_CARRIER_FRAME_BYTES,
                }),
                actions: Vec::new(),
            };
        };
        self.queued_bytes += encoded.len();
        self.frames.push_back(encoded);
        QueueEnqueueResult {
            outcome: CarrierSendOutcome::Accepted,
            actions: Vec::new(),
        }
    }

    /// Borrows the next encoded frame for an attempted handoff.
    pub fn front(&self) -> Option<&[u8]> {
        self.frames.front().map(Vec::as_slice)
    }

    /// Confirms one successful handoff and emits a full-to-ready transition.
    pub fn confirm_handoff(&mut self, counters: &CarrierCounters) -> Vec<CarrierAction> {
        let Some(frame) = self.frames.pop_front() else {
            return Vec::new();
        };
        self.queued_bytes -= frame.len();
        counters.record_frame_sent();

        let became_ready = self
            .blocked_encoded_len
            .is_some_and(|encoded_len| self.can_admit(encoded_len));
        if became_ready {
            self.blocked_encoded_len = None;
            vec![CarrierAction::EmitEvent(CarrierEvent::SendReady {
                key: self.key,
            })]
        } else {
            Vec::new()
        }
    }

    /// Returns the queued frame count.
    pub fn queued_frames(&self) -> usize {
        self.frames.len()
    }

    /// Returns encoded bytes currently owned by the queue.
    pub const fn queued_bytes(&self) -> usize {
        self.queued_bytes
    }

    /// Reports whether the queue is empty.
    pub fn is_empty(&self) -> bool {
        self.frames.is_empty()
    }

    fn can_admit(&self, encoded_len: usize) -> bool {
        self.frames.len() < MAX_QUEUED_CARRIER_FRAMES
            && self
                .queued_bytes
                .checked_add(encoded_len)
                .is_some_and(|total| total <= MAX_QUEUED_CARRIER_BYTES)
    }
}