use crate::identifiers::{ClientIdentifier, GroupIdentifier, ReplicaIdentifier};
use crate::model::{Message, Payload};
use crate::stamps::View;
use std::collections::VecDeque;
#[derive(Copy, Clone, Debug, Eq, PartialEq, Hash)]
pub enum Address {
Replica(ReplicaIdentifier),
Group(GroupIdentifier),
Client(ClientIdentifier),
}
impl From<ReplicaIdentifier> for Address {
fn from(value: ReplicaIdentifier) -> Self {
Self::Replica(value)
}
}
impl From<GroupIdentifier> for Address {
fn from(value: GroupIdentifier) -> Self {
Self::Group(value)
}
}
impl From<ClientIdentifier> for Address {
fn from(value: ClientIdentifier) -> Self {
Self::Client(value)
}
}
#[derive(Debug)]
pub struct Mailbox {
address: Address,
group: Address,
inbound: Vec<Option<Message>>,
outbound: VecDeque<Message>,
}
impl From<ReplicaIdentifier> for Mailbox {
fn from(value: ReplicaIdentifier) -> Self {
Self::new(value.into(), value.group().into())
}
}
impl From<ClientIdentifier> for Mailbox {
fn from(value: ClientIdentifier) -> Self {
Self::new(value.into(), value.into())
}
}
impl Mailbox {
pub fn new(address: Address, group: Address) -> Self {
Self {
address,
group,
inbound: Default::default(),
outbound: Default::default(),
}
}
pub fn is_empty(&self) -> bool {
self.inbound.iter().all(Option::is_none)
}
pub fn deliver(&mut self, message: Message) {
for slot in self.inbound.iter_mut().rev() {
if slot.is_none() {
*slot = Some(message);
return;
}
}
self.inbound.push(Some(message));
}
pub fn drain_inbound(&mut self) -> impl Iterator<Item = Message> + '_ {
self.inbound.drain(..).filter_map(|o| o)
}
pub fn drain_outbound(&mut self) -> impl Iterator<Item = Message> + '_ {
self.outbound.drain(..)
}
pub fn select<F: FnMut(&mut Self, Message) -> Option<Message>>(&mut self, mut f: F) {
for index in 0..self.inbound.len() {
if let Some(message) = self.inbound.get_mut(index).and_then(Option::take) {
self.inbound[index] = f(self, message);
if self.inbound[index].is_none() {
break;
}
}
}
}
pub fn select_all<F: FnMut(&mut Self, Message) -> Option<Message>>(&mut self, mut f: F) {
for index in 0..self.inbound.len() {
if let Some(message) = self.inbound.get_mut(index).and_then(Option::take) {
self.inbound[index] = f(self, message);
}
}
}
pub fn visit<F: FnMut(&Message)>(&mut self, mut f: F) {
for slot in self.inbound.iter() {
if let Some(message) = slot {
f(message);
}
}
}
pub fn send(&mut self, to: impl Into<Address>, view: View, payload: impl Into<Payload>) {
let from = self.address;
let to = to.into();
let payload = payload.into();
let message = Message {
from,
to,
view,
payload,
};
if to == self.address {
self.deliver(message);
} else {
self.outbound.push_back(message);
}
}
pub fn broadcast(&mut self, view: View, payload: impl Into<Payload>) {
let from = self.address;
let to = self.group;
let payload = payload.into();
self.outbound.push_back(Message {
from,
to,
view,
payload,
});
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::client::Client;
use crate::identifiers::{ClientIdentifier, GroupIdentifier};
use crate::model::{Commit, Request};
use crate::stamps::{OpNumber, View};
#[test]
fn mailbox() {
let group = GroupIdentifier::default();
let client = ClientIdentifier::default();
let client_address = Address::from(client);
let replica = Address::from(group.into_iter().next().unwrap());
let view = View::default();
let request = Request {
op: vec![],
c: client,
s: Default::default(),
};
let mut instance = Mailbox::new(replica, Address::from(group));
instance.inbound = vec![
None,
Some(Message {
from: client_address,
to: replica,
view,
payload: request.clone().into(),
}),
Some(Message {
from: client_address,
to: replica,
view,
payload: request.clone().into(),
}),
];
instance.select(|_, m| Some(m));
assert!(!instance.inbound.iter().all(Option::is_none));
instance.select(|_, _| None);
assert!(!instance.inbound.iter().all(Option::is_none));
instance.select(|_, _| None);
assert!(instance.inbound.iter().all(Option::is_none));
}
#[test]
fn mailbox_select_all() {
let group = GroupIdentifier::default();
let client = ClientIdentifier::default();
let client_address = Address::from(client);
let replica = Address::from(group.into_iter().next().unwrap());
let view = View::default();
let request = Request {
op: vec![],
c: client,
s: Default::default(),
};
let message = Message {
from: client_address,
to: replica,
view,
payload: request.clone().into(),
};
let mut instance = Mailbox::new(replica, Address::from(group));
instance.inbound = vec![
None,
Some(message.clone()),
Some(message.clone()),
Some(Message {
from: client_address,
to: client_address,
view,
payload: request.clone().into(),
}),
];
instance.select_all(|_, m| Some(m));
assert!(!instance.inbound.iter().all(Option::is_none));
instance.select_all(|_, _| None);
assert!(instance.inbound.iter().all(Option::is_none));
}
#[test]
fn mailbox_visit() {
let group = GroupIdentifier::default();
let client = ClientIdentifier::default();
let client_address = Address::from(client);
let replica = Address::from(group.into_iter().next().unwrap());
let view = View::default();
let request = Request {
op: vec![],
c: client,
s: Default::default(),
};
let message = Message {
from: client_address,
to: replica,
view,
payload: request.clone().into(),
};
let mut instance = Mailbox::new(replica, Address::from(group));
instance.inbound = vec![None, Some(message.clone()), Some(message.clone())];
let mut counter = 0;
instance.visit(|_| counter += 1);
assert_eq!(counter, 2);
}
#[test]
fn deliver() {
let group = GroupIdentifier::default();
let client = ClientIdentifier::default();
let client_address = Address::from(client);
let replica = Address::from(group.into_iter().next().unwrap());
let view = View::default();
let request = Request {
op: vec![],
c: client,
s: Default::default(),
};
let mut instance = Mailbox::new(replica, Address::from(group));
instance.inbound = vec![
None,
Some(Message {
from: replica,
to: replica,
view,
payload: request.clone().into(),
}),
];
instance.deliver(Message {
from: client_address,
to: replica,
view,
payload: request.clone().into(),
});
assert!(instance.inbound.iter().all(Option::is_some));
assert_eq!(instance.inbound.len(), 2);
instance.deliver(Message {
from: client_address,
to: replica,
view,
payload: request.clone().into(),
});
assert!(instance.inbound.iter().all(Option::is_some));
assert_eq!(instance.inbound.len(), 3);
}
#[test]
fn send_self() {
let group = GroupIdentifier::default();
let replica = Address::from(group.into_iter().next().unwrap());
let view = View::default();
let mut instance = Mailbox::new(replica, Address::from(group));
instance.send(
replica,
view,
Commit {
k: OpNumber::default(),
},
);
assert!(instance.inbound.iter().all(Option::is_some));
assert_eq!(instance.inbound.len(), 1);
}
#[test]
fn empty() {
let instance = Mailbox::from(ClientIdentifier::default());
assert!(instance.is_empty());
}
#[test]
fn not_empty() {
let group = GroupIdentifier::new(3);
let client = Client::new(group);
let mut instance = Mailbox::from(client.identifier());
instance.deliver(Message {
from: client.address(),
to: group.into(),
view: client.view(),
payload: Payload::OutdatedView,
});
assert!(!instance.is_empty());
}
#[test]
fn empty_slot() {
let group = GroupIdentifier::new(3);
let client = Client::new(group);
let mut instance = Mailbox::from(client.identifier());
instance.deliver(Message {
from: client.address(),
to: group.into(),
view: client.view(),
payload: Payload::OutdatedView,
});
instance.drain_inbound().count();
assert!(instance.is_empty());
}
}