use std::collections::VecDeque;
use super::{
CarrierAction, CarrierCloseCause, CarrierCounters, CarrierEvent, CarrierRefusal,
CarrierSendOutcome, CarrierSendRefusal, ConnectionKey, MAX_CARRIER_FRAME_BYTES, OpaqueBytes,
encode_frame,
};
pub const MAX_QUEUED_CARRIER_FRAMES: usize = 64;
pub const MAX_QUEUED_CARRIER_BYTES: usize = 33_554_432;
const LENGTH_PREFIX_BYTES: usize = 4;
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct QueueEnqueueResult {
pub outcome: CarrierSendOutcome,
pub actions: Vec<CarrierAction>,
}
#[derive(Debug)]
pub struct BoundedFrameQueue {
key: ConnectionKey,
frames: VecDeque<Vec<u8>>,
queued_bytes: usize,
blocked_encoded_len: Option<usize>,
}
impl BoundedFrameQueue {
pub const fn new(key: ConnectionKey) -> Self {
Self {
key,
frames: VecDeque::new(),
queued_bytes: 0,
blocked_encoded_len: None,
}
}
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(),
}
}
pub fn front(&self) -> Option<&[u8]> {
self.frames.front().map(Vec::as_slice)
}
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()
}
}
pub fn queued_frames(&self) -> usize {
self.frames.len()
}
pub const fn queued_bytes(&self) -> usize {
self.queued_bytes
}
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)
}
}