use super::*;
use crate::{
peer::DisconnectReason,
peer_manager::{conn_notifs_channel, ConnectionRequest},
transport::ConnectionMetadata,
};
use aptos_config::config::{Peer, PeerRole, PeerSet, HANDSHAKE_VERSION};
use aptos_crypto::{test_utils::TEST_SEED, x25519, Uniform};
use aptos_logger::info;
use aptos_time_service::{MockTimeService, TimeService};
use aptos_types::{account_address::AccountAddress, network_address::NetworkAddress};
use channel::{aptos_channel, message_queues::QueueStyle};
use futures::{executor::block_on, future, SinkExt};
use maplit::{hashmap, hashset};
use rand::rngs::StdRng;
use std::{io, str::FromStr};
use tokio_retry::strategy::FixedInterval;
const MAX_TEST_CONNECTIONS: usize = 3;
const CONNECTIVITY_CHECK_INTERVAL: Duration = Duration::from_secs(5);
const CONNECTION_DELAY: Duration = Duration::from_millis(100);
const MAX_CONNECTION_DELAY: Duration = Duration::from_secs(60);
const DEFAULT_BASE_ADDR: &str = "/ip4/127.0.0.1/tcp/9090";
const MAX_DELAY_WITH_JITTER: Duration = Duration::from_millis(
CONNECTION_DELAY.as_millis() as u64 + MAX_CONNECTION_DELAY_JITTER.as_millis() as u64,
);
fn network_address(addr_str: &'static str) -> NetworkAddress {
NetworkAddress::from_str(addr_str).unwrap()
}
fn network_address_with_pubkey(
addr_str: &'static str,
pubkey: x25519::PublicKey,
) -> NetworkAddress {
network_address(addr_str).append_prod_protos(pubkey, HANDSHAKE_VERSION)
}
fn test_peer(index: AccountAddress) -> (PeerId, Peer, x25519::PublicKey, NetworkAddress) {
test_peer_with_address(index, DEFAULT_BASE_ADDR)
}
fn test_peer_with_address(
peer_id: AccountAddress,
addr_str: &'static str,
) -> (PeerId, Peer, x25519::PublicKey, NetworkAddress) {
let pubkey = x25519::PrivateKey::generate_for_testing().public_key();
let pubkeys = hashset! { pubkey };
let addr = network_address_with_pubkey(addr_str, pubkey);
(
peer_id,
Peer::new(vec![addr.clone()], pubkeys, PeerRole::Validator),
pubkey,
addr,
)
}
fn update_peer_with_address(mut peer: Peer, addr_str: &'static str) -> (Peer, NetworkAddress) {
let keys: Vec<_> = peer.keys.iter().collect();
let key = *keys.first().unwrap();
let addr = network_address_with_pubkey(addr_str, *key);
peer.addresses = vec![addr.clone()];
(peer, addr)
}
struct TestHarness {
trusted_peers: Arc<RwLock<PeerSet>>,
mock_time: MockTimeService,
connection_reqs_rx: aptos_channel::Receiver<PeerId, ConnectionRequest>,
connection_notifs_tx: conn_notifs_channel::Sender,
conn_mgr_reqs_tx: channel::Sender<ConnectivityRequest>,
}
impl TestHarness {
fn new(seeds: PeerSet) -> (Self, ConnectivityManager<FixedInterval>) {
let network_context = NetworkContext::mock();
let time_service = TimeService::mock();
let (connection_reqs_tx, connection_reqs_rx) =
aptos_channel::new(QueueStyle::FIFO, 1, None);
let (connection_notifs_tx, connection_notifs_rx) = conn_notifs_channel::new();
let (conn_mgr_reqs_tx, conn_mgr_reqs_rx) = channel::new_test(0);
let trusted_peers = Arc::new(RwLock::new(HashMap::new()));
let conn_mgr = ConnectivityManager::new(
network_context,
time_service.clone(),
trusted_peers.clone(),
seeds,
ConnectionRequestSender::new(connection_reqs_tx),
connection_notifs_rx,
conn_mgr_reqs_rx,
CONNECTIVITY_CHECK_INTERVAL,
FixedInterval::new(CONNECTION_DELAY),
MAX_CONNECTION_DELAY,
Some(MAX_TEST_CONNECTIONS),
true,
);
let mock = Self {
trusted_peers,
mock_time: time_service.into_mock(),
connection_reqs_rx,
connection_notifs_tx,
conn_mgr_reqs_tx,
};
(mock, conn_mgr)
}
async fn trigger_connectivity_check(&self) {
info!("Advance time to trigger connectivity check");
self.mock_time
.advance_async(CONNECTIVITY_CHECK_INTERVAL)
.await;
}
async fn trigger_pending_dials(&self) {
info!("Advance time to trigger dial");
self.mock_time.advance_async(MAX_DELAY_WITH_JITTER).await;
}
async fn get_connected_size(&mut self) -> usize {
info!("Sending ConnectivityRequest::GetConnectedSize");
let (queue_size_tx, queue_size_rx) = oneshot::channel();
self.conn_mgr_reqs_tx
.send(ConnectivityRequest::GetConnectedSize(queue_size_tx))
.await
.unwrap();
queue_size_rx.await.unwrap()
}
async fn get_dial_queue_size(&mut self) -> usize {
info!("Sending ConnectivityRequest::GetDialQueueSize");
let (queue_size_tx, queue_size_rx) = oneshot::channel();
self.conn_mgr_reqs_tx
.send(ConnectivityRequest::GetDialQueueSize(queue_size_tx))
.await
.unwrap();
queue_size_rx.await.unwrap()
}
async fn send_new_peer_await_delivery(
&mut self,
peer_id: PeerId,
notif_peer_id: PeerId,
address: NetworkAddress,
) {
info!(
"Sending NewPeer notification for peer: {}",
peer_id.short_str()
);
let mut metadata = ConnectionMetadata::mock_with_role_and_origin(
notif_peer_id,
PeerRole::Unknown,
ConnectionOrigin::Outbound,
);
metadata.addr = address;
let notif = peer_manager::ConnectionNotification::NewPeer(metadata, NetworkContext::mock());
self.send_notification_await_delivery(peer_id, notif).await;
}
async fn send_lost_peer_await_delivery(&mut self, peer_id: PeerId, address: NetworkAddress) {
info!(
"Sending LostPeer notification for peer: {}",
peer_id.short_str()
);
let mut metadata = ConnectionMetadata::mock_with_role_and_origin(
peer_id,
PeerRole::Unknown,
ConnectionOrigin::Outbound,
);
metadata.addr = address;
let notif = peer_manager::ConnectionNotification::LostPeer(
metadata,
NetworkContext::mock(),
DisconnectReason::ConnectionLost,
);
self.send_notification_await_delivery(peer_id, notif).await;
}
async fn send_notification_await_delivery(
&mut self,
peer_id: PeerId,
notif: peer_manager::ConnectionNotification,
) {
let (delivered_tx, delivered_rx) = oneshot::channel();
self.connection_notifs_tx
.push_with_feedback(peer_id, notif, Some(delivered_tx))
.unwrap();
delivered_rx.await.unwrap();
}
async fn expect_disconnect_inner(
&mut self,
peer_id: PeerId,
address: NetworkAddress,
result: Result<(), PeerManagerError>,
) {
info!("Waiting to receive disconnect request");
let success = result.is_ok();
match self.connection_reqs_rx.next().await.unwrap() {
ConnectionRequest::DisconnectPeer(p, result_tx) => {
assert_eq!(peer_id, p);
result_tx.send(result).unwrap();
}
request => panic!(
"Unexpected ConnectionRequest, expected DisconnectPeer: {:?}",
request
),
}
if success {
self.send_lost_peer_await_delivery(peer_id, address).await;
}
}
async fn expect_disconnect_success(&mut self, peer_id: PeerId, address: NetworkAddress) {
self.expect_disconnect_inner(peer_id, address, Ok(())).await;
}
async fn expect_disconnect_fail(&mut self, peer_id: PeerId, address: NetworkAddress) {
let error = PeerManagerError::NotConnected(peer_id);
self.expect_disconnect_inner(peer_id, address, Err(error))
.await;
}
async fn wait_until_empty_dial_queue(&mut self) {
info!("Waiting for dial queue to be empty");
while self.get_dial_queue_size().await > 0 {}
}
async fn expect_one_dial_inner(
&mut self,
result: Result<(), PeerManagerError>,
) -> (PeerId, NetworkAddress) {
info!("Waiting to receive dial request");
let success = result.is_ok();
let (peer_id, address) = match self.connection_reqs_rx.next().await.unwrap() {
ConnectionRequest::DialPeer(peer_id, address, result_tx) => {
result_tx.send(result).unwrap();
(peer_id, address)
}
request => panic!(
"Unexpected ConnectionRequest, expected DialPeer: {:?}",
request
),
};
if success {
self.send_new_peer_await_delivery(peer_id, peer_id, address.clone())
.await;
}
(peer_id, address)
}
async fn expect_one_dial(
&mut self,
expected_peer_id: PeerId,
expected_address: NetworkAddress,
result: Result<(), PeerManagerError>,
) {
let (peer_id, address) = self.expect_one_dial_inner(result).await;
assert_eq!(peer_id, expected_peer_id);
assert_eq!(address, expected_address);
self.wait_until_empty_dial_queue().await;
}
async fn expect_one_dial_success(
&mut self,
expected_peer_id: PeerId,
expected_address: NetworkAddress,
) {
self.expect_one_dial(expected_peer_id, expected_address, Ok(()))
.await;
}
async fn expect_one_dial_fail(
&mut self,
expected_peer_id: PeerId,
expected_address: NetworkAddress,
) {
let error = PeerManagerError::IoError(io::Error::from(io::ErrorKind::ConnectionRefused));
self.expect_one_dial(expected_peer_id, expected_address, Err(error))
.await;
}
async fn expect_num_dials(&mut self, num_expected: usize) {
for _ in 0..num_expected {
let _ = self.expect_one_dial_inner(Ok(())).await;
}
self.wait_until_empty_dial_queue().await;
}
async fn send_update_discovered_peers(&mut self, src: DiscoverySource, peers: PeerSet) {
info!("Sending UpdateDiscoveredPeers");
self.conn_mgr_reqs_tx
.send(ConnectivityRequest::UpdateDiscoveredPeers(src, peers))
.await
.unwrap();
}
}
#[test]
fn connect_to_seeds_on_startup() {
let (seed_peer_id, seed_peer, _, seed_addr) = test_peer(AccountAddress::ONE);
let seeds: PeerSet = hashmap! {seed_peer_id => seed_peer.clone()};
let (mut mock, conn_mgr) = TestHarness::new(seeds.clone());
let test = async move {
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(seed_peer_id, seed_addr.clone())
.await;
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, seeds)
.await;
mock.trigger_connectivity_check().await;
assert_eq!(0, mock.get_dial_queue_size().await);
let (new_seed, new_seed_addr) =
update_peer_with_address(seed_peer, "/ip4/127.0.1.1/tcp/8080");
let update = hashmap! {seed_peer_id => new_seed};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.send_lost_peer_await_delivery(seed_peer_id, seed_addr.clone())
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(seed_peer_id, new_seed_addr).await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(seed_peer_id, seed_addr).await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn addr_change() {
let (other_peer_id, other_peer, _, other_addr) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let update = hashmap! {other_peer_id => other_peer.clone()};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update.clone())
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr.clone())
.await;
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.trigger_connectivity_check().await;
assert_eq!(0, mock.get_dial_queue_size().await);
let (other_peer_new, other_addr_new) =
update_peer_with_address(other_peer, "/ip4/127.0.1.1/tcp/8080");
let update = hashmap! {other_peer_id => other_peer_new};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.trigger_connectivity_check().await;
assert_eq!(1, mock.get_connected_size().await);
mock.send_lost_peer_await_delivery(other_peer_id, other_addr_new.clone())
.await;
assert_eq!(0, mock.get_connected_size().await);
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr_new)
.await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn lost_connection() {
let (other_peer_id, other_peer, _, other_addr) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let update = hashmap! {other_peer_id => other_peer};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr.clone())
.await;
mock.send_lost_peer_await_delivery(other_peer_id, other_addr.clone())
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr.clone())
.await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn disconnect() {
let (other_peer_id, other_peer, _, other_addr) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let peers = hashmap! {other_peer_id => other_peer.clone()};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr.clone())
.await;
let mut peer = other_peer;
peer.keys = HashSet::new();
peer.addresses = vec![network_address(DEFAULT_BASE_ADDR)];
let peers = hashmap! {other_peer_id => peer};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers)
.await;
mock.trigger_connectivity_check().await;
mock.expect_disconnect_success(other_peer_id, other_addr)
.await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn retry_on_failure() {
let (other_peer_id, peer, _, other_addr) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let peers = hashmap! {other_peer_id => peer.clone()};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(other_peer_id, other_addr.clone())
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr.clone())
.await;
let mut peer = peer;
peer.keys = HashSet::new();
peer.addresses = vec![network_address(DEFAULT_BASE_ADDR)];
let peers = hashmap! {other_peer_id => peer};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers)
.await;
mock.trigger_connectivity_check().await;
mock.expect_disconnect_fail(other_peer_id, other_addr.clone())
.await;
mock.trigger_connectivity_check().await;
mock.expect_disconnect_success(other_peer_id, other_addr.clone())
.await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn no_op_requests() {
let (other_peer_id, peer, _, other_addr) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let peers = hashmap! {other_peer_id => peer.clone()};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(other_peer_id, other_addr.clone())
.await;
mock.send_new_peer_await_delivery(other_peer_id, other_peer_id, other_addr.clone())
.await;
mock.trigger_connectivity_check().await;
let mut peer = peer;
peer.keys = HashSet::new();
peer.addresses = vec![network_address(DEFAULT_BASE_ADDR)];
let peers = hashmap! {other_peer_id => peer};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers)
.await;
mock.trigger_connectivity_check().await;
mock.expect_disconnect_fail(other_peer_id, other_addr.clone())
.await;
mock.send_lost_peer_await_delivery(other_peer_id, other_addr)
.await;
mock.trigger_connectivity_check().await;
assert_eq!(0, mock.get_connected_size().await);
assert_eq!(0, mock.get_dial_queue_size().await);
};
block_on(future::join(conn_mgr.start(), test));
}
fn generate_account_address(val: usize) -> AccountAddress {
let mut addr = [0u8; AccountAddress::LENGTH];
let array = val.to_be_bytes();
addr[AccountAddress::LENGTH - array.len()..].copy_from_slice(&array);
AccountAddress::new(addr)
}
#[test]
fn backoff_on_failure() {
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let (peer_id_a, peer_a, _, peer_a_addr) = test_peer(AccountAddress::ONE);
let (peer_id_b, peer_b, _, peer_b_addr) = test_peer(generate_account_address(2));
let peers = hashmap! {peer_id_a => peer_a, peer_id_b => peer_b};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers)
.await;
mock.send_new_peer_await_delivery(peer_id_b, peer_id_b, peer_b_addr)
.await;
for _ in 0..10 {
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(peer_id_a, peer_a_addr.clone())
.await;
}
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(peer_id_a, peer_a_addr).await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn multiple_addrs_basic() {
let (other_peer_id, mut peer, pubkey, _) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let other_addr_1 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9091", pubkey);
let other_addr_2 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9092", pubkey);
peer.addresses = vec![other_addr_1.clone(), other_addr_2.clone()];
let update = hashmap! {other_peer_id => peer};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(other_peer_id, other_addr_1).await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr_2)
.await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn multiple_addrs_wrapping() {
let (other_peer_id, mut peer, pubkey, _) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let other_addr_1 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9091", pubkey);
let other_addr_2 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9092", pubkey);
peer.addresses = vec![other_addr_1.clone(), other_addr_2.clone()];
let update = hashmap! {other_peer_id => peer};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(other_peer_id, other_addr_1.clone())
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(other_peer_id, other_addr_2).await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr_1)
.await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn multiple_addrs_shrinking() {
let (other_peer_id, mut peer, pubkey, _) = test_peer(AccountAddress::ZERO);
let (mut mock, conn_mgr) = TestHarness::new(HashMap::new());
let test = async move {
let other_addr_1 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9091", pubkey);
let other_addr_2 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9092", pubkey);
let other_addr_3 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9093", pubkey);
peer.addresses = vec![other_addr_1.clone(), other_addr_2, other_addr_3];
let update = hashmap! {other_peer_id => peer.clone()};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_fail(other_peer_id, other_addr_1).await;
let other_addr_4 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9094", pubkey);
let other_addr_5 = network_address_with_pubkey("/ip4/127.0.0.1/tcp/9095", pubkey);
peer.addresses = vec![other_addr_4.clone(), other_addr_5];
let update = hashmap! {other_peer_id => peer};
mock.send_update_discovered_peers(DiscoverySource::OnChainValidatorSet, update)
.await;
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_one_dial_success(other_peer_id, other_addr_4)
.await;
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn public_connection_limit() {
let mut seeds = HashMap::new();
for i in 0..=MAX_TEST_CONNECTIONS {
let (peer_id, peer, _, _) = test_peer(generate_account_address(i));
seeds.insert(peer_id, peer);
}
let (mut mock, conn_mgr) = TestHarness::new(seeds);
let test = async move {
mock.trigger_connectivity_check().await;
mock.trigger_pending_dials().await;
mock.expect_num_dials(MAX_TEST_CONNECTIONS).await;
assert_eq!(MAX_TEST_CONNECTIONS, mock.get_connected_size().await);
mock.trigger_connectivity_check().await;
assert_eq!(0, mock.get_dial_queue_size().await);
};
block_on(future::join(conn_mgr.start(), test));
}
#[test]
fn basic_update_discovered_peers() {
let mut rng = StdRng::from_seed(TEST_SEED);
let (mock, mut conn_mgr) = TestHarness::new(HashMap::new());
let trusted_peers = mock.trusted_peers;
let peer_id_a = AccountAddress::ZERO;
let peer_id_b = AccountAddress::ONE;
let addr_a = network_address("/ip4/127.0.0.1/tcp/9090");
let addr_b = network_address("/ip4/127.0.0.1/tcp/9091");
let pubkey_1 = x25519::PrivateKey::generate(&mut rng).public_key();
let pubkey_2 = x25519::PrivateKey::generate(&mut rng).public_key();
let pubkeys_1 = hashset! {pubkey_1};
let pubkeys_2 = hashset! {pubkey_2};
let pubkeys_1_2 = hashset! {pubkey_1, pubkey_2};
let peer_a1 = Peer::new(vec![addr_a.clone()], pubkeys_1.clone(), PeerRole::Validator);
let peer_a2 = Peer::new(vec![addr_a.clone()], pubkeys_2, PeerRole::Validator);
let peer_b1 = Peer::new(vec![addr_b], pubkeys_1, PeerRole::Validator);
let peer_a_1_2 = Peer::new(vec![addr_a], pubkeys_1_2, PeerRole::Validator);
let peers_empty = PeerSet::new();
let peers_1 = hashmap! {peer_id_a => peer_a1};
let peers_2 = hashmap! {peer_id_a => peer_a2, peer_id_b => peer_b1.clone()};
let peers_1_2 = hashmap! {peer_id_a => peer_a_1_2, peer_id_b => peer_b1};
conn_mgr.handle_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers_1.clone());
assert_eq!(*trusted_peers.read(), peers_1);
conn_mgr.handle_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers_1.clone());
assert_eq!(*trusted_peers.read(), peers_1);
conn_mgr
.handle_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers_empty.clone());
assert_eq!(*trusted_peers.read(), peers_empty);
conn_mgr.handle_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers_1.clone());
assert_eq!(*trusted_peers.read(), peers_1);
conn_mgr.handle_update_discovered_peers(DiscoverySource::Config, peers_2);
assert_eq!(*trusted_peers.read(), peers_1_2);
conn_mgr
.handle_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers_1_2.clone());
assert_eq!(*trusted_peers.read(), peers_1_2);
conn_mgr.handle_update_discovered_peers(DiscoverySource::Config, peers_1_2.clone());
assert_eq!(*trusted_peers.read(), peers_1_2);
conn_mgr.handle_update_discovered_peers(DiscoverySource::Config, peers_empty.clone());
assert_eq!(*trusted_peers.read(), peers_1_2);
conn_mgr
.handle_update_discovered_peers(DiscoverySource::OnChainValidatorSet, peers_empty.clone());
assert_eq!(*trusted_peers.read(), peers_empty);
conn_mgr.handle_update_discovered_peers(DiscoverySource::Config, peers_empty.clone());
assert_eq!(*trusted_peers.read(), peers_empty);
}