use std::collections::VecDeque;
use std::sync::Mutex;
use crate::bytecode::Value;
use super::error::RuntimeError;
use super::process::Flow;
use super::sync_lock;
#[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) => value
.as_message()
.map(|m| m.tag == expected_tag)
.unwrap_or(false),
Self::Correlation {
expect_request_id,
expect_sender,
} => match value.as_message() {
Some(m) => {
m.request_id == expect_request_id
&& expect_sender.map(|s| m.sender == s).unwrap_or(true)
}
None => false,
},
}
}
}
pub struct Mailbox {
inner: Mutex<MailboxInner>,
}
struct MailboxInner {
queue: VecDeque<Value>,
parked: Option<Box<Flow>>,
parked_filter: WaitFilter,
}
impl Default for MailboxInner {
fn default() -> Self {
Self {
queue: VecDeque::new(),
parked: None,
parked_filter: WaitFilter::Any,
}
}
}
pub enum Delivery {
Queued,
Handoff(Box<Flow>),
}
impl Mailbox {
pub fn new() -> Self {
Mailbox {
inner: Mutex::new(MailboxInner::default()),
}
}
pub fn push(&self, value: Value) -> Result<Delivery, 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;
return Ok(Delivery::Handoff(flow));
}
inner.parked = Some(flow);
inner.queue.push_back(value);
return Ok(Delivery::Queued);
}
inner.queue.push_back(value);
Ok(Delivery::Queued)
}
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")?;
Ok(take_with_filter(&mut inner.queue, filter))
}
pub fn park(&self, flow: Box<Flow>) -> Result<Result<(), Box<Flow>>, RuntimeError> {
self.park_filter(flow, WaitFilter::Any)
}
pub fn park_match(
&self,
flow: Box<Flow>,
tag: u16,
) -> Result<Result<(), Box<Flow>>, RuntimeError> {
self.park_filter(flow, WaitFilter::Tag(tag))
}
pub(crate) fn park_filter(
&self,
flow: Box<Flow>,
filter: WaitFilter,
) -> Result<Result<(), Box<Flow>>, RuntimeError> {
let mut inner = sync_lock::lock(&self.inner, "Mailbox::park_filter")?;
if let Some(value) = take_with_filter(&mut inner.queue, filter) {
drop(inner);
return Ok(Err(with_pending(flow, value)));
}
inner.parked_filter = filter;
inner.parked = Some(flow);
Ok(Ok(()))
}
pub 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())
}
}
impl Default for Mailbox {
fn default() -> Self {
Self::new()
}
}
fn take_with_filter(queue: &mut VecDeque<Value>, filter: WaitFilter) -> Option<Value> {
match filter {
WaitFilter::Any => queue.pop_front(),
other => {
let idx = queue.iter().position(|v| other.matches(v))?;
queue.remove(idx)
}
}
}
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::ChunkBuilder;
use std::sync::Arc;
fn dummy_flow() -> Box<Flow> {
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, &[]).expect("vm");
let (tx, _rx) = oneshot::channel();
Box::new(Flow::new(
next_flow_id(),
vm,
Arc::new(Mailbox::new()),
RestartPolicy::Never,
tx,
))
}
fn hop(sender: u64, request_id: u64, tag: u16, payload: u64) -> Value {
Value::Message(Message::new(sender, request_id, tag, payload))
}
fn msg(tag: u16, payload: u64) -> Value {
hop(1, 1, tag, payload)
}
#[test]
fn try_pop_match_skips_non_matching_fifo() {
let mb = Mailbox::new();
mb.push(msg(9, 1)).unwrap();
mb.push(msg(1, 42)).unwrap();
mb.push(msg(9, 2)).unwrap();
let got = mb.try_pop_match(1).unwrap().expect("match");
assert_eq!(got.as_message().unwrap().payload, 42);
assert_eq!(mb.try_pop().unwrap().unwrap().as_message().unwrap().tag, 9);
assert_eq!(mb.try_pop().unwrap().unwrap().as_message().unwrap().payload, 2);
}
#[test]
fn push_while_park_match_queues_junk_keeps_waiter() {
let mb = Mailbox::new();
let flow = dummy_flow();
assert!(mb.park_match(flow, 1).unwrap().is_ok());
assert!(matches!(mb.push(msg(9, 0)).unwrap(), Delivery::Queued));
assert!(matches!(mb.push(msg(1, 7)).unwrap(), Delivery::Handoff(_)));
assert_eq!(mb.try_pop().unwrap().unwrap().as_message().unwrap().tag, 9);
}
#[test]
fn ask_does_not_consume_reply_for_another_request() {
let mb = Mailbox::new();
mb.push(hop(10, 2, 2, 99)).unwrap();
mb.push(hop(10, 1, 2, 42)).unwrap();
let filter = WaitFilter::Correlation {
expect_request_id: 1,
expect_sender: Some(10),
};
let got = mb.try_pop_filter(filter).unwrap().expect("id=1");
assert_eq!(got.as_message().unwrap().payload, 42);
let left = mb.try_pop().unwrap().unwrap();
assert_eq!(left.as_message().unwrap().request_id, 2);
}
#[test]
fn ask_requires_reply_from_target() {
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).unwrap().is_ok());
assert!(matches!(
mb.push(hop(99, 1, 2, 0)).unwrap(),
Delivery::Queued
));
assert!(matches!(
mb.push(hop(10, 1, 2, 42)).unwrap(),
Delivery::Handoff(_)
));
assert_eq!(mb.try_pop().unwrap().unwrap().as_message().unwrap().sender, 99);
}
}