use std::str::FromStr;
use async_trait::async_trait;
use crate::dht::Chord;
use crate::dht::ChordStorage;
use crate::dht::PeerRingAction;
use crate::dht::PeerRingRemoteAction;
use crate::err::Error;
use crate::err::Result;
use crate::message::types::AlreadyConnected;
use crate::message::types::ConnectNodeReport;
use crate::message::types::ConnectNodeSend;
use crate::message::types::FindSuccessorReport;
use crate::message::types::FindSuccessorSend;
use crate::message::types::JoinDHT;
use crate::message::types::Message;
use crate::message::types::SyncVNodeWithSuccessor;
use crate::message::FindSuccessorReportHandler;
use crate::message::FindSuccessorThen;
use crate::message::HandleMsg;
use crate::message::LeaveDHT;
use crate::message::MessageHandler;
use crate::message::MessagePayload;
use crate::message::PayloadSender;
use crate::prelude::RTCSdpType;
use crate::transports::manager::TransportManager;
use crate::types::ice_transport::IceTrickleScheme;
#[cfg_attr(feature = "wasm", async_trait(?Send))]
#[cfg_attr(not(feature = "wasm"), async_trait)]
impl HandleMsg<LeaveDHT> for MessageHandler {
async fn handle(&self, _ctx: &MessagePayload<Message>, msg: &LeaveDHT) -> Result<()> {
self.swarm.disconnect(msg.did).await
}
}
#[cfg_attr(feature = "wasm", async_trait(?Send))]
#[cfg_attr(not(feature = "wasm"), async_trait)]
impl HandleMsg<JoinDHT> for MessageHandler {
async fn handle(&self, ctx: &MessagePayload<Message>, msg: &JoinDHT) -> Result<()> {
match self.dht.join(msg.did)? {
PeerRingAction::None => Ok(()),
PeerRingAction::RemoteAction(
next,
PeerRingRemoteAction::FindSuccessorForConnect(did),
) => {
if next != ctx.addr {
self.send_direct_message(
Message::FindSuccessorSend(FindSuccessorSend {
did,
strict: false,
then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect),
}),
next,
)
.await?;
}
Ok(())
}
_ => unreachable!(),
}
}
}
#[cfg_attr(feature = "wasm", async_trait(?Send))]
#[cfg_attr(not(feature = "wasm"), async_trait)]
impl HandleMsg<ConnectNodeSend> for MessageHandler {
async fn handle(&self, ctx: &MessagePayload<Message>, msg: &ConnectNodeSend) -> Result<()> {
let mut relay = ctx.relay.clone();
if self.dht.did != relay.destination {
if self
.swarm
.get_and_check_transport(relay.destination)
.await
.is_some()
{
relay.relay(self.dht.did, Some(relay.destination))?;
return self.forward_payload(ctx, relay).await;
} else {
let next_node = match self.dht.find_successor(relay.destination)? {
PeerRingAction::Some(node) => Some(node),
PeerRingAction::RemoteAction(node, _) => Some(node),
_ => None,
}
.ok_or(Error::MessageHandlerMissNextNode)?;
relay.relay(self.dht.did, Some(next_node))?;
return self.forward_payload(ctx, relay).await;
}
} else {
relay.relay(self.dht.did, None)?;
match self.swarm.get_and_check_transport(relay.sender()).await {
None => {
let trans = self.swarm.new_transport().await?;
trans
.register_remote_info(msg.handshake_info.to_owned().into())
.await?;
let handshake_info = trans
.get_handshake_info(self.swarm.session_manager(), RTCSdpType::Answer)
.await?
.to_string();
self.send_report_message(
Message::ConnectNodeReport(ConnectNodeReport {
transport_uuid: msg.transport_uuid.clone(),
handshake_info,
}),
ctx.tx_id,
relay,
)
.await?;
self.swarm.push_pending_transport(&trans)?;
Ok(())
}
_ => {
self.send_report_message(
Message::AlreadyConnected(AlreadyConnected),
ctx.tx_id,
relay,
)
.await
}
}
}
}
}
#[cfg_attr(feature = "wasm", async_trait(?Send))]
#[cfg_attr(not(feature = "wasm"), async_trait)]
impl HandleMsg<ConnectNodeReport> for MessageHandler {
async fn handle(&self, ctx: &MessagePayload<Message>, msg: &ConnectNodeReport) -> Result<()> {
let mut relay = ctx.relay.clone();
relay.relay(self.dht.did, None)?;
if relay.next_hop.is_some() {
self.forward_payload(ctx, relay).await
} else {
let transport = self
.swarm
.find_pending_transport(
uuid::Uuid::from_str(&msg.transport_uuid)
.map_err(|_| Error::InvalidTransportUuid)?,
)?
.ok_or(Error::MessageHandlerMissTransportConnectedNode)?;
transport
.register_remote_info(msg.handshake_info.clone().into())
.await?;
Ok(())
}
}
}
#[cfg_attr(feature = "wasm", async_trait(?Send))]
#[cfg_attr(not(feature = "wasm"), async_trait)]
impl HandleMsg<AlreadyConnected> for MessageHandler {
async fn handle(&self, ctx: &MessagePayload<Message>, _msg: &AlreadyConnected) -> Result<()> {
let mut relay = ctx.relay.clone();
relay.relay(self.dht.did, None)?;
if relay.next_hop.is_some() {
self.forward_payload(ctx, relay).await
} else {
self.swarm
.get_and_check_transport(relay.sender())
.await
.map(|_| ())
.ok_or(Error::MessageHandlerMissTransportAlreadyConnected)
}
}
}
#[cfg_attr(feature = "wasm", async_trait(?Send))]
#[cfg_attr(not(feature = "wasm"), async_trait)]
impl HandleMsg<FindSuccessorSend> for MessageHandler {
async fn handle(&self, ctx: &MessagePayload<Message>, msg: &FindSuccessorSend) -> Result<()> {
let mut relay = ctx.relay.clone();
match self.dht.find_successor(msg.did)? {
PeerRingAction::Some(did) => {
if !msg.strict || self.dht.did == msg.did {
match &msg.then {
FindSuccessorThen::Report(handler) => {
relay.relay(self.dht.did, None)?;
self.send_report_message(
Message::FindSuccessorReport(FindSuccessorReport {
did,
handler: handler.clone(),
}),
ctx.tx_id,
relay,
)
.await
}
}
} else if self.swarm.get_and_check_transport(msg.did).await.is_some() {
relay.relay(self.dht.did, Some(relay.destination))?;
return self.forward_payload(ctx, relay).await;
} else {
return Err(Error::MessageHandlerMissNextNode);
}
}
PeerRingAction::RemoteAction(next, _) => {
relay.relay(self.dht.did, Some(next))?;
relay.reset_destination(next)?;
self.forward_payload(ctx, relay).await
}
act => Err(Error::PeerRingUnexpectedAction(act)),
}
}
}
#[cfg_attr(feature = "wasm", async_trait(?Send))]
#[cfg_attr(not(feature = "wasm"), async_trait)]
impl HandleMsg<FindSuccessorReport> for MessageHandler {
async fn handle(&self, ctx: &MessagePayload<Message>, msg: &FindSuccessorReport) -> Result<()> {
let mut relay = ctx.relay.clone();
relay.relay(self.dht.did, None)?;
if relay.next_hop.is_some() {
return self.forward_payload(ctx, relay).await;
}
match &msg.handler {
FindSuccessorReportHandler::FixFingerTable => self.dht.lock_finger()?.set_fix(msg.did),
FindSuccessorReportHandler::Connect => {
if self.swarm.get_and_check_transport(msg.did).await.is_none()
&& msg.did != self.swarm.did()
{
self.swarm.connect(msg.did).await?;
}
}
FindSuccessorReportHandler::SyncStorage => {
self.dht.lock_successor()?.update(msg.did);
if let Ok(PeerRingAction::RemoteAction(
next,
PeerRingRemoteAction::SyncVNodeWithSuccessor(data),
)) = self.dht.sync_vnode_with_successor(msg.did).await
{
self.send_direct_message(
Message::SyncVNodeWithSuccessor(SyncVNodeWithSuccessor { data }),
next,
)
.await?;
return Ok(());
}
}
_ => {}
}
Ok(())
}
}
#[cfg(not(feature = "wasm"))]
#[cfg(test)]
pub mod tests {
use std::matches;
use std::sync::Arc;
use tokio::time::sleep;
use tokio::time::Duration;
use super::*;
use crate::dht::Did;
use crate::ecc::tests::gen_ordered_keys;
use crate::ecc::SecretKey;
use crate::message::handlers::tests::assert_no_more_msg;
use crate::message::MessageHandler;
use crate::swarm::Swarm;
use crate::tests::default::prepare_node;
use crate::tests::manually_establish_connection;
use crate::transports::manager::TransportManager;
use crate::types::ice_transport::IceTransportInterface;
#[tokio::test]
async fn test_triple_nodes_connection_1_2_3() -> Result<()> {
let keys = gen_ordered_keys(3);
let (key1, key2, key3) = (keys[0], keys[1], keys[2]);
test_triple_ordered_nodes_connection(key1, key2, key3).await?;
Ok(())
}
#[tokio::test]
async fn test_triple_nodes_connection_2_3_1() -> Result<()> {
let keys = gen_ordered_keys(3);
let (key1, key2, key3) = (keys[0], keys[1], keys[2]);
test_triple_ordered_nodes_connection(key2, key3, key1).await?;
Ok(())
}
#[tokio::test]
async fn test_triple_nodes_connection_3_1_2() -> Result<()> {
let keys = gen_ordered_keys(3);
let (key1, key2, key3) = (keys[0], keys[1], keys[2]);
test_triple_ordered_nodes_connection(key3, key1, key2).await?;
Ok(())
}
#[tokio::test]
async fn test_triple_nodes_connection_3_2_1() -> Result<()> {
let keys = gen_ordered_keys(3);
let (key1, key2, key3) = (keys[0], keys[1], keys[2]);
test_triple_desc_ordered_nodes_connection(key3, key2, key1).await?;
Ok(())
}
#[tokio::test]
async fn test_triple_nodes_connection_2_1_3() -> Result<()> {
let keys = gen_ordered_keys(3);
let (key1, key2, key3) = (keys[0], keys[1], keys[2]);
test_triple_desc_ordered_nodes_connection(key2, key1, key3).await?;
Ok(())
}
#[tokio::test]
async fn test_triple_nodes_connection_1_3_2() -> Result<()> {
let keys = gen_ordered_keys(3);
let (key1, key2, key3) = (keys[0], keys[1], keys[2]);
test_triple_desc_ordered_nodes_connection(key1, key3, key2).await?;
Ok(())
}
async fn test_triple_ordered_nodes_connection(
key1: SecretKey,
key2: SecretKey,
key3: SecretKey,
) -> Result<(MessageHandler, MessageHandler, MessageHandler)> {
let (did1, dht1, swarm1, node1, _path1) = prepare_node(key1).await;
let (did2, dht2, swarm2, node2, _path2) = prepare_node(key2).await;
let (did3, dht3, swarm3, node3, _path3) = prepare_node(key3).await;
println!("========================================");
println!("|| now we connect node1 and node2 ||");
println!("========================================");
test_only_two_nodes_establish_connection(&node1, &node2).await?;
assert_eq!(dht1.lock_successor()?.list(), vec![did2]);
assert_eq!(dht2.lock_successor()?.list(), vec![did1]);
assert_eq!(dht3.lock_successor()?.list(), vec![]);
println!("========================================");
println!("|| now we start join node3 to node2 ||");
println!("========================================");
manually_establish_connection(&swarm3, &swarm2).await?;
test_listen_join_and_init_find_succeesor(&node3, &node2).await?;
assert_eq!(dht1.lock_successor()?.list(), vec![did2]);
assert_eq!(dht2.lock_successor()?.list(), vec![did3, did1]);
assert_eq!(dht3.lock_successor()?.list(), vec![did2]);
let ev_3 = node3.listen_once().await.unwrap();
assert_eq!(ev_3.addr, did2);
assert_eq!(ev_3.relay.path, vec![did3, did2]);
assert!(matches!(
ev_3.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did3
));
assert!(!dht3.lock_successor()?.list().contains(&did3));
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did3);
assert_eq!(ev_2.relay.path, vec![did2, did3]);
assert!(matches!(
ev_2.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did2
));
assert!(!dht2.lock_successor()?.list().contains(&did2));
println!("=== Check state before connect via DHT ===");
assert_transports(swarm1.clone(), vec![did2]);
assert_transports(swarm2.clone(), vec![did1, did3]);
assert_transports(swarm3.clone(), vec![did2]);
assert_eq!(dht1.lock_successor()?.list(), vec![did2]);
assert_eq!(dht2.lock_successor()?.list(), vec![did3, did1]);
assert_eq!(dht3.lock_successor()?.list(), vec![did2]);
println!("=============================================");
println!("|| now we connect node1 to node3 via DHT ||");
println!("=============================================");
test_connect_via_dht_and_init_find_succeesor(&node1, &node2, &node3).await?;
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did1);
assert_eq!(ev_2.relay.path, vec![did3, did1]);
assert!(matches!(
ev_2.data,
Message::FindSuccessorSend(FindSuccessorSend{did, strict: false, then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect)}) if did == did3
));
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did3);
assert_eq!(ev_1.relay.path, vec![did1, did3]);
assert!(matches!(
ev_1.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did1
));
assert!(!dht1.lock_successor()?.list().contains(&did1));
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did2);
assert_eq!(ev_1.relay.path, vec![did3, did1, did2]);
let ev_3 = node3.listen_once().await.unwrap();
assert_eq!(ev_3.addr, did1);
assert_eq!(ev_3.relay.path, vec![did3, did1, did2]);
assert!(matches!(
ev_3.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did3
));
assert!(!dht3.lock_successor()?.list().contains(&did3));
assert_no_more_msg(&node1, &node2, &node3).await;
println!("=== Check state after connect via DHT ===");
assert_transports(swarm1, vec![did2, did3]);
assert_transports(swarm2, vec![did1, did3]);
assert_transports(swarm3, vec![did1, did2]);
assert_eq!(dht1.lock_successor()?.list(), vec![did2, did3]);
assert_eq!(dht2.lock_successor()?.list(), vec![did3, did1]);
assert_eq!(dht3.lock_successor()?.list(), vec![did1, did2]);
tokio::fs::remove_dir_all("./tmp").await.ok();
Ok((node1, node2, node3))
}
async fn test_triple_desc_ordered_nodes_connection(
key1: SecretKey,
key2: SecretKey,
key3: SecretKey,
) -> Result<(MessageHandler, MessageHandler, MessageHandler)> {
let (did1, dht1, swarm1, node1, _path1) = prepare_node(key1).await;
let (did2, dht2, swarm2, node2, _path2) = prepare_node(key2).await;
let (did3, dht3, swarm3, node3, _path3) = prepare_node(key3).await;
println!("========================================");
println!("|| now we connect node1 and node2 ||");
println!("========================================");
test_only_two_nodes_establish_connection(&node1, &node2).await?;
assert_eq!(dht1.lock_successor()?.list(), vec![did2]);
assert_eq!(dht2.lock_successor()?.list(), vec![did1]);
assert_eq!(dht3.lock_successor()?.list(), vec![]);
println!("========================================");
println!("|| now we start join node3 to node2 ||");
println!("========================================");
manually_establish_connection(&swarm3, &swarm2).await?;
test_listen_join_and_init_find_succeesor(&node3, &node2).await?;
assert_eq!(dht1.lock_successor()?.list(), vec![did2]);
assert_eq!(dht2.lock_successor()?.list(), vec![did1, did3]);
assert_eq!(dht3.lock_successor()?.list(), vec![did2]);
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did2);
assert_eq!(ev_1.relay.path, vec![did3, did2]);
assert!(matches!(
ev_1.data,
Message::FindSuccessorSend(FindSuccessorSend{did, strict: false, then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect)}) if did == did3
));
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did3);
assert_eq!(ev_2.relay.path, vec![did2, did3]);
assert!(matches!(
ev_2.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did2
));
assert!(!dht2.lock_successor()?.list().contains(&did2));
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did1);
assert_eq!(ev_2.relay.path, vec![did3, did2, did1]);
assert_eq!(ev_2.relay.path_end_cursor, 0);
assert!(matches!(
ev_2.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did2
));
let ev_3 = node3.listen_once().await.unwrap();
assert_eq!(ev_3.addr, did2);
assert_eq!(ev_3.relay.path, vec![did3, did2, did1]);
assert_eq!(ev_3.relay.path_end_cursor, 1);
assert!(matches!(
ev_3.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did2
));
println!("=== Check state before connect via DHT ===");
assert_transports(swarm1.clone(), vec![did2]);
assert_transports(swarm2.clone(), vec![did1, did3]);
assert_transports(swarm3.clone(), vec![did2]);
assert_eq!(dht1.lock_successor()?.list(), vec![did2]);
assert_eq!(dht2.lock_successor()?.list(), vec![did1, did3]);
assert_eq!(dht3.lock_successor()?.list(), vec![did2]);
println!("=============================================");
println!("|| now we connect node1 to node3 via DHT ||");
println!("=============================================");
test_connect_via_dht_and_init_find_succeesor(&node1, &node2, &node3).await?;
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did3);
assert_eq!(ev_2.relay.path, vec![did1, did3]);
assert!(matches!(
ev_2.data,
Message::FindSuccessorSend(FindSuccessorSend{did, strict: false, then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect)}) if did == did1
));
let ev_3 = node3.listen_once().await.unwrap();
assert_eq!(ev_3.addr, did1);
assert_eq!(ev_3.relay.path, vec![did3, did1]);
assert!(matches!(
ev_3.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did3
));
assert!(!node3.dht.lock_successor()?.list().contains(&did3));
let ev_3 = node3.listen_once().await.unwrap();
assert_eq!(ev_3.addr, did2);
assert_eq!(ev_3.relay.path, vec![did1, did3, did2]);
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did3);
assert_eq!(ev_1.relay.path, vec![did1, did3, did2]);
assert!(matches!(
ev_1.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did1
));
assert!(!node1.dht.lock_successor()?.list().contains(&did1));
assert_no_more_msg(&node1, &node2, &node3).await;
println!("=== Check state after connect via DHT ===");
assert_transports(swarm1, vec![did2, did3]);
assert_transports(swarm2, vec![did1, did3]);
assert_transports(swarm3, vec![did1, did2]);
assert_eq!(dht1.lock_successor()?.list(), vec![did3, did2]);
assert_eq!(dht2.lock_successor()?.list(), vec![did1, did3]);
assert_eq!(dht3.lock_successor()?.list(), vec![did2, did1]);
tokio::fs::remove_dir_all("./tmp").await.ok();
Ok((node1, node2, node3))
}
pub async fn test_listen_join_and_init_find_succeesor(
node1: &MessageHandler,
node2: &MessageHandler,
) -> Result<()> {
let did1 = node1.swarm.did();
let did2 = node2.swarm.did();
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did1);
assert_eq!(ev_1.relay.path, vec![did1]);
assert!(matches!(ev_1.data, Message::JoinDHT(JoinDHT{did, ..}) if did == did2));
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did2);
assert_eq!(ev_2.relay.path, vec![did2]);
assert!(matches!(ev_2.data, Message::JoinDHT(JoinDHT{did, ..}) if did == did1));
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did2);
assert_eq!(ev_1.relay.path, vec![did2]);
assert!(matches!(
ev_1.data,
Message::FindSuccessorSend(FindSuccessorSend{did, then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect), strict: false}) if did == did2
));
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did1);
assert_eq!(ev_2.relay.path, vec![did1]);
assert!(matches!(
ev_2.data,
Message::FindSuccessorSend(FindSuccessorSend{did, then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect), strict: false}) if did == did1
));
Ok(())
}
pub async fn test_only_two_nodes_establish_connection(
node1: &MessageHandler,
node2: &MessageHandler,
) -> Result<()> {
let did1 = node1.swarm.did();
let did2 = node2.swarm.did();
let dht1 = node1.swarm.dht();
let dht2 = node2.swarm.dht();
manually_establish_connection(&node1.swarm, &node2.swarm).await?;
test_listen_join_and_init_find_succeesor(node1, node2).await?;
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did2);
assert_eq!(ev_1.relay.path, vec![did1, did2]);
assert!(matches!(
ev_1.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did1
));
assert!(!dht1.lock_successor()?.list().contains(&did1));
let ev_2 = node2.listen_once().await.unwrap();
assert_eq!(ev_2.addr, did1);
assert_eq!(ev_2.relay.path, vec![did2, did1]);
assert!(matches!(
ev_2.data,
Message::FindSuccessorReport(FindSuccessorReport{did, handler: FindSuccessorReportHandler::Connect}) if did == did2
));
assert!(!dht2.lock_successor()?.list().contains(&did2));
Ok(())
}
async fn test_connect_via_dht_and_init_find_succeesor(
node1: &MessageHandler,
node2: &MessageHandler,
node3: &MessageHandler,
) -> Result<()> {
let did1 = node1.swarm.did();
let did2 = node2.swarm.did();
let did3 = node3.swarm.did();
assert!(node1.swarm.get_transport(did3).is_none());
assert_eq!(node1.dht.lock_successor()?.max(), did2);
node1.swarm.connect(did3).await.unwrap();
let ev2 = node2.listen_once().await.unwrap();
assert_eq!(ev2.addr, did1);
assert_eq!(ev2.relay.path, vec![did1]);
assert!(matches!(ev2.data, Message::ConnectNodeSend(_)));
let ev3 = node3.listen_once().await.unwrap();
assert_eq!(ev3.addr, did2);
assert_eq!(ev3.relay.path, vec![did1, did2]);
assert!(matches!(ev3.data, Message::ConnectNodeSend(_)));
let ev2 = node2.listen_once().await.unwrap();
assert_eq!(ev2.addr, did3);
assert_eq!(ev2.relay.path, vec![did1, did2, did3]);
assert!(matches!(ev2.data, Message::ConnectNodeReport(_)));
let ev1 = node1.listen_once().await.unwrap();
assert_eq!(ev1.addr, did2);
assert_eq!(ev1.relay.path, vec![did1, did2, did3]);
assert!(matches!(ev1.data, Message::ConnectNodeReport(_)));
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did1);
assert_eq!(ev_1.relay.path, vec![did1]);
assert!(matches!(ev_1.data, Message::JoinDHT(JoinDHT{did, ..}) if did == did3));
let ev_3 = node3.listen_once().await.unwrap();
assert_eq!(ev_3.addr, did3);
assert_eq!(ev_3.relay.path, vec![did3]);
assert!(matches!(ev_3.data, Message::JoinDHT(JoinDHT{did, ..}) if did == did1));
let ev_1 = node1.listen_once().await.unwrap();
assert_eq!(ev_1.addr, did3);
assert_eq!(ev_1.relay.path, vec![did3]);
assert!(matches!(
ev_1.data,
Message::FindSuccessorSend(FindSuccessorSend{did, strict: false, then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect)}) if did == did3
));
let ev_3 = node3.listen_once().await.unwrap();
assert_eq!(ev_3.addr, did1);
assert_eq!(ev_3.relay.path, vec![did1]);
assert!(matches!(
ev_3.data,
Message::FindSuccessorSend(FindSuccessorSend{did, strict: false, then: FindSuccessorThen::Report(FindSuccessorReportHandler::Connect)}) if did == did1
));
Ok(())
}
fn assert_transports(swarm: Arc<Swarm>, addresses: Vec<Did>) {
println!(
"Check transport of {:?}: {:?} for addresses {:?}",
swarm.did(),
swarm.get_dids(),
addresses
);
assert_eq!(swarm.get_transports().len(), addresses.len());
for addr in addresses {
assert!(swarm.get_transport(addr).is_some());
}
}
#[tokio::test]
async fn test_quadra_desc_node_connection() -> Result<()> {
let keys = gen_ordered_keys(4);
let (key1, key2, key3, key4) = (keys[0], keys[1], keys[2], keys[3]);
let (node1, node2, node3) = test_triple_ordered_nodes_connection(key1, key2, key3).await?;
let did1 = node1.swarm.did();
let did2 = node2.swarm.did();
let did3 = node3.swarm.did();
let (did4, _, swarm4, node4, _path4) = prepare_node(key4).await;
manually_establish_connection(&swarm4, &node2.swarm).await?;
test_listen_join_and_init_find_succeesor(&node4, &node2).await?;
let _ = node2.listen_once().await.unwrap();
let _ = node3.listen_once().await.unwrap();
let _ = node2.listen_once().await.unwrap();
let _ = node4.listen_once().await.unwrap();
let _ = node2.listen_once().await.unwrap();
println!("==================================================");
println!("| test connect node 4 from node 1 via node 2 |");
println!("==================================================");
println!(
"did1: {:?}, did2: {:?}, did3: {:?}, did4: {:?}",
did1, did2, did3, did4,
);
println!("==================================================");
swarm4.connect(did1).await?;
assert!(matches!(
node2.listen_once().await.unwrap().data,
Message::ConnectNodeSend(_)
));
assert!(matches!(
node1.listen_once().await.unwrap().data,
Message::ConnectNodeSend(_)
));
assert!(matches!(
node2.listen_once().await.unwrap().data,
Message::ConnectNodeReport(_)
));
assert!(matches!(
node4.listen_once().await.unwrap().data,
Message::ConnectNodeReport(_)
));
println!("=================Finish handshake here=================");
Ok(())
}
#[tokio::test]
async fn test_finger_when_disconnect() -> Result<()> {
let key1 = SecretKey::random();
let key2 = SecretKey::random();
let key3 = SecretKey::random();
let (did1, dht1, swarm1, node1, _path1) = prepare_node(key1).await;
let (did2, dht2, swarm2, node2, _path2) = prepare_node(key2).await;
let (_did3, _dht3, _swarm3, node3, _path3) = prepare_node(key3).await;
{
assert!(dht1.lock_finger()?.is_empty());
assert!(dht1.lock_finger()?.is_empty());
}
test_only_two_nodes_establish_connection(&node1, &node2).await?;
assert_no_more_msg(&node1, &node2, &node3).await;
assert_transports(swarm1.clone(), vec![did2]);
assert_transports(swarm2.clone(), vec![did1]);
{
let finger1 = dht1.lock_finger()?.clone().clone_finger();
let finger2 = dht2.lock_finger()?.clone().clone_finger();
assert!(finger1.into_iter().any(|x| x == Some(did2)));
assert!(finger2.into_iter().any(|x| x == Some(did1)));
}
println!("===================================");
println!("| test disconnect node1 and node2 |");
println!("===================================");
swarm1.disconnect(did2).await?;
let ev1 = node1.listen_once().await;
assert!(ev1.is_none());
#[cfg(not(feature = "wasm"))]
swarm2.get_transport(did1).unwrap().close().await.unwrap();
for _ in 1..10 {
println!("wait 3 seconds for node2's transport 2to1 closing");
sleep(Duration::from_secs(3)).await;
if let Some(t) = swarm2.get_transport(did1) {
if t.is_disconnected().await {
println!("transport 2to1 is disconnected!!!!");
break;
}
} else {
println!("transport 2to1 is disappeared!!!!");
break;
}
}
let ev2 = node2.listen_once().await.unwrap();
assert_eq!(ev2.addr, did2);
assert!(matches!(ev2.data, Message::LeaveDHT(LeaveDHT{did}) if did == did1));
assert_no_more_msg(&node1, &node2, &node3).await;
assert_transports(swarm1.clone(), vec![]);
assert_transports(swarm2.clone(), vec![]);
{
let finger1 = dht1.lock_finger()?.clone().clone_finger();
let finger2 = dht2.lock_finger()?.clone().clone_finger();
assert!(finger1.into_iter().all(|x| x.is_none()));
assert!(finger2.into_iter().all(|x| x.is_none()));
}
Ok(())
}
#[tokio::test]
async fn test_already_connect_fixture() -> Result<()> {
let keys = gen_ordered_keys(3);
let (key1, key2, key3) = (keys[0], keys[1], keys[2]);
let (did1, _dht1, swarm1, node1, _path1) = prepare_node(key1).await;
let (_did2, _dht2, _swarm2, node2, _path2) = prepare_node(key2).await;
let (did3, _dht3, swarm3, node3, _path3) = prepare_node(key3).await;
test_only_two_nodes_establish_connection(&node1, &node2).await?;
assert_no_more_msg(&node1, &node2, &node3).await;
test_only_two_nodes_establish_connection(&node3, &node2).await?;
assert_no_more_msg(&node1, &node2, &node3).await;
println!("node1 connect node2 twice here");
let _ = swarm1.connect(did3).await.unwrap();
let t_1_3_b = swarm1.connect(did3).await.unwrap();
let _ = node2.listen_once().await.unwrap();
let _ = node2.listen_once().await.unwrap();
let _ = node3.listen_once().await.unwrap();
let _ = node3.listen_once().await.unwrap();
let _ = node2.listen_once().await.unwrap();
let _ = node2.listen_once().await.unwrap();
let _ = node1.listen_once().await.unwrap();
let _ = node1.listen_once().await.unwrap();
println!("wait for handshake here");
sleep(Duration::from_secs(3)).await;
let ev3 = node3.listen_once().await.unwrap();
assert!(matches!(ev3.data, Message::JoinDHT(_)));
let _ = node3.listen_once().await.is_none();
let ev1 = node1.listen_once().await.unwrap();
assert!(matches!(ev1.data, Message::JoinDHT(_)));
let _ = node1.listen_once().await.is_none();
let t1_3 = swarm1.get_transport(did3).unwrap();
println!("transport is replace by second");
assert_eq!(t1_3.id, t_1_3_b.id);
let t3_1 = swarm3.get_transport(did1).unwrap();
assert!(t1_3.is_connected().await);
assert!(t3_1.is_connected().await);
Ok(())
}
}