use std::collections::VecDeque;
use crate::bytecode::Value;
use super::{MailboxFullReason, OverflowPolicy, WaitFilter};
pub(crate) struct MailboxQueue {
inner: VecDeque<Value>,
limit: usize,
bytes: usize,
byte_limit: usize,
}
impl MailboxQueue {
pub(crate) fn new(limit: usize, byte_limit: usize) -> Self {
debug_assert!(limit >= 1);
debug_assert!(byte_limit >= 1);
Self {
inner: VecDeque::new(),
limit,
bytes: 0,
byte_limit,
}
}
#[inline]
pub(crate) fn len(&self) -> usize {
self.inner.len()
}
#[inline]
pub(crate) fn bytes(&self) -> usize {
self.bytes
}
#[inline]
fn is_full(&self) -> bool {
self.inner.len() >= self.limit
}
#[inline]
fn would_exceed_bytes(&self, cost: usize) -> bool {
self.bytes.saturating_add(cost) > self.byte_limit
}
pub(crate) fn take(&mut self, filter: WaitFilter) -> Option<Value> {
let value = match filter {
WaitFilter::Any => self.inner.pop_front(),
other => {
let idx = self.inner.iter().position(|v| other.matches(v))?;
self.inner.remove(idx)
}
};
if let Some(v) = &value {
self.bytes = self.bytes.saturating_sub(v.memory_size());
}
value
}
pub(crate) fn enqueue(
&mut self,
value: Value,
policy: OverflowPolicy,
) -> Result<EnqueueEffect, MailboxFullReason> {
let cost = value.memory_size();
if !self.is_full() && !self.would_exceed_bytes(cost) {
if !reserve_for_push(&mut self.inner, self.limit) {
return Err(MailboxFullReason::MessageLimit);
}
self.bytes = self.bytes.saturating_add(cost);
self.inner.push_back(value);
return Ok(EnqueueEffect::Enqueued);
}
let reason = if self.is_full() {
MailboxFullReason::MessageLimit
} else {
MailboxFullReason::ByteLimit
};
match policy {
OverflowPolicy::Reject => Err(reason),
OverflowPolicy::DropNewest => Ok(EnqueueEffect::DroppedNewest),
OverflowPolicy::DropOldest => {
let mut dropped = false;
while (self.is_full() || self.would_exceed_bytes(cost)) && !self.inner.is_empty() {
if let Some(old) = self.inner.pop_front() {
self.bytes = self.bytes.saturating_sub(old.memory_size());
dropped = true;
}
}
if self.would_exceed_bytes(cost) {
return Err(MailboxFullReason::ByteLimit);
}
self.bytes = self.bytes.saturating_add(cost);
self.inner.push_back(value);
Ok(if dropped {
EnqueueEffect::DroppedOldest
} else {
EnqueueEffect::Enqueued
})
}
}
}
}
fn reserve_for_push(queue: &mut VecDeque<Value>, capacity: usize) -> bool {
if queue.len() < queue.capacity() {
return true;
}
let current = queue.capacity();
let next = current.max(1).saturating_mul(2).min(capacity);
if next <= current {
return queue.len() < capacity;
}
queue.reserve(next - current);
true
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum EnqueueEffect {
Enqueued,
DroppedNewest,
DroppedOldest,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bytecode::{Message, Value};
const ROOMY: usize = 1 << 20;
fn hop(n: u64) -> Value {
Value::Message(Message::new(1, n, 1, n))
}
fn blob(len: usize) -> Value {
Value::bytes(vec![0u8; len])
}
fn hop_payload(q: &mut MailboxQueue) -> Result<u64, &'static str> {
match q.take(WaitFilter::Any) {
Some(v) => match v.as_message() {
Some(m) => Ok(m.payload),
None => Err("expected Message"),
},
None => Err("queue empty"),
}
}
#[test]
fn reject_at_limit() {
let mut q = MailboxQueue::new(2, ROOMY);
assert_eq!(
q.enqueue(hop(1), OverflowPolicy::Reject),
Ok(EnqueueEffect::Enqueued)
);
assert_eq!(
q.enqueue(hop(2), OverflowPolicy::Reject),
Ok(EnqueueEffect::Enqueued)
);
assert_eq!(
q.enqueue(hop(3), OverflowPolicy::Reject),
Err(MailboxFullReason::MessageLimit)
);
assert_eq!(q.len(), 2);
}
#[test]
fn drop_newest_keeps_old() -> Result<(), &'static str> {
let mut q = MailboxQueue::new(1, ROOMY);
let _ = q.enqueue(hop(1), OverflowPolicy::DropNewest);
assert_eq!(
q.enqueue(hop(2), OverflowPolicy::DropNewest),
Ok(EnqueueEffect::DroppedNewest)
);
assert_eq!(hop_payload(&mut q)?, 1);
Ok(())
}
#[test]
fn drop_oldest_slides() -> Result<(), &'static str> {
let mut q = MailboxQueue::new(2, ROOMY);
let _ = q.enqueue(hop(1), OverflowPolicy::DropOldest);
let _ = q.enqueue(hop(2), OverflowPolicy::DropOldest);
assert_eq!(
q.enqueue(hop(3), OverflowPolicy::DropOldest),
Ok(EnqueueEffect::DroppedOldest)
);
assert_eq!(hop_payload(&mut q)?, 2);
assert_eq!(hop_payload(&mut q)?, 3);
Ok(())
}
#[test]
fn byte_limit_rejects_long_before_the_hop_count() {
let mut q = MailboxQueue::new(64, 4096);
assert_eq!(
q.enqueue(blob(3000), OverflowPolicy::Reject),
Ok(EnqueueEffect::Enqueued)
);
assert_eq!(
q.enqueue(blob(3000), OverflowPolicy::Reject),
Err(MailboxFullReason::ByteLimit)
);
assert_eq!(q.len(), 1);
}
#[test]
fn take_refunds_the_charge() {
let mut q = MailboxQueue::new(64, 4096);
let _ = q.enqueue(blob(3000), OverflowPolicy::Reject);
assert!(q.bytes() >= 3000);
assert!(q.take(WaitFilter::Any).is_some());
assert_eq!(q.bytes(), 0);
assert_eq!(
q.enqueue(blob(3000), OverflowPolicy::Reject),
Ok(EnqueueEffect::Enqueued)
);
}
#[test]
fn selective_take_refunds_the_right_charge() {
let mut q = MailboxQueue::new(64, ROOMY);
let _ = q.enqueue(hop(1), OverflowPolicy::Reject);
let _ = q.enqueue(blob(3000), OverflowPolicy::Reject);
let _ = q.enqueue(hop(2), OverflowPolicy::Reject);
let charged = q.bytes();
assert!(q.take(WaitFilter::Tag(1)).is_some());
assert_eq!(q.bytes(), charged - std::mem::size_of::<Value>());
assert_eq!(q.len(), 2);
}
#[test]
fn hop_larger_than_the_whole_budget_is_refused_even_by_drop_oldest() {
let mut q = MailboxQueue::new(64, 2048);
let _ = q.enqueue(hop(1), OverflowPolicy::DropOldest);
assert_eq!(
q.enqueue(blob(4000), OverflowPolicy::DropOldest),
Err(MailboxFullReason::ByteLimit)
);
}
#[test]
fn drop_oldest_evicts_as_many_as_the_byte_budget_needs() {
let mut q = MailboxQueue::new(64, 4096);
for n in 0..3 {
assert_eq!(
q.enqueue(blob(1000), OverflowPolicy::DropOldest),
Ok(EnqueueEffect::Enqueued),
"blob {n} should fit"
);
}
assert_eq!(
q.enqueue(blob(3000), OverflowPolicy::DropOldest),
Ok(EnqueueEffect::DroppedOldest)
);
assert!(q.bytes() <= 4096);
}
}