use std::{
cmp::max,
collections::HashSet,
iter,
net::{IpAddr, Ipv4Addr, SocketAddr},
sync::Arc,
time::Duration,
};
use futures::{stream, FutureExt as _, StreamExt};
use tokio::time::timeout;
use tower::{discover::Change, Service, ServiceExt};
use zakura_chain::{
block,
parameters::{Network, NetworkUpgrade},
};
use crate::{
constants::{CURRENT_NETWORK_PROTOCOL_VERSION, DEFAULT_MAX_CONNS_PER_IP},
peer::{
ClientRequest, ClientTestHarness, ConnectedAddr, LoadTrackedClient, MinimumPeerVersion,
},
peer_set::{inventory_registry::InventoryStatus, stall_tracker::FIND_RESPONSE_STALL_THRESHOLD},
protocol::external::{types::Version, InventoryHash},
BoxError, PeerSocketAddr, Request, Response, SharedPeerError,
};
use indexmap::IndexMap;
use tokio::sync::watch;
use super::{PeerSetBuilder, PeerVersions};
#[test]
fn peer_set_ready_single_connection() {
let peer_versions = PeerVersions {
peer_versions: vec![Version::min_specified_for_upgrade(
&Network::Mainnet,
NetworkUpgrade::Nu6_2,
)],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
let mut client_handle = handles
.into_iter()
.next()
.expect("we always have at least one client");
assert!(client_handle
.try_to_receive_outbound_client_request()
.is_empty());
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.build();
let peer_ready_future = peer_set.ready();
#[allow(unknown_lints)]
#[allow(clippy::drop_non_drop)]
std::mem::drop(peer_ready_future);
let peer_ready1 = peer_set
.ready()
.await
.expect("peer set service is always ready");
assert!(client_handle
.try_to_receive_outbound_client_request()
.is_empty());
let fut = peer_ready1.call(Request::Peers);
assert!(matches!(
client_handle
.try_to_receive_outbound_client_request()
.request(),
Some(ClientRequest {
request: Request::Peers,
..
})
));
std::mem::drop(fut);
let peer_ready2 = peer_set
.ready()
.await
.expect("peer set service is always ready");
let _fut = peer_ready2.call(Request::MempoolTransactionIds);
assert!(matches!(
client_handle
.try_to_receive_outbound_client_request()
.request(),
Some(ClientRequest {
request: Request::MempoolTransactionIds,
..
})
));
});
}
#[test]
fn peer_set_ready_multiple_connections() {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version, peer_version, peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), 3);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.max_conns_per_ip(max(3, DEFAULT_MAX_CONNS_PER_IP))
.build();
let peer_ready = peer_set
.ready()
.await
.expect("peer set service is always ready");
assert_eq!(peer_ready.ready_services.len(), 3);
handles[0].stop_connection_task().await;
handles[1].stop_connection_task().await;
peer_set
.ready()
.await
.expect("peer set service is always ready");
handles[2].stop_connection_task().await;
let peer_ready = peer_set.ready();
assert!(timeout(Duration::from_secs(10), peer_ready).await.is_err());
});
}
#[test]
fn peer_set_rejects_connections_past_per_ip_limit() {
const NUM_PEER_VERSIONS: usize = crate::constants::DEFAULT_MAX_CONNS_PER_IP + 1;
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: [peer_version; NUM_PEER_VERSIONS].into_iter().collect(),
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), NUM_PEER_VERSIONS);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.build();
let peer_ready = peer_set
.ready()
.await
.expect("peer set service is always ready");
assert_eq!(
peer_ready.ready_services.len(),
crate::constants::DEFAULT_MAX_CONNS_PER_IP
);
});
}
#[test]
fn peer_set_route_inv_empty_registry() {
let test_hash = block::Hash([0; 32]);
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version, peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), 2);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.max_conns_per_ip(max(2, DEFAULT_MAX_CONNS_PER_IP))
.build();
let peer_ready = peer_set
.ready()
.await
.expect("peer set service is always ready");
assert_eq!(peer_ready.ready_services.len(), 2);
let sent_request = Request::BlocksByHash(iter::once(test_hash).collect());
let _fut = peer_ready.call(sent_request.clone());
let mut received_count = 0;
for mut handle in handles {
if let Some(ClientRequest { request, .. }) =
handle.try_to_receive_outbound_client_request().request()
{
assert_eq!(sent_request, request);
received_count += 1;
}
}
assert_eq!(received_count, 1);
});
}
#[test]
fn broadcast_all_queued_removes_banned_peers() {
let peer_versions = PeerVersions {
peer_versions: vec![Version::min_specified_for_upgrade(
&Network::Mainnet,
NetworkUpgrade::Nu6_2,
)],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let (discovered_peers, _handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.build();
let banned_ip: std::net::IpAddr = "127.0.0.1".parse().unwrap();
let mut bans_map: IndexMap<std::net::IpAddr, std::time::Instant> = IndexMap::new();
bans_map.insert(banned_ip, std::time::Instant::now());
let (bans_tx, bans_rx) = watch::channel(Arc::new(bans_map));
let _ = bans_tx;
peer_set.bans_receiver = bans_rx;
let banned_addr: PeerSocketAddr = SocketAddr::new(banned_ip, 1).into();
let mut remaining_peers = HashSet::new();
remaining_peers.insert(banned_addr);
let (sender, mut receiver) = tokio::sync::mpsc::channel(1);
peer_set.queued_broadcast_all = Some((Request::Peers, sender, remaining_peers));
peer_set.broadcast_all_queued();
if let Some((_req, _sender, remaining_peers)) = peer_set.queued_broadcast_all.take() {
assert!(remaining_peers.is_empty());
} else {
assert!(receiver.try_recv().is_ok());
}
});
}
#[test]
fn mined_block_gossip_to_unready_peer_is_delivered_not_canceled() {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version, peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), 2);
let mut handles = handles.into_iter();
let handle_1 = handles.next().expect("first peer harness");
let handle_2 = handles.next().expect("second peer harness");
let addr_1: PeerSocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 1).into();
let addr_2: PeerSocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 2).into();
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version)
.max_conns_per_ip(max(2, DEFAULT_MAX_CONNS_PER_IP))
.build();
{
let ready = peer_set.ready().await.expect("peer set is always ready");
assert_eq!(ready.ready_services.len(), 2);
}
for addr in [addr_1, addr_2] {
let svc = peer_set.take_ready_service(&addr).expect("peer is ready");
peer_set.push_unready(addr, svc);
}
let hash = block::Hash([7; 32]);
let broadcast_handle =
tokio::spawn(peer_set.broadcast_all(Request::AdvertiseBlockToAll(hash)));
let mut broadcast_finished = false;
for _ in 0..16 {
{
let _ = peer_set.ready().await.expect("peer set is always ready");
}
tokio::task::yield_now().await;
if broadcast_handle.is_finished() {
broadcast_finished = true;
break;
}
}
assert!(
broadcast_finished,
"the mined-block broadcast future should complete once queued deliveries drain",
);
broadcast_handle
.await
.expect("broadcast task should not panic")
.expect("broadcast_all should succeed");
for mut handle in [handle_1, handle_2] {
let client_request = handle
.try_to_receive_outbound_client_request()
.request()
.expect("each once-unready peer should receive the queued mined-block gossip");
assert!(
matches!(client_request.request, Request::AdvertiseBlockToAll(h) if h == hash),
"expected the mined-block advertisement, got {:?}",
client_request.request,
);
assert!(
!client_request.tx.is_canceled(),
"the queued send future must be spawned, not dropped: a dropped future \
cancels the response channel and the connection skips the block inv",
);
}
});
}
#[test]
fn mined_block_gossip_reaches_ready_peers() {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version, peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), 2);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version)
.max_conns_per_ip(max(2, DEFAULT_MAX_CONNS_PER_IP))
.build();
{
let ready = peer_set.ready().await.expect("peer set is always ready");
assert_eq!(ready.ready_services.len(), 2);
}
let hash = block::Hash([5; 32]);
let _broadcast_fut = peer_set.broadcast_all(Request::AdvertiseBlockToAll(hash));
let mut received = 0;
for mut handle in handles {
if let Some(client_request) = handle.try_to_receive_outbound_client_request().request()
{
assert!(
matches!(client_request.request, Request::AdvertiseBlockToAll(h) if h == hash),
"expected the mined-block advertisement, got {:?}",
client_request.request,
);
received += 1;
}
}
assert_eq!(
received, 2,
"both ready peers should receive the mined block"
);
});
}
#[test]
fn remove_unready_peer_clears_cancel_handle_and_updates_counts() {
let peer_versions = PeerVersions {
peer_versions: vec![Version::min_specified_for_upgrade(
&Network::Mainnet,
NetworkUpgrade::Nu6_2,
)],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let (discovered_peers, _handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.build();
let banned_ip: std::net::IpAddr = "127.0.0.1".parse().unwrap();
let mut bans_map: IndexMap<std::net::IpAddr, std::time::Instant> = IndexMap::new();
bans_map.insert(banned_ip, std::time::Instant::now());
let (_bans_tx, bans_rx) = watch::channel(Arc::new(bans_map));
peer_set.bans_receiver = bans_rx;
let banned_addr: PeerSocketAddr = SocketAddr::new(banned_ip, 1).into();
let (tx, _rx) =
crate::peer_set::set::oneshot::channel::<crate::peer_set::set::CancelClientWork>();
peer_set.cancel_handles.insert(banned_addr, tx);
assert_eq!(peer_set.num_peers_with_ip(banned_ip), 1);
peer_set.remove(&banned_addr);
assert!(!peer_set.cancel_handles.contains_key(&banned_addr));
assert_eq!(peer_set.num_peers_with_ip(banned_ip), 0);
});
}
#[test]
fn peer_set_route_inv_advertised_registry() {
peer_set_route_inv_advertised_registry_order(true);
peer_set_route_inv_advertised_registry_order(false);
}
fn peer_set_route_inv_advertised_registry_order(advertised_first: bool) {
let test_hash = block::Hash([0; 32]);
let test_inv = InventoryHash::Block(test_hash);
let test_peer = if advertised_first {
"127.0.0.1:1"
} else {
"127.0.0.1:2"
}
.parse()
.expect("unexpected invalid peer address");
let test_change = InventoryStatus::new_available(test_inv, test_peer);
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version, peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, mut handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), 2);
runtime.block_on(async move {
let (mut peer_set, mut peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.max_conns_per_ip(max(2, DEFAULT_MAX_CONNS_PER_IP))
.build();
peer_set_guard
.inventory_sender()
.as_mut()
.expect("unexpected missing inv sender")
.send(test_change)
.expect("unexpected dropped receiver");
let peer_ready = peer_set
.ready()
.await
.expect("peer set service is always ready");
assert_eq!(peer_ready.ready_services.len(), 2);
let sent_request = Request::BlocksByHash(iter::once(test_hash).collect());
let _fut = peer_ready.call(sent_request.clone());
let advertised_handle = if advertised_first {
&mut handles[0]
} else {
&mut handles[1]
};
if let Some(ClientRequest { request, .. }) = advertised_handle
.try_to_receive_outbound_client_request()
.request()
{
assert_eq!(sent_request, request);
} else {
panic!("inv request not routed to advertised peer");
}
let other_handle = if advertised_first {
&mut handles[1]
} else {
&mut handles[0]
};
assert!(
other_handle
.try_to_receive_outbound_client_request()
.request()
.is_none(),
"request routed to non-advertised peer",
);
});
}
#[test]
fn peer_set_route_inv_missing_registry() {
peer_set_route_inv_missing_registry_order(true);
peer_set_route_inv_missing_registry_order(false);
}
fn peer_set_route_inv_missing_registry_order(missing_first: bool) {
let test_hash = block::Hash([0; 32]);
let test_inv = InventoryHash::Block(test_hash);
let test_peer = if missing_first {
"127.0.0.1:1"
} else {
"127.0.0.1:2"
}
.parse()
.expect("unexpected invalid peer address");
let test_change = InventoryStatus::new_missing(test_inv, test_peer);
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version, peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, mut handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), 2);
runtime.block_on(async move {
let (mut peer_set, mut peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.max_conns_per_ip(max(2, DEFAULT_MAX_CONNS_PER_IP))
.build();
peer_set_guard
.inventory_sender()
.as_mut()
.expect("unexpected missing inv sender")
.send(test_change)
.expect("unexpected dropped receiver");
let peer_ready = peer_set
.ready()
.await
.expect("peer set service is always ready");
assert_eq!(peer_ready.ready_services.len(), 2);
let sent_request = Request::BlocksByHash(iter::once(test_hash).collect());
let _fut = peer_ready.call(sent_request.clone());
let missing_handle = if missing_first {
&mut handles[0]
} else {
&mut handles[1]
};
assert!(
missing_handle
.try_to_receive_outbound_client_request()
.request()
.is_none(),
"request routed to missing peer",
);
let other_handle = if missing_first {
&mut handles[1]
} else {
&mut handles[0]
};
if let Some(ClientRequest { request, .. }) = other_handle
.try_to_receive_outbound_client_request()
.request()
{
assert_eq!(sent_request, request);
} else {
panic!(
"inv request should have been routed to the only peer not missing the inventory"
);
}
});
}
#[test]
fn peer_set_route_inv_all_missing_fail() {
let test_hash = block::Hash([0; 32]);
let test_inv = InventoryHash::Block(test_hash);
let test_peer = "127.0.0.1:1"
.parse()
.expect("unexpected invalid peer address");
let test_change = InventoryStatus::new_missing(test_inv, test_peer);
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered_peers, mut handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
assert_eq!(handles.len(), 1);
runtime.block_on(async move {
let (mut peer_set, mut peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version.clone())
.build();
peer_set_guard
.inventory_sender()
.as_mut()
.expect("unexpected missing inv sender")
.send(test_change)
.expect("unexpected dropped receiver");
let peer_ready = peer_set
.ready()
.await
.expect("peer set service is always ready");
assert_eq!(peer_ready.ready_services.len(), 1);
let sent_request = Request::BlocksByHash(iter::once(test_hash).collect());
let response_fut = peer_ready.call(sent_request.clone());
let missing_handle = &mut handles[0];
assert!(
missing_handle
.try_to_receive_outbound_client_request()
.request().is_none(),
"request routed to missing peer",
);
let response = response_fut.await;
assert_eq!(
response
.expect_err("peer set should return an error (not a Response)")
.downcast_ref::<SharedPeerError>()
.expect("peer set should return a boxed SharedPeerError")
.inner_debug(),
"NotFoundRegistry([Block(block::Hash(\"0000000000000000000000000000000000000000000000000000000000000000\"))])"
);
});
}
#[test]
fn find_blocks_stall_not_tracked_when_at_tip() {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, best_tip) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
best_tip.send_best_tip_height(Some(block::Height(2_500_000)));
best_tip.send_estimated_distance_to_network_chain_tip(Some(0));
let mut handle = handles.into_iter().next().expect("there is one peer");
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version)
.build();
let request_count = FIND_RESPONSE_STALL_THRESHOLD + 1;
for _ in 0..request_count {
let peer_ready = peer_set.ready().await.expect("peer set is ready");
let response_fut = peer_ready.call(Request::FindBlocks {
known_blocks: vec![],
stop: None,
});
let client_request = handle
.try_to_receive_outbound_client_request()
.request()
.expect("peer received the request");
let _ = client_request.tx.send(Ok(Response::BlockHashes(vec![])));
response_fut.await.expect("response received");
}
assert!(
handle.wants_connection_heartbeats(),
"peer should not be disconnected when at tip"
);
});
}
#[test]
fn find_blocks_stall_not_tracked_for_zcashd_compat() {
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let sidecar_ip = Ipv4Addr::LOCALHOST;
let sidecar_addr: PeerSocketAddr =
SocketAddr::new(IpAddr::V6(sidecar_ip.to_ipv6_mapped()), 1).into();
let (sidecar, mut sidecar_handle) = ClientTestHarness::build()
.with_version(CURRENT_NETWORK_PROTOCOL_VERSION)
.with_connected_addr(ConnectedAddr::new_inbound_direct(sidecar_addr))
.finish();
let discovered_peers = stream::iter([Ok::<_, BoxError>(Change::Insert(
sidecar_addr,
sidecar.into(),
))])
.chain(stream::pending());
let (minimum_peer_version, best_tip) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
best_tip.send_best_tip_height(Some(block::Height(2_490_000)));
best_tip.send_estimated_distance_to_network_chain_tip(Some(10_000));
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_block_gossip_peer_ips(vec![sidecar_ip.into()])
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version)
.build();
for _ in 0..FIND_RESPONSE_STALL_THRESHOLD {
let peer_ready = peer_set.ready().await.expect("peer set is ready");
let response_fut = peer_ready.call(Request::FindBlocks {
known_blocks: vec![],
stop: None,
});
let client_request = sidecar_handle
.try_to_receive_outbound_client_request()
.request()
.expect("sidecar received the request");
let _ = client_request.tx.send(Ok(Response::BlockHashes(vec![])));
response_fut.await.expect("response received");
}
let _ = peer_set.ready().now_or_never();
assert!(
sidecar_handle.wants_connection_heartbeats(),
"zcashd-compat sidecar should not be disconnected by the sync stall detector"
);
});
}
#[test]
fn find_blocks_stall_tracked_when_syncing() {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, best_tip) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
best_tip.send_best_tip_height(Some(block::Height(2_490_000)));
best_tip.send_estimated_distance_to_network_chain_tip(Some(10_000));
let mut handle = handles.into_iter().next().expect("there is one peer");
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version)
.build();
for _ in 0..FIND_RESPONSE_STALL_THRESHOLD {
let peer_ready = peer_set.ready().await.expect("peer set is ready");
let response_fut = peer_ready.call(Request::FindBlocks {
known_blocks: vec![],
stop: None,
});
let client_request = handle
.try_to_receive_outbound_client_request()
.request()
.expect("peer received the request");
let _ = client_request.tx.send(Ok(Response::BlockHashes(vec![])));
response_fut.await.expect("response received");
}
let _ = peer_set.ready().now_or_never();
assert!(
!handle.wants_connection_heartbeats(),
"peer should be disconnected after stall threshold is reached while syncing"
);
});
}
#[test]
fn find_blocks_stall_tracked_when_tip_unknown() {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, _best_tip) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
let mut handle = handles.into_iter().next().expect("there is one peer");
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version)
.build();
for _ in 0..FIND_RESPONSE_STALL_THRESHOLD {
let peer_ready = peer_set.ready().await.expect("peer set is ready");
let response_fut = peer_ready.call(Request::FindBlocks {
known_blocks: vec![],
stop: None,
});
let client_request = handle
.try_to_receive_outbound_client_request()
.request()
.expect("peer received the request");
let _ = client_request.tx.send(Ok(Response::BlockHashes(vec![])));
response_fut.await.expect("response received");
}
let _ = peer_set.ready().now_or_never();
assert!(
!handle.wants_connection_heartbeats(),
"peer should be disconnected when tip is unknown and stall threshold is reached"
);
});
}
#[test]
fn find_blocks_stall_count_preserved_across_tip_transition() {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let peer_versions = PeerVersions {
peer_versions: vec![peer_version],
};
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
let (discovered_peers, handles) = peer_versions.mock_peer_discovery();
let (minimum_peer_version, best_tip) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
best_tip.send_best_tip_height(Some(block::Height(2_490_000)));
best_tip.send_estimated_distance_to_network_chain_tip(Some(10_000));
let mut handle = handles.into_iter().next().expect("there is one peer");
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered_peers)
.with_minimum_peer_version(minimum_peer_version)
.build();
for _ in 0..FIND_RESPONSE_STALL_THRESHOLD - 1 {
let peer_ready = peer_set.ready().await.expect("peer set is ready");
let response_fut = peer_ready.call(Request::FindBlocks {
known_blocks: vec![],
stop: None,
});
let client_request = handle
.try_to_receive_outbound_client_request()
.request()
.expect("peer received the request");
let _ = client_request.tx.send(Ok(Response::BlockHashes(vec![])));
response_fut.await.expect("response received");
}
best_tip.send_best_tip_height(Some(block::Height(2_500_000)));
best_tip.send_estimated_distance_to_network_chain_tip(Some(0));
let peer_ready = peer_set.ready().await.expect("peer set is ready");
let response_fut = peer_ready.call(Request::FindBlocks {
known_blocks: vec![],
stop: None,
});
let client_request = handle
.try_to_receive_outbound_client_request()
.request()
.expect("peer received the request");
let _ = client_request.tx.send(Ok(Response::BlockHashes(vec![])));
response_fut.await.expect("response received");
best_tip.send_estimated_distance_to_network_chain_tip(Some(10_000));
let peer_ready = peer_set.ready().await.expect("peer set is ready");
let response_fut = peer_ready.call(Request::FindBlocks {
known_blocks: vec![],
stop: None,
});
let client_request = handle
.try_to_receive_outbound_client_request()
.request()
.expect("peer received the request");
let _ = client_request.tx.send(Ok(Response::BlockHashes(vec![])));
response_fut.await.expect("response received");
let _ = peer_set.ready().now_or_never();
assert!(
!handle.wants_connection_heartbeats(),
"peer should be disconnected because its syncing stall count was preserved"
);
});
}
fn recv_advertise_block(handle: &mut ClientTestHarness) -> Option<block::Hash> {
match handle.try_to_receive_outbound_client_request().request() {
Some(ClientRequest {
request: Request::AdvertiseBlock(hash, _),
..
}) => Some(hash),
Some(other) => panic!("unexpected outbound request: {:?}", other.request),
None => None,
}
}
fn sidecar_and_ordinary_discovery() -> (
impl futures::Stream<Item = Result<Change<PeerSocketAddr, LoadTrackedClient>, BoxError>>,
PeerSocketAddr,
PeerSocketAddr,
ClientTestHarness,
ClientTestHarness,
) {
let peer_version = Version::min_specified_for_upgrade(&Network::Mainnet, NetworkUpgrade::Nu6_2);
let sidecar_addr: PeerSocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 1).into();
let ordinary_addr: PeerSocketAddr =
SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 2)), 2).into();
let (sidecar_client, sidecar_handle) = ClientTestHarness::build()
.with_version(peer_version)
.with_connected_addr(ConnectedAddr::InboundDirect { addr: sidecar_addr })
.finish();
let (ordinary_client, ordinary_handle) = ClientTestHarness::build()
.with_version(peer_version)
.with_connected_addr(ConnectedAddr::InboundDirect {
addr: ordinary_addr,
})
.finish();
let discovered = stream::iter([
Ok::<_, BoxError>(Change::Insert(
sidecar_addr,
LoadTrackedClient::from(sidecar_client),
)),
Ok::<_, BoxError>(Change::Insert(
ordinary_addr,
LoadTrackedClient::from(ordinary_client),
)),
])
.chain(stream::pending());
(
discovered,
sidecar_addr,
ordinary_addr,
sidecar_handle,
ordinary_handle,
)
}
#[test]
fn unready_sidecar_block_gossip_is_queued_not_dropped() {
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered, sidecar_addr, _ordinary_addr, mut sidecar_handle, mut ordinary_handle) =
sidecar_and_ordinary_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered)
.with_minimum_peer_version(minimum_peer_version)
.max_conns_per_ip(max(2, DEFAULT_MAX_CONNS_PER_IP))
.with_block_gossip_peer_ips(vec![IpAddr::V4(Ipv4Addr::LOCALHOST)])
.build();
{
let ready = peer_set.ready().await.expect("peer set is always ready");
assert_eq!(ready.ready_services.len(), 2);
}
let sidecar_svc = peer_set
.take_ready_service(&sidecar_addr)
.expect("sidecar is ready");
peer_set.push_unready(sidecar_addr, sidecar_svc);
assert!(peer_set.cancel_handles.contains_key(&sidecar_addr));
assert!(!peer_set.ready_services.contains_key(&sidecar_addr));
let hash = block::Hash([7; 32]);
let _fut = peer_set.route_block_broadcast(Request::AdvertiseBlock(hash, None));
assert_eq!(
recv_advertise_block(&mut ordinary_handle),
Some(hash),
"an ordinary ready peer should receive the block gossip immediately",
);
assert_eq!(
recv_advertise_block(&mut sidecar_handle),
None,
"an unready sidecar cannot be sent to synchronously",
);
let (queued_req, queued_peers) = peer_set
.queued_sidecar_block_gossip
.as_ref()
.expect("block gossip should be queued for the unready sidecar");
assert_eq!(*queued_req, Request::AdvertiseBlock(hash, None));
assert!(
queued_peers.contains(&sidecar_addr),
"the unready sidecar should be queued for redelivery"
);
});
}
#[test]
fn queued_sidecar_block_gossip_delivered_once_ready() {
let (runtime, _init_guard) = zakura_test::init_async();
let _guard = runtime.enter();
tokio::time::pause();
let (discovered, sidecar_addr, _ordinary_addr, mut sidecar_handle, _ordinary_handle) =
sidecar_and_ordinary_discovery();
let (minimum_peer_version, _best_tip_height) =
MinimumPeerVersion::with_mock_chain_tip(&Network::Mainnet);
runtime.block_on(async move {
let (mut peer_set, _peer_set_guard) = PeerSetBuilder::new()
.with_discover(discovered)
.with_minimum_peer_version(minimum_peer_version)
.max_conns_per_ip(max(2, DEFAULT_MAX_CONNS_PER_IP))
.with_block_gossip_peer_ips(vec![IpAddr::V4(Ipv4Addr::LOCALHOST)])
.build();
{
let ready = peer_set.ready().await.expect("peer set is always ready");
assert_eq!(ready.ready_services.len(), 2);
}
let sidecar_svc = peer_set
.take_ready_service(&sidecar_addr)
.expect("sidecar is ready");
peer_set.push_unready(sidecar_addr, sidecar_svc);
let hash = block::Hash([9; 32]);
let _fut = peer_set.route_block_broadcast(Request::AdvertiseBlock(hash, None));
assert!(peer_set.queued_sidecar_block_gossip.is_some());
assert_eq!(recv_advertise_block(&mut sidecar_handle), None);
let mut delivered = None;
for _ in 0..8 {
{
let _ = peer_set.ready().await.expect("peer set is always ready");
}
tokio::task::yield_now().await;
if let Some(received_hash) = recv_advertise_block(&mut sidecar_handle) {
delivered = Some(received_hash);
break;
}
}
assert_eq!(
delivered,
Some(hash),
"the sidecar should receive the queued block gossip once it is ready again",
);
assert!(
peer_set.queued_sidecar_block_gossip.is_none(),
"the queue should be cleared after delivery",
);
});
}