Skip to content

Commit b05badb

Browse files
authored
Gate sync peer selection on per-protocol concurrent-request limit (sigp#9456)
Co-Authored-By: dapplion <35266934+dapplion@users.noreply.github.com>
1 parent 10568b1 commit b05badb

3 files changed

Lines changed: 59 additions & 77 deletions

File tree

beacon_node/lighthouse_network/src/rpc/mod.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ mod response_limiter;
4646
mod self_limiter;
4747

4848
// Maximum number of concurrent requests per protocol ID that a client may issue.
49-
const MAX_CONCURRENT_REQUESTS: usize = 2;
49+
pub const MAX_CONCURRENT_REQUESTS: usize = 2;
5050

5151
/// Composite trait for a request id.
5252
pub trait ReqId: Send + 'static + std::fmt::Debug + Copy + Clone {}

beacon_node/network/src/sync/network_context.rs

Lines changed: 49 additions & 70 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,9 @@ use fnv::FnvHashMap;
2727
use lighthouse_network::rpc::methods::{
2828
BlobsByRangeRequest, DataColumnsByRangeRequest, PayloadEnvelopesByRangeRequest,
2929
};
30-
use lighthouse_network::rpc::{BlocksByRangeRequest, GoodbyeReason, RPCError, RequestType};
30+
use lighthouse_network::rpc::{
31+
BlocksByRangeRequest, GoodbyeReason, MAX_CONCURRENT_REQUESTS, RPCError, RequestType,
32+
};
3133
pub use lighthouse_network::service::api_types::RangeRequestId;
3234
use lighthouse_network::service::api_types::{
3335
AppRequestId, BlobsByRangeRequestId, BlocksByRangeRequestId, ComponentsByRangeRequestId,
@@ -40,8 +42,8 @@ use lighthouse_network::{Client, NetworkGlobals, PeerAction, PeerId, ReportSourc
4042
use parking_lot::RwLock;
4143
pub use requests::LookupVerifyError;
4244
use requests::{
43-
ActiveRequests, BlobsByRangeRequestItems, BlocksByRangeRequestItems, BlocksByRootRequestItems,
44-
DataColumnsByRangeRequestItems, DataColumnsByRootRequestItems,
45+
ActiveRequestItems, ActiveRequests, BlobsByRangeRequestItems, BlocksByRangeRequestItems,
46+
BlocksByRootRequestItems, DataColumnsByRangeRequestItems, DataColumnsByRootRequestItems,
4547
PayloadEnvelopesByRangeRequestItems, PayloadEnvelopesByRootRequestItems,
4648
};
4749
#[cfg(test)]
@@ -100,6 +102,30 @@ pub type RpcResponseResult<T> = Result<(T, Duration), RpcResponseError>;
100102
pub type CustodyByRootResult<T> =
101103
Result<DownloadResult<DataColumnSidecarList<T>>, RpcResponseError>;
102104

105+
/// Per-peer count of active requests for a single protocol, to keep peer selection within
106+
/// `MAX_CONCURRENT_REQUESTS` concurrent requests per protocol ID.
107+
struct ActiveRequestsPerPeer {
108+
count_by_peer: HashMap<PeerId, usize>,
109+
}
110+
111+
impl ActiveRequestsPerPeer {
112+
fn new<K, T>(requests: &ActiveRequests<K, T>) -> Self
113+
where
114+
K: Copy + Eq + std::hash::Hash + std::fmt::Display,
115+
T: ActiveRequestItems,
116+
{
117+
let mut count_by_peer = HashMap::<PeerId, usize>::new();
118+
for peer_id in requests.iter_request_peers() {
119+
*count_by_peer.entry(peer_id).or_default() += 1;
120+
}
121+
Self { count_by_peer }
122+
}
123+
124+
fn at_concurrency_limit(&self, peer_id: &PeerId) -> bool {
125+
self.count_by_peer.get(peer_id).copied().unwrap_or(0) >= MAX_CONCURRENT_REQUESTS
126+
}
127+
}
128+
103129
#[derive(Debug)]
104130
#[allow(private_interfaces)]
105131
pub enum RpcResponseError {
@@ -440,47 +466,6 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
440466
}
441467
}
442468

443-
fn active_request_count_by_peer(&self) -> HashMap<PeerId, usize> {
444-
let Self {
445-
network_send: _,
446-
request_id: _,
447-
blocks_by_root_requests,
448-
payload_envelopes_by_root_requests,
449-
data_columns_by_root_requests,
450-
blocks_by_range_requests,
451-
blobs_by_range_requests,
452-
data_columns_by_range_requests,
453-
payload_envelopes_by_range_requests,
454-
// custody_by_root_requests is a meta request of data_columns_by_root_requests
455-
custody_by_root_requests: _,
456-
// components_by_range_requests is a meta request of various _by_range requests
457-
components_by_range_requests: _,
458-
custody_backfill_data_column_batch_requests: _,
459-
execution_engine_state: _,
460-
network_beacon_processor: _,
461-
chain: _,
462-
fork_context: _,
463-
// Don't use a fallback match. We want to be sure that all requests are considered when
464-
// adding new ones
465-
} = self;
466-
467-
let mut active_request_count_by_peer = HashMap::<PeerId, usize>::new();
468-
469-
for peer_id in blocks_by_root_requests
470-
.iter_request_peers()
471-
.chain(payload_envelopes_by_root_requests.iter_request_peers())
472-
.chain(data_columns_by_root_requests.iter_request_peers())
473-
.chain(blocks_by_range_requests.iter_request_peers())
474-
.chain(blobs_by_range_requests.iter_request_peers())
475-
.chain(data_columns_by_range_requests.iter_request_peers())
476-
.chain(payload_envelopes_by_range_requests.iter_request_peers())
477-
{
478-
*active_request_count_by_peer.entry(peer_id).or_default() += 1;
479-
}
480-
481-
active_request_count_by_peer
482-
}
483-
484469
/// Retries only the specified failed columns by requesting them again.
485470
///
486471
/// Note: This function doesn't retry the whole batch, but retries specific requests within
@@ -507,8 +492,6 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
507492
return Err("request id not present".to_string());
508493
};
509494

510-
let active_request_count_by_peer = self.active_request_count_by_peer();
511-
512495
debug!(
513496
?failed_columns,
514497
?id,
@@ -518,12 +501,7 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
518501

519502
// Attempt to find all required custody peers to request the failed columns from
520503
let columns_by_range_peers_to_request = self
521-
.select_columns_by_range_peers_to_request(
522-
failed_columns,
523-
peers,
524-
active_request_count_by_peer,
525-
peers_to_deprioritize,
526-
)
504+
.select_columns_by_range_peers_to_request(failed_columns, peers, peers_to_deprioritize)
527505
.map_err(|e| format!("{:?}", e))?;
528506

529507
// Reuse the id for the request that received partially correct responses
@@ -581,16 +559,16 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
581559
column_peers = column_peers.len()
582560
);
583561
let _guard = range_request_span.clone().entered();
584-
let active_request_count_by_peer = self.active_request_count_by_peer();
562+
let blocks_by_range_per_peer = ActiveRequestsPerPeer::new(&self.blocks_by_range_requests);
585563

586564
let Some(block_peer) = block_peers
587565
.iter()
588566
.map(|peer| {
589567
(
590568
// If contains -> 1 (order after), not contains -> 0 (order first)
591569
peers_to_deprioritize.contains(peer),
592-
// Prefer peers with less overall requests
593-
active_request_count_by_peer.get(peer).copied().unwrap_or(0),
570+
// Strictly de-prioritize peers already at the per-protocol concurrency limit
571+
blocks_by_range_per_peer.at_concurrency_limit(peer),
594572
// Random factor to break ties, otherwise the PeerID breaks ties
595573
rand::random::<u32>(),
596574
peer,
@@ -620,7 +598,6 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
620598
Some(self.select_columns_by_range_peers_to_request(
621599
&column_indexes,
622600
column_peers,
623-
active_request_count_by_peer,
624601
peers_to_deprioritize,
625602
)?)
626603
} else {
@@ -692,6 +669,9 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
692669
let payloads_req_id =
693670
if matches!(batch_type, ByRangeRequestType::BlocksAndEnvelopesAndColumns) {
694671
Some(self.send_payload_envelopes_by_range_request(
672+
// Peer selection: for a given peer, the count of sent blocks_by_range requests
673+
// equals the count of sent payloads_by_range requests. So we are under the
674+
// concurrency limit for payloads_by_range requests
695675
block_peer,
696676
PayloadEnvelopesByRangeRequest {
697677
start_slot: *request.start_slot(),
@@ -731,10 +711,11 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
731711
&self,
732712
custody_indexes: &HashSet<ColumnIndex>,
733713
peers: &HashSet<PeerId>,
734-
active_request_count_by_peer: HashMap<PeerId, usize>,
735714
peers_to_deprioritize: &HashSet<PeerId>,
736715
) -> Result<HashMap<PeerId, Vec<ColumnIndex>>, RpcRequestSendError> {
737716
let mut columns_to_request_by_peer = HashMap::<PeerId, Vec<ColumnIndex>>::new();
717+
let data_columns_by_range_per_peer =
718+
ActiveRequestsPerPeer::new(&self.data_columns_by_range_requests);
738719

739720
for column_index in custody_indexes {
740721
// Strictly consider peers that are custodials of this column AND are part of this
@@ -750,12 +731,10 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
750731
(
751732
// If contains -> 1 (order after), not contains -> 0 (order first)
752733
peers_to_deprioritize.contains(peer),
753-
// Prefer peers with less overall requests
754-
// Also account for requests that are not yet issued tracked in peer_id_to_request_map
755-
// We batch requests to the same peer, so count existance in the
756-
// `columns_to_request_by_peer` as a single 1 request.
757-
active_request_count_by_peer.get(peer).copied().unwrap_or(0)
758-
+ columns_to_request_by_peer.get(peer).map(|_| 1).unwrap_or(0),
734+
// Strictly de-prioritize peers already at the per-protocol concurrency limit
735+
// Note: do not account for to-be-sent requests on
736+
// `data_columns_by_range_by_peer` as we always send at most one request
737+
data_columns_by_range_per_peer.at_concurrency_limit(peer),
759738
// Random factor to break ties, otherwise the PeerID breaks ties
760739
rand::random::<u32>(),
761740
peer,
@@ -881,14 +860,14 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
881860
lookup_peers: Arc<RwLock<HashSet<PeerId>>>,
882861
block_root: Hash256,
883862
) -> Result<LookupRequestResult<Arc<SignedBeaconBlock<T::EthSpec>>>, RpcRequestSendError> {
884-
let active_request_count_by_peer = self.active_request_count_by_peer();
863+
let blocks_by_root_per_peer = ActiveRequestsPerPeer::new(&self.blocks_by_root_requests);
885864
let Some(peer_id) = lookup_peers
886865
.read()
887866
.iter()
888867
.map(|peer| {
889868
(
890-
// Prefer peers with less overall requests
891-
active_request_count_by_peer.get(peer).copied().unwrap_or(0),
869+
// Strictly de-prioritize peers already at the per-protocol concurrency limit
870+
blocks_by_root_per_peer.at_concurrency_limit(peer),
892871
// Random factor to break ties, otherwise the PeerID breaks ties
893872
rand::random::<u32>(),
894873
peer,
@@ -1001,13 +980,15 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
1001980
));
1002981
}
1003982

1004-
let active_request_count_by_peer = self.active_request_count_by_peer();
983+
let payload_envelopes_by_root_per_peer =
984+
ActiveRequestsPerPeer::new(&self.payload_envelopes_by_root_requests);
1005985
let Some(peer_id) = lookup_peers
1006986
.read()
1007987
.iter()
1008988
.map(|peer| {
1009989
(
1010-
active_request_count_by_peer.get(peer).copied().unwrap_or(0),
990+
// Strictly de-prioritize peers already at the per-protocol concurrency limit
991+
payload_envelopes_by_root_per_peer.at_concurrency_limit(peer),
1011992
rand::random::<u32>(),
1012993
peer,
1013994
)
@@ -1757,7 +1738,6 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
17571738
peers: &HashSet<PeerId>,
17581739
peers_to_deprioritize: &HashSet<PeerId>,
17591740
) -> Result<CustodyBackFillBatchRequestId, RpcRequestSendError> {
1760-
let active_request_count_by_peer = self.active_request_count_by_peer();
17611741
// Attempt to find all required custody peers before sending any request or creating an ID
17621742
let columns_by_range_peers_to_request = {
17631743
let column_indexes = self
@@ -1770,7 +1750,6 @@ impl<T: BeaconChainTypes> SyncNetworkContext<T> {
17701750
self.select_columns_by_range_peers_to_request(
17711751
&column_indexes,
17721752
peers,
1773-
active_request_count_by_peer,
17741753
peers_to_deprioritize,
17751754
)?
17761755
};

beacon_node/network/src/sync/network_context/custody.rs

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,9 @@ use tracing::{Span, debug, debug_span, warn};
1616
use types::{DataColumnSidecar, Hash256, Slot, data::ColumnIndex};
1717
use types::{DataColumnSidecarList, EthSpec};
1818

19-
use super::{LookupRequestResult, PeerGroup, RpcResponseResult, SyncNetworkContext};
19+
use super::{
20+
ActiveRequestsPerPeer, LookupRequestResult, PeerGroup, RpcResponseResult, SyncNetworkContext,
21+
};
2022

2123
const MAX_STALE_NO_PEERS_DURATION: Duration = Duration::from_secs(30);
2224

@@ -237,7 +239,8 @@ impl<T: BeaconChainTypes> ActiveCustodyRequest<T> {
237239
)));
238240
}
239241

240-
let active_request_count_by_peer = cx.active_request_count_by_peer();
242+
let data_columns_by_root_per_peer =
243+
ActiveRequestsPerPeer::new(&cx.data_columns_by_root_requests);
241244
let mut columns_to_request_by_peer = HashMap::<PeerId, Vec<ColumnIndex>>::new();
242245
let mut columns_without_peers = vec![];
243246
let lookup_peers = self.lookup_peers.read();
@@ -255,7 +258,7 @@ impl<T: BeaconChainTypes> ActiveCustodyRequest<T> {
255258

256259
let peer_to_request = self.select_column_peer(
257260
cx,
258-
&active_request_count_by_peer,
261+
&data_columns_by_root_per_peer,
259262
&lookup_peers,
260263
*column_index,
261264
&random_state,
@@ -360,7 +363,7 @@ impl<T: BeaconChainTypes> ActiveCustodyRequest<T> {
360363
fn select_column_peer(
361364
&self,
362365
cx: &mut SyncNetworkContext<T>,
363-
active_request_count_by_peer: &HashMap<PeerId, usize>,
366+
data_columns_by_root_per_peer: &ActiveRequestsPerPeer,
364367
lookup_peers: &HashSet<PeerId>,
365368
column_index: ColumnIndex,
366369
random_state: &RandomState,
@@ -377,12 +380,12 @@ impl<T: BeaconChainTypes> ActiveCustodyRequest<T> {
377380
})
378381
.map(|peer| {
379382
(
383+
// Strictly de-prioritize peers already at the per-protocol concurrency limit
384+
data_columns_by_root_per_peer.at_concurrency_limit(peer),
380385
// Prioritize peers that claim to know have imported this block
381386
if lookup_peers.contains(peer) { 0 } else { 1 },
382387
// De-prioritize peers that we have already attempted to download from
383388
self.peer_attempts.get(peer).copied().unwrap_or(0),
384-
// Prefer peers with fewer requests to load balance across peers.
385-
active_request_count_by_peer.get(peer).copied().unwrap_or(0),
386389
// The hash ensures consistent peer ordering within this request
387390
// to avoid fragmentation while varying selection across different requests.
388391
random_state.hash_one(peer),

0 commit comments

Comments
 (0)