use std::sync::Arc;
use crate::{CommsConfig, Inbox, InboxSender, Keypair, Router, TrustedPeers};
use super::types::CommsMessage;
pub struct CommsManagerConfig {
pub keypair: Keypair,
pub trusted_peers: TrustedPeers,
pub comms_config: CommsConfig,
}
impl CommsManagerConfig {
pub fn new() -> Self {
Self {
keypair: Keypair::generate(),
trusted_peers: TrustedPeers::new(),
comms_config: CommsConfig::default(),
}
}
pub fn with_keypair(keypair: Keypair) -> Self {
Self {
keypair,
trusted_peers: TrustedPeers::new(),
comms_config: CommsConfig::default(),
}
}
pub fn trusted_peers(mut self, trusted_peers: TrustedPeers) -> Self {
self.trusted_peers = trusted_peers;
self
}
pub fn comms_config(mut self, config: CommsConfig) -> Self {
self.comms_config = config;
self
}
}
impl Default for CommsManagerConfig {
fn default() -> Self {
Self::new()
}
}
pub struct CommsManager {
keypair: Arc<Keypair>,
trusted_peers: Arc<TrustedPeers>,
inbox: Inbox,
inbox_sender: InboxSender,
router: Arc<Router>,
}
impl CommsManager {
pub fn new(config: CommsManagerConfig) -> std::io::Result<Self> {
let (inbox, inbox_sender) = Inbox::new();
let trusted_peers = Arc::new(config.trusted_peers.clone());
let router = Router::new(
config.keypair,
config.trusted_peers,
config.comms_config,
inbox_sender.clone(),
true,
);
let keypair = router.keypair_arc();
let router = Arc::new(router);
Ok(Self {
keypair,
trusted_peers,
inbox,
inbox_sender,
router,
})
}
pub fn keypair(&self) -> &Keypair {
self.keypair.as_ref()
}
pub fn keypair_arc(&self) -> Arc<Keypair> {
self.keypair.clone()
}
pub fn public_key(&self) -> crate::PubKey {
self.keypair.public_key()
}
pub fn trusted_peers(&self) -> &Arc<TrustedPeers> {
&self.trusted_peers
}
pub fn inbox_sender(&self) -> &InboxSender {
&self.inbox_sender
}
pub fn router(&self) -> &Arc<Router> {
&self.router
}
pub fn drain_messages(&mut self) -> Vec<CommsMessage> {
let items = self.inbox.try_drain();
items
.iter()
.filter_map(|item| CommsMessage::from_inbox_item(item, &self.trusted_peers, true))
.collect()
}
pub async fn recv_message(&mut self) -> Option<CommsMessage> {
loop {
let item = self.inbox.recv().await?;
if let Some(msg) = CommsMessage::from_inbox_item(&item, &self.trusted_peers, true) {
return Some(msg);
}
}
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
use crate::{Envelope, InboxItem, MessageKind, Signature, TrustedPeer};
use uuid::Uuid;
fn make_keypair() -> Keypair {
Keypair::generate()
}
fn make_trusted_peers(name: &str, pubkey: &crate::PubKey) -> TrustedPeers {
TrustedPeers {
peers: vec![TrustedPeer {
name: name.to_string(),
pubkey: *pubkey,
addr: "tcp://127.0.0.1:4200".to_string(),
meta: crate::PeerMeta::default(),
}],
}
}
#[test]
fn test_comms_manager_struct() {
let config = CommsManagerConfig::new();
let manager = CommsManager::new(config).unwrap();
let _ = manager.keypair();
let _ = manager.public_key();
let _ = manager.trusted_peers();
let _ = manager.inbox_sender();
let _ = manager.router();
}
#[test]
fn test_comms_manager_new() {
let keypair = make_keypair();
let pubkey = keypair.public_key();
let trusted = make_trusted_peers("test-peer", &pubkey);
let config = CommsManagerConfig::with_keypair(keypair).trusted_peers(trusted);
let manager = CommsManager::new(config).unwrap();
assert_eq!(manager.trusted_peers().peers.len(), 1);
assert_eq!(manager.trusted_peers().peers[0].name, "test-peer");
}
#[test]
fn test_comms_manager_drain_empty() {
let config = CommsManagerConfig::new();
let mut manager = CommsManager::new(config).unwrap();
let messages = manager.drain_messages();
assert!(messages.is_empty());
}
#[tokio::test]
async fn test_comms_manager_drain() {
let sender = make_keypair();
let sender_pubkey = sender.public_key();
let our_keypair = make_keypair();
let our_pubkey = our_keypair.public_key();
let trusted = make_trusted_peers("sender-agent", &sender_pubkey);
let config = CommsManagerConfig::with_keypair(our_keypair).trusted_peers(trusted);
let mut manager = CommsManager::new(config).unwrap();
let mut envelope = Envelope {
id: Uuid::new_v4(),
from: sender_pubkey,
to: our_pubkey,
kind: MessageKind::Message {
body: "hello from sender".to_string(),
},
sig: Signature::new([0u8; 64]),
};
envelope.sign(&sender);
manager
.inbox_sender()
.send(InboxItem::External { envelope })
.unwrap();
tokio::task::yield_now().await;
let messages = manager.drain_messages();
assert_eq!(messages.len(), 1);
assert_eq!(messages[0].from_peer, "sender-agent");
}
#[test]
fn test_comms_manager_router_access() {
let config = CommsManagerConfig::new();
let manager = CommsManager::new(config).unwrap();
let router = manager.router();
let _ = Arc::strong_count(router);
}
#[test]
fn test_comms_manager_config_builder() {
let keypair = make_keypair();
let keypair_pubkey = keypair.public_key();
let trusted = TrustedPeers::new();
let comms_config = CommsConfig {
ack_timeout_secs: 60,
max_message_bytes: 2_000_000,
};
let config = CommsManagerConfig::with_keypair(keypair)
.trusted_peers(trusted)
.comms_config(comms_config);
let manager = CommsManager::new(config).unwrap();
assert_eq!(manager.public_key(), keypair_pubkey);
}
}