1mod capacity;
33mod metrics;
34mod policy;
35mod queue;
36
37use std::collections::VecDeque;
38use std::sync::Mutex;
39
40use crate::bytecode::Value;
41
42use super::error::RuntimeError;
43use super::process::{Flow, FlowId};
44use super::sync_lock;
45
46pub use capacity::{MailboxBytes, MailboxCapacity};
47pub use metrics::MailboxStats;
48pub use policy::{MailboxConfig, OverflowPolicy};
49
50use queue::{EnqueueEffect, MailboxQueue};
51
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub(crate) enum WaitFilter {
58 Any,
60 Tag(u16),
62 Correlation {
70 expect_request_id: u64,
71 expect_sender: Option<u64>,
72 },
73}
74
75impl WaitFilter {
76 #[inline]
77 pub(crate) fn matches(&self, value: &Value) -> bool {
78 match *self {
79 Self::Any => true,
80 Self::Tag(expected_tag) => match value.as_message() {
81 Some(m) => m.tag == expected_tag,
82 None => false,
83 },
84 Self::Correlation {
85 expect_request_id,
86 expect_sender,
87 } => match value.as_message() {
88 Some(m) => {
89 let id_ok = m.request_id == expect_request_id;
90 let sender_ok = match expect_sender {
91 Some(s) => m.sender == s,
92 None => true,
93 };
94 id_ok && sender_ok
95 }
96 None => false,
97 },
98 }
99 }
100}
101
102pub struct Mailbox {
151 inner: Mutex<MailboxInner>,
152 config: MailboxConfig,
153}
154
155struct MailboxInner {
156 queue: MailboxQueue,
157 parked: Option<Box<Flow>>,
158 parked_filter: WaitFilter,
160 wait_epoch: u64,
163 stats: MailboxStats,
164 waiting_senders: VecDeque<WaitingSender>,
167 closed: bool,
170}
171
172struct WaitingSender {
173 flow: Box<Flow>,
174 message: Value,
175}
176
177pub enum Delivery {
185 Queued,
186 QueuedDropOldest,
187 DroppedNewest,
188 Handoff(Box<Flow>),
189}
190
191#[derive(Debug, Clone, Copy, PartialEq, Eq)]
199pub enum MailboxFullReason {
200 MessageLimit,
202 ByteLimit,
206}
207
208impl std::fmt::Display for MailboxFullReason {
209 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
210 match self {
211 MailboxFullReason::MessageLimit => write!(f, "hop count limit"),
212 MailboxFullReason::ByteLimit => write!(f, "byte budget"),
213 }
214 }
215}
216
217#[derive(Debug, Clone, Copy, PartialEq, Eq)]
219pub struct MailboxFull {
220 reason: MailboxFullReason,
221}
222
223impl MailboxFull {
224 #[inline]
225 pub const fn reason(self) -> MailboxFullReason {
226 self.reason
227 }
228}
229
230impl std::fmt::Display for MailboxFull {
231 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
232 write!(f, "mailbox full ({})", self.reason)
233 }
234}
235
236impl std::error::Error for MailboxFull {}
237
238#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
264pub struct WaitEpoch(u64);
265
266impl WaitEpoch {
267 #[inline]
270 pub const fn get(self) -> u64 {
271 self.0
272 }
273}
274
275impl std::fmt::Display for WaitEpoch {
276 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
277 write!(f, "wait#{}", self.0)
278 }
279}
280
281impl Mailbox {
282 pub fn new() -> Self {
283 Self::with_config(MailboxConfig::DEFAULT)
284 }
285
286 pub fn with_config(config: MailboxConfig) -> Self {
287 Mailbox {
288 inner: Mutex::new(MailboxInner {
289 queue: MailboxQueue::new(config.capacity().get(), config.bytes().get()),
290 parked: None,
291 parked_filter: WaitFilter::Any,
292 wait_epoch: 0,
293 stats: MailboxStats::default(),
294 waiting_senders: VecDeque::new(),
295 closed: false,
296 }),
297 config,
298 }
299 }
300
301 #[inline]
302 pub fn config(&self) -> MailboxConfig {
303 self.config
304 }
305
306 pub fn stats(&self) -> Result<MailboxStats, RuntimeError> {
310 let inner = sync_lock::lock(&self.inner, "Mailbox::stats")?;
311 let mut stats = inner.stats;
312 stats.queued_messages = inner.queue.len();
313 stats.queued_bytes = inner.queue.bytes();
314 Ok(stats)
315 }
316
317 pub fn push(&self, value: Value) -> Result<Result<Delivery, MailboxFull>, RuntimeError> {
329 let mut inner = sync_lock::lock(&self.inner, "Mailbox::push")?;
330 if let Some(flow) = inner.parked.take() {
331 if inner.parked_filter.matches(&value) {
332 inner.parked_filter = WaitFilter::Any;
333 inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
334 inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
335 return Ok(Ok(Delivery::Handoff(flow)));
336 }
337 inner.parked = Some(flow);
338 return Ok(enqueue_locked(
339 &mut inner,
340 value,
341 self.config.overflow(),
342 ));
343 }
344 Ok(enqueue_locked(
345 &mut inner,
346 value,
347 self.config.overflow(),
348 ))
349 }
350
351 pub fn try_pop(&self) -> Result<Option<Value>, RuntimeError> {
353 self.try_pop_filter(WaitFilter::Any)
354 }
355
356 pub fn try_pop_match(&self, tag: u16) -> Result<Option<Value>, RuntimeError> {
358 self.try_pop_filter(WaitFilter::Tag(tag))
359 }
360
361 pub(crate) fn try_pop_filter(
363 &self,
364 filter: WaitFilter,
365 ) -> Result<Option<Value>, RuntimeError> {
366 let mut inner = sync_lock::lock(&self.inner, "Mailbox::try_pop_filter")?;
367 let got = inner.queue.take(filter);
368 if got.is_some() {
369 inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
370 }
371 Ok(got)
372 }
373
374 pub fn park(&self, flow: Box<Flow>) -> Result<Result<WaitEpoch, Box<Flow>>, RuntimeError> {
387 self.park_filter(flow, WaitFilter::Any)
388 }
389
390 pub fn park_match(
394 &self,
395 flow: Box<Flow>,
396 tag: u16,
397 ) -> Result<Result<WaitEpoch, Box<Flow>>, RuntimeError> {
398 self.park_filter(flow, WaitFilter::Tag(tag))
399 }
400
401 pub(crate) fn park_filter(
404 &self,
405 flow: Box<Flow>,
406 filter: WaitFilter,
407 ) -> Result<Result<WaitEpoch, Box<Flow>>, RuntimeError> {
408 let mut inner = sync_lock::lock(&self.inner, "Mailbox::park_filter")?;
409 if let Some(value) = inner.queue.take(filter) {
410 inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
411 drop(inner);
412 return Ok(Err(with_pending(flow, value)));
413 }
414 inner.wait_epoch = inner.wait_epoch.wrapping_add(1);
418 let epoch = WaitEpoch(inner.wait_epoch);
419 inner.parked_filter = filter;
420 inner.parked = Some(flow);
421 Ok(Ok(epoch))
422 }
423
424 pub fn take_parked_at(&self, epoch: WaitEpoch) -> Result<Option<Box<Flow>>, RuntimeError> {
437 let mut inner = sync_lock::lock(&self.inner, "Mailbox::take_parked_at")?;
438 if inner.wait_epoch != epoch.0 {
439 return Ok(None);
440 }
441 match inner.parked.take() {
442 Some(flow) => {
443 inner.parked_filter = WaitFilter::Any;
444 Ok(Some(flow))
445 }
446 None => Ok(None),
447 }
448 }
449
450 pub(crate) fn take_parked(&self) -> Result<Option<Box<Flow>>, RuntimeError> {
452 let mut inner = sync_lock::lock(&self.inner, "Mailbox::take_parked")?;
453 inner.parked_filter = WaitFilter::Any;
454 Ok(inner.parked.take())
455 }
456
457 pub(crate) fn push_system(&self, value: Value) -> Result<Delivery, RuntimeError> {
459 let mut inner = sync_lock::lock(&self.inner, "Mailbox::push_system")?;
460 if let Some(flow) = inner.parked.take() {
461 if inner.parked_filter.matches(&value) {
462 inner.parked_filter = WaitFilter::Any;
463 inner.stats.dequeued = inner.stats.dequeued.saturating_add(1);
464 inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
465 return Ok(Delivery::Handoff(flow));
466 }
467 inner.parked = Some(flow);
468 }
469 inner.queue.force_push(value);
470 inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
471 Ok(Delivery::Queued)
472 }
473
474 pub(crate) fn close(&self) -> Result<Vec<Flow>, RuntimeError> {
479 let mut inner = sync_lock::lock(&self.inner, "Mailbox::close")?;
480 inner.closed = true;
481 Ok(inner.waiting_senders.drain(..).map(|w| *w.flow).collect())
482 }
483
484 pub(crate) fn park_sender(&self, flow: Box<Flow>, message: Value) -> ParkSender {
489 match sync_lock::lock(&self.inner, "Mailbox::park_sender") {
490 Ok(inner) if inner.closed => ParkSender::Closed(flow),
491 Ok(mut inner) => {
492 inner.waiting_senders.push_back(WaitingSender { flow, message });
493 ParkSender::Parked
494 }
495 Err(e) => {
496 super::error::report_fault(e);
497 ParkSender::Closed(flow)
498 }
499 }
500 }
501
502 pub(crate) fn admit_waiting_sender(&self) -> Result<Option<Box<Flow>>, RuntimeError> {
504 let mut inner = sync_lock::lock(&self.inner, "Mailbox::admit_waiting_sender")?;
505 if inner.closed {
506 return Ok(None);
507 }
508 let Some(waiter) = inner.waiting_senders.pop_front() else {
509 return Ok(None);
510 };
511 match enqueue_locked(&mut inner, waiter.message.clone(), OverflowPolicy::Reject) {
512 Ok(_) => Ok(Some(waiter.flow)),
513 Err(_) => {
514 inner.waiting_senders.push_front(waiter);
515 Ok(None)
516 }
517 }
518 }
519
520 pub(crate) fn take_waiting_sender(
522 &self,
523 sender: FlowId,
524 ) -> Result<Option<Box<Flow>>, RuntimeError> {
525 let mut inner = sync_lock::lock(&self.inner, "Mailbox::take_waiting_sender")?;
526 if let Some(pos) = inner.waiting_senders.iter().position(|w| w.flow.id == sender) {
527 return Ok(inner.waiting_senders.remove(pos).map(|w| w.flow));
528 }
529 Ok(None)
530 }
531
532}
533
534pub(crate) enum ParkSender {
536 Parked,
537 Closed(Box<Flow>),
539}
540
541impl Default for Mailbox {
542 fn default() -> Self {
543 Self::new()
544 }
545}
546
547fn enqueue_locked(
548 inner: &mut MailboxInner,
549 value: Value,
550 policy: OverflowPolicy,
551) -> Result<Delivery, MailboxFull> {
552 match inner.queue.enqueue(value, policy) {
553 Ok(EnqueueEffect::Enqueued) => {
554 inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
555 Ok(Delivery::Queued)
556 }
557 Ok(EnqueueEffect::DroppedOldest) => {
558 inner.stats.dropped_oldest = inner.stats.dropped_oldest.saturating_add(1);
559 inner.stats.enqueued = inner.stats.enqueued.saturating_add(1);
560 Ok(Delivery::QueuedDropOldest)
561 }
562 Ok(EnqueueEffect::DroppedNewest) => {
563 inner.stats.dropped_newest = inner.stats.dropped_newest.saturating_add(1);
564 Ok(Delivery::DroppedNewest)
565 }
566 Err(reason) => {
567 inner.stats.rejected = inner.stats.rejected.saturating_add(1);
568 if reason == MailboxFullReason::ByteLimit {
569 inner.stats.rejected_byte_limit =
570 inner.stats.rejected_byte_limit.saturating_add(1);
571 }
572 Err(MailboxFull { reason })
573 }
574 }
575}
576
577fn with_pending(mut flow: Box<Flow>, value: Value) -> Box<Flow> {
581 flow.pending_message = Some(value);
582 flow
583}
584
585#[cfg(test)]
586mod tests {
587 use super::*;
588 use crate::bytecode::Message;
589 use crate::scheduler::oneshot;
590 use crate::scheduler::process::{next_flow_id, RestartPolicy};
591 use crate::vm::{NativeTable, Vm};
592 use crate::bytecode::builder::ChunkBuilder;
593 use std::sync::Arc;
594
595 type TestResult = Result<(), Box<dyn std::error::Error>>;
596
597 fn dummy_flow() -> Result<Box<Flow>, Box<dyn std::error::Error>> {
598 let mut b = ChunkBuilder::new("mb");
599 b.begin_function("main", 0, 1);
600 b.emit_return(0);
601 let chunk = b.finish();
602 let vm = Vm::new(Arc::new(chunk), NativeTable::empty(), 0, &[])?;
603 let (tx, _rx) = oneshot::channel();
604 Ok(Box::new(Flow::new(
605 next_flow_id(),
606 vm,
607 Arc::new(Mailbox::new()),
608 RestartPolicy::Never,
609 tx,
610 )))
611 }
612
613 fn hop(sender: u64, request_id: u64, tag: u16, payload: u64) -> Value {
614 Value::Message(Message::new(sender, request_id, tag, payload))
615 }
616
617 fn msg(tag: u16, payload: u64) -> Value {
618 hop(1, 1, tag, payload)
619 }
620
621 fn tiny_reject(n: u32) -> Result<Mailbox, Box<dyn std::error::Error>> {
622 let cap = MailboxCapacity::new(n).ok_or("invalid mailbox capacity")?;
623 Ok(Mailbox::with_config(MailboxConfig::new(
624 cap,
625 OverflowPolicy::Reject,
626 )))
627 }
628
629 fn hop_msg(value: &Value) -> Result<Message, Box<dyn std::error::Error>> {
630 value.as_message().ok_or("expected Message hop".into())
631 }
632
633 #[test]
634 fn try_pop_match_skips_non_matching_fifo() -> TestResult {
635 let mb = Mailbox::new();
636 mb.push(msg(9, 1))??;
637 mb.push(msg(1, 42))??;
638 mb.push(msg(9, 2))??;
639 let got = mb.try_pop_match(1)?.ok_or("match")?;
640 assert_eq!(hop_msg(&got)?.payload, 42);
641 assert_eq!(
642 hop_msg(&mb.try_pop()?.ok_or("first leftover")?)?.tag,
643 9
644 );
645 assert_eq!(
646 hop_msg(&mb.try_pop()?.ok_or("second leftover")?)?.payload,
647 2
648 );
649 Ok(())
650 }
651
652 #[test]
653 fn push_while_park_match_queues_junk_keeps_waiter() -> TestResult {
654 let mb = Mailbox::new();
655 let flow = dummy_flow()?;
656 assert!(mb.park_match(flow, 1)?.is_ok());
657 assert!(matches!(mb.push(msg(9, 0))??, Delivery::Queued));
658 assert!(matches!(mb.push(msg(1, 7))??, Delivery::Handoff(_)));
659 assert_eq!(hop_msg(&mb.try_pop()?.ok_or("queued junk")?)?.tag, 9);
660 Ok(())
661 }
662
663 #[test]
664 fn ask_does_not_consume_reply_for_another_request() -> TestResult {
665 let mb = Mailbox::new();
666 mb.push(hop(10, 2, 2, 99))??;
667 mb.push(hop(10, 1, 2, 42))??;
668 let filter = WaitFilter::Correlation {
669 expect_request_id: 1,
670 expect_sender: Some(10),
671 };
672 let got = mb.try_pop_filter(filter)?.ok_or("id=1")?;
673 assert_eq!(hop_msg(&got)?.payload, 42);
674 let left = mb.try_pop()?.ok_or("leftover")?;
675 assert_eq!(hop_msg(&left)?.request_id, 2);
676 Ok(())
677 }
678
679 #[test]
680 fn ask_requires_reply_from_target() -> TestResult {
681 let mb = Mailbox::new();
682 let flow = dummy_flow()?;
683 let filter = WaitFilter::Correlation {
684 expect_request_id: 1,
685 expect_sender: Some(10),
686 };
687 assert!(mb.park_filter(flow, filter)?.is_ok());
688 assert!(matches!(mb.push(hop(99, 1, 2, 0))??, Delivery::Queued));
689 assert!(matches!(mb.push(hop(10, 1, 2, 42))??, Delivery::Handoff(_)));
690 assert_eq!(
691 hop_msg(&mb.try_pop()?.ok_or("non-matching queued")?)?.sender,
692 99
693 );
694 Ok(())
695 }
696
697 #[test]
698 fn reject_when_full_without_waiter() -> TestResult {
699 let mb = tiny_reject(1)?;
700 assert!(matches!(mb.push(msg(1, 1))??, Delivery::Queued));
701 match mb.push(msg(1, 2))? {
705 Err(full) => assert_eq!(full.reason(), MailboxFullReason::MessageLimit),
706 Ok(_) => return Err("expected the hop count bound to refuse".into()),
707 }
708 let s = mb.stats()?;
709 assert_eq!(s.enqueued, 1);
710 assert_eq!(s.rejected, 1);
711 assert_eq!(s.rejected_byte_limit, 0);
712 assert_eq!(s.queued_messages, 1);
713 Ok(())
714 }
715
716 fn park_now(mb: &Mailbox, flow: Box<Flow>) -> Result<WaitEpoch, Box<dyn std::error::Error>> {
719 match mb.park(flow)? {
720 Ok(epoch) => Ok(epoch),
721 Err(_) => Err("an empty mailbox should have parked the flow".into()),
722 }
723 }
724
725 fn handoff(mb: &Mailbox, value: Value) -> Result<Box<Flow>, Box<dyn std::error::Error>> {
726 match mb.push(value)?? {
727 Delivery::Handoff(flow) => Ok(flow),
728 other => Err(format!("expected a handoff, got {}", delivery_name(&other)).into()),
729 }
730 }
731
732 fn delivery_name(d: &Delivery) -> &'static str {
733 match d {
734 Delivery::Queued => "Queued",
735 Delivery::QueuedDropOldest => "QueuedDropOldest",
736 Delivery::DroppedNewest => "DroppedNewest",
737 Delivery::Handoff(_) => "Handoff",
738 }
739 }
740
741 #[test]
742 fn a_stale_deadline_cannot_steal_a_later_wait() -> TestResult {
743 let mb = Mailbox::new();
744 let first = park_now(&mb, dummy_flow()?)?;
745 let woken = handoff(&mb, msg(1, 1))?;
747 let second = park_now(&mb, woken)?;
749 assert_ne!(first, second, "each park must get its own epoch");
750
751 assert!(mb.take_parked_at(first)?.is_none());
755 assert!(mb.take_parked_at(second)?.is_some());
757 Ok(())
758 }
759
760 #[test]
761 fn a_stale_deadline_does_not_downgrade_a_selective_waiter() -> TestResult {
762 let mb = Mailbox::new();
763 let first = park_now(&mb, dummy_flow()?)?;
764 let woken = handoff(&mb, msg(1, 1))?;
765 let second = match mb.park_match(woken, 7)? {
767 Ok(epoch) => epoch,
768 Err(_) => return Err("empty mailbox should have parked the flow".into()),
769 };
770 assert_ne!(first, second);
771
772 assert!(mb.take_parked_at(first)?.is_none());
773 assert!(matches!(mb.push(msg(9, 0))??, Delivery::Queued));
776 assert!(matches!(mb.push(msg(7, 0))??, Delivery::Handoff(_)));
778 Ok(())
779 }
780
781 #[test]
782 fn a_deadline_for_a_wait_that_a_hop_ended_does_nothing() -> TestResult {
783 let mb = Mailbox::new();
784 let epoch = park_now(&mb, dummy_flow()?)?;
785 let _woken = handoff(&mb, msg(1, 1))?;
786 assert!(mb.take_parked_at(epoch)?.is_none());
789 Ok(())
790 }
791
792 #[test]
793 fn byte_budget_refuses_before_the_hop_count_and_says_so() -> TestResult {
794 let cap = MailboxCapacity::new(64).ok_or("cap")?;
796 let budget = MailboxBytes::new(MailboxBytes::MIN).ok_or("bytes")?;
797 let mb = Mailbox::with_config(
798 MailboxConfig::new(cap, OverflowPolicy::Reject).with_bytes(budget),
799 );
800 assert!(matches!(mb.push(Value::bytes(vec![0u8; 900]))??, Delivery::Queued));
801 match mb.push(Value::bytes(vec![0u8; 900]))? {
802 Err(full) => assert_eq!(full.reason(), MailboxFullReason::ByteLimit),
803 Ok(_) => return Err("expected the byte budget to refuse".into()),
804 }
805 let s = mb.stats()?;
806 assert_eq!(s.queued_messages, 1);
807 assert!(s.queued_bytes >= 900);
808 assert_eq!(s.rejected, 1);
809 assert_eq!(s.rejected_byte_limit, 1);
810 Ok(())
811 }
812
813 #[test]
814 fn draining_a_hop_frees_its_byte_charge() -> TestResult {
815 let cap = MailboxCapacity::new(64).ok_or("cap")?;
816 let budget = MailboxBytes::new(MailboxBytes::MIN).ok_or("bytes")?;
817 let mb = Mailbox::with_config(
818 MailboxConfig::new(cap, OverflowPolicy::Reject).with_bytes(budget),
819 );
820 mb.push(Value::bytes(vec![0u8; 900]))??;
821 mb.try_pop()?.ok_or("queued blob")?;
822 assert_eq!(mb.stats()?.queued_bytes, 0);
823 assert!(matches!(mb.push(Value::bytes(vec![0u8; 900]))??, Delivery::Queued));
826 Ok(())
827 }
828
829 #[test]
830 fn matching_handoff_does_not_count_as_full() -> TestResult {
831 let mb = tiny_reject(1)?;
832 mb.push(msg(9, 0))??;
833 let flow = dummy_flow()?;
834 assert!(mb.park_match(flow, 1)?.is_ok());
835 assert!(matches!(mb.push(msg(1, 7))??, Delivery::Handoff(_)));
836 assert_eq!(hop_msg(&mb.try_pop()?.ok_or("queued")?)?.tag, 9);
837 Ok(())
838 }
839
840 #[test]
841 fn waiting_send_admits_one_after_pop() -> TestResult {
842 let mb = tiny_reject(1)?;
843 assert!(matches!(mb.push(msg(1, 1))??, Delivery::Queued));
844 assert!(matches!(
845 mb.park_sender(dummy_flow()?, msg(1, 2)),
846 ParkSender::Parked
847 ));
848 assert!(mb.try_pop()?.is_some());
849 let woken = mb.admit_waiting_sender()?.ok_or("admitted")?;
850 drop(woken);
851 assert_eq!(hop_msg(&mb.try_pop()?.ok_or("second hop")?)?.payload, 2);
852 assert!(mb.admit_waiting_sender()?.is_none());
853 Ok(())
854 }
855
856 #[test]
857 fn close_rejects_late_park_sender() -> TestResult {
858 let mb = tiny_reject(1)?;
859 let leftover = mb.close()?;
860 assert!(leftover.is_empty());
861 match mb.park_sender(dummy_flow()?, msg(1, 1)) {
862 ParkSender::Closed(_) => {}
863 ParkSender::Parked => return Err("closed mailbox must not park a sender".into()),
864 }
865 Ok(())
866 }
867
868 #[test]
869 fn drop_oldest_still_wakes_on_match() -> TestResult {
870 let cap = MailboxCapacity::new(1).ok_or("cap")?;
871 let mb = Mailbox::with_config(MailboxConfig::new(cap, OverflowPolicy::DropOldest));
872 let flow = dummy_flow()?;
873 assert!(mb.park_match(flow, 1)?.is_ok());
874 mb.push(msg(9, 1))??;
875 assert!(matches!(mb.push(msg(1, 2))??, Delivery::Handoff(_)));
876 Ok(())
877 }
878}