mod capacity;
mod metrics;
mod policy;
mod queue;
use std::collections::VecDeque;
use std::sync::Mutex;
use crate::bytecode::Value;
use super::error::RuntimeError;
use super::process::{Flow, FlowId};
use super::sync_lock;
pub use capacity::{MailboxBytes, MailboxCapacity};
pub use metrics::MailboxStats;
pub use policy::{MailboxConfig, OverflowPolicy};
use queue::{EnqueueEffect, MailboxQueue};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum WaitFilter {
Any,
Tag(u16),
Correlation {
expect_request_id: u64,
expect_sender: Option<u64>,
},
}
impl WaitFilter {
#[inline]
pub(crate) fn matches(&self, value: &Value) -> bool {
match *self {
Self::Any => true,
Self::Tag(expected_tag) => match value.as_message() {
Some(m) => m.tag == expected_tag,
None => false,
},
Self::Correlation {
expect_request_id,
expect_sender,
} => match value.as_message() {
Some(m) => {
let id_ok = m.request_id == expect_request_id;
let sender_ok = match expect_sender {
Some(s) => m.sender == s,
None => true,
};
id_ok && sender_ok
}
None => false,
},
}
}
}
pub struct Mailbox {
inner: Mutex<MailboxInner>,
config: MailboxConfig,
}
struct MailboxInner {
queue: MailboxQueue,
parked: Option<Box<Flow>>,
parked_filter: WaitFilter,
wait_epoch: u64,
stats: MailboxStats,
waiting_senders: VecDeque<WaitingSender>,
closed: bool,
}
struct WaitingSender {
flow: Box<Flow>,
message: Value,
}
pub enum Delivery {
Queued,
QueuedDropOldest,
DroppedNewest,
Handoff(Box<Flow>),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MailboxFullReason {
MessageLimit,
ByteLimit,
}
impl std::fmt::Display for MailboxFullReason {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
MailboxFullReason::MessageLimit => write!(f, "hop count limit"),
MailboxFullReason::ByteLimit => write!(f, "byte budget"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MailboxFull {
reason: MailboxFullReason,
}
impl MailboxFull {
#[inline]
pub const fn reason(self) -> MailboxFullReason {
self.reason
}
}
impl std::fmt::Display for MailboxFull {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "mailbox full ({})", self.reason)
}
}
impl std::error::Error for MailboxFull {}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub struct WaitEpoch(u64);
impl WaitEpoch {
#[inline]
pub const fn get(self) -> u64 {
self.0
}
}
impl std::fmt::Display for WaitEpoch {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "wait#{}", self.0)
}
}
impl Mailbox {
pub fn new() -> Self {
Self::with_config(MailboxConfig::DEFAULT)
}
pub fn with_config(config: MailboxConfig) -> Self {
Mailbox {
inner: Mutex::new(MailboxInner {
queue: MailboxQueue::new(config.capacity().get(), config.bytes().get()),
parked: None,
parked_filter: WaitFilter::Any,
wait_epoch: 0,
stats: MailboxStats::default(),
waiting_senders: VecDeque::new(),
closed: false,
}),
config,
}
}
#[inline]
pub fn config(&self) -> MailboxConfig {
self.config
}
pub fn stats(&self) -> Result<MailboxStats, RuntimeError> {
let inner = sync_lock::lock(&self.inner, "Mailbox::stats")?;
let mut stats = inner.stats;
stats.queued_messages = inner.queue.len();
stats.queued_bytes = inner.queue.bytes();
Ok(stats)
}
pub fn push(&self, value: Value) -> Result<Result<Delivery, MailboxFull>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::push")?;
if let Some(flow) = inner.parked.take() {
if inner.parked_filter.matches(&value) {
inner.parked_filter = WaitFilter::Any;
inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
return Ok(Ok(Delivery::Handoff(flow)));
}
inner.parked = Some(flow);
return Ok(enqueue_locked(
&mut inner,
value,
self.config.overflow(),
));
}
Ok(enqueue_locked(
&mut inner,
value,
self.config.overflow(),
))
}
pub fn try_pop(&self) -> Result<Option<Value>, RuntimeError> {
self.try_pop_filter(WaitFilter::Any)
}
pub fn try_pop_match(&self, tag: u16) -> Result<Option<Value>, RuntimeError> {
self.try_pop_filter(WaitFilter::Tag(tag))
}
pub(crate) fn try_pop_filter(
&self,
filter: WaitFilter,
) -> Result<Option<Value>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::try_pop_filter")?;
let got = inner.queue.take(filter);
if got.is_some() {
inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
}
Ok(got)
}
pub fn park(&self, flow: Box<Flow>) -> Result<Result<WaitEpoch, Box<Flow>>, RuntimeError> {
self.park_filter(flow, WaitFilter::Any)
}
pub fn park_match(
&self,
flow: Box<Flow>,
tag: u16,
) -> Result<Result<WaitEpoch, Box<Flow>>, RuntimeError> {
self.park_filter(flow, WaitFilter::Tag(tag))
}
pub(crate) fn park_filter(
&self,
flow: Box<Flow>,
filter: WaitFilter,
) -> Result<Result<WaitEpoch, Box<Flow>>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::park_filter")?;
if let Some(value) = inner.queue.take(filter) {
inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
drop(inner);
return Ok(Err(with_pending(flow, value)));
}
inner.wait_epoch = inner.wait_epoch.wrapping_add(1);
let epoch = WaitEpoch(inner.wait_epoch);
inner.parked_filter = filter;
inner.parked = Some(flow);
Ok(Ok(epoch))
}
pub fn take_parked_at(&self, epoch: WaitEpoch) -> Result<Option<Box<Flow>>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::take_parked_at")?;
if inner.wait_epoch != epoch.0 {
return Ok(None);
}
match inner.parked.take() {
Some(flow) => {
inner.parked_filter = WaitFilter::Any;
Ok(Some(flow))
}
None => Ok(None),
}
}
pub(crate) fn take_parked(&self) -> Result<Option<Box<Flow>>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::take_parked")?;
inner.parked_filter = WaitFilter::Any;
Ok(inner.parked.take())
}
pub(crate) fn push_system(&self, value: Value) -> Result<Delivery, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::push_system")?;
if let Some(flow) = inner.parked.take() {
if inner.parked_filter.matches(&value) {
inner.parked_filter = WaitFilter::Any;
inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
return Ok(Delivery::Handoff(flow));
}
inner.parked = Some(flow);
}
inner.queue.force_push(value);
inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
Ok(Delivery::Queued)
}
pub(crate) fn close(&self) -> Result<Vec<Flow>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::close")?;
inner.closed = true;
Ok(inner.waiting_senders.drain(..).map(|w| *w.flow).collect())
}
pub(crate) fn park_sender(&self, flow: Box<Flow>, message: Value) -> ParkSender {
match sync_lock::lock(&self.inner, "Mailbox::park_sender") {
Ok(inner) if inner.closed => ParkSender::Closed(flow),
Ok(mut inner) => {
inner.waiting_senders.push_back(WaitingSender { flow, message });
ParkSender::Parked
}
Err(e) => {
super::error::report_fault(e);
ParkSender::Closed(flow)
}
}
}
pub(crate) fn admit_waiting_sender(&self) -> Result<Option<Box<Flow>>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::admit_waiting_sender")?;
if inner.closed {
return Ok(None);
}
let Some(waiter) = inner.waiting_senders.pop_front() else {
return Ok(None);
};
match enqueue_locked(&mut inner, waiter.message.clone(), OverflowPolicy::Reject) {
Ok(_) => Ok(Some(waiter.flow)),
Err(_) => {
inner.waiting_senders.push_front(waiter);
Ok(None)
}
}
}
pub(crate) fn take_waiting_sender(
&self,
sender: FlowId,
) -> Result<Option<Box<Flow>>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::take_waiting_sender")?;
if let Some(pos) = inner.waiting_senders.iter().position(|w| w.flow.id == sender) {
return Ok(inner.waiting_senders.remove(pos).map(|w| w.flow));
}
Ok(None)
}
}
pub(crate) enum ParkSender {
Parked,
Closed(Box<Flow>),
}
impl Default for Mailbox {
fn default() -> Self {
Self::new()
}
}
fn enqueue_locked(
inner: &mut MailboxInner,
value: Value,
policy: OverflowPolicy,
) -> Result<Delivery, MailboxFull> {
match inner.queue.enqueue(value, policy) {
Ok(EnqueueEffect::Enqueued) => {
inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
Ok(Delivery::Queued)
}
Ok(EnqueueEffect::DroppedOldest) => {
inner.stats.dropped_oldest = inner.stats.dropped_oldest.saturating_add(1);
inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
Ok(Delivery::QueuedDropOldest)
}
Ok(EnqueueEffect::DroppedNewest) => {
inner.stats.dropped_newest = inner.stats.dropped_newest.saturating_add(1);
Ok(Delivery::DroppedNewest)
}
Err(reason) => {
inner.stats.rejected = inner.stats.rejected.saturating_add(1);
if reason == MailboxFullReason::ByteLimit {
inner.stats.rejected_byte_limit =
inner.stats.rejected_byte_limit.saturating_add(1);
}
Err(MailboxFull { reason })
}
}
}
fn with_pending(mut flow: Box<Flow>, value: Value) -> Box<Flow> {
flow.pending_message = Some(value);
flow
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bytecode::Message;
use crate::scheduler::oneshot;
use crate::scheduler::process::{next_flow_id, RestartPolicy};
use crate::vm::{NativeTable, Vm};
use crate::bytecode::builder::ChunkBuilder;
use std::sync::Arc;
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn dummy_flow() -> Result<Box<Flow>, Box<dyn std::error::Error>> {
let mut b = ChunkBuilder::new("mb");
b.begin_function("main", 0, 1);
b.emit_return(0);
let chunk = b.finish();
let vm = Vm::new(Arc::new(chunk), NativeTable::empty(), 0, &[])?;
let (tx, _rx) = oneshot::channel();
Ok(Box::new(Flow::new(
next_flow_id(),
vm,
Arc::new(Mailbox::new()),
RestartPolicy::Never,
tx,
)))
}
fn hop_payload_int(msg: &Message) -> i64 {
msg.payload.as_int().unwrap_or(0)
}
fn hop(sender: u64, request_id: u64, tag: u16, payload: impl Into<Value>) -> Value {
Value::Message(Message::new(sender, request_id, tag, payload))
}
fn msg(tag: u16, payload: impl Into<Value>) -> Value {
hop(1, 1, tag, payload)
}
fn tiny_reject(n: u32) -> Result<Mailbox, Box<dyn std::error::Error>> {
let cap = MailboxCapacity::new(n).ok_or("invalid mailbox capacity")?;
Ok(Mailbox::with_config(MailboxConfig::new(
cap,
OverflowPolicy::Reject,
)))
}
fn hop_msg(value: &Value) -> Result<&Message, Box<dyn std::error::Error>> {
value.as_message().ok_or_else(|| "expected Message hop".into())
}
#[test]
fn try_pop_match_skips_non_matching_fifo() -> TestResult {
let mb = Mailbox::new();
mb.push(msg(9, 1))??;
mb.push(msg(1, 42))??;
mb.push(msg(9, 2))??;
let got = mb.try_pop_match(1)?.ok_or("match")?;
assert_eq!(hop_payload_int(hop_msg(&got)?), 42);
assert_eq!(
hop_msg(&mb.try_pop()?.ok_or("first leftover")?)?.tag,
9
);
assert_eq!(
hop_payload_int(hop_msg(&mb.try_pop()?.ok_or("second leftover")?)?),
2
);
Ok(())
}
#[test]
fn push_while_park_match_queues_junk_keeps_waiter() -> TestResult {
let mb = Mailbox::new();
let flow = dummy_flow()?;
assert!(mb.park_match(flow, 1)?.is_ok());
assert!(matches!(mb.push(msg(9, 0))??, Delivery::Queued));
assert!(matches!(mb.push(msg(1, 7))??, Delivery::Handoff(_)));
assert_eq!(hop_msg(&mb.try_pop()?.ok_or("queued junk")?)?.tag, 9);
Ok(())
}
#[test]
fn ask_does_not_consume_reply_for_another_request() -> TestResult {
let mb = Mailbox::new();
mb.push(hop(10, 2, 2, 99))??;
mb.push(hop(10, 1, 2, 42))??;
let filter = WaitFilter::Correlation {
expect_request_id: 1,
expect_sender: Some(10),
};
let got = mb.try_pop_filter(filter)?.ok_or("id=1")?;
assert_eq!(hop_payload_int(hop_msg(&got)?), 42);
let left = mb.try_pop()?.ok_or("leftover")?;
assert_eq!(hop_msg(&left)?.request_id, 2);
Ok(())
}
#[test]
fn ask_requires_reply_from_target() -> TestResult {
let mb = Mailbox::new();
let flow = dummy_flow()?;
let filter = WaitFilter::Correlation {
expect_request_id: 1,
expect_sender: Some(10),
};
assert!(mb.park_filter(flow, filter)?.is_ok());
assert!(matches!(mb.push(hop(99, 1, 2, 0))??, Delivery::Queued));
assert!(matches!(mb.push(hop(10, 1, 2, 42))??, Delivery::Handoff(_)));
assert_eq!(
hop_msg(&mb.try_pop()?.ok_or("non-matching queued")?)?.sender,
99
);
Ok(())
}
#[test]
fn reject_when_full_without_waiter() -> TestResult {
let mb = tiny_reject(1)?;
assert!(matches!(mb.push(msg(1, 1))??, Delivery::Queued));
match mb.push(msg(1, 2))? {
Err(full) => assert_eq!(full.reason(), MailboxFullReason::MessageLimit),
Ok(_) => return Err("expected the hop count bound to refuse".into()),
}
let s = mb.stats()?;
assert_eq!(s.enqueued, 1);
assert_eq!(s.rejected, 1);
assert_eq!(s.rejected_byte_limit, 0);
assert_eq!(s.queued_messages, 1);
Ok(())
}
fn park_now(mb: &Mailbox, flow: Box<Flow>) -> Result<WaitEpoch, Box<dyn std::error::Error>> {
match mb.park(flow)? {
Ok(epoch) => Ok(epoch),
Err(_) => Err("an empty mailbox should have parked the flow".into()),
}
}
fn handoff(mb: &Mailbox, value: Value) -> Result<Box<Flow>, Box<dyn std::error::Error>> {
match mb.push(value)?? {
Delivery::Handoff(flow) => Ok(flow),
other => Err(format!("expected a handoff, got {}", delivery_name(&other)).into()),
}
}
fn delivery_name(d: &Delivery) -> &'static str {
match d {
Delivery::Queued => "Queued",
Delivery::QueuedDropOldest => "QueuedDropOldest",
Delivery::DroppedNewest => "DroppedNewest",
Delivery::Handoff(_) => "Handoff",
}
}
#[test]
fn a_stale_deadline_cannot_steal_a_later_wait() -> TestResult {
let mb = Mailbox::new();
let first = park_now(&mb, dummy_flow()?)?;
let woken = handoff(&mb, msg(1, 1))?;
let second = park_now(&mb, woken)?;
assert_ne!(first, second, "each park must get its own epoch");
assert!(mb.take_parked_at(first)?.is_none());
assert!(mb.take_parked_at(second)?.is_some());
Ok(())
}
#[test]
fn a_stale_deadline_does_not_downgrade_a_selective_waiter() -> TestResult {
let mb = Mailbox::new();
let first = park_now(&mb, dummy_flow()?)?;
let woken = handoff(&mb, msg(1, 1))?;
let second = match mb.park_match(woken, 7)? {
Ok(epoch) => epoch,
Err(_) => return Err("empty mailbox should have parked the flow".into()),
};
assert_ne!(first, second);
assert!(mb.take_parked_at(first)?.is_none());
assert!(matches!(mb.push(msg(9, 0))??, Delivery::Queued));
assert!(matches!(mb.push(msg(7, 0))??, Delivery::Handoff(_)));
Ok(())
}
#[test]
fn a_deadline_for_a_wait_that_a_hop_ended_does_nothing() -> TestResult {
let mb = Mailbox::new();
let epoch = park_now(&mb, dummy_flow()?)?;
let _woken = handoff(&mb, msg(1, 1))?;
assert!(mb.take_parked_at(epoch)?.is_none());
Ok(())
}
#[test]
fn byte_budget_refuses_before_the_hop_count_and_says_so() -> TestResult {
let cap = MailboxCapacity::new(64).ok_or("cap")?;
let budget = MailboxBytes::new(MailboxBytes::MIN).ok_or("bytes")?;
let mb = Mailbox::with_config(
MailboxConfig::new(cap, OverflowPolicy::Reject).with_bytes(budget),
);
assert!(matches!(mb.push(Value::bytes(vec![0u8; 900]))??, Delivery::Queued));
match mb.push(Value::bytes(vec![0u8; 900]))? {
Err(full) => assert_eq!(full.reason(), MailboxFullReason::ByteLimit),
Ok(_) => return Err("expected the byte budget to refuse".into()),
}
let s = mb.stats()?;
assert_eq!(s.queued_messages, 1);
assert!(s.queued_bytes >= 900);
assert_eq!(s.rejected, 1);
assert_eq!(s.rejected_byte_limit, 1);
Ok(())
}
#[test]
fn draining_a_hop_frees_its_byte_charge() -> TestResult {
let cap = MailboxCapacity::new(64).ok_or("cap")?;
let budget = MailboxBytes::new(MailboxBytes::MIN).ok_or("bytes")?;
let mb = Mailbox::with_config(
MailboxConfig::new(cap, OverflowPolicy::Reject).with_bytes(budget),
);
mb.push(Value::bytes(vec![0u8; 900]))??;
mb.try_pop()?.ok_or("queued blob")?;
assert_eq!(mb.stats()?.queued_bytes, 0);
assert!(matches!(mb.push(Value::bytes(vec![0u8; 900]))??, Delivery::Queued));
Ok(())
}
#[test]
fn matching_handoff_does_not_count_as_full() -> TestResult {
let mb = tiny_reject(1)?;
mb.push(msg(9, 0))??;
let flow = dummy_flow()?;
assert!(mb.park_match(flow, 1)?.is_ok());
assert!(matches!(mb.push(msg(1, 7))??, Delivery::Handoff(_)));
assert_eq!(hop_msg(&mb.try_pop()?.ok_or("queued")?)?.tag, 9);
Ok(())
}
#[test]
fn waiting_send_admits_one_after_pop() -> TestResult {
let mb = tiny_reject(1)?;
assert!(matches!(mb.push(msg(1, 1))??, Delivery::Queued));
assert!(matches!(
mb.park_sender(dummy_flow()?, msg(1, 2)),
ParkSender::Parked
));
assert!(mb.try_pop()?.is_some());
let woken = mb.admit_waiting_sender()?.ok_or("admitted")?;
drop(woken);
assert_eq!(hop_payload_int(hop_msg(&mb.try_pop()?.ok_or("second hop")?)?), 2);
assert!(mb.admit_waiting_sender()?.is_none());
Ok(())
}
#[test]
fn close_rejects_late_park_sender() -> TestResult {
let mb = tiny_reject(1)?;
let leftover = mb.close()?;
assert!(leftover.is_empty());
match mb.park_sender(dummy_flow()?, msg(1, 1)) {
ParkSender::Closed(_) => {}
ParkSender::Parked => return Err("closed mailbox must not park a sender".into()),
}
Ok(())
}
#[test]
fn drop_oldest_still_wakes_on_match() -> TestResult {
let cap = MailboxCapacity::new(1).ok_or("cap")?;
let mb = Mailbox::with_config(MailboxConfig::new(cap, OverflowPolicy::DropOldest));
let flow = dummy_flow()?;
assert!(mb.park_match(flow, 1)?.is_ok());
mb.push(msg(9, 1))??;
assert!(matches!(mb.push(msg(1, 2))??, Delivery::Handoff(_)));
Ok(())
}
}