use anyhow::Result;
use std::net::TcpListener;
use std::sync::Arc;
use tokio::time::Duration;
use velo::Messenger;
use velo::backend::Transport;
use velo::backend::tcp::TcpTransportBuilder;
pub async fn create_messenger_tcp() -> Result<Arc<Messenger>> {
let listener = TcpListener::bind("127.0.0.1:0")?;
let transport: Arc<dyn Transport> = Arc::new(
TcpTransportBuilder::new()
.from_listener(listener)?
.build()?,
);
let messenger = Messenger::builder()
.add_transport(transport)
.build()
.await?;
tokio::time::sleep(Duration::from_millis(100)).await;
Ok(messenger)
}
pub struct MessengerPair {
pub messenger_a: Arc<Messenger>,
pub messenger_b: Arc<Messenger>,
}
pub async fn create_messenger_pair_tcp() -> Result<MessengerPair> {
let messenger_a = create_messenger_tcp().await?;
let messenger_b = create_messenger_tcp().await?;
messenger_a.register_peer(messenger_b.peer_info())?;
messenger_b.register_peer(messenger_a.peer_info())?;
tokio::time::sleep(Duration::from_millis(200)).await;
Ok(MessengerPair {
messenger_a,
messenger_b,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_create_messenger_instance() {
let messenger = create_messenger_tcp()
.await
.expect("Should create Messenger");
let peer_info = messenger.peer_info();
assert_eq!(
peer_info.instance_id().worker_id(),
messenger.instance_id().worker_id()
);
assert!(!peer_info.worker_address().as_bytes().is_empty());
let handlers = messenger.list_local_handlers();
assert!(
handlers.contains(&"_list_handlers".to_string()),
"Expected _list_handlers in local handler list: {:?}",
handlers
);
assert!(
handlers.contains(&"_hello".to_string()),
"Expected _hello in local handler list: {:?}",
handlers
);
}
#[tokio::test]
async fn test_create_messenger_pair() {
let pair = create_messenger_pair_tcp()
.await
.expect("Should create pair");
assert_ne!(
pair.messenger_a.instance_id(),
pair.messenger_b.instance_id()
);
assert_ne!(
pair.messenger_a.peer_info().worker_address().checksum(),
pair.messenger_b.peer_info().worker_address().checksum()
);
let handlers_from_a = pair
.messenger_a
.available_handlers(pair.messenger_b.instance_id())
.await
.expect("Handlers from messenger_b should be available");
assert!(
handlers_from_a.contains(&"_list_handlers".to_string()),
"messenger_a should see _list_handlers on messenger_b: {:?}",
handlers_from_a
);
assert!(
handlers_from_a.contains(&"_hello".to_string()),
"messenger_a should see _hello on messenger_b: {:?}",
handlers_from_a
);
let handlers_from_b = pair
.messenger_b
.available_handlers(pair.messenger_a.instance_id())
.await
.expect("Handlers from messenger_a should be available");
assert!(
handlers_from_b.contains(&"_list_handlers".to_string()),
"messenger_b should see _list_handlers on messenger_a: {:?}",
handlers_from_b
);
assert!(
handlers_from_b.contains(&"_hello".to_string()),
"messenger_b should see _hello on messenger_a: {:?}",
handlers_from_b
);
}
}